From d4aa44d77431dde7dbb071f7fbef1f97c5a6574f Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 26 Aug 2025 16:27:35 -0300 Subject: [PATCH] SIENTIAPDE-1205 Refactor Ingestor and OPC Manager for Asynchronous Operations - Updated main function to be asynchronous and integrated asyncio for better concurrency. - Refactored Ingestor methods to support async operations, including prepare_ingestor, loop, and shutdown. - Enhanced IngestorManager and OpcManager with async methods for improved performance and responsiveness. - Replaced blocking calls with await statements to ensure non-blocking behavior during operations. - Added a new run_async_main function to handle the async event loop setup. --- ingestor/app.py | 30 ++++++++++---- ingestor/ingestor.py | 37 ++++++++--------- ingestor/managers/ingestor_manager.py | 35 ++++++++-------- ingestor/managers/opc_manager.py | 58 +++++++++++++++------------ run_local.sh | 18 +++++++++ 5 files changed, 111 insertions(+), 67 deletions(-) create mode 100755 run_local.sh diff --git a/ingestor/app.py b/ingestor/app.py index 8fc3672..825b35e 100644 --- a/ingestor/app.py +++ b/ingestor/app.py @@ -1,4 +1,5 @@ import os +import asyncio import signal import traceback from threading import Event @@ -13,11 +14,11 @@ exit_signal = Event() POD_ID = os.getenv("HOSTNAME", "localhost") -def main(): +async def main(): start_prometheus_server() ingestor = Ingestor() try: - ingestor.prepare_ingestor() + await ingestor.prepare_ingestor() except Exception as e: metrics.APP_ERRORS_TOTAL.labels( pod_id=POD_ID).inc() # Increment errors @@ -28,11 +29,12 @@ def main(): while not exit_signal.is_set(): start_time = time() # Start loop timer try: - ingestor.loop() + await ingestor.loop() metrics.APP_LOOP_COUNT.labels( pod_id=POD_ID).inc() # Increment loop counter - exit_signal.wait(ingestor.poll_interval) + # Use asyncio.sleep instead of exit_signal.wait for better async compatibility + await asyncio.sleep(ingestor.poll_interval) except KeyboardInterrupt: # Handle Ctrl+C gracefully print("KeyboardInterrupt received. Setting exit_signal flag.") @@ -48,13 +50,13 @@ def main(): duration = time() - start_time metrics.APP_LOOP_DURATION.labels(pod_id=POD_ID).observe(duration) - ingestor.shutdown() + await ingestor.shutdown() metrics.APP_UP.labels(pod_id=POD_ID).set(0) # Mark app as DOWN ingestor.logger.info("Main loop exit_signaled.") # Give Prometheus a chance to scrape one last time before exiting (optional) - sleep(5) + await asyncio.sleep(5) os._exit(0) @@ -75,8 +77,22 @@ def start_prometheus_server(): os._exit(1) +def run_async_main(): + """Run the async main function with proper event loop setup""" + try: + loop = asyncio.new_event_loop() + asyncio.set_event_loop(loop) + loop.run_until_complete(main()) + except KeyboardInterrupt: + print("KeyboardInterrupt received in main thread.") + exit_signal.set() + finally: + loop.close() + + if __name__ == "__main__": signal.signal(signal.SIGINT, signal_handler) signal.signal(signal.SIGTERM, signal_handler) signal.signal(signal.SIGHUP, signal_handler) - main() + + run_async_main() diff --git a/ingestor/ingestor.py b/ingestor/ingestor.py index 0d68d2f..55d7543 100644 --- a/ingestor/ingestor.py +++ b/ingestor/ingestor.py @@ -1,3 +1,4 @@ +import asyncio from os import getenv from copy import deepcopy from typing import Dict, Any @@ -78,14 +79,14 @@ class Ingestor: } self.ingestor_manager = None - def shutdown(self): + async def shutdown(self): if self.ingestor_manager: - self.ingestor_manager.shutdown() + await self.ingestor_manager.shutdown() def __del__(self): - self.shutdown() + asyncio.run(self.shutdown()) - def handle_acquired_tags(self, acquired): + async def handle_acquired_tags(self, acquired): """ Handles the acquired tags by subscribing to them if available. This method checks if there are any acquired tags. If no tags are acquired, @@ -103,9 +104,9 @@ class Ingestor: else: # Subscribe to acquired slots self.ingestor_manager.update_opc_servers() - self.ingestor_manager.subscribe_to_tags(acquired) + await self.ingestor_manager.subscribe_to_tags(acquired) - def prepare_ingestor(self): + async def prepare_ingestor(self): """ Prepares the ingestor by initializing the IngestorManager, declaring the ingestor as active, acquiring slot leases, and handling the acquired tags. @@ -153,7 +154,7 @@ class Ingestor: acquired = self.ingestor_manager.get_slot_leases() self.logger.info(f"Acquired slots: {acquired}") - self.handle_acquired_tags(acquired) + await self.handle_acquired_tags(acquired) metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set( len(self.ingestor_manager.managed_tags) @@ -176,7 +177,7 @@ class Ingestor: # Get slot lease self.ingestor_manager.get_slot_leases(1) - def manage_leases( + async def manage_leases( self, available_slots: int, lacking_ingestors: int, slot_diff: int ): """ @@ -214,7 +215,7 @@ class Ingestor: self.ingestor_manager.drop_slot_leases(overleases) for lease in overleases: - self.ingestor_manager.unsubscribe_slot(lease) + await self.ingestor_manager.unsubscribe_slot(lease) self.ingestor_manager.managed_tags.pop(lease) # Update metric after removal @@ -222,7 +223,7 @@ class Ingestor: len(self.ingestor_manager.managed_tags) ) - def update_ingestor_manager(self, old_managed_tags: Dict[str, Any]): + async def update_ingestor_manager(self, old_managed_tags: Dict[str, Any]): """ Updates the ingestor manager with the new managed tags. Args: @@ -232,7 +233,7 @@ class Ingestor: self.logger.debug( f"Current managed tags: {self.ingestor_manager.managed_tags}") - self.ingestor_manager.update_opc_servers() + await self.ingestor_manager.update_opc_servers() new_managed_tags = deepcopy(self.ingestor_manager.managed_tags) self.logger.debug( @@ -249,25 +250,25 @@ class Ingestor: for slot, config in new_managed_tags.items(): if slot not in old_managed_tags: self.logger.info(f"Subscribing to new slot {slot}") - self.ingestor_manager.subscribe_to_tags({slot: config}) + await self.ingestor_manager.subscribe_to_tags({slot: config}) continue if config != old_managed_tags[slot]: self.logger.info(f"Resubscribing to slot {slot}") - self.ingestor_manager.unsubscribe_slot(slot) - self.ingestor_manager.subscribe_to_tags({slot: config}) + await self.ingestor_manager.unsubscribe_slot(slot) + await self.ingestor_manager.subscribe_to_tags({slot: config}) for slot in old_managed_tags.keys(): if slot not in new_managed_tags: self.logger.info(f"Unsubscribing from slot {slot}") - self.ingestor_manager.unsubscribe_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) ) - def loop(self): + async def loop(self): """ Executes the main loop for managing ingestors and slots. This method performs the following tasks: @@ -310,7 +311,7 @@ class Ingestor: slot_diff = len(self.ingestor_manager.managed_tags) - 1 self.logger.info("Managing leases...") - self.manage_leases(available_slots, lacking_ingestors, slot_diff) + await self.manage_leases(available_slots, lacking_ingestors, slot_diff) # Update managed slots gauge metrics.SLOTS_MANAGED.labels(pod_id=self.pod_id).set( @@ -337,4 +338,4 @@ class Ingestor: self.ingestor_manager.check_opc_servers_integrity() self.logger.info("Updating managed tags...") - self.update_ingestor_manager(current_managed_tags) + await self.update_ingestor_manager(current_managed_tags) diff --git a/ingestor/managers/ingestor_manager.py b/ingestor/managers/ingestor_manager.py index 860321a..70b0451 100644 --- a/ingestor/managers/ingestor_manager.py +++ b/ingestor/managers/ingestor_manager.py @@ -1,3 +1,4 @@ +import asyncio import traceback from typing import Dict, List from copy import deepcopy @@ -57,7 +58,7 @@ class IngestorManager(BaseActivity): notification_handler=notification_handler, set_error_counter=True) - def initialize_opc_from_config(self, server_config: dict) -> OpcManager | None: + async def initialize_opc_from_config(self, server_config: dict) -> OpcManager | None: """ Initializes an OPC Manager instance using the provided server configuration. Args: @@ -93,7 +94,7 @@ class IngestorManager(BaseActivity): ) manager.config = server_config - manager.connect() + await manager.connect() except Exception as e: trace = traceback.format_exc() @@ -111,15 +112,15 @@ class IngestorManager(BaseActivity): return manager - def shutdown(self): + async def shutdown(self): for _server_name, server in self.opc_managers.items(): - server.disconnect() + await server.disconnect() self.data_manager.shutdown() def __del__(self): - self.shutdown() + asyncio.run(self.shutdown()) - def update_opc_servers(self): + async def update_opc_servers(self): """ Updates the OPC (OLE for Process Control) server connections managed by the ingestor. This method ensures that the OPC servers defined in `self.managed_tags` are properly @@ -156,7 +157,7 @@ class IngestorManager(BaseActivity): self.logger.info( f"Initializing OPC manager for server {server}" ) - server_instance = self.initialize_opc_from_config( + server_instance = await self.initialize_opc_from_config( server_config ) @@ -166,7 +167,7 @@ class IngestorManager(BaseActivity): ) server_instance.disconnect() del self.opc_managers[server] - server_instance = self.initialize_opc_from_config( + server_instance = await self.initialize_opc_from_config( server_config ) else: @@ -191,7 +192,7 @@ class IngestorManager(BaseActivity): f"Server {server} not found in managed tags. " f"Desconnecting from server." ) - self.opc_managers[server].disconnect() + await self.opc_managers[server].disconnect() self.opc_managers.pop(server, None) metrics.OPC_MANAGERS_ACTIVE.labels( @@ -312,7 +313,7 @@ class IngestorManager(BaseActivity): pod_id=self.pod_id).set(len(self.managed_tags)) return acquired - def unsubscribe_slot(self, slot: str): + async def unsubscribe_slot(self, slot: str): """ Unsubscribes a specific slot from all associated OPC servers. Args: @@ -323,7 +324,7 @@ class IngestorManager(BaseActivity): for server in self.managed_tags[slot].keys(): if server in self.opc_managers: - self.opc_managers[server].unsubscribe(slot) + await self.opc_managers[server].unsubscribe(slot) def update_slot_config(self): """ @@ -378,7 +379,7 @@ class IngestorManager(BaseActivity): self.resource_manager.drop_tag_lease(lease_id) metrics.SLOTS_RELEASED.labels(pod_id=self.pod_id).inc() - def manage_server(self, slot: str, server: str, server_config: dict, tags: dict) -> int: + async def manage_server(self, slot: str, server: str, server_config: dict, tags: dict) -> int: """ Manages the subscription of tags to a specified OPC server and slot. This method ensures that the specified server and slot have an active subscription @@ -416,7 +417,7 @@ class IngestorManager(BaseActivity): return 1 if slot not in self.opc_managers[server].subscriptions: try: - self.opc_managers[server].create_subscription( + await self.opc_managers[server].create_subscription( slot ) except Exception as e: @@ -428,7 +429,7 @@ class IngestorManager(BaseActivity): self.logger.info( tags_to_sub ) - self.opc_managers[server].subscribe( + await self.opc_managers[server].subscribe( slot, deepcopy(tags_to_sub), self.poll_interval ) self.logger.info( @@ -452,11 +453,11 @@ class IngestorManager(BaseActivity): "Removing subscription from server " f"{server} for slot {slot}" ) - self.opc_managers[server].unsubscribe(slot) + await self.opc_managers[server].unsubscribe(slot) return 2 return 0 - def subscribe_to_tags(self, tags: Dict) -> None: + async def subscribe_to_tags(self, tags: Dict) -> None: """ Subscribes to a set of tags and manages their configurations. This method processes a dictionary of tags, iterating through each slot and server @@ -478,7 +479,7 @@ class IngestorManager(BaseActivity): self.logger.info(tags) for slot, slot_config in tags.items(): for server, server_config in slot_config.items(): - response = self.manage_server( + response = await self.manage_server( slot, server, server_config, tags ) if response == 2: diff --git a/ingestor/managers/opc_manager.py b/ingestor/managers/opc_manager.py index 179c23b..7f846c4 100644 --- a/ingestor/managers/opc_manager.py +++ b/ingestor/managers/opc_manager.py @@ -1,7 +1,8 @@ import json +import asyncio from pathlib import Path from asyncua.crypto.security_policies import SecurityPolicyBasic256 -from asyncua.sync import Client +from asyncua import Client from sientia_do.notifications.models import NotificationLevel from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.temporal.activities.base import BaseActivity @@ -42,7 +43,18 @@ class OpcManager(BaseActivity): return f"OpcManager(name={self.name}, url={self.url}, server_uri={self.server_uri})\n" \ f"nodes={self.nodes}, subscriptions={self.subscriptions}" - def set_security(self): + async def shutdown(self): + """Comprehensive cleanup method""" + try: + + await self.disconnect() + except Exception as e: + self.logger.error(f"Error during cleanup: {e}") + + def __del__(self): + asyncio.run(self.shutdown()) + + async def set_security(self): """ Configures the security settings for the OPC UA client. This method sets up the security policy, certificates, and timeouts @@ -70,18 +82,18 @@ class OpcManager(BaseActivity): server_cert = Path( self.server_cert_path) if self.server_cert_path else None - self.client.application_uri = self.server_uri + await self.client.set_application_uri(self.server_uri) self.logger.info('Setting security...') - self.client.set_security( + await self.client.set_security( SecurityPolicyBasic256, certificate=str(cert), private_key=str(private_key), server_certificate=str(server_cert) ) - self.client.secure_channel_timeout = 10000000 - self.client.session_timeout = 10000000 + await self.client.set_secure_channel_timeout(10000000) + await self.client.set_session_timeout(10000000) - def connect(self): + async def connect(self): """ Establishes a connection to the OPC server. This method initializes the OPC client using the provided URL and @@ -96,9 +108,9 @@ class OpcManager(BaseActivity): try: self.client = Client(self.url) if self.cert_path: - self.set_security() + await self.set_security() self.logger.info(f'Starting connection to {self.name}...') - self.client.connect() + await self.client.connect() metrics.OPC_CONNECTION_STATUS.labels( pod_id=self.pod_id, server_name=self.name, server_url=self.url).set(1) self.logger.info(f'Connection to {self.name} successful.') @@ -110,7 +122,7 @@ class OpcManager(BaseActivity): self.logger.error(f"Failed to connect to {self.name}: {e}") raise - def create_subscription(self, name: str, period: int = 500): + async def create_subscription(self, name: str, period: int = 500): """ Creates a subscription with the specified monitoring period. This method establishes a subscription to monitor data changes or events @@ -129,7 +141,7 @@ class OpcManager(BaseActivity): raise ValueError("Client not connected. Call connect first.") try: p = period if period is not None else 500 - self.subscriptions[name] = self.client.create_subscription(p, self) + self.subscriptions[name] = await self.client.create_subscription(p, 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() @@ -138,7 +150,7 @@ class OpcManager(BaseActivity): f"Failed to create subscription {name} on {self.name}: {e}") raise - def subscribe(self, subscription: str, nodes: dict, collect_period: int): + async def subscribe(self, subscription: str, nodes: dict, collect_period: int): """ Subscribes to a set of OPC UA nodes for data change notifications. This method adds the specified nodes to the subscription and configures @@ -161,8 +173,7 @@ class OpcManager(BaseActivity): self.logger.info(f"Subscribing to {subscription} on {self.name}...") self.logger.info(f"Subscribing to nodes: {nodes}") - addr_nodes = [self.client.get_node( - n) for n in nodes] + addr_nodes = [self.client.get_node(n) for n in nodes] self.logger.debug(f"Addr nodes: {addr_nodes}") self.nodes.update(nodes) self.logger.debug(f"Nodes: {self.nodes}") @@ -176,9 +187,9 @@ class OpcManager(BaseActivity): 'cycle_count': 0 } - self.subscriptions[subscription].subscribe_data_change(addr_nodes) + await self.subscriptions[subscription].subscribe_data_change(addr_nodes) - def unsubscribe(self, subscription: str): + async def unsubscribe(self, subscription: str): """ Unsubscribes from a given subscription. Args: @@ -196,14 +207,11 @@ class OpcManager(BaseActivity): self.logger.warning( f"Subscription '{subscription}' not found. Cannot unsubscribe.") return - self.subscriptions[subscription].delete() + await self.subscriptions[subscription].delete() del self.subscriptions[subscription] self.logger.info(f"Unsubscribed from {subscription}.") - def __del__(self): - self.disconnect() - - def disconnect(self): + async def disconnect(self): """ Disconnects from the OPC UA server. This method handles the disconnection process by deleting the subscription @@ -219,14 +227,14 @@ class OpcManager(BaseActivity): self.logger.warning("Client already disconnected.") return try: - _a = [self.subscriptions[sub].delete() - for sub in self.subscriptions] + for sub in self.subscriptions: + await self.subscriptions[sub].delete() self.logger.warning("Deleted all subscriptions.") except Exception as sub_error: self.logger.error(f"Failed to clean up subscription: {sub_error}") try: - self.client.disconnect() + await self.client.disconnect() except Exception as conn_error: self.logger.error( f"Failed to disconnect from OPC UA server: {conn_error}") @@ -239,7 +247,7 @@ class OpcManager(BaseActivity): pod_id=self.pod_id, server_name=self.name).set(0) self.logger.warning("Disconnected from OPC UA server.") - def datachange_notification(self, node, _val, data): + async def datachange_notification(self, node, _val, data): """ Handles data change notifications for monitored OPC UA nodes. This method is triggered when a monitored node's value changes. It processes diff --git a/run_local.sh b/run_local.sh new file mode 100755 index 0000000..de3ad55 --- /dev/null +++ b/run_local.sh @@ -0,0 +1,18 @@ +#!/bin/bash + +# Exit on any error +set -e + +echo "Activating virtual environment..." +source ./venv/bin/activate + +echo "Loading environment variables from .env..." +if [ -f .env ]; then + export $(cat .env | grep -v '^#' | xargs) + echo "Environment variables loaded from .env" +else + echo "Warning: .env file not found. Continuing without environment variables." +fi + +echo "Starting ingestor application..." +python -m ingestor.app