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.
This commit is contained in:
vitor-aignosi
2025-11-03 10:33:05 -03:00
parent dfd08e883f
commit d93c497914
4 changed files with 62 additions and 24 deletions

View File

@@ -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(

View File

@@ -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',

View File

@@ -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()

View File

@@ -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():