From 85a371a38b51de2cc23f23712433f59a847d2e04 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Fri, 31 Oct 2025 16:32:45 -0300 Subject: [PATCH] SIENTIAPDE-1325 SIENTIAPDE-1325: Refactor Ingestor and Manager Classes for Enhanced Asynchronous Operations - Introduced asynchronous methods across Ingestor, IngestorManager, DataManager, and OpcManager to improve performance and responsiveness. - Integrated MetricsController into various classes for better observability and monitoring. - Updated Redis and MongoDB interactions to support asynchronous operations, enhancing data handling efficiency. - Removed deprecated Redis metrics and streamlined resource management logic. - Adjusted unit tests to accommodate the new asynchronous behavior and ensure proper mocking of async methods. --- ingestor/ingestor.py | 21 +-- ingestor/managers/data_manager.py | 86 ++++++++----- ingestor/managers/ingestor_manager.py | 134 +++++++++++++++----- ingestor/managers/opc_manager.py | 127 +++++++++++++++---- ingestor/managers/resource_manager.py | 155 ++++++----------------- ingestor/metrics.py | 20 +-- tests/unit/managers/test_data_manager.py | 61 +++++---- 7 files changed, 336 insertions(+), 268 deletions(-) diff --git a/ingestor/ingestor.py b/ingestor/ingestor.py index 5f27759..1028e8c 100644 --- a/ingestor/ingestor.py +++ b/ingestor/ingestor.py @@ -4,6 +4,7 @@ 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 import ingestor.metrics as metrics from ingestor.managers.ingestor_manager import IngestorManager @@ -97,6 +98,7 @@ class Ingestor: logger=self.logger, project_name='opc_ingestor', ) + self.metrics_controller = MetricsController(logger=self.logger) self.metadata = { 'model_id': '-', @@ -184,6 +186,7 @@ class Ingestor: metadata=self.metadata, logger=self.logger, notification_handler=self.notification_handler, + metrics_controller=self.metrics_controller, export_to_kafka=self.export_to_kafka, ) assert self.ingestor_manager is not None @@ -201,7 +204,7 @@ class Ingestor: len(self.ingestor_manager.managed_tags) ) # Set initial - def manage_no_slots(self, number_of_slots: int): + async def manage_no_slots(self, number_of_slots: int): """ Manages the scenario where there are no slots assigned to the ingestor. @@ -222,7 +225,7 @@ class Ingestor: # This ingestor is active and has no slots, so we need to try to # Get slot lease - self.ingestor_manager.get_slot_leases(1) + await self.ingestor_manager.get_slot_leases(1) async def manage_leases(self, available_slots: int, lacking_ingestors: int, slot_diff: int): """ @@ -255,7 +258,7 @@ class Ingestor: self.logger.info(f'Slots available: {available_slots}') # Get slot lease - self.ingestor_manager.get_slot_leases(available_slots) + await self.ingestor_manager.get_slot_leases(available_slots) elif lacking_ingestors <= 0 and slot_diff > 0: self.logger.info(f'Extra slots available: {slot_diff}') @@ -264,7 +267,7 @@ class Ingestor: overleases = list(self.ingestor_manager.managed_tags.keys())[1:] - self.ingestor_manager.drop_slot_leases(overleases) + await self.ingestor_manager.drop_slot_leases(overleases) for lease in overleases: await self.ingestor_manager.unsubscribe_slot(lease) @@ -373,17 +376,17 @@ class Ingestor: current_managed_tags = deepcopy(self.ingestor_manager.managed_tags) - ingestors = self.ingestor_manager.get_active_ingestors() + ingestors = await self.ingestor_manager.get_active_ingestors() number_of_ingestors = len(ingestors) - number_of_leases = self.ingestor_manager.get_number_of_leases() - number_of_slots = self.ingestor_manager.get_number_of_slots() + number_of_leases = await self.ingestor_manager.get_number_of_leases() + number_of_slots = await self.ingestor_manager.get_number_of_slots() # Update active ingestors gauge metrics.ACTIVE_INGESTORS.set(number_of_ingestors) # Handle no slots self.logger.info('Managing no slots...') - self.manage_no_slots(number_of_slots) + await self.manage_no_slots(number_of_slots) available_slots = number_of_slots - number_of_leases lacking_ingestors = number_of_slots - number_of_ingestors @@ -410,7 +413,7 @@ class Ingestor: # Update opc servers self.logger.info('Updating slot config...') - self.ingestor_manager.update_slot_config() + await self.ingestor_manager.update_slot_config() # Check OPC cycles self.logger.info('Checking OPC servers integrity...') diff --git a/ingestor/managers/data_manager.py b/ingestor/managers/data_manager.py index 8a9d970..a493aa2 100644 --- a/ingestor/managers/data_manager.py +++ b/ingestor/managers/data_manager.py @@ -5,17 +5,17 @@ from time import sleep from kafka import KafkaProducer from kafka.errors import NoBrokersAvailable -from pymongo import MongoClient from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.notifications.models import NotificationLevel from sientia_do.observability.logger import Logger -from sientia_do.temporal.activities.base import BaseActivity -from sientia_do.temporal.constants import now +from sientia_do.observability.metrics_controller import MetricsController +from sientia_do.observability.sientia_monitoring import SientiaMonitoring +from sientia_do.repository.mongodb_repository import MongoDBRepository import ingestor.metrics as metrics -class DataManager(BaseActivity): +class DataManager(SientiaMonitoring): """ Manages data persistence and export operations for the OPC Ingestor. @@ -57,6 +57,7 @@ class DataManager(BaseActivity): metadata: dict, logger: Logger, notification_handler: NotificationHandler, + metrics_controller: MetricsController, ) -> None: """ Initializes the DataManager instance with Kafka and MongoDB connections. @@ -91,6 +92,13 @@ class DataManager(BaseActivity): self.kafka_producer = None self.export_to_kafka = export_to_kafka + SientiaMonitoring.__init__( + self, + logger=logger, + metrics_controller=metrics_controller, + notification_handler=notification_handler, + ) + if self.export_to_kafka: for i in range(0, 3): logger.info( @@ -129,19 +137,18 @@ class DataManager(BaseActivity): self.connection_string = mongo_connection_string self.database = mongo_database - self.mongo_client: MongoClient = MongoClient(self.connection_string) - self.mongo_client.server_info() + self.mongo_repository = MongoDBRepository( + connection_string=self.connection_string, + database_name=self.database, + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + ) self.metadata = metadata - self.mongo_db = self.mongo_client[self.database] - logger.info(f'DataManager initialized with MongoDB servers: {self.connection_string}') - BaseActivity.__init__( - self, logger=logger, notification_handler=notification_handler, set_error_counter=True - ) - def shutdown(self): """ Gracefully shuts down the DataManager and closes all connections. @@ -165,13 +172,10 @@ class DataManager(BaseActivity): else: self.logger.warning('Kafka producer is already closed or not initialized.') - if self.mongo_client: - try: - self.mongo_client.close() - except Exception as e: - self.logger.error(f'Error closing MongoDB client: {e}') - else: - self.logger.warning('MongoDB client is already closed or not initialized.') + try: + self.mongo_repository.close() + except Exception as e: + self.logger.error(f'Error closing MongoDB client: {e}') def __del__(self): self.shutdown() @@ -203,7 +207,7 @@ class DataManager(BaseActivity): """ self.logger.error(f'Delivery failed for record : {err}') - def publish(self, topic: str, data: dict) -> None: + async def publish(self, topic: str, data: dict) -> None: """ Publishes a message to a specified Kafka topic. @@ -226,12 +230,24 @@ class DataManager(BaseActivity): ).add_errback(self.delivery_error) self.kafka_producer.flush(timeout=10) - metrics.KAFKA_MESSAGES_SENT.labels(pod_id=self.pod_id, topic=topic).inc() + await self.emit_metric( + metric_object=metrics.KAFKA_MESSAGES_SENT, + tags={ + 'pod_id': self.pod_id, + 'topic': topic, + }, + ) except Exception as e: - metrics.KAFKA_MESSAGES_ERRORS.labels(pod_id=self.pod_id, topic=topic).inc() + await self.emit_metric( + metric_object=metrics.KAFKA_MESSAGES_ERRORS, + tags={ + 'pod_id': self.pod_id, + 'topic': topic, + }, + ) trace = traceback.format_exc() - self.send_notification( + await self.send_notification_async( metadata=self.metadata, notification_id=f'KAFKA_PRODUCER_ERROR_{topic}', message=f'Error publishing message to topic {topic}: {e}', @@ -242,23 +258,25 @@ class DataManager(BaseActivity): self.logger.error(trace) try: - collection = self.mongo_db[topic] - - collection.insert_one( - { - **data, - 'inserted_at': now(), - } + await self.mongo_repository.insert( + collection_name=topic, + document=data, + metadata=self.metadata, ) self.logger.debug(f'Message inserted into MongoDB collection {topic}: {data}') - metrics.TAG_WRITTEN_COUNT.labels( - pod_id=self.pod_id, tag_name=data['name'], collection_name=topic - ).inc() + await self.emit_metric( + metric_object=metrics.TAG_WRITTEN_COUNT, + tags={ + 'pod_id': self.pod_id, + 'tag_name': data['name'], + 'collection_name': topic, + }, + ) except Exception as e: trace = traceback.format_exc() - self.send_notification( + await self.send_notification_async( metadata=self.metadata, notification_id=f'MONGO_PRODUCER_ERROR_{topic}', message=f'Error inserting message to MongoDB: {e}', diff --git a/ingestor/managers/ingestor_manager.py b/ingestor/managers/ingestor_manager.py index fc0a534..57e192b 100644 --- a/ingestor/managers/ingestor_manager.py +++ b/ingestor/managers/ingestor_manager.py @@ -5,7 +5,8 @@ from copy import deepcopy from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.notifications.models import NotificationLevel from sientia_do.observability.logger import Logger -from sientia_do.temporal.activities.base import BaseActivity +from sientia_do.observability.metrics_controller import MetricsController +from sientia_do.observability.sientia_monitoring import SientiaMonitoring import ingestor.metrics as metrics from ingestor.managers.data_manager import DataManager @@ -13,7 +14,7 @@ from ingestor.managers.opc_manager import OpcManager from ingestor.managers.resource_manager import ResourceManager -class IngestorManager(BaseActivity): +class IngestorManager(SientiaMonitoring): """ Central coordinator for managing OPC data ingestion operations. @@ -70,6 +71,7 @@ class IngestorManager(BaseActivity): metadata: dict, logger: Logger, notification_handler: NotificationHandler, + metrics_controller: MetricsController, export_to_kafka: bool = False, ): redis_host: str = redis_data['host'] @@ -85,6 +87,7 @@ class IngestorManager(BaseActivity): metadata=metadata, logger=logger, notification_handler=notification_handler, + metrics_controller=metrics_controller, ) self.opc_managers: dict = {} self.resource_manager = ResourceManager( @@ -97,6 +100,7 @@ class IngestorManager(BaseActivity): notification_handler=notification_handler, username=redis_username, password=redis_password, + metrics_controller=metrics_controller, ) self.number_of_slots = 0 self.poll_interval = poll_interval @@ -105,8 +109,12 @@ class IngestorManager(BaseActivity): self.metadata = metadata - BaseActivity.__init__( - self, logger=logger, notification_handler=notification_handler, set_error_counter=True + 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: @@ -150,6 +158,7 @@ class IngestorManager(BaseActivity): cert_path=server_config.get('cert_path'), private_key_path=server_config.get('private_key_path'), server_cert_path=server_config.get('server_cert_path'), + metrics_controller=self.metrics_controller, ) manager.config = server_config @@ -264,7 +273,14 @@ class IngestorManager(BaseActivity): ) await self.remove_server(server) - metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers)) + await self.emit_metric( + metric_object=metrics.OPC_MANAGERS_ACTIVE, + method='set', + value=len(self.opc_managers), + tags={ + 'pod_id': self.pod_id, + }, + ) async def check_opc_servers_integrity(self): """ @@ -294,9 +310,16 @@ class IngestorManager(BaseActivity): for server in to_disconnect: await self.remove_server(server) - metrics.OPC_MANAGERS_ACTIVE.labels(pod_id=self.pod_id).set(len(self.opc_managers)) + await self.emit_metric( + metric_object=metrics.OPC_MANAGERS_ACTIVE, + method='set', + value=len(self.opc_managers), + tags={ + 'pod_id': self.pod_id, + }, + ) - def declare_active(self): + async def declare_active(self): """ Declares the ingestor as active by sending a heartbeat signal to the resource manager. @@ -309,9 +332,9 @@ class IngestorManager(BaseActivity): - Enables load balancing and health monitoring """ - self.resource_manager.ingestor_heartbeat() + await self.resource_manager.ingestor_heartbeat() - def get_active_ingestors(self) -> list[str]: + async def get_active_ingestors(self) -> list[str]: """ Retrieve a list of active ingestors. @@ -325,10 +348,10 @@ class IngestorManager(BaseActivity): the pod identifiers for load balancing and coordination purposes. """ - ingestors = self.resource_manager.get_all_ingestors() + ingestors = await self.resource_manager.get_all_ingestors() return ingestors if ingestors else [] - def get_number_of_leases(self) -> int: + async def get_number_of_leases(self) -> int: """ Retrieves the number of leases managed by the resource manager. @@ -343,12 +366,19 @@ class IngestorManager(BaseActivity): - Updates Prometheus metrics for total leases """ - leases = self.resource_manager.get_all_leases() + leases = await self.resource_manager.get_all_leases() self.number_of_slots = len(leases) if leases else 0 - metrics.LEASES_TOTAL.set(self.number_of_slots) + await self.emit_metric( + metric_object=metrics.LEASES_TOTAL, + method='set', + value=self.number_of_slots, + tags={ + 'pod_id': self.pod_id, + }, + ) return self.number_of_slots - def get_number_of_slots(self) -> int: + async def get_number_of_slots(self) -> int: """ Retrieves the number of slots managed by the resource manager. @@ -363,12 +393,19 @@ class IngestorManager(BaseActivity): - Updates Prometheus metrics for total slots """ - slots = self.resource_manager.get_all_slots() + slots = await self.resource_manager.get_all_slots() self.number_of_slots = len(slots) if slots else 0 - metrics.SLOTS_TOTAL.set(self.number_of_slots) + await self.emit_metric( + metric_object=metrics.SLOTS_TOTAL, + method='set', + value=self.number_of_slots, + tags={ + 'pod_id': self.pod_id, + }, + ) return self.number_of_slots - def get_slot_leases(self, max_slots: int = 1) -> dict: + async def get_slot_leases(self, max_slots: int = 1) -> dict: """ Acquires a specified number of resource slots by leasing them from the resource manager. @@ -398,24 +435,45 @@ class IngestorManager(BaseActivity): acquired = {} for i in range(1, self.number_of_slots + 1): - if self.resource_manager.lease_tag(str(i)): + if await self.resource_manager.lease_tag(str(i)): self.logger.info(f'Leased slot {i}') - slots = self.resource_manager.get_tag_slot(str(i)) + slots = await self.resource_manager.get_tag_slot(str(i)) if slots is None: continue acquired[str(i)] = slots - metrics.SLOTS_ACQUIRED.labels(pod_id=self.pod_id).inc() + await self.emit_metric( + metric_object=metrics.SLOTS_ACQUIRED, + method='inc', + value=1, + tags={ + 'pod_id': self.pod_id, + }, + ) if len(acquired) >= max_slots: self.managed_tags.update(acquired) - metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(len(self.managed_tags)) + await self.emit_metric( + metric_object=metrics.SLOTS_MANAGED, + method='set', + value=len(self.managed_tags), + tags={ + 'pod_id': self.pod_id, + }, + ) return acquired self.logger.warning( f'Unable to acquire {max_slots} slots. Only {acquired} slots were leased.' ) self.managed_tags.update(acquired) - metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set(len(self.managed_tags)) + await self.emit_metric( + metric_object=metrics.SLOTS_MANAGED, + method='set', + value=len(self.managed_tags), + tags={ + 'pod_id': self.pod_id, + }, + ) return acquired async def unsubscribe_slot(self, slot: str): @@ -441,7 +499,7 @@ class IngestorManager(BaseActivity): if server in self.opc_managers: await self.opc_managers[server].unsubscribe(slot) - def update_slot_config(self): + async def update_slot_config(self): """ Updates the configuration of managed slots by renewing their leases, fetching the latest configurations, and handling any changes or removals. @@ -470,8 +528,8 @@ class IngestorManager(BaseActivity): removed_slots: list[str] = [] for slot, _slot_config in self.managed_tags.items(): - self.resource_manager.renew_tag_lease(slot) - update = self.resource_manager.get_tag_slot(slot) + await self.resource_manager.renew_tag_lease(slot) + update = await self.resource_manager.get_tag_slot(slot) if update is None: removed_slots.append(slot) continue @@ -481,7 +539,7 @@ class IngestorManager(BaseActivity): for slot in removed_slots: self.managed_tags.pop(slot, None) - def drop_slot_leases(self, ids: list[str]) -> None: + async def drop_slot_leases(self, ids: list[str]) -> None: """ Releases the leases associated with the specified slot IDs. @@ -502,8 +560,15 @@ class IngestorManager(BaseActivity): - Updates metrics for released slots count """ for lease_id in ids: - self.resource_manager.drop_tag_lease(lease_id) - metrics.SLOTS_RELEASED.labels(pod_id=self.pod_id).inc() + await self.resource_manager.drop_tag_lease(lease_id) + await self.emit_metric( + metric_object=metrics.SLOTS_RELEASED, + method='inc', + value=1, + tags={ + 'pod_id': self.pod_id, + }, + ) async def manage_server(self, slot: str, server: str, server_config: dict, tags: dict) -> int: """ @@ -560,11 +625,16 @@ class IngestorManager(BaseActivity): ) self.logger.info(tags_to_sub) except Exception as e: - metrics.OPC_SUBSCRIPTION_ERRORS.labels( - pod_id=self.pod_id, server=server, slot=slot - ).inc() + await self.emit_metric( + metric_object=metrics.OPC_SUBSCRIPTION_ERRORS, + tags={ + 'pod_id': self.pod_id, + 'server': server, + 'slot': slot, + }, + ) trace = traceback.format_exc() - self.send_notification( + await self.send_notification_async( metadata=self.metadata, notification_id=f'OPC_SUBSCRIPTION_ERROR_{slot}:{server}', message=f'Failed to subscribe to tags from {slot}:{server}\n{tags}: {e}', diff --git a/ingestor/managers/opc_manager.py b/ingestor/managers/opc_manager.py index 2c0ba7e..0558d57 100644 --- a/ingestor/managers/opc_manager.py +++ b/ingestor/managers/opc_manager.py @@ -8,14 +8,15 @@ from asyncua.crypto.security_policies import SecurityPolicyBasic256 from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.notifications.models import NotificationLevel from sientia_do.observability.logger import Logger -from sientia_do.temporal.activities.base import BaseActivity +from sientia_do.observability.metrics_controller import MetricsController +from sientia_do.observability.sientia_monitoring import SientiaMonitoring from sientia_do.temporal.constants import DATETIME_FORMAT_WITH_TZ, OPC_TIMEZONE import ingestor.metrics as metrics from ingestor.managers.data_manager import DataManager -class OpcManager(BaseActivity): +class OpcManager(SientiaMonitoring): """ Manages OPC UA server connections and tag subscriptions. @@ -66,6 +67,7 @@ class OpcManager(BaseActivity): logger: Logger, server_uri: str, notification_handler: NotificationHandler, + metrics_controller: MetricsController, metadata: dict, cert_path: str | None = None, private_key_path: str | None = None, @@ -86,8 +88,12 @@ class OpcManager(BaseActivity): self.data_manager = data_manager self.metadata = metadata - BaseActivity.__init__( - self, logger=logger, notification_handler=notification_handler, set_error_counter=True + SientiaMonitoring.__init__( + self, + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + set_error_counter=True, ) metrics.OPC_CONNECTION_STATUS.labels( @@ -192,7 +198,15 @@ class OpcManager(BaseActivity): - OPC_CONNECTION_STATUS: Set to 1 on successful connection """ - metrics.OPC_CONNECTIONS_TOTAL.labels(pod_id=self.pod_id, server_name=self.name).inc() + await self.emit_metric( + metric_object=metrics.OPC_CONNECTIONS_TOTAL, + method='inc', + value=1, + tags={ + 'pod_id': self.pod_id, + 'server_name': self.name, + }, + ) try: self.client = Client(self.url, timeout=10, watchdog_intervall=3600000) assert self.client is not None # Informa ao mypy que client não é None @@ -207,9 +221,16 @@ class OpcManager(BaseActivity): await self.set_security() self.logger.info(f'Starting connection to {self.name}...') await self.client.connect() - metrics.OPC_CONNECTION_STATUS.labels( - pod_id=self.pod_id, server_name=self.name, server_url=self.url - ).set(1) + await self.emit_metric( + metric_object=metrics.OPC_CONNECTION_STATUS, + method='set', + value=1, + tags={ + 'pod_id': self.pod_id, + 'server_name': self.name, + 'server_url': self.url, + }, + ) self.logger.info(f'Connection to {self.name} successful.') except Exception: await self.disconnect() @@ -244,9 +265,16 @@ class OpcManager(BaseActivity): self.subscription_period_ms, self ) self.logger.info(f'Subscription {name} created on {self.name}.') - metrics.OPC_SUBSCRIPTIONS_CREATED.labels( - pod_id=self.pod_id, server_name=self.name, slot_name=name - ).inc() + await self.emit_metric( + metric_object=metrics.OPC_SUBSCRIPTIONS_CREATED, + method='inc', + value=1, + tags={ + 'pod_id': self.pod_id, + 'server_name': self.name, + 'slot_name': name, + }, + ) except Exception as e: self.logger.error(f'Failed to create subscription {name} on {self.name}: {e}') raise @@ -297,8 +325,14 @@ class OpcManager(BaseActivity): await self.subscriptions[subscription].subscribe_data_change(addr_nodes) - metrics.OPC_TAGS_SUBSCRIBED.labels(pod_id=self.pod_id, server_name=self.name).set( - len(self.nodes) + await self.emit_metric( + metric_object=metrics.OPC_TAGS_SUBSCRIBED, + method='set', + value=len(self.nodes), + tags={ + 'pod_id': self.pod_id, + 'server_name': self.name, + }, ) async def unsubscribe(self, subscription: str): @@ -387,7 +421,7 @@ class OpcManager(BaseActivity): errors = await self.disconnection_fallback() if errors: - self.send_notification( + await self.send_notification_async( metadata=self.metadata, notification_id=f'OPC_DISCONNECTION_ERROR_{self.name}', message=f'Failed to disconnect from OPC UA server {self.name} after 5 attempts', @@ -400,10 +434,25 @@ class OpcManager(BaseActivity): del self.client self.client = None - metrics.OPC_CONNECTION_STATUS.labels( - pod_id=self.pod_id, server_name=self.name, server_url=self.url - ).set(0) - metrics.OPC_TAGS_SUBSCRIBED.labels(pod_id=self.pod_id, server_name=self.name).set(0) + await self.emit_metric( + metric_object=metrics.OPC_CONNECTION_STATUS, + method='set', + value=0, + tags={ + 'pod_id': self.pod_id, + 'server_name': self.name, + 'server_url': self.url, + }, + ) + await self.emit_metric( + metric_object=metrics.OPC_TAGS_SUBSCRIBED, + method='set', + value=0, + tags={ + 'pod_id': self.pod_id, + 'server_name': self.name, + }, + ) async def datachange_notification(self, node, _val, data): """ @@ -447,13 +496,21 @@ class OpcManager(BaseActivity): } for topic in self.nodes[tag]['topics']: - self.data_manager.publish(topic, data) + await self.data_manager.publish(topic, data) self.nodes[tag]['cycle_rule']['cycle_count'] = 0 self.non_receive_count = 0 - metrics.OPC_CYCLES_WITHOUT_DATA.labels(pod_id=self.pod_id, server_name=self.name).set(0) + await self.emit_metric( + metric_object=metrics.OPC_CYCLES_WITHOUT_DATA, + method='set', + value=0, + tags={ + 'pod_id': self.pod_id, + 'server_name': self.name, + }, + ) - def check_cycles(self): + async def check_cycles(self): """ Checks the cycle counts for all monitored nodes and sends notifications if thresholds are exceeded. @@ -471,7 +528,7 @@ class OpcManager(BaseActivity): if self.nodes[node]['cycle_rule']['cycle_count'] >= 5: name = config['tag_name'] cycles = self.nodes[node]['cycle_rule']['cycle_count'] - self.send_notification( + await self.send_notification_async( metadata=self.metadata, notification_id=f'TAG_{node}:{name}_LISTENNING_STOPPED', message=f'{cycles} cycles without receive from {node}:{name}', @@ -479,7 +536,7 @@ class OpcManager(BaseActivity): level=NotificationLevel.WARNING, ) - def check_opc_listenning(self) -> bool: + async def check_opc_listenning(self) -> bool: """ Checks the OPC connection and triggers notifications if the connection is lost. @@ -499,11 +556,17 @@ class OpcManager(BaseActivity): """ self.non_receive_count += 1 - metrics.OPC_CYCLES_WITHOUT_DATA.labels(pod_id=self.pod_id, server_name=self.name).set( - self.non_receive_count + await self.emit_metric( + metric_object=metrics.OPC_CYCLES_WITHOUT_DATA, + method='set', + value=self.non_receive_count, + tags={ + 'pod_id': self.pod_id, + 'server_name': self.name, + }, ) if self.non_receive_count >= 5: - self.send_notification( + await self.send_notification_async( metadata=self.metadata, notification_id=f'OPC_LISTENNING_STOPPED__{self.name}', message=f'{self.non_receive_count} cycles without ' @@ -512,8 +575,16 @@ class OpcManager(BaseActivity): level=NotificationLevel.ERROR, ) if self.non_receive_count >= 15: - metrics.OPC_RECONNECTIONS_TOTAL.labels(pod_id=self.pod_id, server_name=self.name).inc() - self.send_notification( + await self.emit_metric( + metric_object=metrics.OPC_RECONNECTIONS_TOTAL, + method='inc', + value=1, + tags={ + 'pod_id': self.pod_id, + 'server_name': self.name, + }, + ) + await self.send_notification_async( metadata=self.metadata, notification_id=f'OPC_CONNECTION_RETRY__{self.name}', message=f'Retrying to connect to server {self.name}', diff --git a/ingestor/managers/resource_manager.py b/ingestor/managers/resource_manager.py index 30502c7..b467b00 100644 --- a/ingestor/managers/resource_manager.py +++ b/ingestor/managers/resource_manager.py @@ -1,16 +1,11 @@ -import json -from time import time - -from redis import Redis from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler -from sientia_do.notifications.models import NotificationLevel from sientia_do.observability.logger import Logger -from sientia_do.temporal.activities.base import BaseActivity - -import ingestor.metrics as metrics +from sientia_do.observability.metrics_controller import MetricsController +from sientia_do.observability.sientia_monitoring import SientiaMonitoring +from sientia_do.repository.redis_repository import RedisRepository -class ResourceManager(BaseActivity): +class ResourceManager(SientiaMonitoring): """ Manages Redis-based resource coordination and slot leasing for the OPC Ingestor. @@ -53,6 +48,7 @@ class ResourceManager(BaseActivity): metadata: dict, logger: Logger, notification_handler: NotificationHandler, + metrics_controller: MetricsController, username: str | None = None, password: str | None = None, ) -> None: @@ -81,102 +77,34 @@ class ResourceManager(BaseActivity): Metrics: - REDIS_CONNECTION_STATUS: Set to 1 on successful connection, 0 on failure """ - BaseActivity.__init__( - self, logger=logger, notification_handler=notification_handler, set_error_counter=True + SientiaMonitoring.__init__( + self, + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + set_error_counter=True, ) + try: - self.redis = Redis( + self.redis_repository = RedisRepository( host=host, port=port, - decode_responses=True, username=username, password=password, + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, ) - self.redis.ping() - metrics.REDIS_CONNECTION_STATUS.labels(pod_id=self.pod_id).set(1) + self.redis_repository.redis_client.ping() except Exception as e: self.logger.error(f'Failed to connect to Redis: {e}') - metrics.REDIS_CONNECTION_STATUS.labels(pod_id=self.pod_id).set(0) raise self.lease_ttl = lease_ttl self.heartbeat_ttl = heartbeat_ttl self.metadata = metadata - def _execute_redis_op(self, operation_name: str, func, *args, **kwargs): - """ - Wrapper to execute Redis operations and record metrics. - - This method provides a unified interface for Redis operations that: - - Records operation timing and success/failure metrics - - Handles error notifications consistently - - Ensures all Redis operations are properly monitored - - Args: - operation_name (str): Name of the Redis operation for metrics labeling - func: The Redis function to execute - *args: Positional arguments for the Redis function - **kwargs: Keyword arguments for the Redis function - - Returns: - The result of the Redis operation - - Raises: - Exception: Re-raises any exception from the Redis operation after - recording error metrics and sending notifications. - - Metrics: - - REDIS_OPERATIONS_TOTAL: Incremented on successful operations - - REDIS_OPERATIONS_DURATION: Records operation timing - - REDIS_OPERATIONS_ERRORS: Incremented on operation failures - """ - start_time = time() - try: - result = func(*args, **kwargs) - metrics.REDIS_OPERATIONS_TOTAL.labels( - pod_id=self.pod_id, operation=operation_name - ).inc() - duration = time() - start_time - metrics.REDIS_OPERATIONS_DURATION.labels( - pod_id=self.pod_id, operation=operation_name - ).observe(duration) - return result - except Exception as e: - metrics.REDIS_OPERATIONS_ERRORS.labels( - pod_id=self.pod_id, operation=operation_name - ).inc() - self.send_notification( - metadata=self.metadata, - notification_id=f'REDIS_OPERATION_ERROR_{operation_name}', - message=f"Error in Redis operation '{operation_name}': {e}", - block='redis_manager', - level=NotificationLevel.ERROR, - ) - raise - - def get(self, key: str) -> dict | None: - """ - Retrieve a value from Redis by its key and return it as a dictionary. - - This method fetches a value from Redis and attempts to parse it as JSON. - If the key doesn't exist or the value is empty, it returns None. - - Args: - key (str): The key to look up in Redis. - - Returns: - dict: The value associated with the key, parsed as a dictionary, - or None if the key does not exist or the value is empty. - - Metrics: - - REDIS_OPERATIONS_TOTAL: Incremented with operation="get" - - REDIS_OPERATIONS_DURATION: Records timing for get operations - """ - - history = self._execute_redis_op('get', self.redis.get, key) - return json.loads(history) if history else None - - def get_tag_slot(self, tag_id: str) -> dict | None: + async def get_tag_slot(self, tag_id: str) -> dict | None: """ Retrieve the tag slot information for a given ID. @@ -194,9 +122,9 @@ class ResourceManager(BaseActivity): and delegates to the get() method for the actual Redis operation. """ - return self.get(f'slot:opc_tags:{tag_id}') + return await self.redis_repository.get(f'slot:opc_tags:{tag_id}', metadata=self.metadata) - def ingestor_heartbeat(self) -> None: + async def ingestor_heartbeat(self) -> None: """ Sends a heartbeat signal to Redis to indicate that the ingestor is active. @@ -215,15 +143,11 @@ class ResourceManager(BaseActivity): - REDIS_OPERATIONS_DURATION: Records timing for heartbeat operations """ - self._execute_redis_op( - 'set', - self.redis.set, - f'heartbeat:ingestor:{self.pod_id}', - 1, - ex=self.heartbeat_ttl, + await self.redis_repository.set( + f'heartbeat:ingestor:{self.pod_id}', 1, ttl=self.heartbeat_ttl, metadata=self.metadata ) - def lease_tag(self, tag_id: str) -> bool: + async def lease_tag(self, tag_id: str) -> bool: """ Attempts to lease a tag by setting a key in Redis with a specified TTL. @@ -249,16 +173,15 @@ class ResourceManager(BaseActivity): - REDIS_OPERATIONS_DURATION: Records timing for lease operations """ - return self._execute_redis_op( - 'set_nx', - self.redis.set, + return await self.redis_repository.set( f'lease:opc_tags:{tag_id}', self.pod_id, + ttl=self.lease_ttl, nx=True, - ex=self.lease_ttl, + metadata=self.metadata, ) - def renew_tag_lease(self, tag_id: str) -> bool: + async def renew_tag_lease(self, tag_id: str) -> bool: """ Renews the lease for a specific OPC tag if the current pod holds the lease. @@ -283,15 +206,17 @@ class ResourceManager(BaseActivity): - REDIS_OPERATIONS_DURATION: Records timing for renewal operations """ - current = self._execute_redis_op('get', self.redis.get, f'lease:opc_tags:{tag_id}') + current = await self.redis_repository.get( + f'lease:opc_tags:{tag_id}', metadata=self.metadata + ) if current == self.pod_id: - self._execute_redis_op( - 'expire', self.redis.expire, f'lease:opc_tags:{tag_id}', self.lease_ttl + await self.redis_repository.expire( + f'lease:opc_tags:{tag_id}', self.lease_ttl, metadata=self.metadata ) return True return False - def drop_tag_lease(self, tag_id: str) -> None: + async def drop_tag_lease(self, tag_id: str) -> None: """ Drops the lease for a specific OPC tag. @@ -312,9 +237,9 @@ class ResourceManager(BaseActivity): - REDIS_OPERATIONS_DURATION: Records timing for lease dropping operations """ - self._execute_redis_op('delete', self.redis.delete, f'lease:opc_tags:{tag_id}') + await self.redis_repository.delete(f'lease:opc_tags:{tag_id}', metadata=self.metadata) - def get_all_ingestors(self) -> list[str]: + async def get_all_ingestors(self) -> list[str]: """ Retrieves all active ingestors from Redis. @@ -336,9 +261,9 @@ class ResourceManager(BaseActivity): - REDIS_OPERATIONS_DURATION: Records timing for ingestor discovery """ - return self._execute_redis_op('keys', self.redis.keys, 'heartbeat:ingestor:*') + return await self.redis_repository.keys('heartbeat:ingestor:*', metadata=self.metadata) - def get_all_slots(self) -> list[str]: + async def get_all_slots(self) -> list[str]: """ Retrieves all available slots from Redis. @@ -360,9 +285,9 @@ class ResourceManager(BaseActivity): - REDIS_OPERATIONS_DURATION: Records timing for slot discovery """ - return self._execute_redis_op('keys', self.redis.keys, 'slot:opc_tags:*') + return await self.redis_repository.keys('slot:opc_tags:*', metadata=self.metadata) - def get_all_leases(self) -> list[str]: + async def get_all_leases(self) -> list[str]: """ Retrieves all active leases from Redis. @@ -384,4 +309,4 @@ class ResourceManager(BaseActivity): - REDIS_OPERATIONS_DURATION: Records timing for lease discovery """ - return self._execute_redis_op('keys', self.redis.keys, 'lease:opc_tags:*') + return await self.redis_repository.keys('lease:opc_tags:*', metadata=self.metadata) diff --git a/ingestor/metrics.py b/ingestor/metrics.py index 13c121e..10181de 100644 --- a/ingestor/metrics.py +++ b/ingestor/metrics.py @@ -16,6 +16,7 @@ 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'] @@ -145,25 +146,6 @@ KAFKA_CONNECTION_STATUS = Gauge( POD_ID_LABEL, ) -# --- Resource Manager (Redis) Metrics --- -REDIS_OPERATIONS_TOTAL = Counter( - 'redis_operations_total', 'Total number of Redis operations performed', REDIS_LABELS -) -REDIS_OPERATIONS_ERRORS = Counter( - 'redis_operations_errors_total', - 'Total number of errors in Redis operations', - REDIS_LABELS, -) -REDIS_OPERATIONS_DURATION = Histogram( - 'redis_operations_duration_seconds', - 'Duration of Redis operations in seconds', - REDIS_LABELS, -) -REDIS_CONNECTION_STATUS = Gauge( - 'redis_connection_status', - 'Connection status with Redis (1=connected, 0=disconnected)', - POD_ID_LABEL, -) # --- Notification Metrics --- NOTIFICATIONS_SENT = Counter( diff --git a/tests/unit/managers/test_data_manager.py b/tests/unit/managers/test_data_manager.py index 037579e..1eba5fe 100644 --- a/tests/unit/managers/test_data_manager.py +++ b/tests/unit/managers/test_data_manager.py @@ -1,7 +1,7 @@ -from unittest.mock import ANY, MagicMock, patch +from unittest.mock import ANY, AsyncMock, MagicMock, patch from kafka.errors import NoBrokersAvailable -from pytest import fixture +from pytest import fixture, mark from sientia_do.notifications.models import NotificationLevel from ingestor.managers.data_manager import DataManager @@ -19,8 +19,8 @@ metadata = { @fixture @patch('ingestor.managers.data_manager.KafkaProducer') -@patch('ingestor.managers.data_manager.MongoClient') -def data_manager(mongo, kafka): +@patch('ingestor.managers.data_manager.MongoDBRepository') +def data_manager(mongodb_repository, kafka): data_manager = DataManager( kafka_servers='localhost:9092', mongo_connection_string='mongodb://localhost:27017', @@ -29,16 +29,18 @@ def data_manager(mongo, kafka): logger=MagicMock(), notification_handler=MagicMock(), metadata=metadata['metadata'], + metrics_controller=MagicMock(), ) data_manager.send_notification = MagicMock() - + data_manager.send_notification_async = AsyncMock() + data_manager.emit_metric = AsyncMock() return data_manager @patch('ingestor.managers.data_manager.KafkaProducer') -@patch('ingestor.managers.data_manager.MongoClient') -def test___init___success(mongo, kafka): +@patch('ingestor.managers.data_manager.MongoDBRepository') +def test___init___success(mongodb_repository, kafka): logger_mock = MagicMock() data_manager = DataManager( @@ -49,6 +51,7 @@ def test___init___success(mongo, kafka): export_to_kafka=True, logger=logger_mock, notification_handler=MagicMock(), + metrics_controller=MagicMock(), ) kafka.assert_called_once_with( @@ -64,8 +67,8 @@ def test___init___success(mongo, kafka): @patch('ingestor.managers.data_manager.KafkaProducer') -@patch('ingestor.managers.data_manager.MongoClient') -def test___init___second_attempt(mongo, kafka): +@patch('ingestor.managers.data_manager.MongoDBRepository') +def test___init___second_attempt(mongodb_repository, kafka): kafka.side_effect = [NoBrokersAvailable, MagicMock()] logger_mock = MagicMock() @@ -77,6 +80,7 @@ def test___init___second_attempt(mongo, kafka): logger=logger_mock, notification_handler=MagicMock(), metadata=metadata['metadata'], + metrics_controller=MagicMock(), ) kafka.assert_any_call( @@ -98,8 +102,8 @@ def test___init___second_attempt(mongo, kafka): @patch('ingestor.managers.data_manager.KafkaProducer') -@patch('ingestor.managers.data_manager.MongoClient') -def test___init___failure_max_attempts(mongo, kafka): +@patch('ingestor.managers.data_manager.MongoDBRepository') +def test___init___failure_max_attempts(mongodb_repository, kafka): kafka.side_effect = NoBrokersAvailable logger_mock = MagicMock() @@ -112,6 +116,7 @@ def test___init___failure_max_attempts(mongo, kafka): logger=logger_mock, notification_handler=MagicMock(), metadata=metadata['metadata'], + metrics_controller=MagicMock(), ) except NoBrokersAvailable as e: assert ( @@ -160,15 +165,6 @@ def test_shutdown_no_producer(data_manager): ) -def test_shutdown_no_mongo_client(data_manager): - data_manager.mongo_client = None - - data_manager.shutdown() - - data_manager.logger.warning.assert_any_call( - 'MongoDB client is already closed or not initialized.' - ) - def test_shutdown_exception(data_manager): data_manager.kafka_producer.flush = MagicMock(side_effect=Exception('Test error')) @@ -179,7 +175,7 @@ def test_shutdown_exception(data_manager): def test_shutdown_exception_mongo(data_manager): - data_manager.mongo_client.close = MagicMock(side_effect=Exception('Test error')) + data_manager.mongo_repository.close = MagicMock(side_effect=Exception('Test error')) data_manager.shutdown() @@ -210,9 +206,10 @@ 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}') + - -def test_publish(data_manager): +@mark.asyncio +async def test_publish(data_manager): topic = 'test_topic' data = {'key': 'value'} @@ -221,7 +218,7 @@ def test_publish(data_manager): data_manager.kafka_producer.send = send_mock # Call the publish method - data_manager.publish(topic, data) + await data_manager.publish(topic, data) # Check if the send method was called with the correct arguments send_mock.assert_called_once_with(topic=topic, value=data) @@ -242,7 +239,8 @@ def test_publish_no_kafka(data_manager): @patch('ingestor.managers.data_manager.traceback') -def test_publish_error(traceback, data_manager): +@mark.asyncio +async def test_publish_error(traceback, data_manager): topic = 'test_topic' data = {'key': 'value', 'name': 'test_tag'} @@ -251,13 +249,13 @@ def test_publish_error(traceback, data_manager): data_manager.kafka_producer.send = send_mock # Call the publish method - data_manager.publish(topic, data) + await data_manager.publish(topic, data) # Check if the send method was called with the correct arguments send_mock.assert_called_once_with(topic=topic, value=data) # Check if the error was logged - data_manager.send_notification.assert_called_once_with( + data_manager.send_notification_async.assert_called_once_with( notification_id=f'KAFKA_PRODUCER_ERROR_{topic}', message=f'Error publishing message to topic {topic}: Test error', block='kafka_producer', @@ -267,15 +265,16 @@ def test_publish_error(traceback, data_manager): ) -def test_publish_error_mongo(data_manager): +@mark.asyncio +async def test_publish_error_mongo(data_manager): data_manager.export_to_kafka = False - data_manager.mongo_db.__getitem__.return_value.insert_one = MagicMock( + data_manager.mongo_repository.insert = AsyncMock( side_effect=Exception('Test error') ) - data_manager.publish('test_topic', {'key': 'value'}) + await data_manager.publish('test_topic', {'key': 'value'}) - data_manager.send_notification.assert_called_once_with( + data_manager.send_notification_async.assert_called_once_with( notification_id='MONGO_PRODUCER_ERROR_test_topic', message='Error inserting message to MongoDB: Test error', block='mongo_producer',