diff --git a/ingestor/managers/ingestor_manager.py b/ingestor/managers/ingestor_manager.py index f1258b3..c63039c 100644 --- a/ingestor/managers/ingestor_manager.py +++ b/ingestor/managers/ingestor_manager.py @@ -189,7 +189,17 @@ class IngestorManager(BaseActivity): def __del__(self): asyncio.run(self.shutdown()) - async def update_opc_servers(self): # NOSONAR + async def remove_server(self, server: str): + """ + Removes an OPC server from the ingestor. + """ + if server in self.opc_managers: + await self.opc_managers[server].shutdown() + del self.opc_managers[server] + for slot, _config in self.managed_tags.items(): + self.managed_tags[slot].pop(server, None) + + async def update_opc_servers(self): """ Updates the OPC (OLE for Process Control) server connections managed by the ingestor. This method ensures that the OPC servers defined in `self.managed_tags` are properly @@ -244,11 +254,7 @@ class IngestorManager(BaseActivity): f'removing server from managed tags.' ) - for slot, _value in current_managed_tags.items(): - if server in self.opc_managers: - await self.opc_managers[server].shutdown() - del self.opc_managers[server] - self.managed_tags[slot].pop(server, None) + await self.remove_server(server) servers = list(self.opc_managers.keys()) for server in servers: @@ -256,8 +262,7 @@ class IngestorManager(BaseActivity): self.logger.warning( f'Server {server} not found in managed tags. Desconnecting from server.' ) - await self.opc_managers[server].shutdown() - self.opc_managers.pop(server, None) + await self.remove_server(server) metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers)) @@ -276,6 +281,7 @@ class IngestorManager(BaseActivity): - Removes lost servers from managed tags - Updates OPC manager metrics """ + to_disconnect: list[str] = [] for server, opc_manager in self.opc_managers.items(): opc_manager.check_cycles() @@ -284,11 +290,14 @@ class IngestorManager(BaseActivity): self.logger.warning(f'OPC server {server} is lost. Server will be disconnected.') for slot, _config in self.managed_tags.items(): - if server in self.opc_managers: - await self.opc_managers[server].shutdown() - del self.opc_managers[server] + to_disconnect.append(server) self.managed_tags[slot].pop(server, None) + for server in to_disconnect: + if server in self.opc_managers: + await self.opc_managers[server].shutdown() + del self.opc_managers[server] + metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers)) def declare_active(self): diff --git a/ingestor/managers/opc_manager.py b/ingestor/managers/opc_manager.py index d1232dd..2c0ba7e 100644 --- a/ingestor/managers/opc_manager.py +++ b/ingestor/managers/opc_manager.py @@ -154,21 +154,21 @@ class OpcManager(BaseActivity): raise ValueError( 'Certificate and private key paths must be provided for secure connection.' ) - cert = Path(self.cert_path) if self.cert_path else None - private_key = Path(self.private_key_path) if self.private_key_path else None - server_cert = Path(self.server_cert_path) if self.server_cert_path else None + cert = str(Path(self.cert_path)) if self.cert_path else None + private_key = str(Path(self.private_key_path)) if self.private_key_path else None + server_cert = str(Path(self.server_cert_path)) if self.server_cert_path else None if self.client: - await self.client.set_application_uri(self.server_uri) + self.client.application_uri = self.server_uri self.logger.info('Setting security...') await self.client.set_security( SecurityPolicyBasic256, - certificate=str(cert), - private_key=str(private_key), - server_certificate=str(server_cert), + certificate=cert, + private_key=private_key, + server_certificate=server_cert, ) - await self.client.set_secure_channel_timeout(10000000) - await self.client.set_session_timeout(10000000) + self.client.secure_channel_timeout = 10000000 + self.client.session_timeout = 10000000 async def connect(self): """ @@ -286,9 +286,7 @@ class OpcManager(BaseActivity): self.logger.debug(f'Addr nodes: {addr_nodes}') self.nodes.update(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(): @@ -299,6 +297,10 @@ class OpcManager(BaseActivity): await self.subscriptions[subscription].subscribe_data_change(addr_nodes) + metrics.OPC_TAGS_SUBSCRIBED.labels(pod_id=self.pod_id, server_name=self.name).set( + len(self.nodes) + ) + async def unsubscribe(self, subscription: str): """ Unsubscribes from a given subscription. @@ -437,10 +439,6 @@ class OpcManager(BaseActivity): f'{tag} after {self.nodes[tag]["cycle_rule"]["cycle_count"]} cycles' ) - 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) - data = { 'tag': tag, 'name': self.nodes[str(node)]['tag_name'], @@ -451,6 +449,10 @@ class OpcManager(BaseActivity): for topic in self.nodes[tag]['topics']: self.data_manager.publish(topic, data) + 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) + def check_cycles(self): """ Checks the cycle counts for all monitored nodes and sends diff --git a/tests/unit/managers/test_ingestor_manager.py b/tests/unit/managers/test_ingestor_manager.py index b218e4e..ba4c437 100644 --- a/tests/unit/managers/test_ingestor_manager.py +++ b/tests/unit/managers/test_ingestor_manager.py @@ -451,7 +451,11 @@ async def test_subscribe_to_tags(ingestor_manager): ingestor_manager.manage_server = AsyncMock(side_effect=[0, 1, 2]) ingestor_manager.managed_tags = {'slot1': MagicMock(), 'slot2': MagicMock()} - ingestor_manager.opc_managers = {'server1': AsyncMock(), 'server2': AsyncMock()} + ingestor_manager.opc_managers = { + 'server1': AsyncMock(), + 'server2': AsyncMock(), + 'server3': AsyncMock(), + } ingestor_manager.subscriptions = {'server1': AsyncMock()} tags = { 'slot1': { @@ -509,21 +513,22 @@ async def test_check_opc_servers_integrity_all_healthy(metrics, ingestor_manager ) -def test_check_opc_servers_integrity_server_lost(ingestor_manager): +@mark.asyncio +async def test_check_opc_servers_integrity_server_lost(ingestor_manager): # Setup mock OPC manager that will be lost - opc_manager = MagicMock() + opc_manager = AsyncMock() opc_manager.check_cycles.return_value = None opc_manager.check_opc_listenning.return_value = True # Server is lost opc_manager.config = {'config': 'config1'} ingestor_manager.opc_managers = {'server1': opc_manager} - - # Mock the initialize_opc_from_config method to return a new manager - new_manager = MagicMock() - ingestor_manager.initialize_opc_from_config = MagicMock(return_value=new_manager) + ingestor_manager.managed_tags = {'slot1': MagicMock()} # Call the method - ingestor_manager.check_opc_servers_integrity() + await ingestor_manager.check_opc_servers_integrity() + + assert 'server1' not in ingestor_manager.opc_managers + ingestor_manager.managed_tags['slot1'].pop.assert_called_once_with('server1', None) def test_check_opc_servers_integrity_server_lost_with_tags(ingestor_manager): diff --git a/tests/unit/managers/test_opc_manager.py b/tests/unit/managers/test_opc_manager.py index f440e43..282328b 100644 --- a/tests/unit/managers/test_opc_manager.py +++ b/tests/unit/managers/test_opc_manager.py @@ -105,7 +105,7 @@ async def test_shutdown_error(opc_manager): async def test_set_security_success(opc_manager): await opc_manager.set_security() - opc_manager.client.set_application_uri.assert_called_once_with(opc_manager.server_uri) + assert opc_manager.client.application_uri == opc_manager.server_uri opc_manager.client.set_security.assert_called_once_with( SecurityPolicyBasic256, @@ -114,8 +114,8 @@ async def test_set_security_success(opc_manager): server_certificate=opc_manager.server_cert_path, ) - opc_manager.client.set_secure_channel_timeout.assert_called_once_with(10000000) - opc_manager.client.set_session_timeout.assert_called_once_with(10000000) + assert opc_manager.client.secure_channel_timeout == 10000000 + assert opc_manager.client.session_timeout == 10000000 @mark.asyncio