From 7f88d89849095a8dcc1fb7ba8d289ce5c113424c Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 8 Jul 2025 11:35:18 -0300 Subject: [PATCH 1/7] SIENTIAPDE-1110 Enhance logging in Ingestor class for subscription management - Added debug logging for subscribing to new slots, resubscribing to modified slots, and unsubscribing from removed slots in the Ingestor class. --- ingestor/ingestor.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/ingestor/ingestor.py b/ingestor/ingestor.py index 681e64d..572d900 100644 --- a/ingestor/ingestor.py +++ b/ingestor/ingestor.py @@ -260,15 +260,18 @@ class Ingestor: for slot, config in new_managed_tags.items(): if slot not in old_managed_tags: + self.logger.debug("Subscribing to new slot %s", slot) self.ingestor_manager.subscribe_to_tags({slot: config}) continue if config != old_managed_tags[slot]: + self.logger.debug("Resubscribing to slot %s", slot) self.ingestor_manager.unsubscribe_slot(slot) self.ingestor_manager.subscribe_to_tags({slot: config}) for slot in old_managed_tags.keys(): if slot not in new_managed_tags: + self.logger.debug("Unsubscribing from slot %s", slot) self.ingestor_manager.unsubscribe_slot(slot) # Ensure the gauge is updated after any potential changes here From 22ed2c8cdeac8bf12b8786d178da2a118f25ce1e Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 8 Jul 2025 11:38:32 -0300 Subject: [PATCH 2/7] SIENTIAPDE-1110 Update GITHUB_BRANCH in values.yaml and remove debug log in DataManager - Changed GITHUB_BRANCH value to "SIENTIAPDE-1110-criar-testes-e-2-e" in values.yaml. - Removed unnecessary debug log for skipping messages in DataManager class. --- ingestor/managers/data_manager.py | 2 -- values.yaml | 2 +- 2 files changed, 1 insertion(+), 3 deletions(-) diff --git a/ingestor/managers/data_manager.py b/ingestor/managers/data_manager.py index ad97a1a..9c689ed 100644 --- a/ingestor/managers/data_manager.py +++ b/ingestor/managers/data_manager.py @@ -171,8 +171,6 @@ class DataManager: attachment_content=trace, ) self.logger.error(trace) - else: - self.logger.debug(f"Skipping message to topic {topic}: {data}") try: collection = self.mongo_db[topic] diff --git a/values.yaml b/values.yaml index 5aa0321..cd916b5 100644 --- a/values.yaml +++ b/values.yaml @@ -140,7 +140,7 @@ env: - name: GITHUB_REPO_URL value: "git@github.com:Aignosi/sientia-dataops-opc-ingestor.git" - name: GITHUB_BRANCH - value: "error-test" + value: "SIENTIAPDE-1110-criar-testes-e-2-e" - name: PYTHON_APP value: "ingestor.app" From 9d510c593694c4e60b063299e5b348f4704c383f Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 8 Jul 2025 11:54:19 -0300 Subject: [PATCH 3/7] SIENTIAPDE-1110 Enhance debug logging in Ingestor class to track changes in managed tags - Added debug logging to capture changes between new and old managed tags in the Ingestor class, improving traceability of subscription management. --- ingestor/ingestor.py | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/ingestor/ingestor.py b/ingestor/ingestor.py index 572d900..528bf0f 100644 --- a/ingestor/ingestor.py +++ b/ingestor/ingestor.py @@ -258,6 +258,13 @@ class Ingestor: 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(): if slot not in old_managed_tags: self.logger.debug("Subscribing to new slot %s", slot) From 7391ed3ba574a20fc1421fac069069e1ab2ebdc4 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 8 Jul 2025 12:25:11 -0300 Subject: [PATCH 4/7] SIENTIAPDE-1110 Refactor OPC Manager metrics logging for improved readability - Reformatted metrics logging statements in the OpcManager class for better readability and consistency. - Enhanced debug logging to provide clearer insights into subscription management and connection status. --- ingestor/managers/opc_manager.py | 50 ++++++++++++++++++++++---------- 1 file changed, 34 insertions(+), 16 deletions(-) diff --git a/ingestor/managers/opc_manager.py b/ingestor/managers/opc_manager.py index 73a2b12..9e1cff9 100644 --- a/ingestor/managers/opc_manager.py +++ b/ingestor/managers/opc_manager.py @@ -29,8 +29,10 @@ class OpcManager(): self.data_manager = data_manager self.notification_handler = notification_handler 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_TAGS_SUBSCRIBED.labels(pod_id=self.pod_id, server_name=self.name).set(0) + metrics.OPC_CONNECTION_STATUS.labels( + 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): 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. """ - 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: self.client = Client(self.url) if self.cert_path: self.set_security() self.logger.info(f'Starting connection to {self.name}...') 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.') 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_CONNECTIONS_FAILED.labels(pod_id=self.pod_id, server_name=self.name).inc() + metrics.OPC_CONNECTION_STATUS.labels( + 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}") raise @@ -121,9 +127,11 @@ class OpcManager(): p = period if period is not None else 500 self.subscriptions[name] = self.client.create_subscription(p, self) 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: - 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 def subscribe(self, subscription: str, nodes: dict, collect_period: int): @@ -144,13 +152,18 @@ class OpcManager(): """ 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 nodes: {nodes}") - self.addr_nodes = [self.client.get_node(n) for n in nodes if n not in self.nodes] + self.addr_nodes = [self.client.get_node( + n) for n in nodes if n not in self.nodes] + self.logger.debug(f"Addr nodes: {self.addr_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 for node, config in self.nodes.items(): @@ -216,8 +229,10 @@ class OpcManager(): finally: del self.client 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_TAGS_SUBSCRIBED.labels(pod_id=self.pod_id, server_name=self.name).set(0) + metrics.OPC_CONNECTION_STATUS.labels( + 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.") def datachange_notification(self, node, _val, data): @@ -252,7 +267,8 @@ class OpcManager(): self.nodes[tag]['cycle_rule']['cycle_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 = { 'tag': tag, @@ -295,7 +311,8 @@ class OpcManager(): """ 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: self.notification_handler.build_and_send_notification( notification_id=f'OPC_LISTENNING_STOPPED__{self.name}', @@ -305,7 +322,8 @@ class OpcManager(): level=NotificationLevel.ERROR ) 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( notification_id=f'OPC_CONNECTION_RETRY__{self.name}', message=f'Retrying to connect to server {self.name}', From 1512ca98bdc9b97e41fd5c25d9904f33fc21493d Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 8 Jul 2025 12:32:21 -0300 Subject: [PATCH 5/7] SIENTIAPDE-1110 Refactor OpcManager to improve node subscription handling - Simplified the address node assignment in the OpcManager class by directly using a local variable instead of an instance variable. - Enhanced debug logging to provide clearer insights into the nodes being subscribed to. --- ingestor/managers/opc_manager.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/ingestor/managers/opc_manager.py b/ingestor/managers/opc_manager.py index 9e1cff9..cba0286 100644 --- a/ingestor/managers/opc_manager.py +++ b/ingestor/managers/opc_manager.py @@ -157,9 +157,9 @@ class OpcManager(): self.logger.info(f"Subscribing to {subscription} on {self.name}...") 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] - self.logger.debug(f"Addr nodes: {self.addr_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.logger.debug(f"Nodes: {self.nodes}") metrics.OPC_TAGS_SUBSCRIBED.labels( @@ -172,7 +172,7 @@ class OpcManager(): '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): """ From 8a5989dddc4082371cf02e953640dac0bac01938 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 8 Jul 2025 15:21:07 -0300 Subject: [PATCH 6/7] SIENTIAPDE-1110 shortcut --- README.md | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index e1b09a7..22facb3 100644 --- a/README.md +++ b/README.md @@ -71,4 +71,11 @@ pytest --cov=ingestor ### Generate complete report ``` pytest --cov=ingestor --cov-report=html -``` \ No newline at end of file +``` + +#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 \ No newline at end of file From 1a96ada0e0131957abe028e5d8e7ff2e1c6b71e6 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 8 Jul 2025 16:24:11 -0300 Subject: [PATCH 7/7] SIENTIAPDE-1110 Refactor test assertions for improved readability and consistency - Updated assertions in test cases for DataManager and OpcManager to enhance readability by formatting long lines. - Ensured that the expected behavior of metrics logging and message publishing is clearly defined in the tests. --- tests/unit/managers/test_data_manager.py | 4 +-- tests/unit/managers/test_opc_manager.py | 44 +++++++++++++++--------- 2 files changed, 29 insertions(+), 19 deletions(-) diff --git a/tests/unit/managers/test_data_manager.py b/tests/unit/managers/test_data_manager.py index 8a79566..da78e26 100644 --- a/tests/unit/managers/test_data_manager.py +++ b/tests/unit/managers/test_data_manager.py @@ -234,9 +234,7 @@ def test_publish_no_kafka(data_manager): data_manager.publish(topic, data) - data_manager.logger.debug.assert_any_call( - f"Skipping message to topic {topic}: {data}" - ) + data_manager.kafka_producer.send.assert_not_called() @patch("ingestor.managers.data_manager.traceback") diff --git a/tests/unit/managers/test_opc_manager.py b/tests/unit/managers/test_opc_manager.py index 90c25b9..f9be3dc 100644 --- a/tests/unit/managers/test_opc_manager.py +++ b/tests/unit/managers/test_opc_manager.py @@ -124,7 +124,8 @@ def test_connect_no_security(client, mock_metrics, raw_opc_manager): server_name=raw_opc_manager.name, 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() @@ -149,7 +150,8 @@ def test_connect_exception_handling_and_metrics( ): mock_client_instance = mock_opc_client_class.return_value 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.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): 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 def test_create_subscription_success_no_period(opc_manager): 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 @@ -250,7 +254,8 @@ def test_subscribe_no_subscription(metrics, opc_manager): try: opc_manager.subscribe("sub1", tags, 1000) 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: assert False, "ValueError not raised" 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) assert opc_manager_subscribed.nodes == tags - assert opc_manager_subscribed.addr_nodes == [ - opc_manager_subscribed.client.get_node(n) for n in tags if n != "ns=3;i=1001" - ] + opc_manager_subscribed.subscriptions["sub1"].subscribe_data_change.assert_called_once_with( + [opc_manager_subscribed.client.get_node(n) for n in tags] + ) 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_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( pod_id=raw_opc_manager.pod_id, 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') @@ -395,7 +402,8 @@ def test_datachange_notification(metrics, opc_manager_subscribed): pod_id=opc_manager_subscribed.pod_id, 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): @@ -458,7 +466,8 @@ def test_check_opc_listenning_no_notification(metrics, opc_manager): pod_id=opc_manager.pod_id, 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() @@ -511,8 +520,9 @@ def test_check_opc_listenning_error_notification_and_retry(metrics, opc_manager) pod_id=opc_manager.pod_id, 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( pod_id=opc_manager.pod_id, server_name=opc_manager.name @@ -538,11 +548,13 @@ def test_init_metrics_calls_correct_metric_methods(mock_metrics): 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( pod_id=opc_manager.pod_id, 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)