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 diff --git a/ingestor/ingestor.py b/ingestor/ingestor.py index 681e64d..528bf0f 100644 --- a/ingestor/ingestor.py +++ b/ingestor/ingestor.py @@ -258,17 +258,27 @@ 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) 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 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/ingestor/managers/opc_manager.py b/ingestor/managers/opc_manager.py index 73a2b12..cba0286 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] + addr_nodes = [self.client.get_node( + n) for n in nodes] + self.logger.debug(f"Addr nodes: {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(): @@ -159,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): """ @@ -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}', 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) 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"