From 47b3159e7177d9163dbc29e2e38012308983762b Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 20 Oct 2025 12:41:05 -0300 Subject: [PATCH] SIENTIAPDE-1312 Enhance OpcManager to accept subscription period from configuration - Updated OpcManager to include subscription_period_ms as a parameter for better configurability. - Refactored create_subscription method to utilize the new subscription period parameter. - Adjusted unit tests to reflect changes in subscription handling and ensure correct functionality. --- ingestor/managers/ingestor_manager.py | 1 + ingestor/managers/opc_manager.py | 26 ++++++++++++++++++++++--- tests/unit/managers/test_opc_manager.py | 14 ++++++++----- 3 files changed, 33 insertions(+), 8 deletions(-) diff --git a/ingestor/managers/ingestor_manager.py b/ingestor/managers/ingestor_manager.py index ea0b3a1..bf964c0 100644 --- a/ingestor/managers/ingestor_manager.py +++ b/ingestor/managers/ingestor_manager.py @@ -141,6 +141,7 @@ class IngestorManager(BaseActivity): manager = OpcManager( name=server_config['name'], url=server_config['url'], + subscription_period_ms=server_config['subscription_period_ms'], data_manager=self.data_manager, logger=self.logger, server_uri=server_config['server_uri'], diff --git a/ingestor/managers/opc_manager.py b/ingestor/managers/opc_manager.py index bf66ef1..abd81c1 100644 --- a/ingestor/managers/opc_manager.py +++ b/ingestor/managers/opc_manager.py @@ -59,6 +59,7 @@ class OpcManager(BaseActivity): self, name: str, url: str, + subscription_period_ms: int, data_manager: DataManager, logger: Logger, server_uri: str, @@ -74,6 +75,7 @@ class OpcManager(BaseActivity): self.data_queue: dict = {} self.non_receive_count = 0 self.client: Client | None = None + self.subscription_period_ms = subscription_period_ms self.cert_path = cert_path self.private_key_path = private_key_path self.server_cert_path = server_cert_path @@ -82,6 +84,7 @@ class OpcManager(BaseActivity): self.data_manager = data_manager self.metadata = metadata + BaseActivity.__init__( self, logger=logger, notification_handler=notification_handler, set_error_counter=True ) @@ -189,7 +192,13 @@ class OpcManager(BaseActivity): try: self.client = Client(self.url, watchdog_intervall=3600000) assert self.client is not None # Informa ao mypy que client não é None + self.client.name = self.pod_id + self.client.application_name = self.pod_id + pod_uri = f'urn:sientia-do:pod:{self.pod_id}' + self.client.application_uri = pod_uri + self.client.product_uri = pod_uri + if self.cert_path: await self.set_security() self.logger.info(f'Starting connection to {self.name}...') @@ -199,6 +208,15 @@ class OpcManager(BaseActivity): ).set(1) self.logger.info(f'Connection to {self.name} successful.') except Exception as e: + if self.client: + try: + await self.client.disconnect() + except Exception as internal_e: + self.logger.error(f'And error occurred while creating connection from {self.name}, and an exception occurred while disconnecting: {internal_e}') + + 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) @@ -206,7 +224,7 @@ class OpcManager(BaseActivity): self.logger.error(f'Failed to connect to {self.name}: {e}') raise - async def create_subscription(self, name: str, period: int = 500): + async def create_subscription(self, name: str): """ Creates a subscription with the specified monitoring period. @@ -230,8 +248,8 @@ class OpcManager(BaseActivity): if not self.client: raise ValueError('Client not connected. Call connect first.') try: - p = period if period is not None else 500 - self.subscriptions[name] = await self.client.create_subscription(p, self) + self.subscriptions[name] = await self.client.create_subscription( + self.subscription_period_ms, 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 @@ -345,10 +363,12 @@ class OpcManager(BaseActivity): except Exception as sub_error: self.logger.error(f'Failed to clean up subscription: {sub_error}') + # Inserir backoff para disconnect, emitindo métrica e alerta try: await self.client.disconnect() except Exception as conn_error: self.logger.error(f'Failed to disconnect from OPC UA server: {conn_error}') + finally: del self.client self.client = None diff --git a/tests/unit/managers/test_opc_manager.py b/tests/unit/managers/test_opc_manager.py index 0c4e625..d7a2eac 100644 --- a/tests/unit/managers/test_opc_manager.py +++ b/tests/unit/managers/test_opc_manager.py @@ -50,6 +50,7 @@ def raw_opc_manager(mock_metrics): name='TestConnector', url='opc.tcp://localhost:4840', data_manager=MagicMock(), + subscription_period_ms=1000, logger=MagicMock(), server_uri='opc.tcp://localhost:4840', notification_handler=MagicMock(), @@ -223,24 +224,26 @@ async def test_create_subscription_no_client(raw_opc_manager): @mark.asyncio async def test_create_subscription_success_has_period(opc_manager): - await opc_manager.create_subscription('sub1', 1000) + await opc_manager.create_subscription('sub1') - opc_manager.client.create_subscription.assert_called_once_with(1000, opc_manager) + opc_manager.client.create_subscription.assert_called_once_with( + opc_manager.subscription_period_ms, opc_manager) assert opc_manager.subscriptions['sub1'] is not None @mark.asyncio async def test_create_subscription_success_no_period(opc_manager): - await opc_manager.create_subscription('sub1', None) + await opc_manager.create_subscription('sub1') - opc_manager.client.create_subscription.assert_called_once_with(500, opc_manager) + opc_manager.client.create_subscription.assert_called_once_with( + opc_manager.subscription_period_ms, opc_manager) assert opc_manager.subscriptions['sub1'] is not None @patch('ingestor.managers.opc_manager.metrics') @mark.asyncio async def test_create_subscription_with_metrics(metrics, opc_manager): - await opc_manager.create_subscription('sub1', 1000) + await opc_manager.create_subscription('sub1') metrics.OPC_SUBSCRIPTIONS_CREATED.labels.assert_called_once_with( pod_id=opc_manager.pod_id, server_name=opc_manager.name, slot_name='sub1' @@ -572,6 +575,7 @@ def test_init_metrics_calls_correct_metric_methods(metrics): name='TestInitConnector', url='opc.tcp://init.test:4840', data_manager=MagicMock(), + subscription_period_ms=1000, logger=MagicMock(), server_uri='opc.tcp://init.test:4840/uri', notification_handler=MagicMock(),