diff --git a/ingestor/managers/ingestor_manager.py b/ingestor/managers/ingestor_manager.py index 85d5602..667d928 100644 --- a/ingestor/managers/ingestor_manager.py +++ b/ingestor/managers/ingestor_manager.py @@ -7,6 +7,7 @@ from sientia_do.notifications.models import NotificationLevel from ingestor.managers.data_manager import DataManager from ingestor.managers.opc_manager import OpcManager from ingestor.managers.resource_manager import ResourceManager +import ingestor.metrics as metrics class IngestorManager(): @@ -162,6 +163,8 @@ class IngestorManager(): self.opc_managers[server].disconnect() 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): """ Checks the integrity of the OPC servers and updates the OPC servers if necessary. @@ -179,6 +182,8 @@ class IngestorManager(): for slot, _config in self.managed_tags.items(): self.managed_tags[slot].pop(server, None) + metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers)) + def declare_active(self): """ Declares the ingestor as active by sending a heartbeat signal to the resource manager. @@ -211,6 +216,7 @@ class IngestorManager(): leases = self.resource_manager.get_all_leases() self.number_of_slots = len(leases) if leases else 0 + metrics.LEASES_TOTAL.set(self.number_of_slots) return self.number_of_slots def get_number_of_slots(self) -> int: @@ -224,6 +230,7 @@ class IngestorManager(): slots = self.resource_manager.get_all_slots() self.number_of_slots = len(slots) if slots else 0 + metrics.SLOTS_TOTAL.set(self.number_of_slots) return self.number_of_slots def get_slot_leases(self, max_slots: int = 1) -> Dict: @@ -255,9 +262,11 @@ class IngestorManager(): if slots is None: continue acquired[str(i)] = slots + metrics.SLOTS_ACQUIRED.labels(pod_id=self.pod_id).inc() if len(acquired) >= max_slots: self.managed_tags.update(acquired) + metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(len(self.managed_tags)) return acquired self.logger.warning( @@ -265,6 +274,7 @@ class IngestorManager(): f"Only {acquired} slots were leased." ) self.managed_tags.update(acquired) + metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(len(self.managed_tags)) return acquired def unsubscribe_slot(self, slot: str): @@ -332,6 +342,7 @@ class IngestorManager(): for lease_id in ids: self.resource_manager.drop_tag_lease(lease_id) + metrics.SLOTS_RELEASED.labels(pod_id=self.pod_id).inc() def manage_server(self, slot: str, server: str, server_config: dict, tags: dict) -> int: """ @@ -390,6 +401,7 @@ class IngestorManager(): tags_to_sub ) except Exception as e: + metrics.OPC_SUBSCRIPTION_ERRORS.labels(pod_id=self.pod_id, server=server, slot=slot).inc() trace = traceback.format_exc() self.notification_handler.build_and_send_notification( notification_id=f'OPC_SUBSCRIPTION_ERROR_{slot}:{server}', diff --git a/tests/unit/managers/test_ingestor_manager.py b/tests/unit/managers/test_ingestor_manager.py index f1bcbf3..a1f2d12 100644 --- a/tests/unit/managers/test_ingestor_manager.py +++ b/tests/unit/managers/test_ingestor_manager.py @@ -111,7 +111,8 @@ def test_initialize_opc_from_config_exception(traceback_mock, opc_manager, inges @patch('ingestor.managers.ingestor_manager.OpcManager') -def test_update_opc_servers(opc_manager, ingestor_manager): +@patch('ingestor.managers.ingestor_manager.metrics') +def test_update_opc_servers(metrics, opc_manager, ingestor_manager): manager1 = MagicMock( config={"config": "config1"}) manager2 = MagicMock( @@ -175,6 +176,9 @@ def test_update_opc_servers(opc_manager, ingestor_manager): 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)) + def test_declare_active(ingestor_manager): ingestor_manager.resource_manager.ingestor_heartbeat = MagicMock() @@ -196,36 +200,44 @@ def test_get_active_ingestors_empty(ingestor_manager): ingestor_manager.resource_manager.get_all_ingestors.assert_called_once() -def test_get_number_of_leases_success(ingestor_manager): +@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() assert result == 2 ingestor_manager.resource_manager.get_all_leases.assert_called_once() + metrics.LEASES_TOTAL.set.assert_called_once_with(2) -def test_get_number_of_leases_empty(ingestor_manager): +@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) -def test_get_number_of_slots_success(ingestor_manager): +@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() assert result == 2 ingestor_manager.resource_manager.get_all_slots.assert_called_once() + metrics.SLOTS_TOTAL.set.assert_called_once_with(2) -def test_get_number_of_slots_empty(ingestor_manager): +@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() assert result == 0 ingestor_manager.resource_manager.get_all_slots.assert_called_once() + metrics.SLOTS_TOTAL.set.assert_called_once_with(0) def test_get_slot_leases_1_success(ingestor_manager): @@ -339,13 +351,17 @@ def test_update_slot_config(ingestor_manager): assert ingestor_manager.resource_manager.renew_tag_lease.call_count == 3 -def test_drop_slot_leases(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"]) 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() + def test_manage_server_no_server(ingestor_manager): ingestor_manager.opc_managers = { @@ -491,7 +507,8 @@ def test_subscribe_to_tags(ingestor_manager): 'server3', None) -def test_check_opc_servers_integrity_all_healthy(ingestor_manager): +@patch('ingestor.managers.ingestor_manager.metrics') +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 @@ -523,6 +540,9 @@ def test_check_opc_servers_integrity_all_healthy(ingestor_manager): # 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)) + def test_check_opc_servers_integrity_server_lost(ingestor_manager): # Setup mock OPC manager that will be lost