diff --git a/ingestor/managers/ingestor_manager.py b/ingestor/managers/ingestor_manager.py index 57e192b..bb04c63 100644 --- a/ingestor/managers/ingestor_manager.py +++ b/ingestor/managers/ingestor_manager.py @@ -79,6 +79,13 @@ class IngestorManager(SientiaMonitoring): redis_username: str | None = redis_data.get('username', None) redis_password: str | None = redis_data.get('password', None) + SientiaMonitoring.__init__( + self, + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + ) + self.data_manager = DataManager( kafka_servers=kafka_servers, mongo_connection_string=mongo_connection_string, @@ -109,14 +116,6 @@ class IngestorManager(SientiaMonitoring): self.metadata = metadata - SientiaMonitoring.__init__( - self, - logger=logger, - notification_handler=notification_handler, - metrics_controller=metrics_controller, - set_error_counter=True, - ) - async def initialize_opc_from_config(self, server_config: dict) -> OpcManager | None: """ Initializes an OPC Manager instance using the provided server configuration. diff --git a/ingestor/managers/opc_manager.py b/ingestor/managers/opc_manager.py index 0558d57..ed16b87 100644 --- a/ingestor/managers/opc_manager.py +++ b/ingestor/managers/opc_manager.py @@ -93,7 +93,6 @@ class OpcManager(SientiaMonitoring): logger=logger, notification_handler=notification_handler, metrics_controller=metrics_controller, - set_error_counter=True, ) metrics.OPC_CONNECTION_STATUS.labels( diff --git a/ingestor/managers/resource_manager.py b/ingestor/managers/resource_manager.py index b467b00..d6273e4 100644 --- a/ingestor/managers/resource_manager.py +++ b/ingestor/managers/resource_manager.py @@ -82,7 +82,6 @@ class ResourceManager(SientiaMonitoring): logger=logger, notification_handler=notification_handler, metrics_controller=metrics_controller, - set_error_counter=True, ) try: diff --git a/ingestor/metrics.py b/ingestor/metrics.py index 10181de..8c71814 100644 --- a/ingestor/metrics.py +++ b/ingestor/metrics.py @@ -16,7 +16,6 @@ labels for multi-dimensional analysis and alerting. from prometheus_client import Counter, Gauge, Histogram - # Metric label definitions for consistent labeling across all metrics POD_ID_LABEL = ['pod_id'] SERVER_LABELS = ['pod_id', 'server_name', 'server_url'] diff --git a/pyproject.toml b/pyproject.toml index ae13d36..e4841ad 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -113,11 +113,7 @@ python_classes = ["Test*"] python_functions = ["test_*"] addopts = [ "-v", - "--strict-markers", - "--cov=model_manager", - "--cov-report=term-missing", - "--cov-report=html", - "--cov-report=xml", + "--strict-markers" ] markers = [ "asyncio: marks tests as async", diff --git a/tests/unit/managers/test_data_manager.py b/tests/unit/managers/test_data_manager.py index 1eba5fe..2cb4ae9 100644 --- a/tests/unit/managers/test_data_manager.py +++ b/tests/unit/managers/test_data_manager.py @@ -165,7 +165,6 @@ def test_shutdown_no_producer(data_manager): ) - def test_shutdown_exception(data_manager): data_manager.kafka_producer.flush = MagicMock(side_effect=Exception('Test error')) data_manager.kafka_producer.close = MagicMock() @@ -206,7 +205,7 @@ def test_delivery_error(data_manager): data_manager.delivery_error(err) data_manager.logger.error.assert_called_once_with(f'Delivery failed for record : {err}') - + @mark.asyncio async def test_publish(data_manager): @@ -248,6 +247,8 @@ async def test_publish_error(traceback, data_manager): send_mock = MagicMock(side_effect=Exception('Test error')) data_manager.kafka_producer.send = send_mock + data_manager.mongo_repository = AsyncMock() + # Call the publish method await data_manager.publish(topic, data) @@ -268,9 +269,7 @@ async def test_publish_error(traceback, data_manager): @mark.asyncio async def test_publish_error_mongo(data_manager): data_manager.export_to_kafka = False - data_manager.mongo_repository.insert = AsyncMock( - side_effect=Exception('Test error') - ) + data_manager.mongo_repository.insert = AsyncMock(side_effect=Exception('Test error')) await data_manager.publish('test_topic', {'key': 'value'}) diff --git a/tests/unit/managers/test_ingestor_manager.py b/tests/unit/managers/test_ingestor_manager.py index ba4c437..3c40e3e 100644 --- a/tests/unit/managers/test_ingestor_manager.py +++ b/tests/unit/managers/test_ingestor_manager.py @@ -1,4 +1,4 @@ -from unittest.mock import AsyncMock, MagicMock, patch +from unittest.mock import AsyncMock, MagicMock, call, patch from pytest import fixture, mark from sientia_do.notifications.models import NotificationLevel @@ -31,9 +31,12 @@ def ingestor_manager(data_manager_mock, resource_manager_mock): logger=MagicMock(), notification_handler=MagicMock(), metadata=metadata['metadata'], + metrics_controller=MagicMock(), ) ingestor.send_notification = MagicMock() + ingestor.send_notification_async = AsyncMock() + ingestor.emit_metric = AsyncMock() return ingestor @@ -57,6 +60,7 @@ def test___init__( logger=MagicMock(), notification_handler=MagicMock(), metadata=metadata['metadata'], + metrics_controller=MagicMock(), ) opc_manager_mock.assert_not_called() @@ -68,6 +72,7 @@ def test___init__( metadata=metadata['metadata'], logger=ingestor.logger, notification_handler=ingestor.notification_handler, + metrics_controller=ingestor.metrics_controller, ) resource_manager_mock.assert_called_once_with( host='localhost', @@ -79,6 +84,7 @@ def test___init__( notification_handler=ingestor.notification_handler, username=None, password=None, + metrics_controller=ingestor.metrics_controller, ) assert ingestor.poll_interval == 5 assert ingestor.managed_tags == {} @@ -117,6 +123,7 @@ async def test_initialize_opc_from_config(opc_manager, ingestor_manager): private_key_path=server_config['private_key_path'], server_cert_path=server_config['server_cert_path'], metadata=metadata['metadata'], + metrics_controller=ingestor_manager.metrics_controller, ) assert result == opc_manager.return_value @@ -155,6 +162,61 @@ async def test_initialize_opc_from_config_exception(traceback_mock, opc_manager, ) +@mark.asyncio +async def test_shutdown(ingestor_manager): + ingestor_manager.opc_managers = {'server1': AsyncMock(), 'server2': AsyncMock()} + + ingestor_manager.data_manager.shutdown = MagicMock() + + await ingestor_manager.shutdown() + + ingestor_manager.opc_managers['server1'].shutdown.assert_called_once() + ingestor_manager.opc_managers['server2'].shutdown.assert_called_once() + + ingestor_manager.data_manager.shutdown.assert_called_once() + + +@patch('ingestor.managers.ingestor_manager.asyncio') +def test___del__(asyncio_mock, ingestor_manager): + ingestor_manager.shutdown = MagicMock() + ingestor_manager.__del__() + asyncio_mock.run.assert_called_once_with(ingestor_manager.shutdown.return_value) + + +@mark.asyncio +async def test_remove_server(ingestor_manager): + server1 = AsyncMock() + server2 = AsyncMock() + ingestor_manager.opc_managers = {'server1': server1, 'server2': server2} + ingestor_manager.managed_tags = { + 'slot1': {'server1': {'config': 'config1'}, 'server2': {'config': 'config2'}} + } + await ingestor_manager.remove_server('server1') + server1.shutdown.assert_called_once() + server2.shutdown.assert_not_called() + assert 'server1' not in ingestor_manager.opc_managers + assert 'server2' in ingestor_manager.opc_managers + assert ingestor_manager.managed_tags == {'slot1': {'server2': {'config': 'config2'}}} + + +@mark.asyncio +async def test_remove_server_not_found(ingestor_manager): + server1 = AsyncMock() + server2 = AsyncMock() + ingestor_manager.opc_managers = {'server1': server1, 'server2': server2} + ingestor_manager.managed_tags = { + 'slot1': {'server1': {'config': 'config1'}, 'server2': {'config': 'config2'}} + } + await ingestor_manager.remove_server('server3') + server1.shutdown.assert_not_called() + server2.shutdown.assert_not_called() + assert 'server1' in ingestor_manager.opc_managers + assert 'server2' in ingestor_manager.opc_managers + assert ingestor_manager.managed_tags == { + 'slot1': {'server1': {'config': 'config1'}, 'server2': {'config': 'config2'}} + } + + @mark.asyncio @patch('ingestor.managers.ingestor_manager.OpcManager') @patch('ingestor.managers.ingestor_manager.metrics') @@ -174,6 +236,7 @@ async def test_update_opc_servers(metrics, opc_manager, ingestor_manager): return None ingestor_manager.initialize_opc_from_config = AsyncMock(side_effect=mock_initialize_from_config) + ingestor_manager.remove_server = AsyncMock() ingestor_manager.managed_tags = { 'slot1': { @@ -181,7 +244,10 @@ async def test_update_opc_servers(metrics, opc_manager, ingestor_manager): 'server2': {'config': 'config2'}, 'server5': {'config': 'config5'}, }, - 'slot2': {'server3': {'config': 'config3'}, 'server1': {'config': 'config1'}}, + 'slot2': { + 'server3': {'config': 'config3'}, + 'server1': {'config': 'config1'}, + }, } mock = AsyncMock(config={'config': 'old_config2'}) @@ -191,121 +257,220 @@ async def test_update_opc_servers(metrics, opc_manager, ingestor_manager): await ingestor_manager.update_opc_servers() - assert len(ingestor_manager.opc_managers) == 3 - ingestor_manager.initialize_opc_from_config.assert_any_call({'config': 'config1'}) ingestor_manager.initialize_opc_from_config.assert_any_call({'config': 'config2'}) ingestor_manager.initialize_opc_from_config.assert_any_call({'config': 'config5'}) + assert ingestor_manager.initialize_opc_from_config.call_count == 3 assert ingestor_manager.opc_managers['server1'].config == {'config': 'config1'} assert ingestor_manager.opc_managers['server2'].config == {'config': 'config2'} assert ingestor_manager.opc_managers['server3'].config == {'config': 'config3'} - assert 'server4' not in ingestor_manager.opc_managers - assert 'server5' not in ingestor_manager.opc_managers + + ingestor_manager.remove_server.assert_any_call('server4') + ingestor_manager.remove_server.assert_any_call('server5') + assert ingestor_manager.remove_server.call_count == 2 assert ingestor_manager.opc_managers['server2'] != mock - metrics.OPC_MANAGERS_ACTIVE.labels.assert_called_once_with(pod_id=ingestor_manager.pod_id) - metrics.OPC_MANAGERS_ACTIVE.labels.return_value.set.assert_called_once_with( - len(ingestor_manager.opc_managers) + ingestor_manager.emit_metric.assert_called_once_with( + metric_object=metrics.OPC_MANAGERS_ACTIVE, + method='set', + value=len(ingestor_manager.opc_managers), + tags={'pod_id': ingestor_manager.pod_id}, ) -def test_declare_active(ingestor_manager): - ingestor_manager.resource_manager.ingestor_heartbeat = MagicMock() - ingestor_manager.declare_active() +@mark.asyncio +@patch('ingestor.managers.ingestor_manager.metrics') +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 + opc_manager1.check_opc_listenning.return_value = False + opc_manager1.config = {'config': 'config1'} + + opc_manager2 = MagicMock() + opc_manager2.check_cycles.return_value = None + opc_manager2.check_opc_listenning.return_value = False + opc_manager2.config = {'config': 'config2'} + + ingestor_manager.opc_managers = {'server1': opc_manager1, 'server2': opc_manager2} + + # Mock the initialize_opc_from_config method + ingestor_manager.initialize_opc_from_config = MagicMock() + + # Call the method + 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() + opc_manager1.check_opc_listenning.assert_called_once() + opc_manager2.check_cycles.assert_called_once() + opc_manager2.check_opc_listenning.assert_called_once() + + # Verify that no reinitialization was needed + ingestor_manager.initialize_opc_from_config.assert_not_called() + + ingestor_manager.emit_metric.assert_called_once_with( + metric_object=metrics.OPC_MANAGERS_ACTIVE, + method='set', + value=len(ingestor_manager.opc_managers), + tags={'pod_id': ingestor_manager.pod_id}, + ) + + +@mark.asyncio +async def test_check_opc_servers_integrity_server_lost(ingestor_manager): + # Setup mock OPC manager that will be lost + 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} + ingestor_manager.managed_tags = {'slot1': MagicMock()} + + # Call the method + 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): + # Setup mock OPC manager that will be lost + opc_manager = MagicMock() + 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} + + # Setup managed tags + ingestor_manager.managed_tags = { + 'slot1': {'server1': {'config': 'config1', 'tags': {'tag1': 'value1'}}} + } + + # 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) + + # Call the method + ingestor_manager.check_opc_servers_integrity() + + +@mark.asyncio +async def test_declare_active(ingestor_manager): + ingestor_manager.resource_manager.ingestor_heartbeat = AsyncMock() + await ingestor_manager.declare_active() ingestor_manager.resource_manager.ingestor_heartbeat.assert_called_once() -def test_get_active_ingestors(ingestor_manager): - ingestor_manager.resource_manager.get_all_ingestors = MagicMock() - ingestor_manager.get_active_ingestors() +@mark.asyncio +async def test_get_active_ingestors(ingestor_manager): + ingestor_manager.resource_manager.get_all_ingestors = AsyncMock() + await ingestor_manager.get_active_ingestors() ingestor_manager.resource_manager.get_all_ingestors.assert_called_once() -def test_get_active_ingestors_empty(ingestor_manager): - ingestor_manager.resource_manager.get_all_ingestors = MagicMock(return_value=None) - result = ingestor_manager.get_active_ingestors() +@mark.asyncio +async def test_get_active_ingestors_empty(ingestor_manager): + ingestor_manager.resource_manager.get_all_ingestors = AsyncMock(return_value=None) + result = await ingestor_manager.get_active_ingestors() assert result == [] ingestor_manager.resource_manager.get_all_ingestors.assert_called_once() @patch('ingestor.managers.ingestor_manager.metrics') -def test_get_number_of_leases_success(metrics, ingestor_manager): - ingestor_manager.resource_manager.get_all_leases = MagicMock(return_value=['lease1', 'lease2']) - result = ingestor_manager.get_number_of_leases() +@mark.asyncio +async def test_get_number_of_leases_success(metrics, ingestor_manager): + ingestor_manager.resource_manager.get_all_leases = AsyncMock(return_value=['lease1', 'lease2']) + result = await ingestor_manager.get_number_of_leases() assert result == 2 + + ingestor_manager.emit_metric.assert_called_once_with( + metric_object=metrics.LEASES_TOTAL, + method='set', + value=2, + tags={'pod_id': ingestor_manager.pod_id}, + ) ingestor_manager.resource_manager.get_all_leases.assert_called_once() - metrics.LEASES_TOTAL.set.assert_called_once_with(2) @patch('ingestor.managers.ingestor_manager.metrics') -def test_get_number_of_leases_empty(metrics, ingestor_manager): - ingestor_manager.resource_manager.get_all_leases = MagicMock(return_value=None) - result = ingestor_manager.get_number_of_leases() - assert result == 0 - ingestor_manager.resource_manager.get_all_leases.assert_called_once() - metrics.LEASES_TOTAL.set.assert_called_once_with(0) - - -@patch('ingestor.managers.ingestor_manager.metrics') -def test_get_number_of_slots_success(metrics, ingestor_manager): - ingestor_manager.resource_manager.get_all_slots = MagicMock(return_value=['slot1', 'slot2']) - result = ingestor_manager.get_number_of_slots() +@mark.asyncio +async def test_get_number_of_slots_success(metrics, ingestor_manager): + ingestor_manager.resource_manager.get_all_slots = AsyncMock(return_value=['slot1', 'slot2']) + result = await ingestor_manager.get_number_of_slots() assert result == 2 ingestor_manager.resource_manager.get_all_slots.assert_called_once() - metrics.SLOTS_TOTAL.set.assert_called_once_with(2) + ingestor_manager.emit_metric.assert_called_once_with( + metric_object=metrics.SLOTS_TOTAL, + method='set', + value=2, + tags={'pod_id': ingestor_manager.pod_id}, + ) @patch('ingestor.managers.ingestor_manager.metrics') -def test_get_number_of_slots_empty(metrics, ingestor_manager): - ingestor_manager.resource_manager.get_all_slots = MagicMock(return_value=None) - result = ingestor_manager.get_number_of_slots() +@mark.asyncio +async def test_get_number_of_slots_empty(metrics, ingestor_manager): + ingestor_manager.resource_manager.get_all_slots = AsyncMock(return_value=None) + result = await ingestor_manager.get_number_of_slots() assert result == 0 ingestor_manager.resource_manager.get_all_slots.assert_called_once() - metrics.SLOTS_TOTAL.set.assert_called_once_with(0) + ingestor_manager.emit_metric.assert_called_once_with( + metric_object=metrics.SLOTS_TOTAL, + method='set', + value=0, + tags={'pod_id': ingestor_manager.pod_id}, + ) -def test_get_slot_leases_1_success(ingestor_manager): - ingestor_manager.resource_manager.lease_tag = MagicMock(return_value=True) - ingestor_manager.resource_manager.get_tag_slot = MagicMock(return_value={'tags': ['tag1']}) +@mark.asyncio +async def test_get_slot_leases_1_success(ingestor_manager): + ingestor_manager.resource_manager.lease_tag = AsyncMock(return_value=True) + ingestor_manager.resource_manager.get_tag_slot = AsyncMock(return_value={'tags': ['tag1']}) ingestor_manager.number_of_slots = 1 - result = ingestor_manager.get_slot_leases() + result = await ingestor_manager.get_slot_leases() assert result == {'1': {'tags': ['tag1']}} -def test_get_slot_leases_2_success(ingestor_manager): - ingestor_manager.resource_manager.lease_tag = MagicMock(side_effect=[True, True]) - ingestor_manager.resource_manager.get_tag_slot = MagicMock( +@mark.asyncio +async def test_get_slot_leases_2_success(ingestor_manager): + ingestor_manager.resource_manager.lease_tag = AsyncMock(side_effect=[True, True]) + ingestor_manager.resource_manager.get_tag_slot = AsyncMock( side_effect=[{'tags': ['tag1']}, {'tags': ['tag2']}] ) ingestor_manager.number_of_slots = 2 - result = ingestor_manager.get_slot_leases(max_slots=2) + result = await ingestor_manager.get_slot_leases(max_slots=2) assert result == {'1': {'tags': ['tag1']}, '2': {'tags': ['tag2']}} -def test_get_slot_leases_2_1_none(ingestor_manager): - ingestor_manager.resource_manager.lease_tag = MagicMock(side_effect=[True, True]) - ingestor_manager.resource_manager.get_tag_slot = MagicMock( +@mark.asyncio +async def test_get_slot_leases_2_1_none(ingestor_manager): + ingestor_manager.resource_manager.lease_tag = AsyncMock(side_effect=[True, True]) + ingestor_manager.resource_manager.get_tag_slot = AsyncMock( side_effect=[None, {'tags': ['tag1']}] ) ingestor_manager.number_of_slots = 1 - result = ingestor_manager.get_slot_leases(max_slots=1) + result = await ingestor_manager.get_slot_leases(max_slots=1) assert result == {} -def test_get_slot_leases_1_failure(ingestor_manager): - ingestor_manager.resource_manager.lease_tag = MagicMock(return_value=False) - ingestor_manager.resource_manager.get_tag_slot = MagicMock(return_value={'tags': ['tag1']}) +@mark.asyncio +async def test_get_slot_leases_1_failure(ingestor_manager): + ingestor_manager.resource_manager.lease_tag = AsyncMock(return_value=False) + ingestor_manager.resource_manager.get_tag_slot = AsyncMock(return_value={'tags': ['tag1']}) - result = ingestor_manager.get_slot_leases() + result = await ingestor_manager.get_slot_leases() ingestor_manager.resource_manager.get_tag_slot.assert_not_called() assert result == {} @@ -329,22 +494,20 @@ async def test_unsubscribe_slot(ingestor_manager): ingestor_manager.opc_managers['server3'].unsubscribe.assert_not_called() -def test_update_slot_config(ingestor_manager): +@mark.asyncio +async def test_update_slot_config(ingestor_manager): ingestor_manager.managed_tags = { 'slot1': {'config': 'old_config'}, 'slot2': {'config': 'new_config'}, 'slot3': {'config': 'old_config'}, } - ingestor_manager.resource_manager.get_tag_slot = MagicMock( + ingestor_manager.resource_manager.get_tag_slot = AsyncMock( side_effect=[{'config': 'updated_config'}, {'config': 'new_config'}, None] ) + ingestor_manager.resource_manager.renew_tag_lease = AsyncMock() - ingestor_manager.update_opc_servers = MagicMock() - ingestor_manager.subscribe_to_tags = MagicMock() - ingestor_manager.unsubscribe_slot = MagicMock() - - ingestor_manager.update_slot_config() + await ingestor_manager.update_slot_config() assert ingestor_manager.managed_tags['slot1'] == {'config': 'updated_config'} assert ingestor_manager.managed_tags['slot2'] == {'config': 'new_config'} @@ -357,15 +520,30 @@ def test_update_slot_config(ingestor_manager): @patch('ingestor.managers.ingestor_manager.metrics') -def test_drop_slot_leases(metrics, ingestor_manager): - ingestor_manager.resource_manager.drop_tag_lease = MagicMock() - ingestor_manager.drop_slot_leases(['1', '2']) +@mark.asyncio +async def test_drop_slot_leases(metrics, ingestor_manager): + ingestor_manager.resource_manager.drop_tag_lease = AsyncMock() + await ingestor_manager.drop_slot_leases(['1', '2']) ingestor_manager.resource_manager.drop_tag_lease.assert_any_call('1') ingestor_manager.resource_manager.drop_tag_lease.assert_any_call('2') - metrics.SLOTS_RELEASED.labels.assert_any_call(pod_id=ingestor_manager.pod_id) - metrics.SLOTS_RELEASED.labels.return_value.inc.assert_any_call() + ingestor_manager.emit_metric.assert_has_calls( + [ + call( + metric_object=metrics.SLOTS_RELEASED, + method='inc', + value=1, + tags={'pod_id': ingestor_manager.pod_id}, + ), + call( + metric_object=metrics.SLOTS_RELEASED, + method='inc', + value=1, + tags={'pod_id': ingestor_manager.pod_id}, + ), + ] + ) @mark.asyncio @@ -432,7 +610,7 @@ async def test_manage_server_subscribe_failure(traceback_mock, ingestor_manager) traceback_mock.format_exc.assert_called_once() - ingestor_manager.send_notification.assert_called_once_with( + ingestor_manager.send_notification_async.assert_called_once_with( metadata=metadata['metadata'], notification_id='OPC_SUBSCRIPTION_ERROR_slot1:server1', message="Failed to subscribe to tags from slot1:server1\n{'tags': 'config1'}: Subscription error", @@ -474,80 +652,3 @@ async def test_subscribe_to_tags(ingestor_manager): assert ingestor_manager.manage_server.call_count == 3 ingestor_manager.managed_tags['slot1'].pop.assert_called_once_with('server3', None) - - -@mark.asyncio -@patch('ingestor.managers.ingestor_manager.metrics') -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 - opc_manager1.check_opc_listenning.return_value = False - opc_manager1.config = {'config': 'config1'} - - opc_manager2 = MagicMock() - opc_manager2.check_cycles.return_value = None - opc_manager2.check_opc_listenning.return_value = False - opc_manager2.config = {'config': 'config2'} - - ingestor_manager.opc_managers = {'server1': opc_manager1, 'server2': opc_manager2} - - # Mock the initialize_opc_from_config method - ingestor_manager.initialize_opc_from_config = MagicMock() - - # Call the method - 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() - opc_manager1.check_opc_listenning.assert_called_once() - opc_manager2.check_cycles.assert_called_once() - opc_manager2.check_opc_listenning.assert_called_once() - - # Verify that no reinitialization was needed - ingestor_manager.initialize_opc_from_config.assert_not_called() - - metrics.OPC_MANAGERS_ACTIVE.labels.assert_called_once_with(pod_id=ingestor_manager.pod_id) - metrics.OPC_MANAGERS_ACTIVE.labels.return_value.set.assert_called_once_with( - len(ingestor_manager.opc_managers) - ) - - -@mark.asyncio -async def test_check_opc_servers_integrity_server_lost(ingestor_manager): - # Setup mock OPC manager that will be lost - 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} - ingestor_manager.managed_tags = {'slot1': MagicMock()} - - # Call the method - 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): - # Setup mock OPC manager that will be lost - opc_manager = MagicMock() - 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} - - # Setup managed tags - ingestor_manager.managed_tags = { - 'slot1': {'server1': {'config': 'config1', 'tags': {'tag1': 'value1'}}} - } - - # 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) - - # Call the method - ingestor_manager.check_opc_servers_integrity() diff --git a/tests/unit/managers/test_opc_manager.py b/tests/unit/managers/test_opc_manager.py index 282328b..85a1448 100644 --- a/tests/unit/managers/test_opc_manager.py +++ b/tests/unit/managers/test_opc_manager.py @@ -46,17 +46,24 @@ metadata = { @fixture @patch('ingestor.managers.opc_manager.metrics') def raw_opc_manager(mock_metrics): - return OpcManager( + opc_manager = OpcManager( name='TestConnector', url='opc.tcp://localhost:4840', - data_manager=MagicMock(), + data_manager=AsyncMock(), subscription_period_ms=1000, logger=MagicMock(), server_uri='opc.tcp://localhost:4840', notification_handler=MagicMock(), metadata=metadata['metadata'], + metrics_controller=AsyncMock(), ) + opc_manager.emit_metric = AsyncMock() + opc_manager.send_notification_async = AsyncMock() + opc_manager.send_notification = MagicMock() + + return opc_manager + @fixture def opc_manager(raw_opc_manager): @@ -64,7 +71,6 @@ def opc_manager(raw_opc_manager): raw_opc_manager.cert_path = 'cert.pem' raw_opc_manager.private_key_path = 'private_key.pem' raw_opc_manager.server_cert_path = 'server_cert.pem' - raw_opc_manager.send_notification = MagicMock() return raw_opc_manager @@ -145,17 +151,30 @@ async def test_connect_no_security(client, mock_metrics, raw_opc_manager): client.assert_called_once_with(raw_opc_manager.url, timeout=10, watchdog_intervall=3600000) raw_opc_manager.client.connect.assert_called_once() raw_opc_manager.set_security.assert_not_called() - mock_metrics.OPC_CONNECTIONS_TOTAL.labels.assert_called_once_with( - pod_id=raw_opc_manager.pod_id, server_name=raw_opc_manager.name + raw_opc_manager.emit_metric.assert_has_calls( + [ + call( + metric_object=mock_metrics.OPC_CONNECTIONS_TOTAL, + method='inc', + value=1, + tags={ + 'pod_id': raw_opc_manager.pod_id, + 'server_name': raw_opc_manager.name, + }, + ), + call( + metric_object=mock_metrics.OPC_CONNECTION_STATUS, + method='set', + value=1, + tags={ + 'pod_id': raw_opc_manager.pod_id, + 'server_name': raw_opc_manager.name, + 'server_url': raw_opc_manager.url, + }, + ), + ], + any_order=True, ) - mock_metrics.OPC_CONNECTIONS_TOTAL.labels.return_value.inc.assert_called_once() - mock_metrics.OPC_CONNECTION_STATUS.labels.assert_called_once_with( - pod_id=raw_opc_manager.pod_id, - 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_CONNECTIONS_FAILED.labels.assert_not_called() @mark.asyncio @@ -230,10 +249,16 @@ async def test_create_subscription_success_no_period(opc_manager): async def test_create_subscription_with_metrics(metrics, opc_manager): 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' + opc_manager.emit_metric.assert_called_once_with( + metric_object=metrics.OPC_SUBSCRIPTIONS_CREATED, + method='inc', + value=1, + tags={ + 'pod_id': opc_manager.pod_id, + 'server_name': opc_manager.name, + 'slot_name': 'sub1', + }, ) - metrics.OPC_SUBSCRIPTIONS_CREATED.labels.return_value.inc.assert_called_once() @patch('ingestor.managers.opc_manager.metrics') @@ -390,17 +415,30 @@ async def test_disconnect_metrics_on_successful_path(mock_metrics_module, raw_op await raw_opc_manager.disconnect() - mock_metrics_module.OPC_CONNECTION_STATUS.labels.assert_called_once_with( - pod_id=raw_opc_manager.pod_id, - server_name=raw_opc_manager.name, - server_url=raw_opc_manager.url, + raw_opc_manager.emit_metric.assert_has_calls( + [ + call( + metric_object=mock_metrics_module.OPC_CONNECTION_STATUS, + method='set', + value=0, + tags={ + 'pod_id': raw_opc_manager.pod_id, + 'server_name': raw_opc_manager.name, + 'server_url': raw_opc_manager.url, + }, + ), + call( + metric_object=mock_metrics_module.OPC_TAGS_SUBSCRIBED, + method='set', + value=0, + tags={ + 'pod_id': raw_opc_manager.pod_id, + 'server_name': raw_opc_manager.name, + }, + ), + ], + any_order=True, ) - 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) @patch('ingestor.managers.opc_manager.metrics') @@ -422,8 +460,6 @@ async def test_datachange_notification(metrics, opc_manager_subscribed): } } - metrics.OPC_CYCLES_WITHOUT_DATA.reset_mock() - await opc_manager_subscribed.datachange_notification('ns=3;i=1001', None, data) opc_manager_subscribed.data_manager.publish.assert_any_call( @@ -446,13 +482,19 @@ async def test_datachange_notification(metrics, opc_manager_subscribed): ) assert opc_manager_subscribed.nodes['ns=3;i=1001']['cycle_rule']['cycle_count'] == 0 - metrics.OPC_CYCLES_WITHOUT_DATA.labels.assert_called_once_with( - pod_id=opc_manager_subscribed.pod_id, server_name=opc_manager_subscribed.name + opc_manager_subscribed.emit_metric.assert_called_once_with( + metric_object=metrics.OPC_CYCLES_WITHOUT_DATA, + method='set', + value=0, + tags={ + '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) -def test_check_cycles_no_notification(opc_manager): +@mark.asyncio +async def test_check_cycles_no_notification(opc_manager): # Setup: node with cycle_count just below threshold opc_manager.nodes = { 'ns=3;i=1001': { @@ -461,14 +503,15 @@ def test_check_cycles_no_notification(opc_manager): } } - opc_manager.check_cycles() + await opc_manager.check_cycles() # After one increment, cycle_count = 4.0, still below threshold assert opc_manager.nodes['ns=3;i=1001']['cycle_rule']['cycle_count'] == pytest.approx(4.0) - opc_manager.send_notification.assert_not_called() + opc_manager.send_notification_async.assert_not_called() -def test_check_cycles_triggers_notification(opc_manager): +@mark.asyncio +async def test_check_cycles_triggers_notification(opc_manager): # Setup: node with cycle_count just below threshold, increment will cross threshold opc_manager.nodes = { 'ns=3;i=1001': { @@ -478,11 +521,11 @@ def test_check_cycles_triggers_notification(opc_manager): } opc_manager.notification_handler.build_and_send_notification = MagicMock() - opc_manager.check_cycles() + await opc_manager.check_cycles() # After increment, cycle_count = 5.5, should trigger notification assert opc_manager.nodes['ns=3;i=1001']['cycle_rule']['cycle_count'] == pytest.approx(5.5) - opc_manager.send_notification.assert_called_once_with( + opc_manager.send_notification_async.assert_called_once_with( notification_id='TAG_ns=3;i=1001:Counter_LISTENNING_STOPPED', message='5.5 cycles without receive from ns=3;i=1001:Counter', block='opc_manager', @@ -492,34 +535,35 @@ def test_check_cycles_triggers_notification(opc_manager): @patch('ingestor.managers.opc_manager.metrics') -def test_check_opc_listenning_no_notification(metrics, opc_manager): +@mark.asyncio +async def test_check_opc_listenning_no_notification(metrics, opc_manager): opc_manager.non_receive_count = 3 - opc_manager.notification_handler.build_and_send_notification = MagicMock() - metrics.OPC_CYCLES_WITHOUT_DATA.reset_mock() - metrics.OPC_RECONNECTIONS_TOTAL.reset_mock() - result = opc_manager.check_opc_listenning() + result = await opc_manager.check_opc_listenning() assert opc_manager.non_receive_count == 4 - opc_manager.notification_handler.build_and_send_notification.assert_not_called() + opc_manager.send_notification_async.assert_not_called() assert result is False - metrics.OPC_CYCLES_WITHOUT_DATA.labels.assert_called_once_with( - pod_id=opc_manager.pod_id, server_name=opc_manager.name + opc_manager.emit_metric.assert_called_once_with( + metric_object=metrics.OPC_CYCLES_WITHOUT_DATA, + method='set', + value=opc_manager.non_receive_count, + tags={ + '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_RECONNECTIONS_TOTAL.labels.assert_not_called() -def test_check_opc_listenning_warning_notification(opc_manager): +@mark.asyncio +async def test_check_opc_listenning_warning_notification(opc_manager): opc_manager.non_receive_count = 4 - result = opc_manager.check_opc_listenning() + result = await opc_manager.check_opc_listenning() assert opc_manager.non_receive_count == 5 - opc_manager.send_notification.assert_called_once_with( + opc_manager.send_notification_async.assert_called_once_with( notification_id=f'OPC_LISTENNING_STOPPED__{opc_manager.name}', message=f'5 cycles without receive from OPC {opc_manager.name}. Tags: {json.dumps(opc_manager.nodes)}', block='opc_manager', @@ -530,18 +574,16 @@ def test_check_opc_listenning_warning_notification(opc_manager): @patch('ingestor.managers.opc_manager.metrics') -def test_check_opc_listenning_error_notification_and_retry(metrics, opc_manager): +@mark.asyncio +async def test_check_opc_listenning_error_notification_and_retry(metrics, opc_manager): opc_manager.non_receive_count = 14 - opc_manager.notification_handler.build_and_send_notification = MagicMock() - metrics.OPC_CYCLES_WITHOUT_DATA.reset_mock() - metrics.OPC_RECONNECTIONS_TOTAL.reset_mock() - result = opc_manager.check_opc_listenning() + result = await opc_manager.check_opc_listenning() assert opc_manager.non_receive_count == 15 # Should be called twice: once for 5, once for 15 - assert opc_manager.send_notification.call_count == 2 - calls = opc_manager.send_notification.call_args_list + assert opc_manager.send_notification_async.call_count == 2 + calls = opc_manager.send_notification_async.call_args_list # First call: 5 cycles warning assert calls[0].kwargs == { 'notification_id': f'OPC_LISTENNING_STOPPED__{opc_manager.name}', @@ -560,42 +602,30 @@ def test_check_opc_listenning_error_notification_and_retry(metrics, opc_manager) } assert result is True - metrics.OPC_CYCLES_WITHOUT_DATA.labels.assert_called_once_with( - 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 + opc_manager.emit_metric.assert_has_calls( + [ + call( + metric_object=metrics.OPC_CYCLES_WITHOUT_DATA, + method='set', + value=opc_manager.non_receive_count, + tags={ + 'pod_id': opc_manager.pod_id, + 'server_name': opc_manager.name, + }, + ), + ] ) - metrics.OPC_RECONNECTIONS_TOTAL.labels.assert_called_once_with( - pod_id=opc_manager.pod_id, server_name=opc_manager.name + opc_manager.emit_metric.assert_has_calls( + [ + call( + metric_object=metrics.OPC_RECONNECTIONS_TOTAL, + method='inc', + value=1, + tags={ + 'pod_id': opc_manager.pod_id, + 'server_name': opc_manager.name, + }, + ), + ] ) - metrics.OPC_RECONNECTIONS_TOTAL.labels.return_value.inc.assert_called_once() - - -@patch('ingestor.managers.opc_manager.metrics') -def test_init_metrics_calls_correct_metric_methods(metrics): - opc_manager = OpcManager( - 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(), - metadata=metadata['metadata'], - ) - - metrics.OPC_CONNECTION_STATUS.labels.assert_called_with( - pod_id=opc_manager.pod_id, - server_name=opc_manager.name, - server_url=opc_manager.url, - ) - - metrics.OPC_CONNECTION_STATUS.labels.return_value.set.assert_called_once_with(0) - - metrics.OPC_TAGS_SUBSCRIBED.labels.assert_called_once_with( - pod_id=opc_manager.pod_id, server_name=opc_manager.name - ) - - metrics.OPC_TAGS_SUBSCRIBED.labels.return_value.set.assert_called_once_with(0) diff --git a/tests/unit/managers/test_resource_manager.py b/tests/unit/managers/test_resource_manager.py index 9479672..d7e236c 100644 --- a/tests/unit/managers/test_resource_manager.py +++ b/tests/unit/managers/test_resource_manager.py @@ -1,7 +1,6 @@ -from unittest.mock import MagicMock, patch +from unittest.mock import AsyncMock, MagicMock, patch -from pytest import fixture, raises -from sientia_do.notifications.models import NotificationLevel +from pytest import fixture, mark, raises from ingestor.managers.resource_manager import ResourceManager @@ -15,9 +14,58 @@ metadata = { } +@patch('ingestor.managers.resource_manager.RedisRepository') +def test___init__(redis_repository): + logger = MagicMock() + notification_handler = MagicMock() + metrics_controller = MagicMock() + resource_manager = ResourceManager( + host='localhost', + port=6379, + lease_ttl=10, + heartbeat_ttl=10, + metadata=metadata['metadata'], + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + ) + assert resource_manager.redis_repository == redis_repository.return_value + redis_repository.assert_called_once_with( + host='localhost', + port=6379, + username=None, + password=None, + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + ) + redis_repository.return_value.redis_client.ping.assert_called_once() + + +@patch('ingestor.managers.resource_manager.RedisRepository') +def test___init__connection_failure(redis_repository): + redis_repository.return_value.redis_client.ping.side_effect = Exception('Connection failed') + with raises(Exception, match='Connection failed'): + ResourceManager( + host='localhost', + port=6379, + lease_ttl=10, + heartbeat_ttl=10, + metadata=metadata['metadata'], + logger=MagicMock(), + notification_handler=MagicMock(), + metrics_controller=MagicMock(), + ) + redis_repository.return_value.redis_client.ping.assert_called_once() + redis_repository.logger.error.assert_called_once_with( + 'Failed to connect to Redis: Connection failed' + ) + redis_repository.return_value.redis_client.ping.assert_called_once() + + @fixture -@patch('ingestor.managers.resource_manager.Redis') -def resource_manager(redis): +@patch('ingestor.managers.resource_manager.RedisRepository') +def resource_manager(redis_repository): resource_manager = ResourceManager( host='localhost', port=6379, @@ -26,199 +74,105 @@ def resource_manager(redis): metadata=metadata['metadata'], logger=MagicMock(), notification_handler=MagicMock(), + metrics_controller=MagicMock(), ) resource_manager.send_notification = MagicMock() + resource_manager.send_notification_async = AsyncMock() + resource_manager.emit_metric = AsyncMock() + + resource_manager.redis_repository = AsyncMock() return resource_manager -def test_get_success(resource_manager): - resource_manager.redis.get.return_value = '{"key": "value"}' - result = resource_manager.get('key') - assert result == {'key': 'value'} - resource_manager.redis.get.assert_called_once_with('key') +@mark.asyncio +async def test_get_tag_slot(resource_manager): + result = await resource_manager.get_tag_slot('id') + assert result == resource_manager.redis_repository.get.return_value - -def test_get_failure(resource_manager): - resource_manager.redis.get.return_value = None - result = resource_manager.get('key') - assert result is None - resource_manager.redis.get.assert_called_once_with('key') - - -def test_get_tag_slot(resource_manager): - resource_manager.get = MagicMock(return_value={'tag': 'slot'}) - result = resource_manager.get_tag_slot('id') - assert result == {'tag': 'slot'} - resource_manager.get.assert_called_once_with('slot:opc_tags:id') - - -def test_ingestor_heartbeat(resource_manager): - resource_manager.redis.set.return_value = True - resource_manager.ingestor_heartbeat() - resource_manager.redis.set.assert_called_once_with('heartbeat:ingestor:localhost', 1, ex=10) - - -def test_lease_tag(resource_manager): - resource_manager.redis.set.return_value = True - output = resource_manager.lease_tag('tag_id') - assert output is True - resource_manager.redis.set.assert_called_once_with( - 'lease:opc_tags:tag_id', 'localhost', nx=True, ex=10 + resource_manager.redis_repository.get.assert_called_once_with( + 'slot:opc_tags:id', metadata=metadata['metadata'] ) -def test_renew_tag_lease_success(resource_manager): - resource_manager.redis.get.return_value = 'localhost' - resource_manager.redis.expire.return_value = True - result = resource_manager.renew_tag_lease('tag_id') - assert result is True - resource_manager.redis.get.assert_called_once_with('lease:opc_tags:tag_id') - resource_manager.redis.expire.assert_called_once_with('lease:opc_tags:tag_id', 10) +@mark.asyncio +async def test_ingestor_heartbeat(resource_manager): + await resource_manager.ingestor_heartbeat() + resource_manager.redis_repository.set.assert_called_once_with( + 'heartbeat:ingestor:localhost', 1, ttl=10, metadata=metadata['metadata'] + ) -def test_renew_tag_lease_failure(resource_manager): - resource_manager.redis.get.return_value = 'other_pod_id' - resource_manager.redis.expire.return_value = False - result = resource_manager.renew_tag_lease('tag_id') +@mark.asyncio +async def test_lease_tag(resource_manager): + output = await resource_manager.lease_tag('tag_id') + assert output is resource_manager.redis_repository.set.return_value + resource_manager.redis_repository.set.assert_called_once_with( + 'lease:opc_tags:tag_id', 'localhost', ttl=10, nx=True, metadata=metadata['metadata'] + ) + + +@mark.asyncio +async def test_renew_tag_lease_success(resource_manager): + resource_manager.redis_repository.get.return_value = 'localhost' + resource_manager.redis_repository.expire.return_value = True + result = await resource_manager.renew_tag_lease('tag_id') + assert result is resource_manager.redis_repository.expire.return_value + resource_manager.redis_repository.get.assert_called_once_with( + 'lease:opc_tags:tag_id', metadata=metadata['metadata'] + ) + resource_manager.redis_repository.expire.assert_called_once_with( + 'lease:opc_tags:tag_id', 10, metadata=metadata['metadata'] + ) + + +@mark.asyncio +async def test_renew_tag_lease_failure(resource_manager): + resource_manager.redis_repository.get.return_value = 'other_pod_id' + + result = await resource_manager.renew_tag_lease('tag_id') assert result is False - resource_manager.redis.get.assert_called_once_with('lease:opc_tags:tag_id') - resource_manager.redis.expire.assert_not_called() + + resource_manager.redis_repository.get.assert_called_once_with( + 'lease:opc_tags:tag_id', metadata=metadata['metadata'] + ) + resource_manager.redis_repository.expire.assert_not_called() -def test_drop_tag_lease(resource_manager): - resource_manager.redis.delete.return_value = True - resource_manager.drop_tag_lease('tag_id') - resource_manager.redis.delete.assert_called_once_with('lease:opc_tags:tag_id') +@mark.asyncio +async def test_drop_tag_lease(resource_manager): + await resource_manager.drop_tag_lease('tag_id') + resource_manager.redis_repository.delete.assert_called_once_with( + 'lease:opc_tags:tag_id', metadata=metadata['metadata'] + ) -def test_get_all_ingestors(resource_manager): - resource_manager.redis.keys.return_value = ['ingestor1', 'ingestor2'] - result = resource_manager.get_all_ingestors() +@mark.asyncio +async def test_get_all_ingestors(resource_manager): + resource_manager.redis_repository.keys.return_value = ['ingestor1', 'ingestor2'] + result = await resource_manager.get_all_ingestors() assert result == ['ingestor1', 'ingestor2'] - resource_manager.redis.keys.assert_called_once_with('heartbeat:ingestor:*') + resource_manager.redis_repository.keys.assert_called_once_with( + 'heartbeat:ingestor:*', metadata=metadata['metadata'] + ) -def test_get_all_slots(resource_manager): - resource_manager.redis.keys.return_value = ['slot1', 'slot2'] - result = resource_manager.get_all_slots() +@mark.asyncio +async def test_get_all_slots(resource_manager): + resource_manager.redis_repository.keys.return_value = ['slot1', 'slot2'] + result = await resource_manager.get_all_slots() assert result == ['slot1', 'slot2'] - resource_manager.redis.keys.assert_called_once_with('slot:opc_tags:*') + resource_manager.redis_repository.keys.assert_called_once_with( + 'slot:opc_tags:*', metadata=metadata['metadata'] + ) -def test_get_all_leases(resource_manager): - resource_manager.redis.keys.return_value = ['lease1', 'lease2'] - result = resource_manager.get_all_leases() +@mark.asyncio +async def test_get_all_leases(resource_manager): + resource_manager.redis_repository.keys.return_value = ['lease1', 'lease2'] + result = await resource_manager.get_all_leases() assert result == ['lease1', 'lease2'] - resource_manager.redis.keys.assert_called_once_with('lease:opc_tags:*') - - -def test_init_connection_failure(monkeypatch): - # Mock Redis to raise an exception during initialization - mock_redis = MagicMock() - mock_redis.side_effect = Exception('Connection failed') - - monkeypatch.setattr('ingestor.managers.resource_manager.Redis', mock_redis) - - # Test that the exception is raised and metrics are set properly - with patch('ingestor.metrics.REDIS_CONNECTION_STATUS') as mock_metrics: - mock_status = MagicMock() - mock_metrics.labels.return_value = mock_status - - with raises(Exception, match='Connection failed'): - ResourceManager( - host='localhost', - port=6379, - lease_ttl=10, - heartbeat_ttl=10, - metadata=metadata['metadata'], - logger=MagicMock(), - notification_handler=MagicMock(), - ) - - mock_metrics.labels.assert_called_once_with(pod_id='localhost') - mock_status.set.assert_called_once_with(0) - - -def test_init_ping_failure(monkeypatch): - # Mock Redis ping to raise an exception - mock_redis_instance = MagicMock() - mock_redis_instance.ping.side_effect = Exception('Ping failed') - - mock_redis_class = MagicMock(return_value=mock_redis_instance) - monkeypatch.setattr('ingestor.managers.resource_manager.Redis', mock_redis_class) - - # Test that the exception is raised and metrics are set properly - with patch('ingestor.metrics.REDIS_CONNECTION_STATUS') as mock_metrics: - mock_status = MagicMock() - mock_metrics.labels.return_value = mock_status - - with raises(Exception, match='Ping failed'): - ResourceManager( - host='localhost', - port=6379, - lease_ttl=10, - heartbeat_ttl=10, - metadata=metadata['metadata'], - logger=MagicMock(), - notification_handler=MagicMock(), - ) - - mock_metrics.labels.assert_called_once_with(pod_id='localhost') - mock_status.set.assert_called_once_with(0) - - -def test_execute_redis_op_success(resource_manager): - # Mock the Redis operation and time function - mock_func = MagicMock(return_value='test_result') - - with patch('ingestor.managers.resource_manager.time', side_effect=[100, 100.5]): - with patch('ingestor.metrics.REDIS_OPERATIONS_TOTAL') as mock_total: - with patch('ingestor.metrics.REDIS_OPERATIONS_DURATION') as mock_duration: - mock_total_labels = MagicMock() - mock_duration_labels = MagicMock() - mock_total.labels.return_value = mock_total_labels - mock_duration.labels.return_value = mock_duration_labels - - # Execute the operation - result = resource_manager._execute_redis_op( - 'test_op', mock_func, 'arg1', kwarg1='value1' - ) - - # Verify the result and metrics - assert result == 'test_result' - mock_func.assert_called_once_with('arg1', kwarg1='value1') - - mock_total.labels.assert_called_once_with(pod_id='localhost', operation='test_op') - mock_total_labels.inc.assert_called_once() - - mock_duration.labels.assert_called_once_with( - pod_id='localhost', operation='test_op' - ) - mock_duration_labels.observe.assert_called_once_with(0.5) - - -def test_execute_redis_op_exception(resource_manager): - # Mock the Redis operation to raise an exception - mock_func = MagicMock(side_effect=Exception('Operation failed')) - - with patch('ingestor.managers.resource_manager.time', return_value=100): - with patch('ingestor.metrics.REDIS_OPERATIONS_ERRORS') as mock_errors: - mock_errors_labels = MagicMock() - mock_errors.labels.return_value = mock_errors_labels - - # Execute the operation and expect an exception - with raises(Exception, match='Operation failed'): - resource_manager._execute_redis_op('test_op', mock_func, 'arg1') - - # Verify metrics and error handling - mock_errors.labels.assert_called_once_with(pod_id='localhost', operation='test_op') - mock_errors_labels.inc.assert_called_once() - resource_manager.send_notification.assert_called_once_with( - metadata=metadata['metadata'], - notification_id='REDIS_OPERATION_ERROR_test_op', - message="Error in Redis operation 'test_op': Operation failed", - block='redis_manager', - level=NotificationLevel.ERROR, - ) + resource_manager.redis_repository.keys.assert_called_once_with( + 'lease:opc_tags:*', metadata=metadata['metadata'] + ) diff --git a/tests/unit/test_ingestor.py b/tests/unit/test_ingestor.py index c98fbf3..ea1afd6 100644 --- a/tests/unit/test_ingestor.py +++ b/tests/unit/test_ingestor.py @@ -74,7 +74,7 @@ def ingestor(_notification_handler, _getenv): @fixture def ingestor_manager_started(ingestor): - ingestor.ingestor_manager = MagicMock( + ingestor.ingestor_manager = AsyncMock( initialize_opc_from_config=AsyncMock(), shutdown=AsyncMock(), update_opc_servers=AsyncMock(), @@ -138,6 +138,7 @@ async def test_prepare_ingestor(ingestor_manager_mock, ingestor): logger=ingestor.logger, notification_handler=ingestor.notification_handler, export_to_kafka=ingestor.export_to_kafka, + metrics_controller=ingestor.metrics_controller, ) ingestor_manager.declare_active.assert_called_once() ingestor_manager.get_slot_leases.assert_called_once() @@ -177,11 +178,12 @@ def test_manage_slots_none_available(ingestor_manager_started): ingestor_manager_started.ingestor_manager.handle_acquired_tags.assert_not_called() -def test_manage_slots_none_available_none_available(ingestor_manager_started): +@mark.asyncio +async def test_manage_slots_none_available_none_available(ingestor_manager_started): ingestor_manager_started.handle_acquired_tags = MagicMock() ingestor_manager_started.ingestor_manager.managed_tags = False - ingestor_manager_started.manage_no_slots(2) + await ingestor_manager_started.manage_no_slots(2) ingestor_manager_started.ingestor_manager.get_slot_leases.assert_called_once_with(1) @@ -235,7 +237,7 @@ async def test_manage_leases_no_available_slots_extra_sltos(ingestor_manager_sta @mark.asyncio async def test_loop(ingestor_manager_started): - ingestor_manager_started.manage_no_slots = MagicMock() + ingestor_manager_started.manage_no_slots = AsyncMock() ingestor_manager_started.manage_leases = AsyncMock() ingestor_manager_started.update_ingestor_manager = AsyncMock() ingestor_manager_started.ingestor_manager.check_opc_servers_integrity = AsyncMock() @@ -244,11 +246,11 @@ async def test_loop(ingestor_manager_started): 'slot2': 'server2', 'slot3': 'server3', } - ingestor_manager_started.ingestor_manager.get_active_ingestors = MagicMock( + ingestor_manager_started.ingestor_manager.get_active_ingestors = AsyncMock( return_value=['ingestor1', 'ingestor2'] ) - ingestor_manager_started.ingestor_manager.get_number_of_slots = MagicMock(return_value=5) - ingestor_manager_started.ingestor_manager.get_number_of_leases = MagicMock(return_value=1) + ingestor_manager_started.ingestor_manager.get_number_of_slots = AsyncMock(return_value=5) + ingestor_manager_started.ingestor_manager.get_number_of_leases = AsyncMock(return_value=1) await ingestor_manager_started.loop() @@ -269,15 +271,15 @@ async def test_loop(ingestor_manager_started): @mark.asyncio async def test_loop_no_managed(ingestor_manager_started): - ingestor_manager_started.manage_no_slots = MagicMock() + ingestor_manager_started.manage_no_slots = AsyncMock() 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( + ingestor_manager_started.ingestor_manager.get_active_ingestors = AsyncMock( return_value=['ingestor1', 'ingestor2'] ) - ingestor_manager_started.ingestor_manager.get_number_of_slots = MagicMock(return_value=5) + ingestor_manager_started.ingestor_manager.get_number_of_slots = AsyncMock(return_value=5) await ingestor_manager_started.loop() diff --git a/tests/unit/test_metrics.py b/tests/unit/test_metrics.py index a60a6c7..9fc545e 100644 --- a/tests/unit/test_metrics.py +++ b/tests/unit/test_metrics.py @@ -207,38 +207,6 @@ def test_kafka_connection_status(): assert set(metrics.KAFKA_CONNECTION_STATUS._labelnames) == {'pod_id'} -def test_redis_operations_total(): - """Verify the definition of REDIS_OPERATIONS_TOTAL.""" - assert metrics.REDIS_OPERATIONS_TOTAL is not None - assert isinstance(metrics.REDIS_OPERATIONS_TOTAL, Counter) - assert metrics.REDIS_OPERATIONS_TOTAL._name == 'redis_operations' # REMOVED _total - assert set(metrics.REDIS_OPERATIONS_TOTAL._labelnames) == {'pod_id', 'operation'} - - -def test_redis_operations_errors(): - """Verify the definition of REDIS_OPERATIONS_ERRORS.""" - assert metrics.REDIS_OPERATIONS_ERRORS is not None - assert isinstance(metrics.REDIS_OPERATIONS_ERRORS, Counter) - assert metrics.REDIS_OPERATIONS_ERRORS._name == 'redis_operations_errors' # REMOVED _total - assert set(metrics.REDIS_OPERATIONS_ERRORS._labelnames) == {'pod_id', 'operation'} - - -def test_redis_operations_duration(): - """Verify the definition of REDIS_OPERATIONS_DURATION.""" - assert metrics.REDIS_OPERATIONS_DURATION is not None - assert isinstance(metrics.REDIS_OPERATIONS_DURATION, Histogram) - assert metrics.REDIS_OPERATIONS_DURATION._name == 'redis_operations_duration_seconds' - assert set(metrics.REDIS_OPERATIONS_DURATION._labelnames) == {'pod_id', 'operation'} - - -def test_redis_connection_status(): - """Verify the definition of REDIS_CONNECTION_STATUS.""" - assert metrics.REDIS_CONNECTION_STATUS is not None - assert isinstance(metrics.REDIS_CONNECTION_STATUS, Gauge) - assert metrics.REDIS_CONNECTION_STATUS._name == 'redis_connection_status' - assert set(metrics.REDIS_CONNECTION_STATUS._labelnames) == {'pod_id'} - - def test_notifications_sent(): """Verify the definition of NOTIFICATIONS_SENT.""" assert metrics.NOTIFICATIONS_SENT is not None