From d93c497914c21cec8c8608bb455e4b7613106619 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 3 Nov 2025 10:33:05 -0300 Subject: [PATCH] SIENTIAPDE-1325 Enhance Ingestor Class with Asynchronous Metrics and Monitoring Integration - Refactored Ingestor class to inherit from SientiaMonitoring for improved observability. - Updated methods to utilize asynchronous operations for declaring active ingestors and managing slots. - Integrated pod_id label into various metrics for better tracking. - Adjusted unit tests to reflect changes in the Ingestor initialization and asynchronous behavior. --- ingestor/ingestor.py | 65 ++++++++++++++++++++++++++++--------- ingestor/metrics.py | 3 ++ tests/unit/test_ingestor.py | 12 ++++--- tests/unit/test_metrics.py | 6 ++-- 4 files changed, 62 insertions(+), 24 deletions(-) diff --git a/ingestor/ingestor.py b/ingestor/ingestor.py index 1028e8c..4b76888 100644 --- a/ingestor/ingestor.py +++ b/ingestor/ingestor.py @@ -5,12 +5,13 @@ from typing import Any from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.observability.logger import get_logger from sientia_do.observability.metrics_controller import MetricsController +from sientia_do.observability.sientia_monitoring import SientiaMonitoring import ingestor.metrics as metrics from ingestor.managers.ingestor_manager import IngestorManager -class Ingestor: +class Ingestor(SientiaMonitoring): """ Main OPC Ingestor class that orchestrates data collection from OPC UA servers. @@ -100,6 +101,13 @@ class Ingestor: ) self.metrics_controller = MetricsController(logger=self.logger) + SientiaMonitoring.__init__( + self, + logger=self.logger, + notification_handler=self.notification_handler, + metrics_controller=self.metrics_controller, + ) + self.metadata = { 'model_id': '-', 'model_name': '-', @@ -192,17 +200,22 @@ class Ingestor: assert self.ingestor_manager is not None # Declare ingestor active - self.ingestor_manager.declare_active() + await self.ingestor_manager.declare_active() # Get slot lease - acquired = self.ingestor_manager.get_slot_leases() + acquired = await self.ingestor_manager.get_slot_leases() self.logger.info(f'Acquired slots: {acquired}') await self.handle_acquired_tags(acquired) - metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set( - len(self.ingestor_manager.managed_tags) - ) # Set initial + await self.emit_metric( + metric_object=metrics.SLOTS_MANAGED, + method='set', + value=len(self.ingestor_manager.managed_tags), + tags={ + 'pod_id': self.pod_id, + }, + ) async def manage_no_slots(self, number_of_slots: int): """ @@ -273,9 +286,13 @@ class Ingestor: await self.ingestor_manager.unsubscribe_slot(lease) self.ingestor_manager.managed_tags.pop(lease) - # Update metric after removal - metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set( - len(self.ingestor_manager.managed_tags) + await self.emit_metric( + metric_object=metrics.SLOTS_MANAGED, + method='set', + value=len(self.ingestor_manager.managed_tags), + tags={ + 'pod_id': self.pod_id, + }, ) async def update_ingestor_manager(self, old_managed_tags: dict[str, Any]): @@ -337,9 +354,13 @@ class Ingestor: self.logger.info(f'Unsubscribing from slot {slot}') await self.ingestor_manager.unsubscribe_slot(slot) - # Ensure the gauge is updated after any potential changes here - metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set( - len(self.ingestor_manager.managed_tags) + await self.emit_metric( + metric_object=metrics.SLOTS_MANAGED, + method='set', + value=len(self.ingestor_manager.managed_tags), + tags={ + 'pod_id': self.pod_id, + }, ) async def loop(self): @@ -369,7 +390,7 @@ class Ingestor: if not self.ingestor_manager: return - self.ingestor_manager.declare_active() + await self.ingestor_manager.declare_active() self.logger.info('Polling for slot updates...') # Get active ingestors @@ -382,7 +403,14 @@ class Ingestor: number_of_slots = await self.ingestor_manager.get_number_of_slots() # Update active ingestors gauge - metrics.ACTIVE_INGESTORS.set(number_of_ingestors) + await self.emit_metric( + metric_object=metrics.ACTIVE_INGESTORS, + method='set', + value=number_of_ingestors, + tags={ + 'pod_id': self.pod_id, + }, + ) # Handle no slots self.logger.info('Managing no slots...') @@ -396,8 +424,13 @@ class Ingestor: await self.manage_leases(available_slots, lacking_ingestors, slot_diff) # Update managed slots gauge - metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set( - len(self.ingestor_manager.managed_tags) + await self.emit_metric( + metric_object=metrics.SLOTS_MANAGED, + method='set', + value=len(self.ingestor_manager.managed_tags), + tags={ + 'pod_id': self.pod_id, + }, ) self.logger.debug( diff --git a/ingestor/metrics.py b/ingestor/metrics.py index 8c71814..ca44e22 100644 --- a/ingestor/metrics.py +++ b/ingestor/metrics.py @@ -58,14 +58,17 @@ APP_UP = Gauge( ACTIVE_INGESTORS = Gauge( 'ingestor_active_total', 'Number of active ingestors reported by Redis', + POD_ID_LABEL, ) SLOTS_TOTAL = Gauge( 'ingestor_slots_total', 'Total number of slots configured in Redis', + POD_ID_LABEL, ) LEASES_TOTAL = Gauge( 'ingestor_leases_total', 'Total number of leases (allocated slots) in Redis', + POD_ID_LABEL, ) SLOTS_MANAGED = Gauge( 'ingestor_slots_managed_current', diff --git a/tests/unit/test_ingestor.py b/tests/unit/test_ingestor.py index ea1afd6..63ea638 100644 --- a/tests/unit/test_ingestor.py +++ b/tests/unit/test_ingestor.py @@ -11,13 +11,13 @@ def test___init__(notification_handler, getenv): getenv.side_effect = [ 'localhost:9092,localhost:35', # KAFKA_SERVERS 'true', # EXPORT_TO_KAFKA - 'localhost1', # REDIS_HOST + 'localhost', # REDIS_HOST '63790', # REDIS_PORT 'user', # REDIS_USERNAME 'password', # REDIS_PASSWORD '100', # LEASE_TTL '200', # HEARTBEAT_TTL - 'localhost1', # HOSTNAME + 'localhost', # HOSTNAME '50', # POLL_INTERVAL 'localhost:27017', # MONGODB_URL 'sientia', # MONGODB_USERNAME @@ -38,20 +38,20 @@ def test___init__(notification_handler, getenv): getenv.assert_any_call('POLL_INTERVAL', '5') assert ingestor.kafka_servers == ['localhost:9092', 'localhost:35'] - assert ingestor.redis_host == 'localhost1' + assert ingestor.redis_host == 'localhost' assert ingestor.redis_port == 63790 assert ingestor.redis_username == 'user' assert ingestor.redis_password == 'password' assert ingestor.lease_ttl == 100 assert ingestor.heartbeat_ttl == 200 - assert ingestor.pod_id == 'localhost1' + assert ingestor.pod_id == 'localhost' assert ingestor.poll_interval == 50 assert ingestor.metadata == { 'model_id': '-', 'model_name': '-', 'workflow_name': 'opc_ingestor', 'schema_name': 'opc_ingestor', - 'pod_id': 'localhost1', + 'pod_id': 'localhost', } notification_handler.assert_called_once_with( @@ -117,6 +117,8 @@ async def test_prepare_ingestor(ingestor_manager_mock, ingestor): ingestor_manager = ingestor_manager_mock.return_value ingestor_manager.get_slot_leases.return_value = True + ingestor_manager.declare_active = AsyncMock() + ingestor_manager.get_slot_leases = AsyncMock() ingestor.handle_acquired_tags = AsyncMock() await ingestor.prepare_ingestor() diff --git a/tests/unit/test_metrics.py b/tests/unit/test_metrics.py index 9fc545e..809abdf 100644 --- a/tests/unit/test_metrics.py +++ b/tests/unit/test_metrics.py @@ -54,7 +54,7 @@ def test_active_ingestors(): assert metrics.ACTIVE_INGESTORS is not None assert isinstance(metrics.ACTIVE_INGESTORS, Gauge) assert metrics.ACTIVE_INGESTORS._name == 'ingestor_active_total' - assert set(metrics.ACTIVE_INGESTORS._labelnames) == set() + assert set(metrics.ACTIVE_INGESTORS._labelnames) == {'pod_id'} def test_slots_total(): @@ -62,7 +62,7 @@ def test_slots_total(): assert metrics.SLOTS_TOTAL is not None assert isinstance(metrics.SLOTS_TOTAL, Gauge) assert metrics.SLOTS_TOTAL._name == 'ingestor_slots_total' - assert set(metrics.SLOTS_TOTAL._labelnames) == set() + assert set(metrics.SLOTS_TOTAL._labelnames) == {'pod_id'} def test_leases_total(): @@ -70,7 +70,7 @@ def test_leases_total(): assert metrics.LEASES_TOTAL is not None assert isinstance(metrics.LEASES_TOTAL, Gauge) assert metrics.LEASES_TOTAL._name == 'ingestor_leases_total' - assert set(metrics.LEASES_TOTAL._labelnames) == set() + assert set(metrics.LEASES_TOTAL._labelnames) == {'pod_id'} def test_slots_managed():