SIENTIAPDE-1083: add prometheus metrics to ingestor_manager.
This commit is contained in:
@@ -7,6 +7,7 @@ from sientia_do.notifications.models import NotificationLevel
|
|||||||
from ingestor.managers.data_manager import DataManager
|
from ingestor.managers.data_manager import DataManager
|
||||||
from ingestor.managers.opc_manager import OpcManager
|
from ingestor.managers.opc_manager import OpcManager
|
||||||
from ingestor.managers.resource_manager import ResourceManager
|
from ingestor.managers.resource_manager import ResourceManager
|
||||||
|
import ingestor.metrics as metrics
|
||||||
|
|
||||||
|
|
||||||
class IngestorManager():
|
class IngestorManager():
|
||||||
@@ -162,6 +163,8 @@ class IngestorManager():
|
|||||||
self.opc_managers[server].disconnect()
|
self.opc_managers[server].disconnect()
|
||||||
self.opc_managers.pop(server, None)
|
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):
|
def check_opc_servers_integrity(self):
|
||||||
"""
|
"""
|
||||||
Checks the integrity of the OPC servers and updates the OPC servers if necessary.
|
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():
|
for slot, _config in self.managed_tags.items():
|
||||||
self.managed_tags[slot].pop(server, None)
|
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):
|
def declare_active(self):
|
||||||
"""
|
"""
|
||||||
Declares the ingestor as active by sending a heartbeat signal to the resource manager.
|
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()
|
leases = self.resource_manager.get_all_leases()
|
||||||
self.number_of_slots = len(leases) if leases else 0
|
self.number_of_slots = len(leases) if leases else 0
|
||||||
|
metrics.LEASES_TOTAL.set(self.number_of_slots)
|
||||||
return self.number_of_slots
|
return self.number_of_slots
|
||||||
|
|
||||||
def get_number_of_slots(self) -> int:
|
def get_number_of_slots(self) -> int:
|
||||||
@@ -224,6 +230,7 @@ class IngestorManager():
|
|||||||
|
|
||||||
slots = self.resource_manager.get_all_slots()
|
slots = self.resource_manager.get_all_slots()
|
||||||
self.number_of_slots = len(slots) if slots else 0
|
self.number_of_slots = len(slots) if slots else 0
|
||||||
|
metrics.SLOTS_TOTAL.set(self.number_of_slots)
|
||||||
return self.number_of_slots
|
return self.number_of_slots
|
||||||
|
|
||||||
def get_slot_leases(self, max_slots: int = 1) -> Dict:
|
def get_slot_leases(self, max_slots: int = 1) -> Dict:
|
||||||
@@ -255,9 +262,11 @@ class IngestorManager():
|
|||||||
if slots is None:
|
if slots is None:
|
||||||
continue
|
continue
|
||||||
acquired[str(i)] = slots
|
acquired[str(i)] = slots
|
||||||
|
metrics.SLOTS_ACQUIRED.labels(pod_id=self.pod_id).inc()
|
||||||
|
|
||||||
if len(acquired) >= max_slots:
|
if len(acquired) >= max_slots:
|
||||||
self.managed_tags.update(acquired)
|
self.managed_tags.update(acquired)
|
||||||
|
metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(len(self.managed_tags))
|
||||||
return acquired
|
return acquired
|
||||||
|
|
||||||
self.logger.warning(
|
self.logger.warning(
|
||||||
@@ -265,6 +274,7 @@ class IngestorManager():
|
|||||||
f"Only {acquired} slots were leased."
|
f"Only {acquired} slots were leased."
|
||||||
)
|
)
|
||||||
self.managed_tags.update(acquired)
|
self.managed_tags.update(acquired)
|
||||||
|
metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(len(self.managed_tags))
|
||||||
return acquired
|
return acquired
|
||||||
|
|
||||||
def unsubscribe_slot(self, slot: str):
|
def unsubscribe_slot(self, slot: str):
|
||||||
@@ -332,6 +342,7 @@ class IngestorManager():
|
|||||||
|
|
||||||
for lease_id in ids:
|
for lease_id in ids:
|
||||||
self.resource_manager.drop_tag_lease(lease_id)
|
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:
|
def manage_server(self, slot: str, server: str, server_config: dict, tags: dict) -> int:
|
||||||
"""
|
"""
|
||||||
@@ -390,6 +401,7 @@ class IngestorManager():
|
|||||||
tags_to_sub
|
tags_to_sub
|
||||||
)
|
)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
|
metrics.OPC_SUBSCRIPTION_ERRORS.labels(pod_id=self.pod_id, server=server, slot=slot).inc()
|
||||||
trace = traceback.format_exc()
|
trace = traceback.format_exc()
|
||||||
self.notification_handler.build_and_send_notification(
|
self.notification_handler.build_and_send_notification(
|
||||||
notification_id=f'OPC_SUBSCRIPTION_ERROR_{slot}:{server}',
|
notification_id=f'OPC_SUBSCRIPTION_ERROR_{slot}:{server}',
|
||||||
|
|||||||
@@ -111,7 +111,8 @@ def test_initialize_opc_from_config_exception(traceback_mock, opc_manager, inges
|
|||||||
|
|
||||||
|
|
||||||
@patch('ingestor.managers.ingestor_manager.OpcManager')
|
@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(
|
manager1 = MagicMock(
|
||||||
config={"config": "config1"})
|
config={"config": "config1"})
|
||||||
manager2 = MagicMock(
|
manager2 = MagicMock(
|
||||||
@@ -175,6 +176,9 @@ def test_update_opc_servers(opc_manager, ingestor_manager):
|
|||||||
|
|
||||||
assert ingestor_manager.opc_managers['server2'] != mock
|
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):
|
def test_declare_active(ingestor_manager):
|
||||||
ingestor_manager.resource_manager.ingestor_heartbeat = MagicMock()
|
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()
|
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(
|
ingestor_manager.resource_manager.get_all_leases = MagicMock(
|
||||||
return_value=["lease1", "lease2"])
|
return_value=["lease1", "lease2"])
|
||||||
result = ingestor_manager.get_number_of_leases()
|
result = ingestor_manager.get_number_of_leases()
|
||||||
assert result == 2
|
assert result == 2
|
||||||
ingestor_manager.resource_manager.get_all_leases.assert_called_once()
|
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(
|
ingestor_manager.resource_manager.get_all_leases = MagicMock(
|
||||||
return_value=None)
|
return_value=None)
|
||||||
result = ingestor_manager.get_number_of_leases()
|
result = ingestor_manager.get_number_of_leases()
|
||||||
assert result == 0
|
assert result == 0
|
||||||
ingestor_manager.resource_manager.get_all_leases.assert_called_once()
|
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(
|
ingestor_manager.resource_manager.get_all_slots = MagicMock(
|
||||||
return_value=["slot1", "slot2"])
|
return_value=["slot1", "slot2"])
|
||||||
result = ingestor_manager.get_number_of_slots()
|
result = ingestor_manager.get_number_of_slots()
|
||||||
assert result == 2
|
assert result == 2
|
||||||
ingestor_manager.resource_manager.get_all_slots.assert_called_once()
|
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(
|
ingestor_manager.resource_manager.get_all_slots = MagicMock(
|
||||||
return_value=None)
|
return_value=None)
|
||||||
result = ingestor_manager.get_number_of_slots()
|
result = ingestor_manager.get_number_of_slots()
|
||||||
assert result == 0
|
assert result == 0
|
||||||
ingestor_manager.resource_manager.get_all_slots.assert_called_once()
|
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):
|
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
|
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.resource_manager.drop_tag_lease = MagicMock()
|
||||||
ingestor_manager.drop_slot_leases(["1", "2"])
|
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("1")
|
||||||
ingestor_manager.resource_manager.drop_tag_lease.assert_any_call("2")
|
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):
|
def test_manage_server_no_server(ingestor_manager):
|
||||||
ingestor_manager.opc_managers = {
|
ingestor_manager.opc_managers = {
|
||||||
@@ -491,7 +507,8 @@ def test_subscribe_to_tags(ingestor_manager):
|
|||||||
'server3', None)
|
'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
|
# Setup mock OPC managers
|
||||||
opc_manager1 = MagicMock()
|
opc_manager1 = MagicMock()
|
||||||
opc_manager1.check_cycles.return_value = None
|
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
|
# Verify that no reinitialization was needed
|
||||||
ingestor_manager.initialize_opc_from_config.assert_not_called()
|
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):
|
def test_check_opc_servers_integrity_server_lost(ingestor_manager):
|
||||||
# Setup mock OPC manager that will be lost
|
# Setup mock OPC manager that will be lost
|
||||||
|
|||||||
Reference in New Issue
Block a user