SIENTIAPDE-988

Enhance README and Docker setup; improve IngestorManager and DataManager functionality, add error handling, and implement unit tests for OPC initialization and subscription management.
This commit is contained in:
vitor-aignosi
2025-04-17 17:13:13 -03:00
parent 4e919a08be
commit fa51d142e8
13 changed files with 546 additions and 80 deletions

View File

@@ -6,14 +6,14 @@ from time import sleep
def main():
# Get os parameters
kafka_servers = getenv("KAFKA_SERVERS")
redis_host = getenv("REDIS_HOST")
redis_port = int(getenv("REDIS_PORT"))
lease_ttl = int(getenv("LEASE_TTL"))
heartbeat_ttl = int(getenv("HEARTBEAT_TTL"))
pod_id = getenv("HOSTNAME")
number_of_ingestors = int(getenv("REPLICA_COUNT"))
poll_interval = int(getenv("POLL_INTERVAL"))
kafka_servers = getenv("KAFKA_SERVERS", "localhost:9092")
redis_host = getenv("REDIS_HOST", "localhost")
redis_port = int(getenv("REDIS_PORT", 6379))
lease_ttl = int(getenv("LEASE_TTL", 20))
heartbeat_ttl = int(getenv("HEARTBEAT_TTL", 10))
pod_id = getenv("HOSTNAME", "localhost")
number_of_ingestors = int(getenv("REPLICA_COUNT", 1))
poll_interval = int(getenv("POLL_INTERVAL", 5))
logger = getLogger(__name__)
logger.setLevel(getenv("LOG_LEVEL", "INFO"))
@@ -29,22 +29,23 @@ def main():
lease_ttl, heartbeat_ttl, pod_id,
number_of_ingestors, poll_interval, logger
)
# Declare ingestor ative
ingestor_manager.declare_active()
# Get slot lease
acquired = ingestor_manager.get_slot_leases()
logger.info(f"Acquired slots: {acquired}")
# Subscribe to acquired slots
ingestor_manager.update_opc_servers()
ingestor_manager.subscribe_to_tags(acquired)
if not acquired:
logger.warning("No slots available")
else:
# Declare ingestor ative
ingestor_manager.declare_active()
# Subscribe to acquired slots
ingestor_manager.update_opc_servers()
ingestor_manager.subscribe_to_tags(acquired)
while True:
# Declare ingestor as active
ingestor_manager.declare_active()
# Update opc servers
ingestor_manager.update_slot_config()
logger.info("Polling for slot updates...")
# Get active ingestors
ingestors = ingestor_manager.get_active_ingestors()
@@ -53,13 +54,18 @@ def main():
if ingestor_diff > 0:
# Some ingestors are innactive, so theres "ingestor_diff" slots available
logger.info(f"Slots available: {ingestor_diff}")
# Get slot lease
acquired = ingestor_manager.get_slot_leases(ingestor_diff)
# Subscribe to acquired slots
ingestor_manager.update_opc_servers()
ingestor_manager.subscribe_to_tags(acquired)
if not acquired:
logger.info("No slots acquired")
else:
# Subscribe to acquired slots
ingestor_manager.update_opc_servers()
ingestor_manager.subscribe_to_tags(acquired)
elif slot_diff > 0:
# Some ingestors are active and without slots, so we need to drop
@@ -68,6 +74,17 @@ def main():
ingestor_manager.drop_slot_leases(overleases)
# Update opc servers
ingestor_manager.update_slot_config()
if not ingestor_manager.managed_tags:
logger.info("No managed tags found")
sleep(poll_interval)
continue
# Declare ingestor as active
ingestor_manager.declare_active()
# Sleep for poll interval
sleep(poll_interval)

View File

@@ -43,9 +43,16 @@ class DataManager():
Exception: If there is an error during message delivery, it will be handled by the `delivery_error` callback.
"""
self.kafka_producer.send(
topic=topic, value=data).add_callback(
self.delivery_report).add_errback(
self.delivery_error)
try:
self.kafka_producer.flush()
self.logger.info(
f"Publishing message to topic {topic}: {data}")
self.kafka_producer.send(
topic=topic, value=data).add_callback(
self.delivery_report).add_errback(
self.delivery_error)
self.kafka_producer.flush(timeout=10)
except Exception as e:
self.logger.error(f"Failed to publish message: {e}")

View File

@@ -1,5 +1,7 @@
from logging import Logger
import traceback
from typing import Dict, List
from copy import deepcopy
from ingestor.managers.data_manager import DataManager
from ingestor.managers.opc_manager import OpcManager
from ingestor.managers.resource_manager import ResourceManager
@@ -23,37 +25,88 @@ class IngestorManager():
self.managed_tags = {}
self.opc_servers = {}
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:
server_config (dict): A dictionary containing the OPC server configuration.
Expected keys include:
- 'name' (str): The name of the OPC server.
- 'url' (str): The URL of the OPC server.
- 'server_uri' (str): The URI of the OPC server.
- 'cert_path' (str, optional): Path to the client certificate file.
- 'private_key_path' (str, optional): Path to the private key file.
- 'server_cert_path' (str, optional): Path to the server certificate file.
data_manager (DataManager): An instance of the DataManager to handle data operations.
logger (Logger): A logger instance for logging messages.
Returns:
OpcManager | None: An initialized OpcManager instance if successful,
otherwise None if an error occurs during initialization.
"""
try:
manager = OpcManager(
server_config['name'], server_config['url'], data_manager, logger, server_config['server_uri'],
server_config.get('cert_path'), server_config.get(
'private_key_path'), server_config.get('server_cert_path')
)
manager.config = server_config
manager.connect()
except Exception as e:
logger.error(f"Failed to initialize OpcManager: {e}")
return None
return manager
def update_opc_servers(self):
registered_servers = []
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)
if server not in self.opc_servers:
self.opc_managers[server] = OpcManager.initialize_from_config(
server_config, self.data_manager, self.logger
)
self.opc_managers[server].connect()
elif self.opc_servers[server] != server_config:
del self.opc_managers[server]
self.opc_managers[server] = OpcManager.initialize_from_config(
server_config, self.data_manager, self.logger
)
self.opc_managers[server].connect()
server_instance = None
if server not in self.opc_managers:
self.opc_servers[server] = server_config
server_instance = self.initialize_opc_from_config(
server_config, self.data_manager, self.logger
)
elif self.opc_managers[server].config != server_config:
del self.opc_managers[server]
server_instance = self.initialize_opc_from_config(
server_config, self.data_manager, self.logger
)
if server_instance is not None:
self.opc_managers[server] = server_instance
for server in list(self.opc_managers.keys()):
if server not in registered_servers:
self.logger.warning(
f"Server {server} not found in managed tags. "
f"Desconnecting from server."
)
self.opc_managers[server].disconnect()
self.opc_managers.pop(server, None)
def declare_active(self):
self.resource_manager.ingestor_heartbeat()
def get_active_ingestors(self) -> List[str]:
self.resource_manager.get_all_ingestors()
ingestors = self.resource_manager.get_all_ingestors()
return ingestors if ingestors else []
def get_slot_leases(self, max_slots: int = 1) -> Dict:
acquired = {}
for i in range(1, self.number_of_ingestors + 1):
if self.resource_manager.lease_tag(str(i)):
self.logger.info(f"Leased slot {i}")
acquired[str(i)] = self.resource_manager.get_tag_slot(str(i))
slots = self.resource_manager.get_tag_slot(str(i))
if slots is None:
continue
acquired[str(i)] = slots
if len(acquired) >= max_slots:
self.managed_tags.update(acquired)
@@ -66,9 +119,26 @@ class IngestorManager():
self.managed_tags.update(acquired)
return acquired
def update_slot_config(self) -> Dict:
def unsubscribe_slot(self, slot: str):
for server in self.managed_tags[slot].keys():
if server in self.opc_managers:
self.opc_managers[server].unsubscribe(slot)
def update_slot_config(self):
removed_slots = []
update = {}
for slot, slot_config in self.managed_tags.items():
self.resource_manager.renew_tag_lease(slot)
update = self.resource_manager.get_tag_slot(slot)
if update is None:
self.logger.warning(
f"Slot {slot} configuration not found. "
f"Removing slot from managed tags."
)
self.unsubscribe_slot(slot)
removed_slots.append(slot)
continue
if update != slot_config:
self.logger.info(
@@ -77,27 +147,64 @@ class IngestorManager():
)
self.managed_tags[slot] = update
self.unsubscribe_slot(slot)
self.update_opc_servers()
self.subscribe_to_tags(update)
self.subscribe_to_tags({slot: update})
return update
for slot in removed_slots:
del self.managed_tags[slot]
self.update_opc_servers()
def drop_slot_leases(self, ids: List[str]) -> None:
for id in ids:
self.resource_manager.drop_tag_lease(id)
def subscribe_to_tags(self, tags: Dict) -> None:
def subscribe_to_tags(self, tags: Dict) -> None: # NOSONAR
to_remove = []
self.logger.info(tags)
for slot, slot_config in tags.items():
for server, server_config in slot_config.items():
self.logger.info(
f"Subscribing to tags from {slot}:{server}"
)
tags_to_sub = server_config.get('tags')
if server not in self.opc_managers:
self.logger.error(
f"Server {server} not found in opc_managers."
)
continue
if slot not in self.opc_managers[server].subscriptions:
self.opc_managers[server].create_subscription(
slot
try:
self.opc_managers[server].create_subscription(
slot
)
except Exception as e:
self.logger.error(
f"Failed to create subscription for slot {slot}: {e}"
)
to_remove.append([slot, server])
continue
try:
self.logger.info(
tags_to_sub
)
self.opc_managers[server].subscribe(
slot, server_config['tags'], self.poll_interval
)
self.opc_managers[server].subscribe(
slot, deepcopy(tags_to_sub), self.poll_interval
)
self.logger.info(
tags_to_sub
)
except Exception as e:
self.logger.error(
f"Failed to subscribe to tags from {slot}:{server}\n{tags}: {e}"
)
self.logger.error(traceback.format_exc())
self.logger.warning(
"Removing subscription from server "
f"{server} for slot {slot}"
)
self.opc_managers[server].unsubscribe(slot)
to_remove.append([slot, server])
for slot, server in to_remove:
self.managed_tags[slot].pop(server, None)

View File

@@ -40,6 +40,7 @@ class OpcManager():
server_config.get('cert_path'), server_config.get(
'private_key_path'), server_config.get('server_cert_path')
)
self.config = server_config
def set_security(self):
"""
@@ -140,6 +141,8 @@ class OpcManager():
raise ValueError(
"Subscription not created. Call create_subscription first.")
self.logger.info(f"Subscribing to {subscription}...")
self.logger.info(f"Subscribing to nodes: {nodes}")
self.addr_nodes = [self.client.get_node(
n) for n in nodes if n not in self.nodes]
self.nodes.update(nodes)
@@ -156,8 +159,9 @@ class OpcManager():
def unsubscribe(self, subscription: str):
if not self.subscriptions.get(subscription):
raise ValueError(
"Subscription not created. Call create_subscription first.")
self.logger.warning(
f"Subscription {subscription} not found. Cannot unsubscribe.")
return
self.subscriptions[subscription].delete()
del self.subscriptions[subscription]
self.logger.info(f"Unsubscribed from {subscription}.")
@@ -177,10 +181,14 @@ class OpcManager():
"""
self.logger.warning('Disconnecting from OPC server')
if self.client is None:
self.logger.warning("Client already disconnected.")
return
try:
[self.subscriptions[sub].delete() for sub in self.subscriptions]
self.logger.warning("Deleted all subscriptions.")
self.client.disconnect()
del self.client
self.client = None
self.logger.warning("Disconnected from OPC UA server.")
except Exception as sub_error:
self.logger.error(f"Failed to clean up subscription: {sub_error}")
@@ -217,7 +225,7 @@ class OpcManager():
data = {
'tag': tag,
'name': self.nodes[str(node)]['tag_name'],
'timestamp': source_timestamp,
'timestamp': source_timestamp.strftime('%Y-%m-%d %H:%M:%S'),
'value': value
}