diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 3edd1c1..45d63bb 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -8,6 +8,7 @@ on: jobs: release: + if: github.event.pull_request.merged == true uses: Aignosi/github_workflow_templates/.github/workflows/dataops-module-release.yml@main permissions: write-all with: diff --git a/ingestor/ingestor.py b/ingestor/ingestor.py index 5f27759..4b76888 100644 --- a/ingestor/ingestor.py +++ b/ingestor/ingestor.py @@ -4,12 +4,14 @@ 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. @@ -97,6 +99,14 @@ class Ingestor: logger=self.logger, project_name='opc_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': '-', @@ -184,24 +194,30 @@ 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 # 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, + }, + ) - 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 +238,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 +271,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,15 +280,19 @@ 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) 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]): @@ -334,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): @@ -366,24 +390,31 @@ 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 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) + 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...') - 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 @@ -393,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( @@ -410,7 +446,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..44d248e 100644 --- a/ingestor/managers/data_manager.py +++ b/ingestor/managers/data_manager.py @@ -5,17 +5,18 @@ 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.observability.metrics_controller import MetricsController +from sientia_do.observability.sientia_monitoring import SientiaMonitoring +from sientia_do.repository.mongodb_repository import MongoDBRepository from sientia_do.temporal.constants import now import ingestor.metrics as metrics -class DataManager(BaseActivity): +class DataManager(SientiaMonitoring): """ Manages data persistence and export operations for the OPC Ingestor. @@ -57,6 +58,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 +93,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 +138,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 +173,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 +208,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 +231,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 +259,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, 'inserted_at': now()}, + 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..030e305 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'] @@ -77,6 +79,13 @@ class IngestorManager(BaseActivity): redis_username: str | None = redis_data.get('username', None) redis_password: str | None = redis_data.get('password', None) + SientiaMonitoring.__init__( + self, + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + ) + self.data_manager = DataManager( kafka_servers=kafka_servers, mongo_connection_string=mongo_connection_string, @@ -85,6 +94,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 +107,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,10 +116,6 @@ class IngestorManager(BaseActivity): self.metadata = metadata - BaseActivity.__init__( - self, logger=logger, notification_handler=notification_handler, set_error_counter=True - ) - async def initialize_opc_from_config(self, server_config: dict) -> OpcManager | None: """ Initializes an OPC Manager instance using the provided server configuration. @@ -150,6 +157,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 +272,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): """ @@ -283,9 +298,9 @@ class IngestorManager(BaseActivity): """ to_disconnect: list[str] = [] for server, opc_manager in self.opc_managers.items(): - opc_manager.check_cycles() + await opc_manager.check_cycles() - is_lost = opc_manager.check_opc_listenning() + is_lost = await opc_manager.check_opc_listenning() if is_lost: self.logger.warning(f'OPC server {server} is lost. Server will be disconnected.') @@ -294,9 +309,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 +331,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 +347,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 +365,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 +392,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 +434,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 +498,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 +527,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 +538,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 +559,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 +624,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..7b446c9 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,11 @@ 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, ) metrics.OPC_CONNECTION_STATUS.labels( @@ -192,7 +197,13 @@ 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, + 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 +218,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 +262,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 +322,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 +418,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 +431,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 +493,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 +525,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 +533,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 +553,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 +572,14 @@ 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, + 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..50cbc17 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,33 @@ 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, ) + 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 +121,13 @@ class ResourceManager(BaseActivity): and delegates to the get() method for the actual Redis operation. """ - return self.get(f'slot:opc_tags:{tag_id}') + self.info(f'Getting tag slot for tag_id: {tag_id}', metadata=self.metadata) - def ingestor_heartbeat(self) -> None: + slot = await self.redis_repository.get(f'slot:opc_tags:{tag_id}', metadata=self.metadata) + self.info(f'Tag slot for tag_id: {tag_id} is: {slot}', metadata=self.metadata) + return slot + + async def ingestor_heartbeat(self) -> None: """ Sends a heartbeat signal to Redis to indicate that the ingestor is active. @@ -215,15 +146,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 +176,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 +209,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 +240,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 +264,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 +288,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 +312,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..503d299 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', @@ -82,6 +85,8 @@ SLOTS_RELEASED = Counter( 'Total number of slots released by this instance', POD_ID_LABEL, ) + +# --- OPC Manager Metrics --- OPC_MANAGERS_ACTIVE = Gauge( 'ingestor_opc_managers_active', 'Number of active OPC Managers in this instance', @@ -93,7 +98,6 @@ OPC_SUBSCRIPTION_ERRORS = Counter( ['pod_id', 'server', 'slot'], ) -# --- OPC Manager Metrics --- OPC_CONNECTIONS_TOTAL = Counter( 'opc_connections_initiated_total', 'Total connection attempts to OPC servers', @@ -145,25 +149,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/pyproject.toml b/pyproject.toml index ae13d36..e4841ad 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -113,11 +113,7 @@ python_classes = ["Test*"] python_functions = ["test_*"] addopts = [ "-v", - "--strict-markers", - "--cov=model_manager", - "--cov-report=term-missing", - "--cov-report=html", - "--cov-report=xml", + "--strict-markers" ] markers = [ "asyncio: marks tests as async", diff --git a/requirements.txt b/requirements.txt index 8d25238..bb176bf 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,5 +1,5 @@ asyncua==1.1.5 redis -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.7 +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.5.3 prometheus_client pymongo \ No newline at end of file diff --git a/tests/unit/managers/test_data_manager.py b/tests/unit/managers/test_data_manager.py index 037579e..2cb4ae9 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,16 +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')) data_manager.kafka_producer.close = MagicMock() @@ -179,7 +174,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() @@ -212,7 +207,8 @@ def test_delivery_error(data_manager): 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 +217,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 +238,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'} @@ -250,14 +247,16 @@ def test_publish_error(traceback, data_manager): send_mock = MagicMock(side_effect=Exception('Test error')) data_manager.kafka_producer.send = send_mock + data_manager.mongo_repository = AsyncMock() + # 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 +266,14 @@ 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( - side_effect=Exception('Test error') - ) + 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', diff --git a/tests/unit/managers/test_ingestor_manager.py b/tests/unit/managers/test_ingestor_manager.py index ba4c437..32ac827 100644 --- a/tests/unit/managers/test_ingestor_manager.py +++ b/tests/unit/managers/test_ingestor_manager.py @@ -1,4 +1,4 @@ -from unittest.mock import AsyncMock, MagicMock, patch +from unittest.mock import AsyncMock, MagicMock, call, patch from pytest import fixture, mark from sientia_do.notifications.models import NotificationLevel @@ -31,9 +31,12 @@ def ingestor_manager(data_manager_mock, resource_manager_mock): logger=MagicMock(), notification_handler=MagicMock(), metadata=metadata['metadata'], + metrics_controller=MagicMock(), ) ingestor.send_notification = MagicMock() + ingestor.send_notification_async = AsyncMock() + ingestor.emit_metric = AsyncMock() return ingestor @@ -57,6 +60,7 @@ def test___init__( logger=MagicMock(), notification_handler=MagicMock(), metadata=metadata['metadata'], + metrics_controller=MagicMock(), ) opc_manager_mock.assert_not_called() @@ -68,6 +72,7 @@ def test___init__( metadata=metadata['metadata'], logger=ingestor.logger, notification_handler=ingestor.notification_handler, + metrics_controller=ingestor.metrics_controller, ) resource_manager_mock.assert_called_once_with( host='localhost', @@ -79,6 +84,7 @@ def test___init__( notification_handler=ingestor.notification_handler, username=None, password=None, + metrics_controller=ingestor.metrics_controller, ) assert ingestor.poll_interval == 5 assert ingestor.managed_tags == {} @@ -117,6 +123,7 @@ async def test_initialize_opc_from_config(opc_manager, ingestor_manager): private_key_path=server_config['private_key_path'], server_cert_path=server_config['server_cert_path'], metadata=metadata['metadata'], + metrics_controller=ingestor_manager.metrics_controller, ) assert result == opc_manager.return_value @@ -155,6 +162,61 @@ async def test_initialize_opc_from_config_exception(traceback_mock, opc_manager, ) +@mark.asyncio +async def test_shutdown(ingestor_manager): + ingestor_manager.opc_managers = {'server1': AsyncMock(), 'server2': AsyncMock()} + + ingestor_manager.data_manager.shutdown = MagicMock() + + await ingestor_manager.shutdown() + + ingestor_manager.opc_managers['server1'].shutdown.assert_called_once() + ingestor_manager.opc_managers['server2'].shutdown.assert_called_once() + + ingestor_manager.data_manager.shutdown.assert_called_once() + + +@patch('ingestor.managers.ingestor_manager.asyncio') +def test___del__(asyncio_mock, ingestor_manager): + ingestor_manager.shutdown = MagicMock() + ingestor_manager.__del__() + asyncio_mock.run.assert_called_once_with(ingestor_manager.shutdown.return_value) + + +@mark.asyncio +async def test_remove_server(ingestor_manager): + server1 = AsyncMock() + server2 = AsyncMock() + ingestor_manager.opc_managers = {'server1': server1, 'server2': server2} + ingestor_manager.managed_tags = { + 'slot1': {'server1': {'config': 'config1'}, 'server2': {'config': 'config2'}} + } + await ingestor_manager.remove_server('server1') + server1.shutdown.assert_called_once() + server2.shutdown.assert_not_called() + assert 'server1' not in ingestor_manager.opc_managers + assert 'server2' in ingestor_manager.opc_managers + assert ingestor_manager.managed_tags == {'slot1': {'server2': {'config': 'config2'}}} + + +@mark.asyncio +async def test_remove_server_not_found(ingestor_manager): + server1 = AsyncMock() + server2 = AsyncMock() + ingestor_manager.opc_managers = {'server1': server1, 'server2': server2} + ingestor_manager.managed_tags = { + 'slot1': {'server1': {'config': 'config1'}, 'server2': {'config': 'config2'}} + } + await ingestor_manager.remove_server('server3') + server1.shutdown.assert_not_called() + server2.shutdown.assert_not_called() + assert 'server1' in ingestor_manager.opc_managers + assert 'server2' in ingestor_manager.opc_managers + assert ingestor_manager.managed_tags == { + 'slot1': {'server1': {'config': 'config1'}, 'server2': {'config': 'config2'}} + } + + @mark.asyncio @patch('ingestor.managers.ingestor_manager.OpcManager') @patch('ingestor.managers.ingestor_manager.metrics') @@ -174,6 +236,7 @@ async def test_update_opc_servers(metrics, opc_manager, ingestor_manager): return None ingestor_manager.initialize_opc_from_config = AsyncMock(side_effect=mock_initialize_from_config) + ingestor_manager.remove_server = AsyncMock() ingestor_manager.managed_tags = { 'slot1': { @@ -181,7 +244,10 @@ async def test_update_opc_servers(metrics, opc_manager, ingestor_manager): 'server2': {'config': 'config2'}, 'server5': {'config': 'config5'}, }, - 'slot2': {'server3': {'config': 'config3'}, 'server1': {'config': 'config1'}}, + 'slot2': { + 'server3': {'config': 'config3'}, + 'server1': {'config': 'config1'}, + }, } mock = AsyncMock(config={'config': 'old_config2'}) @@ -191,121 +257,220 @@ async def test_update_opc_servers(metrics, opc_manager, ingestor_manager): await ingestor_manager.update_opc_servers() - assert len(ingestor_manager.opc_managers) == 3 - ingestor_manager.initialize_opc_from_config.assert_any_call({'config': 'config1'}) ingestor_manager.initialize_opc_from_config.assert_any_call({'config': 'config2'}) ingestor_manager.initialize_opc_from_config.assert_any_call({'config': 'config5'}) + assert ingestor_manager.initialize_opc_from_config.call_count == 3 assert ingestor_manager.opc_managers['server1'].config == {'config': 'config1'} assert ingestor_manager.opc_managers['server2'].config == {'config': 'config2'} assert ingestor_manager.opc_managers['server3'].config == {'config': 'config3'} - assert 'server4' not in ingestor_manager.opc_managers - assert 'server5' not in ingestor_manager.opc_managers + + ingestor_manager.remove_server.assert_any_call('server4') + ingestor_manager.remove_server.assert_any_call('server5') + assert ingestor_manager.remove_server.call_count == 2 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) + ingestor_manager.emit_metric.assert_called_once_with( + metric_object=metrics.OPC_MANAGERS_ACTIVE, + method='set', + value=len(ingestor_manager.opc_managers), + tags={'pod_id': ingestor_manager.pod_id}, ) -def test_declare_active(ingestor_manager): - ingestor_manager.resource_manager.ingestor_heartbeat = MagicMock() - ingestor_manager.declare_active() +@mark.asyncio +@patch('ingestor.managers.ingestor_manager.metrics') +async def test_check_opc_servers_integrity_all_healthy(metrics, ingestor_manager): + # Setup mock OPC managers + opc_manager1 = MagicMock() + opc_manager1.check_cycles = AsyncMock(return_value=None) + opc_manager1.check_opc_listenning = AsyncMock(return_value=False) + opc_manager1.config = {'config': 'config1'} + + opc_manager2 = MagicMock() + opc_manager2.check_cycles = AsyncMock(return_value=None) + opc_manager2.check_opc_listenning = AsyncMock(return_value=False) + opc_manager2.config = {'config': 'config2'} + + ingestor_manager.opc_managers = {'server1': opc_manager1, 'server2': opc_manager2} + + # Mock the initialize_opc_from_config method + ingestor_manager.initialize_opc_from_config = MagicMock() + + # Call the method + await ingestor_manager.check_opc_servers_integrity() + + # Verify that check_cycles and check_opc_listenning were called for each server + opc_manager1.check_cycles.assert_called_once() + opc_manager1.check_opc_listenning.assert_called_once() + opc_manager2.check_cycles.assert_called_once() + opc_manager2.check_opc_listenning.assert_called_once() + + # Verify that no reinitialization was needed + ingestor_manager.initialize_opc_from_config.assert_not_called() + + ingestor_manager.emit_metric.assert_called_once_with( + metric_object=metrics.OPC_MANAGERS_ACTIVE, + method='set', + value=len(ingestor_manager.opc_managers), + tags={'pod_id': ingestor_manager.pod_id}, + ) + + +@mark.asyncio +async def test_check_opc_servers_integrity_server_lost(ingestor_manager): + # Setup mock OPC manager that will be lost + opc_manager = AsyncMock() + opc_manager.check_cycles.return_value = None + opc_manager.check_opc_listenning.return_value = True # Server is lost + opc_manager.config = {'config': 'config1'} + + ingestor_manager.opc_managers = {'server1': opc_manager} + ingestor_manager.managed_tags = {'slot1': MagicMock()} + + # Call the method + await ingestor_manager.check_opc_servers_integrity() + + assert 'server1' not in ingestor_manager.opc_managers + ingestor_manager.managed_tags['slot1'].pop.assert_called_once_with('server1', None) + + +def test_check_opc_servers_integrity_server_lost_with_tags(ingestor_manager): + # Setup mock OPC manager that will be lost + opc_manager = MagicMock() + opc_manager.check_cycles.return_value = None + opc_manager.check_opc_listenning.return_value = True # Server is lost + opc_manager.config = {'config': 'config1'} + + ingestor_manager.opc_managers = {'server1': opc_manager} + + # Setup managed tags + ingestor_manager.managed_tags = { + 'slot1': {'server1': {'config': 'config1', 'tags': {'tag1': 'value1'}}} + } + + # Mock the initialize_opc_from_config method to return a new manager + new_manager = MagicMock() + ingestor_manager.initialize_opc_from_config = MagicMock(return_value=new_manager) + + # Call the method + ingestor_manager.check_opc_servers_integrity() + + +@mark.asyncio +async def test_declare_active(ingestor_manager): + ingestor_manager.resource_manager.ingestor_heartbeat = AsyncMock() + await ingestor_manager.declare_active() ingestor_manager.resource_manager.ingestor_heartbeat.assert_called_once() -def test_get_active_ingestors(ingestor_manager): - ingestor_manager.resource_manager.get_all_ingestors = MagicMock() - ingestor_manager.get_active_ingestors() +@mark.asyncio +async def test_get_active_ingestors(ingestor_manager): + ingestor_manager.resource_manager.get_all_ingestors = AsyncMock() + await ingestor_manager.get_active_ingestors() ingestor_manager.resource_manager.get_all_ingestors.assert_called_once() -def test_get_active_ingestors_empty(ingestor_manager): - ingestor_manager.resource_manager.get_all_ingestors = MagicMock(return_value=None) - result = ingestor_manager.get_active_ingestors() +@mark.asyncio +async def test_get_active_ingestors_empty(ingestor_manager): + ingestor_manager.resource_manager.get_all_ingestors = AsyncMock(return_value=None) + result = await ingestor_manager.get_active_ingestors() assert result == [] ingestor_manager.resource_manager.get_all_ingestors.assert_called_once() @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() +@mark.asyncio +async def test_get_number_of_leases_success(metrics, ingestor_manager): + ingestor_manager.resource_manager.get_all_leases = AsyncMock(return_value=['lease1', 'lease2']) + result = await ingestor_manager.get_number_of_leases() assert result == 2 + + ingestor_manager.emit_metric.assert_called_once_with( + metric_object=metrics.LEASES_TOTAL, + method='set', + value=2, + tags={'pod_id': ingestor_manager.pod_id}, + ) ingestor_manager.resource_manager.get_all_leases.assert_called_once() - metrics.LEASES_TOTAL.set.assert_called_once_with(2) @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) - - -@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() +@mark.asyncio +async def test_get_number_of_slots_success(metrics, ingestor_manager): + ingestor_manager.resource_manager.get_all_slots = AsyncMock(return_value=['slot1', 'slot2']) + result = await 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) + ingestor_manager.emit_metric.assert_called_once_with( + metric_object=metrics.SLOTS_TOTAL, + method='set', + value=2, + tags={'pod_id': ingestor_manager.pod_id}, + ) @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() +@mark.asyncio +async def test_get_number_of_slots_empty(metrics, ingestor_manager): + ingestor_manager.resource_manager.get_all_slots = AsyncMock(return_value=None) + result = await 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) + ingestor_manager.emit_metric.assert_called_once_with( + metric_object=metrics.SLOTS_TOTAL, + method='set', + value=0, + tags={'pod_id': ingestor_manager.pod_id}, + ) -def test_get_slot_leases_1_success(ingestor_manager): - ingestor_manager.resource_manager.lease_tag = MagicMock(return_value=True) - ingestor_manager.resource_manager.get_tag_slot = MagicMock(return_value={'tags': ['tag1']}) +@mark.asyncio +async def test_get_slot_leases_1_success(ingestor_manager): + ingestor_manager.resource_manager.lease_tag = AsyncMock(return_value=True) + ingestor_manager.resource_manager.get_tag_slot = AsyncMock(return_value={'tags': ['tag1']}) ingestor_manager.number_of_slots = 1 - result = ingestor_manager.get_slot_leases() + result = await ingestor_manager.get_slot_leases() assert result == {'1': {'tags': ['tag1']}} -def test_get_slot_leases_2_success(ingestor_manager): - ingestor_manager.resource_manager.lease_tag = MagicMock(side_effect=[True, True]) - ingestor_manager.resource_manager.get_tag_slot = MagicMock( +@mark.asyncio +async def test_get_slot_leases_2_success(ingestor_manager): + ingestor_manager.resource_manager.lease_tag = AsyncMock(side_effect=[True, True]) + ingestor_manager.resource_manager.get_tag_slot = AsyncMock( side_effect=[{'tags': ['tag1']}, {'tags': ['tag2']}] ) ingestor_manager.number_of_slots = 2 - result = ingestor_manager.get_slot_leases(max_slots=2) + result = await ingestor_manager.get_slot_leases(max_slots=2) assert result == {'1': {'tags': ['tag1']}, '2': {'tags': ['tag2']}} -def test_get_slot_leases_2_1_none(ingestor_manager): - ingestor_manager.resource_manager.lease_tag = MagicMock(side_effect=[True, True]) - ingestor_manager.resource_manager.get_tag_slot = MagicMock( +@mark.asyncio +async def test_get_slot_leases_2_1_none(ingestor_manager): + ingestor_manager.resource_manager.lease_tag = AsyncMock(side_effect=[True, True]) + ingestor_manager.resource_manager.get_tag_slot = AsyncMock( side_effect=[None, {'tags': ['tag1']}] ) ingestor_manager.number_of_slots = 1 - result = ingestor_manager.get_slot_leases(max_slots=1) + result = await ingestor_manager.get_slot_leases(max_slots=1) assert result == {} -def test_get_slot_leases_1_failure(ingestor_manager): - ingestor_manager.resource_manager.lease_tag = MagicMock(return_value=False) - ingestor_manager.resource_manager.get_tag_slot = MagicMock(return_value={'tags': ['tag1']}) +@mark.asyncio +async def test_get_slot_leases_1_failure(ingestor_manager): + ingestor_manager.resource_manager.lease_tag = AsyncMock(return_value=False) + ingestor_manager.resource_manager.get_tag_slot = AsyncMock(return_value={'tags': ['tag1']}) - result = ingestor_manager.get_slot_leases() + result = await ingestor_manager.get_slot_leases() ingestor_manager.resource_manager.get_tag_slot.assert_not_called() assert result == {} @@ -329,22 +494,20 @@ async def test_unsubscribe_slot(ingestor_manager): ingestor_manager.opc_managers['server3'].unsubscribe.assert_not_called() -def test_update_slot_config(ingestor_manager): +@mark.asyncio +async def test_update_slot_config(ingestor_manager): ingestor_manager.managed_tags = { 'slot1': {'config': 'old_config'}, 'slot2': {'config': 'new_config'}, 'slot3': {'config': 'old_config'}, } - ingestor_manager.resource_manager.get_tag_slot = MagicMock( + ingestor_manager.resource_manager.get_tag_slot = AsyncMock( side_effect=[{'config': 'updated_config'}, {'config': 'new_config'}, None] ) + ingestor_manager.resource_manager.renew_tag_lease = AsyncMock() - ingestor_manager.update_opc_servers = MagicMock() - ingestor_manager.subscribe_to_tags = MagicMock() - ingestor_manager.unsubscribe_slot = MagicMock() - - ingestor_manager.update_slot_config() + await ingestor_manager.update_slot_config() assert ingestor_manager.managed_tags['slot1'] == {'config': 'updated_config'} assert ingestor_manager.managed_tags['slot2'] == {'config': 'new_config'} @@ -357,15 +520,30 @@ def test_update_slot_config(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']) +@mark.asyncio +async def test_drop_slot_leases(metrics, ingestor_manager): + ingestor_manager.resource_manager.drop_tag_lease = AsyncMock() + await 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() + ingestor_manager.emit_metric.assert_has_calls( + [ + call( + metric_object=metrics.SLOTS_RELEASED, + method='inc', + value=1, + tags={'pod_id': ingestor_manager.pod_id}, + ), + call( + metric_object=metrics.SLOTS_RELEASED, + method='inc', + value=1, + tags={'pod_id': ingestor_manager.pod_id}, + ), + ] + ) @mark.asyncio @@ -432,7 +610,7 @@ async def test_manage_server_subscribe_failure(traceback_mock, ingestor_manager) traceback_mock.format_exc.assert_called_once() - ingestor_manager.send_notification.assert_called_once_with( + ingestor_manager.send_notification_async.assert_called_once_with( metadata=metadata['metadata'], notification_id='OPC_SUBSCRIPTION_ERROR_slot1:server1', message="Failed to subscribe to tags from slot1:server1\n{'tags': 'config1'}: Subscription error", @@ -474,80 +652,3 @@ async def test_subscribe_to_tags(ingestor_manager): assert ingestor_manager.manage_server.call_count == 3 ingestor_manager.managed_tags['slot1'].pop.assert_called_once_with('server3', None) - - -@mark.asyncio -@patch('ingestor.managers.ingestor_manager.metrics') -async 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 - opc_manager1.check_opc_listenning.return_value = False - opc_manager1.config = {'config': 'config1'} - - opc_manager2 = MagicMock() - opc_manager2.check_cycles.return_value = None - opc_manager2.check_opc_listenning.return_value = False - opc_manager2.config = {'config': 'config2'} - - ingestor_manager.opc_managers = {'server1': opc_manager1, 'server2': opc_manager2} - - # Mock the initialize_opc_from_config method - ingestor_manager.initialize_opc_from_config = MagicMock() - - # Call the method - await ingestor_manager.check_opc_servers_integrity() - - # Verify that check_cycles and check_opc_listenning were called for each server - opc_manager1.check_cycles.assert_called_once() - opc_manager1.check_opc_listenning.assert_called_once() - opc_manager2.check_cycles.assert_called_once() - opc_manager2.check_opc_listenning.assert_called_once() - - # 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) - ) - - -@mark.asyncio -async def test_check_opc_servers_integrity_server_lost(ingestor_manager): - # Setup mock OPC manager that will be lost - opc_manager = AsyncMock() - opc_manager.check_cycles.return_value = None - opc_manager.check_opc_listenning.return_value = True # Server is lost - opc_manager.config = {'config': 'config1'} - - ingestor_manager.opc_managers = {'server1': opc_manager} - ingestor_manager.managed_tags = {'slot1': MagicMock()} - - # Call the method - await ingestor_manager.check_opc_servers_integrity() - - assert 'server1' not in ingestor_manager.opc_managers - ingestor_manager.managed_tags['slot1'].pop.assert_called_once_with('server1', None) - - -def test_check_opc_servers_integrity_server_lost_with_tags(ingestor_manager): - # Setup mock OPC manager that will be lost - opc_manager = MagicMock() - opc_manager.check_cycles.return_value = None - opc_manager.check_opc_listenning.return_value = True # Server is lost - opc_manager.config = {'config': 'config1'} - - ingestor_manager.opc_managers = {'server1': opc_manager} - - # Setup managed tags - ingestor_manager.managed_tags = { - 'slot1': {'server1': {'config': 'config1', 'tags': {'tag1': 'value1'}}} - } - - # Mock the initialize_opc_from_config method to return a new manager - new_manager = MagicMock() - ingestor_manager.initialize_opc_from_config = MagicMock(return_value=new_manager) - - # Call the method - ingestor_manager.check_opc_servers_integrity() diff --git a/tests/unit/managers/test_opc_manager.py b/tests/unit/managers/test_opc_manager.py index 282328b..407c166 100644 --- a/tests/unit/managers/test_opc_manager.py +++ b/tests/unit/managers/test_opc_manager.py @@ -46,17 +46,24 @@ metadata = { @fixture @patch('ingestor.managers.opc_manager.metrics') def raw_opc_manager(mock_metrics): - return OpcManager( + opc_manager = OpcManager( name='TestConnector', url='opc.tcp://localhost:4840', - data_manager=MagicMock(), + data_manager=AsyncMock(), subscription_period_ms=1000, logger=MagicMock(), server_uri='opc.tcp://localhost:4840', notification_handler=MagicMock(), metadata=metadata['metadata'], + metrics_controller=AsyncMock(), ) + opc_manager.emit_metric = AsyncMock() + opc_manager.send_notification_async = AsyncMock() + opc_manager.send_notification = MagicMock() + + return opc_manager + @fixture def opc_manager(raw_opc_manager): @@ -64,7 +71,6 @@ def opc_manager(raw_opc_manager): raw_opc_manager.cert_path = 'cert.pem' raw_opc_manager.private_key_path = 'private_key.pem' raw_opc_manager.server_cert_path = 'server_cert.pem' - raw_opc_manager.send_notification = MagicMock() return raw_opc_manager @@ -145,17 +151,28 @@ async def test_connect_no_security(client, mock_metrics, raw_opc_manager): client.assert_called_once_with(raw_opc_manager.url, timeout=10, watchdog_intervall=3600000) raw_opc_manager.client.connect.assert_called_once() raw_opc_manager.set_security.assert_not_called() - mock_metrics.OPC_CONNECTIONS_TOTAL.labels.assert_called_once_with( - pod_id=raw_opc_manager.pod_id, server_name=raw_opc_manager.name + raw_opc_manager.emit_metric.assert_has_calls( + [ + call( + metric_object=mock_metrics.OPC_CONNECTIONS_TOTAL, + tags={ + 'pod_id': raw_opc_manager.pod_id, + 'server_name': raw_opc_manager.name, + }, + ), + call( + metric_object=mock_metrics.OPC_CONNECTION_STATUS, + method='set', + value=1, + tags={ + 'pod_id': raw_opc_manager.pod_id, + 'server_name': raw_opc_manager.name, + 'server_url': raw_opc_manager.url, + }, + ), + ], + any_order=True, ) - mock_metrics.OPC_CONNECTIONS_TOTAL.labels.return_value.inc.assert_called_once() - mock_metrics.OPC_CONNECTION_STATUS.labels.assert_called_once_with( - pod_id=raw_opc_manager.pod_id, - server_name=raw_opc_manager.name, - server_url=raw_opc_manager.url, - ) - mock_metrics.OPC_CONNECTION_STATUS.labels.return_value.set.assert_called_once_with(1) - mock_metrics.OPC_CONNECTIONS_FAILED.labels.assert_not_called() @mark.asyncio @@ -230,10 +247,16 @@ async def test_create_subscription_success_no_period(opc_manager): async def test_create_subscription_with_metrics(metrics, opc_manager): await opc_manager.create_subscription('sub1') - metrics.OPC_SUBSCRIPTIONS_CREATED.labels.assert_called_once_with( - pod_id=opc_manager.pod_id, server_name=opc_manager.name, slot_name='sub1' + opc_manager.emit_metric.assert_called_once_with( + metric_object=metrics.OPC_SUBSCRIPTIONS_CREATED, + method='inc', + value=1, + tags={ + 'pod_id': opc_manager.pod_id, + 'server_name': opc_manager.name, + 'slot_name': 'sub1', + }, ) - metrics.OPC_SUBSCRIPTIONS_CREATED.labels.return_value.inc.assert_called_once() @patch('ingestor.managers.opc_manager.metrics') @@ -390,17 +413,30 @@ async def test_disconnect_metrics_on_successful_path(mock_metrics_module, raw_op await raw_opc_manager.disconnect() - mock_metrics_module.OPC_CONNECTION_STATUS.labels.assert_called_once_with( - pod_id=raw_opc_manager.pod_id, - server_name=raw_opc_manager.name, - server_url=raw_opc_manager.url, + raw_opc_manager.emit_metric.assert_has_calls( + [ + call( + metric_object=mock_metrics_module.OPC_CONNECTION_STATUS, + method='set', + value=0, + tags={ + 'pod_id': raw_opc_manager.pod_id, + 'server_name': raw_opc_manager.name, + 'server_url': raw_opc_manager.url, + }, + ), + call( + metric_object=mock_metrics_module.OPC_TAGS_SUBSCRIBED, + method='set', + value=0, + tags={ + 'pod_id': raw_opc_manager.pod_id, + 'server_name': raw_opc_manager.name, + }, + ), + ], + any_order=True, ) - mock_metrics_module.OPC_CONNECTION_STATUS.labels.return_value.set.assert_called_once_with(0) - - mock_metrics_module.OPC_TAGS_SUBSCRIBED.labels.assert_called_once_with( - pod_id=raw_opc_manager.pod_id, server_name=raw_opc_manager.name - ) - mock_metrics_module.OPC_TAGS_SUBSCRIBED.labels.return_value.set.assert_called_once_with(0) @patch('ingestor.managers.opc_manager.metrics') @@ -422,8 +458,6 @@ async def test_datachange_notification(metrics, opc_manager_subscribed): } } - metrics.OPC_CYCLES_WITHOUT_DATA.reset_mock() - await opc_manager_subscribed.datachange_notification('ns=3;i=1001', None, data) opc_manager_subscribed.data_manager.publish.assert_any_call( @@ -446,13 +480,19 @@ async def test_datachange_notification(metrics, opc_manager_subscribed): ) assert opc_manager_subscribed.nodes['ns=3;i=1001']['cycle_rule']['cycle_count'] == 0 - metrics.OPC_CYCLES_WITHOUT_DATA.labels.assert_called_once_with( - pod_id=opc_manager_subscribed.pod_id, server_name=opc_manager_subscribed.name + opc_manager_subscribed.emit_metric.assert_called_once_with( + metric_object=metrics.OPC_CYCLES_WITHOUT_DATA, + method='set', + value=0, + tags={ + 'pod_id': opc_manager_subscribed.pod_id, + 'server_name': opc_manager_subscribed.name, + }, ) - metrics.OPC_CYCLES_WITHOUT_DATA.labels.return_value.set.assert_called_once_with(0) -def test_check_cycles_no_notification(opc_manager): +@mark.asyncio +async def test_check_cycles_no_notification(opc_manager): # Setup: node with cycle_count just below threshold opc_manager.nodes = { 'ns=3;i=1001': { @@ -461,14 +501,15 @@ def test_check_cycles_no_notification(opc_manager): } } - opc_manager.check_cycles() + await opc_manager.check_cycles() # After one increment, cycle_count = 4.0, still below threshold assert opc_manager.nodes['ns=3;i=1001']['cycle_rule']['cycle_count'] == pytest.approx(4.0) - opc_manager.send_notification.assert_not_called() + opc_manager.send_notification_async.assert_not_called() -def test_check_cycles_triggers_notification(opc_manager): +@mark.asyncio +async def test_check_cycles_triggers_notification(opc_manager): # Setup: node with cycle_count just below threshold, increment will cross threshold opc_manager.nodes = { 'ns=3;i=1001': { @@ -478,11 +519,11 @@ def test_check_cycles_triggers_notification(opc_manager): } opc_manager.notification_handler.build_and_send_notification = MagicMock() - opc_manager.check_cycles() + await opc_manager.check_cycles() # After increment, cycle_count = 5.5, should trigger notification assert opc_manager.nodes['ns=3;i=1001']['cycle_rule']['cycle_count'] == pytest.approx(5.5) - opc_manager.send_notification.assert_called_once_with( + opc_manager.send_notification_async.assert_called_once_with( notification_id='TAG_ns=3;i=1001:Counter_LISTENNING_STOPPED', message='5.5 cycles without receive from ns=3;i=1001:Counter', block='opc_manager', @@ -492,34 +533,35 @@ def test_check_cycles_triggers_notification(opc_manager): @patch('ingestor.managers.opc_manager.metrics') -def test_check_opc_listenning_no_notification(metrics, opc_manager): +@mark.asyncio +async def test_check_opc_listenning_no_notification(metrics, opc_manager): opc_manager.non_receive_count = 3 - opc_manager.notification_handler.build_and_send_notification = MagicMock() - metrics.OPC_CYCLES_WITHOUT_DATA.reset_mock() - metrics.OPC_RECONNECTIONS_TOTAL.reset_mock() - result = opc_manager.check_opc_listenning() + result = await opc_manager.check_opc_listenning() assert opc_manager.non_receive_count == 4 - opc_manager.notification_handler.build_and_send_notification.assert_not_called() + opc_manager.send_notification_async.assert_not_called() assert result is False - metrics.OPC_CYCLES_WITHOUT_DATA.labels.assert_called_once_with( - pod_id=opc_manager.pod_id, server_name=opc_manager.name + opc_manager.emit_metric.assert_called_once_with( + metric_object=metrics.OPC_CYCLES_WITHOUT_DATA, + method='set', + value=opc_manager.non_receive_count, + tags={ + 'pod_id': opc_manager.pod_id, + 'server_name': opc_manager.name, + }, ) - metrics.OPC_CYCLES_WITHOUT_DATA.labels.return_value.set.assert_called_once_with( - opc_manager.non_receive_count - ) - metrics.OPC_RECONNECTIONS_TOTAL.labels.assert_not_called() -def test_check_opc_listenning_warning_notification(opc_manager): +@mark.asyncio +async def test_check_opc_listenning_warning_notification(opc_manager): opc_manager.non_receive_count = 4 - result = opc_manager.check_opc_listenning() + result = await opc_manager.check_opc_listenning() assert opc_manager.non_receive_count == 5 - opc_manager.send_notification.assert_called_once_with( + opc_manager.send_notification_async.assert_called_once_with( notification_id=f'OPC_LISTENNING_STOPPED__{opc_manager.name}', message=f'5 cycles without receive from OPC {opc_manager.name}. Tags: {json.dumps(opc_manager.nodes)}', block='opc_manager', @@ -530,18 +572,16 @@ def test_check_opc_listenning_warning_notification(opc_manager): @patch('ingestor.managers.opc_manager.metrics') -def test_check_opc_listenning_error_notification_and_retry(metrics, opc_manager): +@mark.asyncio +async def test_check_opc_listenning_error_notification_and_retry(metrics, opc_manager): opc_manager.non_receive_count = 14 - opc_manager.notification_handler.build_and_send_notification = MagicMock() - metrics.OPC_CYCLES_WITHOUT_DATA.reset_mock() - metrics.OPC_RECONNECTIONS_TOTAL.reset_mock() - result = opc_manager.check_opc_listenning() + result = await opc_manager.check_opc_listenning() assert opc_manager.non_receive_count == 15 # Should be called twice: once for 5, once for 15 - assert opc_manager.send_notification.call_count == 2 - calls = opc_manager.send_notification.call_args_list + assert opc_manager.send_notification_async.call_count == 2 + calls = opc_manager.send_notification_async.call_args_list # First call: 5 cycles warning assert calls[0].kwargs == { 'notification_id': f'OPC_LISTENNING_STOPPED__{opc_manager.name}', @@ -560,42 +600,28 @@ def test_check_opc_listenning_error_notification_and_retry(metrics, opc_manager) } assert result is True - metrics.OPC_CYCLES_WITHOUT_DATA.labels.assert_called_once_with( - pod_id=opc_manager.pod_id, server_name=opc_manager.name - ) - metrics.OPC_CYCLES_WITHOUT_DATA.labels.return_value.set.assert_called_once_with( - opc_manager.non_receive_count + opc_manager.emit_metric.assert_has_calls( + [ + call( + metric_object=metrics.OPC_CYCLES_WITHOUT_DATA, + method='set', + value=opc_manager.non_receive_count, + tags={ + 'pod_id': opc_manager.pod_id, + 'server_name': opc_manager.name, + }, + ), + ] ) - metrics.OPC_RECONNECTIONS_TOTAL.labels.assert_called_once_with( - pod_id=opc_manager.pod_id, server_name=opc_manager.name + opc_manager.emit_metric.assert_has_calls( + [ + call( + metric_object=metrics.OPC_RECONNECTIONS_TOTAL, + tags={ + 'pod_id': opc_manager.pod_id, + 'server_name': opc_manager.name, + }, + ), + ] ) - metrics.OPC_RECONNECTIONS_TOTAL.labels.return_value.inc.assert_called_once() - - -@patch('ingestor.managers.opc_manager.metrics') -def test_init_metrics_calls_correct_metric_methods(metrics): - opc_manager = OpcManager( - name='TestInitConnector', - url='opc.tcp://init.test:4840', - data_manager=MagicMock(), - subscription_period_ms=1000, - logger=MagicMock(), - server_uri='opc.tcp://init.test:4840/uri', - notification_handler=MagicMock(), - metadata=metadata['metadata'], - ) - - metrics.OPC_CONNECTION_STATUS.labels.assert_called_with( - pod_id=opc_manager.pod_id, - server_name=opc_manager.name, - server_url=opc_manager.url, - ) - - metrics.OPC_CONNECTION_STATUS.labels.return_value.set.assert_called_once_with(0) - - metrics.OPC_TAGS_SUBSCRIBED.labels.assert_called_once_with( - pod_id=opc_manager.pod_id, server_name=opc_manager.name - ) - - metrics.OPC_TAGS_SUBSCRIBED.labels.return_value.set.assert_called_once_with(0) diff --git a/tests/unit/managers/test_resource_manager.py b/tests/unit/managers/test_resource_manager.py index 9479672..d7e236c 100644 --- a/tests/unit/managers/test_resource_manager.py +++ b/tests/unit/managers/test_resource_manager.py @@ -1,7 +1,6 @@ -from unittest.mock import MagicMock, patch +from unittest.mock import AsyncMock, MagicMock, patch -from pytest import fixture, raises -from sientia_do.notifications.models import NotificationLevel +from pytest import fixture, mark, raises from ingestor.managers.resource_manager import ResourceManager @@ -15,9 +14,58 @@ metadata = { } +@patch('ingestor.managers.resource_manager.RedisRepository') +def test___init__(redis_repository): + logger = MagicMock() + notification_handler = MagicMock() + metrics_controller = MagicMock() + resource_manager = ResourceManager( + host='localhost', + port=6379, + lease_ttl=10, + heartbeat_ttl=10, + metadata=metadata['metadata'], + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + ) + assert resource_manager.redis_repository == redis_repository.return_value + redis_repository.assert_called_once_with( + host='localhost', + port=6379, + username=None, + password=None, + logger=logger, + notification_handler=notification_handler, + metrics_controller=metrics_controller, + ) + redis_repository.return_value.redis_client.ping.assert_called_once() + + +@patch('ingestor.managers.resource_manager.RedisRepository') +def test___init__connection_failure(redis_repository): + redis_repository.return_value.redis_client.ping.side_effect = Exception('Connection failed') + with raises(Exception, match='Connection failed'): + ResourceManager( + host='localhost', + port=6379, + lease_ttl=10, + heartbeat_ttl=10, + metadata=metadata['metadata'], + logger=MagicMock(), + notification_handler=MagicMock(), + metrics_controller=MagicMock(), + ) + redis_repository.return_value.redis_client.ping.assert_called_once() + redis_repository.logger.error.assert_called_once_with( + 'Failed to connect to Redis: Connection failed' + ) + redis_repository.return_value.redis_client.ping.assert_called_once() + + @fixture -@patch('ingestor.managers.resource_manager.Redis') -def resource_manager(redis): +@patch('ingestor.managers.resource_manager.RedisRepository') +def resource_manager(redis_repository): resource_manager = ResourceManager( host='localhost', port=6379, @@ -26,199 +74,105 @@ def resource_manager(redis): metadata=metadata['metadata'], logger=MagicMock(), notification_handler=MagicMock(), + metrics_controller=MagicMock(), ) resource_manager.send_notification = MagicMock() + resource_manager.send_notification_async = AsyncMock() + resource_manager.emit_metric = AsyncMock() + + resource_manager.redis_repository = AsyncMock() return resource_manager -def test_get_success(resource_manager): - resource_manager.redis.get.return_value = '{"key": "value"}' - result = resource_manager.get('key') - assert result == {'key': 'value'} - resource_manager.redis.get.assert_called_once_with('key') +@mark.asyncio +async def test_get_tag_slot(resource_manager): + result = await resource_manager.get_tag_slot('id') + assert result == resource_manager.redis_repository.get.return_value - -def test_get_failure(resource_manager): - resource_manager.redis.get.return_value = None - result = resource_manager.get('key') - assert result is None - resource_manager.redis.get.assert_called_once_with('key') - - -def test_get_tag_slot(resource_manager): - resource_manager.get = MagicMock(return_value={'tag': 'slot'}) - result = resource_manager.get_tag_slot('id') - assert result == {'tag': 'slot'} - resource_manager.get.assert_called_once_with('slot:opc_tags:id') - - -def test_ingestor_heartbeat(resource_manager): - resource_manager.redis.set.return_value = True - resource_manager.ingestor_heartbeat() - resource_manager.redis.set.assert_called_once_with('heartbeat:ingestor:localhost', 1, ex=10) - - -def test_lease_tag(resource_manager): - resource_manager.redis.set.return_value = True - output = resource_manager.lease_tag('tag_id') - assert output is True - resource_manager.redis.set.assert_called_once_with( - 'lease:opc_tags:tag_id', 'localhost', nx=True, ex=10 + resource_manager.redis_repository.get.assert_called_once_with( + 'slot:opc_tags:id', metadata=metadata['metadata'] ) -def test_renew_tag_lease_success(resource_manager): - resource_manager.redis.get.return_value = 'localhost' - resource_manager.redis.expire.return_value = True - result = resource_manager.renew_tag_lease('tag_id') - assert result is True - resource_manager.redis.get.assert_called_once_with('lease:opc_tags:tag_id') - resource_manager.redis.expire.assert_called_once_with('lease:opc_tags:tag_id', 10) +@mark.asyncio +async def test_ingestor_heartbeat(resource_manager): + await resource_manager.ingestor_heartbeat() + resource_manager.redis_repository.set.assert_called_once_with( + 'heartbeat:ingestor:localhost', 1, ttl=10, metadata=metadata['metadata'] + ) -def test_renew_tag_lease_failure(resource_manager): - resource_manager.redis.get.return_value = 'other_pod_id' - resource_manager.redis.expire.return_value = False - result = resource_manager.renew_tag_lease('tag_id') +@mark.asyncio +async def test_lease_tag(resource_manager): + output = await resource_manager.lease_tag('tag_id') + assert output is resource_manager.redis_repository.set.return_value + resource_manager.redis_repository.set.assert_called_once_with( + 'lease:opc_tags:tag_id', 'localhost', ttl=10, nx=True, metadata=metadata['metadata'] + ) + + +@mark.asyncio +async def test_renew_tag_lease_success(resource_manager): + resource_manager.redis_repository.get.return_value = 'localhost' + resource_manager.redis_repository.expire.return_value = True + result = await resource_manager.renew_tag_lease('tag_id') + assert result is resource_manager.redis_repository.expire.return_value + resource_manager.redis_repository.get.assert_called_once_with( + 'lease:opc_tags:tag_id', metadata=metadata['metadata'] + ) + resource_manager.redis_repository.expire.assert_called_once_with( + 'lease:opc_tags:tag_id', 10, metadata=metadata['metadata'] + ) + + +@mark.asyncio +async def test_renew_tag_lease_failure(resource_manager): + resource_manager.redis_repository.get.return_value = 'other_pod_id' + + result = await resource_manager.renew_tag_lease('tag_id') assert result is False - resource_manager.redis.get.assert_called_once_with('lease:opc_tags:tag_id') - resource_manager.redis.expire.assert_not_called() + + resource_manager.redis_repository.get.assert_called_once_with( + 'lease:opc_tags:tag_id', metadata=metadata['metadata'] + ) + resource_manager.redis_repository.expire.assert_not_called() -def test_drop_tag_lease(resource_manager): - resource_manager.redis.delete.return_value = True - resource_manager.drop_tag_lease('tag_id') - resource_manager.redis.delete.assert_called_once_with('lease:opc_tags:tag_id') +@mark.asyncio +async def test_drop_tag_lease(resource_manager): + await resource_manager.drop_tag_lease('tag_id') + resource_manager.redis_repository.delete.assert_called_once_with( + 'lease:opc_tags:tag_id', metadata=metadata['metadata'] + ) -def test_get_all_ingestors(resource_manager): - resource_manager.redis.keys.return_value = ['ingestor1', 'ingestor2'] - result = resource_manager.get_all_ingestors() +@mark.asyncio +async def test_get_all_ingestors(resource_manager): + resource_manager.redis_repository.keys.return_value = ['ingestor1', 'ingestor2'] + result = await resource_manager.get_all_ingestors() assert result == ['ingestor1', 'ingestor2'] - resource_manager.redis.keys.assert_called_once_with('heartbeat:ingestor:*') + resource_manager.redis_repository.keys.assert_called_once_with( + 'heartbeat:ingestor:*', metadata=metadata['metadata'] + ) -def test_get_all_slots(resource_manager): - resource_manager.redis.keys.return_value = ['slot1', 'slot2'] - result = resource_manager.get_all_slots() +@mark.asyncio +async def test_get_all_slots(resource_manager): + resource_manager.redis_repository.keys.return_value = ['slot1', 'slot2'] + result = await resource_manager.get_all_slots() assert result == ['slot1', 'slot2'] - resource_manager.redis.keys.assert_called_once_with('slot:opc_tags:*') + resource_manager.redis_repository.keys.assert_called_once_with( + 'slot:opc_tags:*', metadata=metadata['metadata'] + ) -def test_get_all_leases(resource_manager): - resource_manager.redis.keys.return_value = ['lease1', 'lease2'] - result = resource_manager.get_all_leases() +@mark.asyncio +async def test_get_all_leases(resource_manager): + resource_manager.redis_repository.keys.return_value = ['lease1', 'lease2'] + result = await resource_manager.get_all_leases() assert result == ['lease1', 'lease2'] - resource_manager.redis.keys.assert_called_once_with('lease:opc_tags:*') - - -def test_init_connection_failure(monkeypatch): - # Mock Redis to raise an exception during initialization - mock_redis = MagicMock() - mock_redis.side_effect = Exception('Connection failed') - - monkeypatch.setattr('ingestor.managers.resource_manager.Redis', mock_redis) - - # Test that the exception is raised and metrics are set properly - with patch('ingestor.metrics.REDIS_CONNECTION_STATUS') as mock_metrics: - mock_status = MagicMock() - mock_metrics.labels.return_value = mock_status - - with raises(Exception, match='Connection failed'): - ResourceManager( - host='localhost', - port=6379, - lease_ttl=10, - heartbeat_ttl=10, - metadata=metadata['metadata'], - logger=MagicMock(), - notification_handler=MagicMock(), - ) - - mock_metrics.labels.assert_called_once_with(pod_id='localhost') - mock_status.set.assert_called_once_with(0) - - -def test_init_ping_failure(monkeypatch): - # Mock Redis ping to raise an exception - mock_redis_instance = MagicMock() - mock_redis_instance.ping.side_effect = Exception('Ping failed') - - mock_redis_class = MagicMock(return_value=mock_redis_instance) - monkeypatch.setattr('ingestor.managers.resource_manager.Redis', mock_redis_class) - - # Test that the exception is raised and metrics are set properly - with patch('ingestor.metrics.REDIS_CONNECTION_STATUS') as mock_metrics: - mock_status = MagicMock() - mock_metrics.labels.return_value = mock_status - - with raises(Exception, match='Ping failed'): - ResourceManager( - host='localhost', - port=6379, - lease_ttl=10, - heartbeat_ttl=10, - metadata=metadata['metadata'], - logger=MagicMock(), - notification_handler=MagicMock(), - ) - - mock_metrics.labels.assert_called_once_with(pod_id='localhost') - mock_status.set.assert_called_once_with(0) - - -def test_execute_redis_op_success(resource_manager): - # Mock the Redis operation and time function - mock_func = MagicMock(return_value='test_result') - - with patch('ingestor.managers.resource_manager.time', side_effect=[100, 100.5]): - with patch('ingestor.metrics.REDIS_OPERATIONS_TOTAL') as mock_total: - with patch('ingestor.metrics.REDIS_OPERATIONS_DURATION') as mock_duration: - mock_total_labels = MagicMock() - mock_duration_labels = MagicMock() - mock_total.labels.return_value = mock_total_labels - mock_duration.labels.return_value = mock_duration_labels - - # Execute the operation - result = resource_manager._execute_redis_op( - 'test_op', mock_func, 'arg1', kwarg1='value1' - ) - - # Verify the result and metrics - assert result == 'test_result' - mock_func.assert_called_once_with('arg1', kwarg1='value1') - - mock_total.labels.assert_called_once_with(pod_id='localhost', operation='test_op') - mock_total_labels.inc.assert_called_once() - - mock_duration.labels.assert_called_once_with( - pod_id='localhost', operation='test_op' - ) - mock_duration_labels.observe.assert_called_once_with(0.5) - - -def test_execute_redis_op_exception(resource_manager): - # Mock the Redis operation to raise an exception - mock_func = MagicMock(side_effect=Exception('Operation failed')) - - with patch('ingestor.managers.resource_manager.time', return_value=100): - with patch('ingestor.metrics.REDIS_OPERATIONS_ERRORS') as mock_errors: - mock_errors_labels = MagicMock() - mock_errors.labels.return_value = mock_errors_labels - - # Execute the operation and expect an exception - with raises(Exception, match='Operation failed'): - resource_manager._execute_redis_op('test_op', mock_func, 'arg1') - - # Verify metrics and error handling - mock_errors.labels.assert_called_once_with(pod_id='localhost', operation='test_op') - mock_errors_labels.inc.assert_called_once() - resource_manager.send_notification.assert_called_once_with( - metadata=metadata['metadata'], - notification_id='REDIS_OPERATION_ERROR_test_op', - message="Error in Redis operation 'test_op': Operation failed", - block='redis_manager', - level=NotificationLevel.ERROR, - ) + resource_manager.redis_repository.keys.assert_called_once_with( + 'lease:opc_tags:*', metadata=metadata['metadata'] + ) diff --git a/tests/unit/test_ingestor.py b/tests/unit/test_ingestor.py index c98fbf3..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( @@ -74,7 +74,7 @@ def ingestor(_notification_handler, _getenv): @fixture def ingestor_manager_started(ingestor): - ingestor.ingestor_manager = MagicMock( + ingestor.ingestor_manager = AsyncMock( initialize_opc_from_config=AsyncMock(), shutdown=AsyncMock(), update_opc_servers=AsyncMock(), @@ -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() @@ -138,6 +140,7 @@ async def test_prepare_ingestor(ingestor_manager_mock, ingestor): logger=ingestor.logger, notification_handler=ingestor.notification_handler, export_to_kafka=ingestor.export_to_kafka, + metrics_controller=ingestor.metrics_controller, ) ingestor_manager.declare_active.assert_called_once() ingestor_manager.get_slot_leases.assert_called_once() @@ -177,11 +180,12 @@ def test_manage_slots_none_available(ingestor_manager_started): ingestor_manager_started.ingestor_manager.handle_acquired_tags.assert_not_called() -def test_manage_slots_none_available_none_available(ingestor_manager_started): +@mark.asyncio +async def test_manage_slots_none_available_none_available(ingestor_manager_started): ingestor_manager_started.handle_acquired_tags = MagicMock() ingestor_manager_started.ingestor_manager.managed_tags = False - ingestor_manager_started.manage_no_slots(2) + await ingestor_manager_started.manage_no_slots(2) ingestor_manager_started.ingestor_manager.get_slot_leases.assert_called_once_with(1) @@ -235,7 +239,7 @@ async def test_manage_leases_no_available_slots_extra_sltos(ingestor_manager_sta @mark.asyncio async def test_loop(ingestor_manager_started): - ingestor_manager_started.manage_no_slots = MagicMock() + ingestor_manager_started.manage_no_slots = AsyncMock() ingestor_manager_started.manage_leases = AsyncMock() ingestor_manager_started.update_ingestor_manager = AsyncMock() ingestor_manager_started.ingestor_manager.check_opc_servers_integrity = AsyncMock() @@ -244,11 +248,11 @@ async def test_loop(ingestor_manager_started): 'slot2': 'server2', 'slot3': 'server3', } - ingestor_manager_started.ingestor_manager.get_active_ingestors = MagicMock( + ingestor_manager_started.ingestor_manager.get_active_ingestors = AsyncMock( return_value=['ingestor1', 'ingestor2'] ) - ingestor_manager_started.ingestor_manager.get_number_of_slots = MagicMock(return_value=5) - ingestor_manager_started.ingestor_manager.get_number_of_leases = MagicMock(return_value=1) + ingestor_manager_started.ingestor_manager.get_number_of_slots = AsyncMock(return_value=5) + ingestor_manager_started.ingestor_manager.get_number_of_leases = AsyncMock(return_value=1) await ingestor_manager_started.loop() @@ -269,15 +273,15 @@ async def test_loop(ingestor_manager_started): @mark.asyncio async def test_loop_no_managed(ingestor_manager_started): - ingestor_manager_started.manage_no_slots = MagicMock() + ingestor_manager_started.manage_no_slots = AsyncMock() ingestor_manager_started.manage_leases = AsyncMock() ingestor_manager_started.update_ingestor_manager = AsyncMock() ingestor_manager_started.ingestor_manager.check_opc_servers_integrity = AsyncMock() ingestor_manager_started.ingestor_manager.managed_tags = {} - ingestor_manager_started.ingestor_manager.get_active_ingestors = MagicMock( + ingestor_manager_started.ingestor_manager.get_active_ingestors = AsyncMock( return_value=['ingestor1', 'ingestor2'] ) - ingestor_manager_started.ingestor_manager.get_number_of_slots = MagicMock(return_value=5) + ingestor_manager_started.ingestor_manager.get_number_of_slots = AsyncMock(return_value=5) await ingestor_manager_started.loop() diff --git a/tests/unit/test_metrics.py b/tests/unit/test_metrics.py index a60a6c7..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(): @@ -207,38 +207,6 @@ def test_kafka_connection_status(): assert set(metrics.KAFKA_CONNECTION_STATUS._labelnames) == {'pod_id'} -def test_redis_operations_total(): - """Verify the definition of REDIS_OPERATIONS_TOTAL.""" - assert metrics.REDIS_OPERATIONS_TOTAL is not None - assert isinstance(metrics.REDIS_OPERATIONS_TOTAL, Counter) - assert metrics.REDIS_OPERATIONS_TOTAL._name == 'redis_operations' # REMOVED _total - assert set(metrics.REDIS_OPERATIONS_TOTAL._labelnames) == {'pod_id', 'operation'} - - -def test_redis_operations_errors(): - """Verify the definition of REDIS_OPERATIONS_ERRORS.""" - assert metrics.REDIS_OPERATIONS_ERRORS is not None - assert isinstance(metrics.REDIS_OPERATIONS_ERRORS, Counter) - assert metrics.REDIS_OPERATIONS_ERRORS._name == 'redis_operations_errors' # REMOVED _total - assert set(metrics.REDIS_OPERATIONS_ERRORS._labelnames) == {'pod_id', 'operation'} - - -def test_redis_operations_duration(): - """Verify the definition of REDIS_OPERATIONS_DURATION.""" - assert metrics.REDIS_OPERATIONS_DURATION is not None - assert isinstance(metrics.REDIS_OPERATIONS_DURATION, Histogram) - assert metrics.REDIS_OPERATIONS_DURATION._name == 'redis_operations_duration_seconds' - assert set(metrics.REDIS_OPERATIONS_DURATION._labelnames) == {'pod_id', 'operation'} - - -def test_redis_connection_status(): - """Verify the definition of REDIS_CONNECTION_STATUS.""" - assert metrics.REDIS_CONNECTION_STATUS is not None - assert isinstance(metrics.REDIS_CONNECTION_STATUS, Gauge) - assert metrics.REDIS_CONNECTION_STATUS._name == 'redis_connection_status' - assert set(metrics.REDIS_CONNECTION_STATUS._labelnames) == {'pod_id'} - - def test_notifications_sent(): """Verify the definition of NOTIFICATIONS_SENT.""" assert metrics.NOTIFICATIONS_SENT is not None diff --git a/values.yaml b/values.yaml index 9ee918c..b7b5129 100644 --- a/values.yaml +++ b/values.yaml @@ -11,7 +11,7 @@ image: # This sets the pull policy for images. pullPolicy: Always # Overrides the image tag whose default is the chart appVersion. - tag: "0.4.9" + tag: "0.5.0" # This is for the secrets for pulling an image from a private repository more information can be found here: https://kubernetes.io/docs/tasks/configure-pod-container/pull-image-private-registry/ imagePullSecrets: @@ -139,7 +139,7 @@ env: - name: GITHUB_REPO_URL value: "git@github.com:Aignosi/sientia-dataops-opc-ingestor.git" - name: GITHUB_BRANCH - value: "fix/SIENTIAPDE-1314-ajustes-nas-camadas-de-monitoramento-do-sientia" + value: "feature/SIENTIAPDE-1325-adicionar-metricas-especificas-de-operacoes-externas" - name: PYTHON_APP value: "ingestor.app"