Merge pull request #11 from Aignosi/SIENTIAPDE-1110-criar-testes-e-2-e

Sientiapde 1110 criar testes e 2 e
This commit is contained in:
Matheus Demoner
2025-07-08 19:11:58 -03:00
committed by GitHub
7 changed files with 83 additions and 40 deletions

View File

@@ -72,3 +72,10 @@ pytest --cov=ingestor
``` ```
pytest --cov=ingestor --cov-report=html pytest --cov=ingestor --cov-report=html
``` ```
#PR shortcut
```
git log origin/main..HEAD --no-merges > git_log
```
Prompt:
Write a summary of PR changes in markdown. Be objective and direct. Write to file

View File

@@ -258,17 +258,27 @@ class Ingestor:
old_managed_tags, old_managed_tags,
) )
keys = set(new_managed_tags) | set(old_managed_tags)
changes = {k: (new_managed_tags.get(k), old_managed_tags.get(k))
for k in keys if new_managed_tags.get(k) != old_managed_tags.get(k)}
self.logger.debug("Changes: %s", changes)
for slot, config in new_managed_tags.items(): for slot, config in new_managed_tags.items():
if slot not in old_managed_tags: if slot not in old_managed_tags:
self.logger.debug("Subscribing to new slot %s", slot)
self.ingestor_manager.subscribe_to_tags({slot: config}) self.ingestor_manager.subscribe_to_tags({slot: config})
continue continue
if config != old_managed_tags[slot]: if config != old_managed_tags[slot]:
self.logger.debug("Resubscribing to slot %s", slot)
self.ingestor_manager.unsubscribe_slot(slot) self.ingestor_manager.unsubscribe_slot(slot)
self.ingestor_manager.subscribe_to_tags({slot: config}) self.ingestor_manager.subscribe_to_tags({slot: config})
for slot in old_managed_tags.keys(): for slot in old_managed_tags.keys():
if slot not in new_managed_tags: if slot not in new_managed_tags:
self.logger.debug("Unsubscribing from slot %s", slot)
self.ingestor_manager.unsubscribe_slot(slot) self.ingestor_manager.unsubscribe_slot(slot)
# Ensure the gauge is updated after any potential changes here # Ensure the gauge is updated after any potential changes here

View File

@@ -171,8 +171,6 @@ class DataManager:
attachment_content=trace, attachment_content=trace,
) )
self.logger.error(trace) self.logger.error(trace)
else:
self.logger.debug(f"Skipping message to topic {topic}: {data}")
try: try:
collection = self.mongo_db[topic] collection = self.mongo_db[topic]

View File

@@ -29,8 +29,10 @@ class OpcManager():
self.data_manager = data_manager self.data_manager = data_manager
self.notification_handler = notification_handler self.notification_handler = notification_handler
self.pod_id = pod_id self.pod_id = pod_id
metrics.OPC_CONNECTION_STATUS.labels(pod_id=self.pod_id, server_name=self.name, server_url=self.url).set(0) metrics.OPC_CONNECTION_STATUS.labels(
metrics.OPC_TAGS_SUBSCRIBED.labels(pod_id=self.pod_id, server_name=self.name).set(0) pod_id=self.pod_id, server_name=self.name, server_url=self.url).set(0)
metrics.OPC_TAGS_SUBSCRIBED.labels(
pod_id=self.pod_id, server_name=self.name).set(0)
def __str__(self): def __str__(self):
return f"OpcManager(name={self.name}, url={self.url}, server_uri={self.server_uri})\n" \ return f"OpcManager(name={self.name}, url={self.url}, server_uri={self.server_uri})\n" \
@@ -85,18 +87,22 @@ class OpcManager():
Exception: If the connection to the OPC server fails. Exception: If the connection to the OPC server fails.
""" """
metrics.OPC_CONNECTIONS_TOTAL.labels(pod_id=self.pod_id, server_name=self.name).inc() metrics.OPC_CONNECTIONS_TOTAL.labels(
pod_id=self.pod_id, server_name=self.name).inc()
try: try:
self.client = Client(self.url) self.client = Client(self.url)
if self.cert_path: if self.cert_path:
self.set_security() self.set_security()
self.logger.info(f'Starting connection to {self.name}...') self.logger.info(f'Starting connection to {self.name}...')
self.client.connect() self.client.connect()
metrics.OPC_CONNECTION_STATUS.labels(pod_id=self.pod_id, server_name=self.name, server_url=self.url).set(1) metrics.OPC_CONNECTION_STATUS.labels(
pod_id=self.pod_id, server_name=self.name, server_url=self.url).set(1)
self.logger.info(f'Connection to {self.name} successful.') self.logger.info(f'Connection to {self.name} successful.')
except Exception as e: except Exception as e:
metrics.OPC_CONNECTION_STATUS.labels(pod_id=self.pod_id, server_name=self.name, server_url=self.url).set(0) metrics.OPC_CONNECTION_STATUS.labels(
metrics.OPC_CONNECTIONS_FAILED.labels(pod_id=self.pod_id, server_name=self.name).inc() pod_id=self.pod_id, server_name=self.name, server_url=self.url).set(0)
metrics.OPC_CONNECTIONS_FAILED.labels(
pod_id=self.pod_id, server_name=self.name).inc()
self.logger.error(f"Failed to connect to {self.name}: {e}") self.logger.error(f"Failed to connect to {self.name}: {e}")
raise raise
@@ -121,9 +127,11 @@ class OpcManager():
p = period if period is not None else 500 p = period if period is not None else 500
self.subscriptions[name] = self.client.create_subscription(p, self) self.subscriptions[name] = self.client.create_subscription(p, self)
self.logger.info(f'Subscription {name} created on {self.name}.') self.logger.info(f'Subscription {name} created on {self.name}.')
metrics.OPC_SUBSCRIPTIONS_CREATED.labels(pod_id=self.pod_id, server_name=self.name, slot_name=name).inc() metrics.OPC_SUBSCRIPTIONS_CREATED.labels(
pod_id=self.pod_id, server_name=self.name, slot_name=name).inc()
except Exception as e: except Exception as e:
self.logger.error(f"Failed to create subscription {name} on {self.name}: {e}") self.logger.error(
f"Failed to create subscription {name} on {self.name}: {e}")
raise raise
def subscribe(self, subscription: str, nodes: dict, collect_period: int): def subscribe(self, subscription: str, nodes: dict, collect_period: int):
@@ -144,13 +152,18 @@ class OpcManager():
""" """
if not self.subscriptions.get(subscription): if not self.subscriptions.get(subscription):
raise ValueError("Subscription not created. Call create_subscription first.") raise ValueError(
"Subscription not created. Call create_subscription first.")
self.logger.info(f"Subscribing to {subscription} on {self.name}...") self.logger.info(f"Subscribing to {subscription} on {self.name}...")
self.logger.info(f"Subscribing to nodes: {nodes}") self.logger.info(f"Subscribing to nodes: {nodes}")
self.addr_nodes = [self.client.get_node(n) for n in nodes if n not in self.nodes] addr_nodes = [self.client.get_node(
n) for n in nodes]
self.logger.debug(f"Addr nodes: {addr_nodes}")
self.nodes.update(nodes) self.nodes.update(nodes)
metrics.OPC_TAGS_SUBSCRIBED.labels(pod_id=self.pod_id, server_name=self.name).set(len(self.nodes)) self.logger.debug(f"Nodes: {self.nodes}")
metrics.OPC_TAGS_SUBSCRIBED.labels(
pod_id=self.pod_id, server_name=self.name).set(len(self.nodes))
self.collect_period = collect_period self.collect_period = collect_period
for node, config in self.nodes.items(): for node, config in self.nodes.items():
@@ -159,7 +172,7 @@ class OpcManager():
'cycle_count': 0 'cycle_count': 0
} }
self.subscriptions[subscription].subscribe_data_change(self.addr_nodes) self.subscriptions[subscription].subscribe_data_change(addr_nodes)
def unsubscribe(self, subscription: str): def unsubscribe(self, subscription: str):
""" """
@@ -216,8 +229,10 @@ class OpcManager():
finally: finally:
del self.client del self.client
self.client = None self.client = None
metrics.OPC_CONNECTION_STATUS.labels(pod_id=self.pod_id, server_name=self.name, server_url=self.url).set(0) metrics.OPC_CONNECTION_STATUS.labels(
metrics.OPC_TAGS_SUBSCRIBED.labels(pod_id=self.pod_id, server_name=self.name).set(0) pod_id=self.pod_id, server_name=self.name, server_url=self.url).set(0)
metrics.OPC_TAGS_SUBSCRIBED.labels(
pod_id=self.pod_id, server_name=self.name).set(0)
self.logger.warning("Disconnected from OPC UA server.") self.logger.warning("Disconnected from OPC UA server.")
def datachange_notification(self, node, _val, data): def datachange_notification(self, node, _val, data):
@@ -252,7 +267,8 @@ class OpcManager():
self.nodes[tag]['cycle_rule']['cycle_count'] = 0 self.nodes[tag]['cycle_rule']['cycle_count'] = 0
self.non_receive_count = 0 self.non_receive_count = 0
metrics.OPC_CYCLES_WITHOUT_DATA.labels(pod_id=self.pod_id, server_name=self.name).set(0) metrics.OPC_CYCLES_WITHOUT_DATA.labels(
pod_id=self.pod_id, server_name=self.name).set(0)
data = { data = {
'tag': tag, 'tag': tag,
@@ -295,7 +311,8 @@ class OpcManager():
""" """
self.non_receive_count += 1 self.non_receive_count += 1
metrics.OPC_CYCLES_WITHOUT_DATA.labels(pod_id=self.pod_id, server_name=self.name).set(self.non_receive_count) metrics.OPC_CYCLES_WITHOUT_DATA.labels(
pod_id=self.pod_id, server_name=self.name).set(self.non_receive_count)
if self.non_receive_count >= 5: if self.non_receive_count >= 5:
self.notification_handler.build_and_send_notification( self.notification_handler.build_and_send_notification(
notification_id=f'OPC_LISTENNING_STOPPED__{self.name}', notification_id=f'OPC_LISTENNING_STOPPED__{self.name}',
@@ -305,7 +322,8 @@ class OpcManager():
level=NotificationLevel.ERROR level=NotificationLevel.ERROR
) )
if self.non_receive_count >= 15: if self.non_receive_count >= 15:
metrics.OPC_RECONNECTIONS_TOTAL.labels(pod_id=self.pod_id, server_name=self.name).inc() metrics.OPC_RECONNECTIONS_TOTAL.labels(
pod_id=self.pod_id, server_name=self.name).inc()
self.notification_handler.build_and_send_notification( self.notification_handler.build_and_send_notification(
notification_id=f'OPC_CONNECTION_RETRY__{self.name}', notification_id=f'OPC_CONNECTION_RETRY__{self.name}',
message=f'Retrying to connect to server {self.name}', message=f'Retrying to connect to server {self.name}',

View File

@@ -234,9 +234,7 @@ def test_publish_no_kafka(data_manager):
data_manager.publish(topic, data) data_manager.publish(topic, data)
data_manager.logger.debug.assert_any_call( data_manager.kafka_producer.send.assert_not_called()
f"Skipping message to topic {topic}: {data}"
)
@patch("ingestor.managers.data_manager.traceback") @patch("ingestor.managers.data_manager.traceback")

View File

@@ -124,7 +124,8 @@ def test_connect_no_security(client, mock_metrics, raw_opc_manager):
server_name=raw_opc_manager.name, server_name=raw_opc_manager.name,
server_url=raw_opc_manager.url server_url=raw_opc_manager.url
) )
mock_metrics.OPC_CONNECTION_STATUS.labels.return_value.set.assert_called_once_with(1) mock_metrics.OPC_CONNECTION_STATUS.labels.return_value.set.assert_called_once_with(
1)
mock_metrics.OPC_CONNECTIONS_FAILED.labels.assert_not_called() mock_metrics.OPC_CONNECTIONS_FAILED.labels.assert_not_called()
@@ -149,7 +150,8 @@ def test_connect_exception_handling_and_metrics(
): ):
mock_client_instance = mock_opc_client_class.return_value mock_client_instance = mock_opc_client_class.return_value
simulated_error_message = "Erro de conexão simulado" simulated_error_message = "Erro de conexão simulado"
mock_client_instance.connect.side_effect = Exception(simulated_error_message) mock_client_instance.connect.side_effect = Exception(
simulated_error_message)
opc_manager_instance = raw_opc_manager opc_manager_instance = raw_opc_manager
opc_manager_instance.cert_path = None opc_manager_instance.cert_path = None
@@ -193,14 +195,16 @@ def test_create_subscription_no_client(raw_opc_manager):
def test_create_subscription_success_has_period(opc_manager): def test_create_subscription_success_has_period(opc_manager):
opc_manager.create_subscription("sub1", 1000) opc_manager.create_subscription("sub1", 1000)
opc_manager.client.create_subscription.assert_called_once_with(1000, opc_manager) opc_manager.client.create_subscription.assert_called_once_with(
1000, opc_manager)
assert opc_manager.subscriptions["sub1"] is not None assert opc_manager.subscriptions["sub1"] is not None
def test_create_subscription_success_no_period(opc_manager): def test_create_subscription_success_no_period(opc_manager):
opc_manager.create_subscription("sub1", None) opc_manager.create_subscription("sub1", None)
opc_manager.client.create_subscription.assert_called_once_with(500, opc_manager) opc_manager.client.create_subscription.assert_called_once_with(
500, opc_manager)
assert opc_manager.subscriptions["sub1"] is not None assert opc_manager.subscriptions["sub1"] is not None
@@ -250,7 +254,8 @@ def test_subscribe_no_subscription(metrics, opc_manager):
try: try:
opc_manager.subscribe("sub1", tags, 1000) opc_manager.subscribe("sub1", tags, 1000)
except ValueError as e: except ValueError as e:
assert str(e) == "Subscription not created. Call create_subscription first." assert str(
e) == "Subscription not created. Call create_subscription first."
else: else:
assert False, "ValueError not raised" assert False, "ValueError not raised"
metrics.OPC_TAGS_SUBSCRIBED.labels.assert_not_called() metrics.OPC_TAGS_SUBSCRIBED.labels.assert_not_called()
@@ -262,9 +267,9 @@ def test_subscribe_success(opc_manager_subscribed):
opc_manager_subscribed.subscribe("sub1", tags, 1000) opc_manager_subscribed.subscribe("sub1", tags, 1000)
assert opc_manager_subscribed.nodes == tags assert opc_manager_subscribed.nodes == tags
assert opc_manager_subscribed.addr_nodes == [ opc_manager_subscribed.subscriptions["sub1"].subscribe_data_change.assert_called_once_with(
opc_manager_subscribed.client.get_node(n) for n in tags if n != "ns=3;i=1001" [opc_manager_subscribed.client.get_node(n) for n in tags]
] )
def test_unsubscribe_no_subscription(opc_manager): def test_unsubscribe_no_subscription(opc_manager):
@@ -338,13 +343,15 @@ def test_disconnect_metrics_on_successful_path(mock_metrics_module, raw_opc_mana
server_name=raw_opc_manager.name, server_name=raw_opc_manager.name,
server_url=raw_opc_manager.url server_url=raw_opc_manager.url
) )
mock_metrics_module.OPC_CONNECTION_STATUS.labels.return_value.set.assert_called_once_with(0) mock_metrics_module.OPC_CONNECTION_STATUS.labels.return_value.set.assert_called_once_with(
0)
mock_metrics_module.OPC_TAGS_SUBSCRIBED.labels.assert_called_once_with( mock_metrics_module.OPC_TAGS_SUBSCRIBED.labels.assert_called_once_with(
pod_id=raw_opc_manager.pod_id, pod_id=raw_opc_manager.pod_id,
server_name=raw_opc_manager.name server_name=raw_opc_manager.name
) )
mock_metrics_module.OPC_TAGS_SUBSCRIBED.labels.return_value.set.assert_called_once_with(0) mock_metrics_module.OPC_TAGS_SUBSCRIBED.labels.return_value.set.assert_called_once_with(
0)
@patch('ingestor.managers.opc_manager.metrics') @patch('ingestor.managers.opc_manager.metrics')
@@ -395,7 +402,8 @@ def test_datachange_notification(metrics, opc_manager_subscribed):
pod_id=opc_manager_subscribed.pod_id, pod_id=opc_manager_subscribed.pod_id,
server_name=opc_manager_subscribed.name server_name=opc_manager_subscribed.name
) )
metrics.OPC_CYCLES_WITHOUT_DATA.labels.return_value.set.assert_called_once_with(0) metrics.OPC_CYCLES_WITHOUT_DATA.labels.return_value.set.assert_called_once_with(
0)
def test_check_cycles_no_notification(opc_manager): def test_check_cycles_no_notification(opc_manager):
@@ -458,7 +466,8 @@ def test_check_opc_listenning_no_notification(metrics, opc_manager):
pod_id=opc_manager.pod_id, pod_id=opc_manager.pod_id,
server_name=opc_manager.name server_name=opc_manager.name
) )
metrics.OPC_CYCLES_WITHOUT_DATA.labels.return_value.set.assert_called_once_with(opc_manager.non_receive_count) metrics.OPC_CYCLES_WITHOUT_DATA.labels.return_value.set.assert_called_once_with(
opc_manager.non_receive_count)
metrics.OPC_RECONNECTIONS_TOTAL.labels.assert_not_called() metrics.OPC_RECONNECTIONS_TOTAL.labels.assert_not_called()
@@ -511,7 +520,8 @@ def test_check_opc_listenning_error_notification_and_retry(metrics, opc_manager)
pod_id=opc_manager.pod_id, pod_id=opc_manager.pod_id,
server_name=opc_manager.name server_name=opc_manager.name
) )
metrics.OPC_CYCLES_WITHOUT_DATA.labels.return_value.set.assert_called_once_with(opc_manager.non_receive_count) metrics.OPC_CYCLES_WITHOUT_DATA.labels.return_value.set.assert_called_once_with(
opc_manager.non_receive_count)
metrics.OPC_RECONNECTIONS_TOTAL.labels.assert_called_once_with( metrics.OPC_RECONNECTIONS_TOTAL.labels.assert_called_once_with(
pod_id=opc_manager.pod_id, pod_id=opc_manager.pod_id,
@@ -538,11 +548,13 @@ def test_init_metrics_calls_correct_metric_methods(mock_metrics):
server_url=opc_manager.url, server_url=opc_manager.url,
) )
mock_metrics.OPC_CONNECTION_STATUS.labels.return_value.set.assert_called_once_with(0) mock_metrics.OPC_CONNECTION_STATUS.labels.return_value.set.assert_called_once_with(
0)
mock_metrics.OPC_TAGS_SUBSCRIBED.labels.assert_called_once_with( mock_metrics.OPC_TAGS_SUBSCRIBED.labels.assert_called_once_with(
pod_id=opc_manager.pod_id, pod_id=opc_manager.pod_id,
server_name=opc_manager.name server_name=opc_manager.name
) )
mock_metrics.OPC_TAGS_SUBSCRIBED.labels.return_value.set.assert_called_once_with(0) mock_metrics.OPC_TAGS_SUBSCRIBED.labels.return_value.set.assert_called_once_with(
0)

View File

@@ -140,7 +140,7 @@ env:
- name: GITHUB_REPO_URL - name: GITHUB_REPO_URL
value: "git@github.com:Aignosi/sientia-dataops-opc-ingestor.git" value: "git@github.com:Aignosi/sientia-dataops-opc-ingestor.git"
- name: GITHUB_BRANCH - name: GITHUB_BRANCH
value: "error-test" value: "SIENTIAPDE-1110-criar-testes-e-2-e"
- name: PYTHON_APP - name: PYTHON_APP
value: "ingestor.app" value: "ingestor.app"