diff --git a/ingestor/ingestor.py b/ingestor/ingestor.py index 45b33f8..f7fb6f5 100644 --- a/ingestor/ingestor.py +++ b/ingestor/ingestor.py @@ -1,9 +1,10 @@ from logging import Formatter, StreamHandler, getLogger from os import getenv -from ingestor.managers.ingestor_manager import IngestorManager from sientia_do.notifications.handlers import NotificationHandler +from ingestor.managers.ingestor_manager import IngestorManager + class Ingestor: def __init__(self): @@ -29,13 +30,13 @@ class Ingestor: kafka_servers = getenv("KAFKA_SERVERS", "localhost:9092") self.redis_host = getenv("REDIS_HOST", "localhost") - self.redis_port = int(getenv("REDIS_PORT", 6379)) + self.redis_port = int(getenv("REDIS_PORT", '6379')) self.redis_username = getenv("REDIS_USERNAME", None) self.redis_password = getenv("REDIS_PASSWORD", None) - self.lease_ttl = int(getenv("LEASE_TTL", 10)) - self.heartbeat_ttl = int(getenv("HEARTBEAT_TTL", 20)) + self.lease_ttl = int(getenv("LEASE_TTL", '10')) + self.heartbeat_ttl = int(getenv("HEARTBEAT_TTL", '20')) self.pod_id = getenv("HOSTNAME", "localhost") - self.poll_interval = int(getenv("POLL_INTERVAL", 5)) + self.poll_interval = int(getenv("POLL_INTERVAL", '5')) self.kafka_servers = kafka_servers.split(",") self.logger = None @@ -49,9 +50,8 @@ class Ingestor: model_name="-", model="-" ) - # build args for build notificarions components - # call build notifications components + self.ingestor_manager = None def init_logger(self): """ @@ -128,7 +128,7 @@ class Ingestor: # Get slot lease acquired = self.ingestor_manager.get_slot_leases() - self.logger.info(f"Acquired slots: {acquired}") + self.logger.info("Acquired slots: %s", acquired) self.handle_acquired_tags(acquired) @@ -171,7 +171,7 @@ class Ingestor: if available_slots > 0 and lacking_ingestors > 0: # Some ingestors are innactive, so theres "available_slots" slots available - self.logger.info(f"Slots available: {available_slots}") + self.logger.info("Slots available: %s", available_slots) # Get slot lease acquired = self.ingestor_manager.get_slot_leases(available_slots) @@ -180,7 +180,7 @@ class Ingestor: elif lacking_ingestors == 0 and slot_diff > 0: - self.logger.info(f"Extra slots available: {slot_diff}") + self.logger.info("Extra slots available: %s", slot_diff) # There's enough slots for all ingestors, but this ingestor has more than one slot # So we need to drop the extra leases @@ -210,37 +210,35 @@ class Ingestor: self.ingestor_manager.declare_active() self.logger.info("Polling for slot updates...") - # Get active ingestors[ + # Get active ingestors ingestors = self.ingestor_manager.get_active_ingestors() - self.logger.debug(f"Active ingestors: {ingestors}") number_of_ingestors = len(ingestors) - self.logger.debug(f"Number of ingestors: {number_of_ingestors}") number_of_leases = self.ingestor_manager.get_number_of_leases() - self.logger.debug(f"Number of leases: {number_of_leases}") number_of_slots = self.ingestor_manager.get_number_of_slots() - self.logger.debug(f"Number of slots: {number_of_slots}") # Handle no slots self.logger.debug("Managing no slots...") self.manage_no_slots(number_of_slots) available_slots = number_of_slots - number_of_leases - self.logger.debug(f"Available slots: {available_slots}") lacking_ingestors = number_of_slots - number_of_ingestors - self.logger.debug(f"Lacking ingestors: {lacking_ingestors}") slot_diff = len(self.ingestor_manager.managed_tags) - 1 - self.logger.debug(f"Slot diff: {slot_diff}") self.logger.debug("Managing leases...") self.manage_leases(available_slots, lacking_ingestors, slot_diff) self.logger.debug( - f"Active ingestors: {ingestors}, " - f"Number of slots: {number_of_slots}, " - f"Number of leases: {number_of_leases}, " - f"Managed tags: {self.ingestor_manager.managed_tags}" - f"Managed servers: {self.ingestor_manager.opc_managers}" + "Active ingestors: %s, " + "Number of slots: %s, " + "Number of leases: %s, " + "Managed tags: %s, " + "Managed servers: %s", + ingestors, + number_of_slots, + number_of_leases, + self.ingestor_manager.managed_tags, + self.ingestor_manager.opc_managers ) if not self.ingestor_manager.managed_tags: # No slots acquired diff --git a/ingestor/managers/ingestor_manager.py b/ingestor/managers/ingestor_manager.py index 06631c4..187a234 100644 --- a/ingestor/managers/ingestor_manager.py +++ b/ingestor/managers/ingestor_manager.py @@ -2,11 +2,11 @@ from logging import Logger import traceback from typing import Dict, List from copy import deepcopy +from sientia_do.notifications.handlers import NotificationHandler +from sientia_do.notifications.models import NotificationLevel from ingestor.managers.data_manager import DataManager from ingestor.managers.opc_manager import OpcManager from ingestor.managers.resource_manager import ResourceManager -from sientia_do.notifications.models import Notification, NotificationLevel -from sientia_do.notifications.handlers import NotificationHandler class IngestorManager(): @@ -30,7 +30,8 @@ class IngestorManager(): self.notification_handler = notification_handler - def initialize_opc_from_config(self, server_config: dict, data_manager: DataManager, logger: Logger) -> OpcManager | None: + def initialize_opc_from_config(self, server_config: dict, + data_manager: DataManager, logger: Logger) -> OpcManager | None: """ Initializes an OPC Manager instance using the provided server configuration. Args: @@ -53,9 +54,11 @@ class IngestorManager(): self.logger.info( f"Initializing OpcManager at {server_config['url']}") manager = OpcManager( - server_config['name'], server_config['url'], data_manager, logger, server_config['server_uri'], - self.notification_handler, server_config.get('cert_path'), server_config.get( - 'private_key_path'), server_config.get('server_cert_path') + server_config['name'], server_config['url'], + data_manager, logger, server_config['server_uri'], + self.notification_handler, server_config.get('cert_path'), + server_config.get('private_key_path'), + server_config.get('server_cert_path') ) manager.config = server_config @@ -91,7 +94,8 @@ class IngestorManager(): 4. Disconnects and removes OPC servers that are no longer registered. Attributes: self.managed_tags (dict): A nested dictionary containing slot and server configurations. - self.opc_managers (dict): A dictionary mapping server names to their OPC manager instances. + self.opc_managers (dict): A dictionary mapping server names to their + OPC manager instances. self.data_manager: An object responsible for managing data operations. self.logger: A logging object for recording warnings and other messages. Raises: @@ -101,19 +105,24 @@ class IngestorManager(): """ registered_servers = [] - for slot, slot_config in self.managed_tags.items(): + for _slot, slot_config in self.managed_tags.items(): for server, server_config in slot_config.items(): registered_servers.append(server) server_config = server_config.copy() server_config.pop('tags', None) server_instance = None if server not in self.opc_managers: - + self.logger.debug( + f"Initializing OPC manager for server {server}" + ) server_instance = self.initialize_opc_from_config( server_config, self.data_manager, self.logger ) elif self.opc_managers[server].config != server_config: + self.logger.warning( + f"Reinitializing OPC manager for server {server}" + ) self.opc_managers[server].disconnect() del self.opc_managers[server] server_instance = self.initialize_opc_from_config( @@ -150,16 +159,9 @@ class IngestorManager(): self.logger.warning( f"Reconnecting to OPC server {server}" ) - try: - self.opc_managers[server] = self.initialize_opc_from_config( - opc_manager.config, self.data_manager, self.logger - ) - except Exception as e: - self.logger.error( - f"Failed to reconnect to OPC server {server}: {e}" - ) - else: - self.update_opc_servers() + self.opc_managers.pop(server, None) + self.update_opc_servers() + if self.opc_managers.get(server) is not None: for slot, slot_config in self.managed_tags.items(): if server in slot_config: self.manage_server( @@ -222,14 +224,16 @@ class IngestorManager(): Dict: A dictionary where the keys are the slot identifiers (as strings) and the values are the leased slot details. Behavior: - - Iterates through available slots and attempts to lease them using the resource manager. + - Iterates through available slots and attempts to lease + them using the resource manager. - Logs the leasing of each slot. - Updates the `managed_tags` attribute with the acquired slots. - Stops leasing once the specified `max_slots` are acquired. - - If unable to acquire the requested number of slots, logs a warning and returns the slots that were leased. + - If unable to acquire the requested number of slots, + logs a warning and returns the slots that were leased. Notes: - - If a slot is leased but its details cannot be retrieved (i.e., `get_tag_slot` returns None), - that slot is skipped. + - If a slot is leased but its details cannot be retrieved + (i.e., `get_tag_slot` returns None), that slot is skipped. """ acquired = {} @@ -331,8 +335,8 @@ class IngestorManager(): None """ - for id in ids: - self.resource_manager.drop_tag_lease(id) + for lease_id in ids: + self.resource_manager.drop_tag_lease(lease_id) def manage_server(self, slot: str, server: str, server_config: dict, tags: dict) -> int: """