diff --git a/ingestor/ingestor.py b/ingestor/ingestor.py index 662f539..5f27759 100644 --- a/ingestor/ingestor.py +++ b/ingestor/ingestor.py @@ -414,7 +414,7 @@ class Ingestor: # Check OPC cycles self.logger.info('Checking OPC servers integrity...') - self.ingestor_manager.check_opc_servers_integrity() + await self.ingestor_manager.check_opc_servers_integrity() self.logger.info('Updating managed tags...') await self.update_ingestor_manager(current_managed_tags) diff --git a/ingestor/managers/ingestor_manager.py b/ingestor/managers/ingestor_manager.py index bf964c0..f1258b3 100644 --- a/ingestor/managers/ingestor_manager.py +++ b/ingestor/managers/ingestor_manager.py @@ -189,7 +189,7 @@ class IngestorManager(BaseActivity): def __del__(self): asyncio.run(self.shutdown()) - async def update_opc_servers(self): + async def update_opc_servers(self): # NOSONAR """ 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 @@ -228,7 +228,7 @@ class IngestorManager(BaseActivity): elif server_instance.config != server_config: self.logger.warning(f'Reinitializing OPC manager for server {server}') - server_instance.disconnect() + await server_instance.shutdown() del self.opc_managers[server] server_instance = await self.initialize_opc_from_config(server_config) else: @@ -243,10 +243,12 @@ class IngestorManager(BaseActivity): f'Failed to initialize OPC manager for server {server}, ' f'removing server from managed tags.' ) - _a = [ + + 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) - for slot, _value in current_managed_tags.items() - ] servers = list(self.opc_managers.keys()) for server in servers: @@ -254,12 +256,12 @@ class IngestorManager(BaseActivity): self.logger.warning( f'Server {server} not found in managed tags. Desconnecting from server.' ) - await self.opc_managers[server].disconnect() + await self.opc_managers[server].shutdown() self.opc_managers.pop(server, None) metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers)) - def check_opc_servers_integrity(self): + async def check_opc_servers_integrity(self): """ Checks the integrity of the OPC servers and updates the OPC servers if necessary. @@ -282,6 +284,9 @@ 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] self.managed_tags[slot].pop(server, None) metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers)) @@ -604,4 +609,7 @@ class IngestorManager(BaseActivity): to_remove.append([slot, server]) for slot, server in to_remove: + if server in self.opc_managers: + await self.opc_managers[server].shutdown() + del self.opc_managers[server] self.managed_tags[slot].pop(server, None) diff --git a/ingestor/managers/opc_manager.py b/ingestor/managers/opc_manager.py index 0ad077d..d1232dd 100644 --- a/ingestor/managers/opc_manager.py +++ b/ingestor/managers/opc_manager.py @@ -107,6 +107,9 @@ class OpcManager(BaseActivity): f'nodes={self.nodes}, subscriptions={self.subscriptions}' ) + def __del__(self): + asyncio.run(self.shutdown()) + async def shutdown(self): """ Comprehensive cleanup method for graceful shutdown. diff --git a/tests/unit/managers/test_ingestor_manager.py b/tests/unit/managers/test_ingestor_manager.py index d6dc973..b218e4e 100644 --- a/tests/unit/managers/test_ingestor_manager.py +++ b/tests/unit/managers/test_ingestor_manager.py @@ -184,7 +184,7 @@ async def test_update_opc_servers(metrics, opc_manager, ingestor_manager): 'slot2': {'server3': {'config': 'config3'}, 'server1': {'config': 'config1'}}, } - mock = MagicMock(config={'config': 'old_config2'}) + mock = AsyncMock(config={'config': 'old_config2'}) ingestor_manager.opc_managers['server3'] = AsyncMock(config={'config': 'config3'}) ingestor_manager.opc_managers['server2'] = mock ingestor_manager.opc_managers['server4'] = AsyncMock() @@ -472,8 +472,9 @@ async def test_subscribe_to_tags(ingestor_manager): ingestor_manager.managed_tags['slot1'].pop.assert_called_once_with('server3', None) +@mark.asyncio @patch('ingestor.managers.ingestor_manager.metrics') -def test_check_opc_servers_integrity_all_healthy(metrics, ingestor_manager): +async def test_check_opc_servers_integrity_all_healthy(metrics, ingestor_manager): # Setup mock OPC managers opc_manager1 = MagicMock() opc_manager1.check_cycles.return_value = None @@ -491,7 +492,7 @@ def test_check_opc_servers_integrity_all_healthy(metrics, ingestor_manager): ingestor_manager.initialize_opc_from_config = MagicMock() # Call the method - ingestor_manager.check_opc_servers_integrity() + await ingestor_manager.check_opc_servers_integrity() # Verify that check_cycles and check_opc_listenning were called for each server opc_manager1.check_cycles.assert_called_once() diff --git a/tests/unit/test_ingestor.py b/tests/unit/test_ingestor.py index b636f2b..c98fbf3 100644 --- a/tests/unit/test_ingestor.py +++ b/tests/unit/test_ingestor.py @@ -238,6 +238,7 @@ async def test_loop(ingestor_manager_started): ingestor_manager_started.manage_no_slots = MagicMock() ingestor_manager_started.manage_leases = AsyncMock() ingestor_manager_started.update_ingestor_manager = AsyncMock() + ingestor_manager_started.ingestor_manager.check_opc_servers_integrity = AsyncMock() ingestor_manager_started.ingestor_manager.managed_tags = { 'slot1': 'server1', 'slot2': 'server2', @@ -262,12 +263,16 @@ async def test_loop(ingestor_manager_started): ingestor_manager_started.manage_leases.assert_called_once_with(4, 3, 2) ingestor_manager_started.ingestor_manager.update_slot_config.assert_called_once() + ingestor_manager_started.ingestor_manager.check_opc_servers_integrity.assert_called_once() + ingestor_manager_started.update_ingestor_manager.assert_called_once() + @mark.asyncio async def test_loop_no_managed(ingestor_manager_started): ingestor_manager_started.manage_no_slots = MagicMock() ingestor_manager_started.manage_leases = AsyncMock() ingestor_manager_started.update_ingestor_manager = AsyncMock() + ingestor_manager_started.ingestor_manager.check_opc_servers_integrity = AsyncMock() ingestor_manager_started.ingestor_manager.managed_tags = {} ingestor_manager_started.ingestor_manager.get_active_ingestors = MagicMock( return_value=['ingestor1', 'ingestor2'] @@ -288,6 +293,9 @@ async def test_loop_no_managed(ingestor_manager_started): ingestor_manager_started.ingestor_manager.update_slot_config.assert_called_once() ingestor_manager_started.logger.info.assert_any_call('No slots acquired in this loop') + ingestor_manager_started.ingestor_manager.check_opc_servers_integrity.assert_called_once() + ingestor_manager_started.update_ingestor_manager.assert_called_once() + @mark.asyncio async def test_loop_no_ingestor_manager(ingestor_manager_started):