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