Update environment configuration and refactor OPC activities for async handling
- Expanded the .env file with configurations for MongoDB, Postgres, MlFlow, and Temporal. - Refactored OPC class methods to be asynchronous, including init_opc, write_data, manage_output_tags, and shutdown. - Updated the worker to initialize OPC asynchronously and adjusted shutdown handling for activities.
This commit is contained in:
2
.gitignore
vendored
2
.gitignore
vendored
@@ -42,3 +42,5 @@ htmlcov/
|
|||||||
git_key*
|
git_key*
|
||||||
|
|
||||||
git_log
|
git_log
|
||||||
|
|
||||||
|
.env
|
||||||
@@ -44,6 +44,6 @@ class Activities(Postgres, MLFlow, Gates, OPC):
|
|||||||
logger=logger,
|
logger=logger,
|
||||||
notification_handler=notification_handler)
|
notification_handler=notification_handler)
|
||||||
|
|
||||||
def shutdown(self):
|
async def shutdown(self):
|
||||||
Postgres.close(self)
|
Postgres.close(self)
|
||||||
OPC.shutdown(self)
|
await OPC.shutdown(self)
|
||||||
|
|||||||
@@ -26,7 +26,10 @@ class OPC(BaseActivity):
|
|||||||
self, logger, notification_handler, set_error_counter=True)
|
self, logger, notification_handler, set_error_counter=True)
|
||||||
|
|
||||||
self.opc_repository: dict[str, OpcRepository] = {}
|
self.opc_repository: dict[str, OpcRepository] = {}
|
||||||
for id, server in opc_servers.items():
|
self.opc_servers = opc_servers
|
||||||
|
|
||||||
|
async def init_opc(self):
|
||||||
|
for id, server in self.opc_servers.items():
|
||||||
self.opc_repository[id] = OpcRepository(
|
self.opc_repository[id] = OpcRepository(
|
||||||
id=server['id'],
|
id=server['id'],
|
||||||
url=server['url'],
|
url=server['url'],
|
||||||
@@ -39,7 +42,7 @@ class OPC(BaseActivity):
|
|||||||
reconnection_interval=server['reconnection_interval'],
|
reconnection_interval=server['reconnection_interval'],
|
||||||
pod_id=self.pod_id
|
pod_id=self.pod_id
|
||||||
)
|
)
|
||||||
is_connected, error_data = self.opc_repository[id].connect()
|
is_connected, error_data = await self.opc_repository[id].connect()
|
||||||
if not is_connected:
|
if not is_connected:
|
||||||
self.send_notification(
|
self.send_notification(
|
||||||
metadata={
|
metadata={
|
||||||
@@ -56,7 +59,7 @@ class OPC(BaseActivity):
|
|||||||
'attachment_content', None)
|
'attachment_content', None)
|
||||||
)
|
)
|
||||||
|
|
||||||
def write_data(self, server_id: str, tag: str, data: Any,
|
async def write_data(self, server_id: str, tag: str, data: Any,
|
||||||
data_type: str, tag_type: str, metadata: dict[str, Any]) -> bool:
|
data_type: str, tag_type: str, metadata: dict[str, Any]) -> bool:
|
||||||
"""
|
"""
|
||||||
Write data to OPC server.
|
Write data to OPC server.
|
||||||
@@ -73,7 +76,7 @@ class OPC(BaseActivity):
|
|||||||
"""
|
"""
|
||||||
|
|
||||||
try:
|
try:
|
||||||
is_success, error_data = self.opc_repository[server_id].write_data(
|
is_success, error_data = await self.opc_repository[server_id].write_data(
|
||||||
tag, data, data_type, self.logger, metadata)
|
tag, data, data_type, self.logger, metadata)
|
||||||
if not is_success:
|
if not is_success:
|
||||||
self.send_notification(
|
self.send_notification(
|
||||||
@@ -113,14 +116,14 @@ class OPC(BaseActivity):
|
|||||||
return False
|
return False
|
||||||
return True
|
return True
|
||||||
|
|
||||||
def manage_output_tags(
|
async def manage_output_tags(
|
||||||
self, server_id: str, config: dict[str, Any], data: DataFrame,
|
self, server_id: str, config: dict[str, Any], data: DataFrame,
|
||||||
metadata: dict[str, Any], success: bool) -> tuple[bool, int]:
|
metadata: dict[str, Any], success: bool) -> tuple[bool, int]:
|
||||||
|
|
||||||
count = 0
|
count = 0
|
||||||
if 'prediction_tags' in config:
|
if 'prediction_tags' in config:
|
||||||
for tag, tag_config in config['prediction_tags'].items():
|
for tag, tag_config in config['prediction_tags'].items():
|
||||||
local_success = self.write_data(
|
local_success = await self.write_data(
|
||||||
server_id=server_id,
|
server_id=server_id,
|
||||||
tag=tag,
|
tag=tag,
|
||||||
data=data.head(1)['prediction'].values[0],
|
data=data.head(1)['prediction'].values[0],
|
||||||
@@ -136,7 +139,7 @@ class OPC(BaseActivity):
|
|||||||
|
|
||||||
if 'confidence_tags' in config:
|
if 'confidence_tags' in config:
|
||||||
for tag, tag_config in config['confidence_tags'].items():
|
for tag, tag_config in config['confidence_tags'].items():
|
||||||
local_success = self.write_data(
|
local_success = await self.write_data(
|
||||||
server_id=server_id,
|
server_id=server_id,
|
||||||
tag=tag,
|
tag=tag,
|
||||||
data=data.head(1)['prediction_confidence'].values[0],
|
data=data.head(1)['prediction_confidence'].values[0],
|
||||||
@@ -185,7 +188,7 @@ class OPC(BaseActivity):
|
|||||||
success = False
|
success = False
|
||||||
continue
|
continue
|
||||||
|
|
||||||
local_success, local_count = self.manage_output_tags(
|
local_success, local_count = await self.manage_output_tags(
|
||||||
server_id, config, data, metadata, success)
|
server_id, config, data, metadata, success)
|
||||||
success = success and local_success
|
success = success and local_success
|
||||||
|
|
||||||
@@ -221,6 +224,6 @@ class OPC(BaseActivity):
|
|||||||
|
|
||||||
return data.to_dict()
|
return data.to_dict()
|
||||||
|
|
||||||
def shutdown(self):
|
async def shutdown(self):
|
||||||
for opc in self.opc_repository.values():
|
for opc in self.opc_repository.values():
|
||||||
opc.disconnect()
|
await opc.disconnect()
|
||||||
|
|||||||
@@ -1,9 +1,10 @@
|
|||||||
|
import asyncio
|
||||||
import traceback
|
import traceback
|
||||||
import time
|
import time
|
||||||
from datetime import datetime
|
from datetime import datetime
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from typing import Any
|
from typing import Any
|
||||||
from asyncua.sync import Client
|
from asyncua import Client
|
||||||
from asyncua.crypto.security_policies import SecurityPolicyBasic256
|
from asyncua.crypto.security_policies import SecurityPolicyBasic256
|
||||||
from asyncua.ua import DataValue, Variant, VariantType, DateTime
|
from asyncua.ua import DataValue, Variant, VariantType, DateTime
|
||||||
from regex import F
|
from regex import F
|
||||||
@@ -62,7 +63,7 @@ class OpcRepository():
|
|||||||
'schedule_name': '-'
|
'schedule_name': '-'
|
||||||
}
|
}
|
||||||
|
|
||||||
def set_security(self):
|
async def set_security(self):
|
||||||
"""
|
"""
|
||||||
Configures the security settings for the OPC UA client.
|
Configures the security settings for the OPC UA client.
|
||||||
This method sets up the security policy, certificates, and timeouts
|
This method sets up the security policy, certificates, and timeouts
|
||||||
@@ -92,7 +93,7 @@ class OpcRepository():
|
|||||||
|
|
||||||
self.client.application_uri = self.server_uri
|
self.client.application_uri = self.server_uri
|
||||||
self.logger.custom_info('Setting security...', self.metadata)
|
self.logger.custom_info('Setting security...', self.metadata)
|
||||||
self.client.set_security(
|
await self.client.set_security(
|
||||||
SecurityPolicyBasic256,
|
SecurityPolicyBasic256,
|
||||||
certificate=str(cert),
|
certificate=str(cert),
|
||||||
private_key=str(private_key),
|
private_key=str(private_key),
|
||||||
@@ -101,7 +102,7 @@ class OpcRepository():
|
|||||||
self.client.secure_channel_timeout = 10000000
|
self.client.secure_channel_timeout = 10000000
|
||||||
self.client.session_timeout = 10000000
|
self.client.session_timeout = 10000000
|
||||||
|
|
||||||
def connect(self) -> tuple[bool, dict[str, Any]]:
|
async def connect(self) -> tuple[bool, dict[str, Any]]:
|
||||||
"""
|
"""
|
||||||
Establishes a connection to the OPC server.
|
Establishes a connection to the OPC server.
|
||||||
This method initializes the OPC client using the provided URL and
|
This method initializes the OPC client using the provided URL and
|
||||||
@@ -113,12 +114,12 @@ class OpcRepository():
|
|||||||
|
|
||||||
self.client = Client(self.url)
|
self.client = Client(self.url)
|
||||||
if self.cert_path:
|
if self.cert_path:
|
||||||
self.set_security()
|
await self.set_security()
|
||||||
self.logger.custom_info(
|
self.logger.custom_info(
|
||||||
f'Starting connection to OPC server {self.id}...', self.metadata)
|
f'Starting connection to OPC server {self.id}...', self.metadata)
|
||||||
return self.try_connect()
|
return await self.try_connect()
|
||||||
|
|
||||||
def try_connect(self) -> tuple[bool, dict[str, Any]]:
|
async def try_connect(self) -> tuple[bool, dict[str, Any]]:
|
||||||
"""
|
"""
|
||||||
Tries to connect to the OPC server.
|
Tries to connect to the OPC server.
|
||||||
|
|
||||||
@@ -128,7 +129,7 @@ class OpcRepository():
|
|||||||
|
|
||||||
try:
|
try:
|
||||||
self.last_reconnection_time = datetime.now()
|
self.last_reconnection_time = datetime.now()
|
||||||
self.client.connect()
|
await self.client.connect()
|
||||||
return True, {}
|
return True, {}
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
trace = traceback.format_exc()
|
trace = traceback.format_exc()
|
||||||
@@ -142,14 +143,14 @@ class OpcRepository():
|
|||||||
"attachment_content": trace
|
"attachment_content": trace
|
||||||
}
|
}
|
||||||
|
|
||||||
def disconnect(self):
|
async def disconnect(self):
|
||||||
"""
|
"""
|
||||||
Disconnects from the OPC server.
|
Disconnects from the OPC server.
|
||||||
"""
|
"""
|
||||||
if self.client is None:
|
if self.client is None:
|
||||||
return
|
return
|
||||||
try:
|
try:
|
||||||
self.client.disconnect()
|
await self.client.disconnect()
|
||||||
self.logger.custom_info(
|
self.logger.custom_info(
|
||||||
'Disconnected from OPC server', self.metadata)
|
'Disconnected from OPC server', self.metadata)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
@@ -162,12 +163,12 @@ class OpcRepository():
|
|||||||
Disconnects from the OPC server when the object is destroyed.
|
Disconnects from the OPC server when the object is destroyed.
|
||||||
"""
|
"""
|
||||||
try:
|
try:
|
||||||
self.disconnect()
|
asyncio.run(self.disconnect())
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
self.logger.custom_error(
|
self.logger.custom_error(
|
||||||
f"Error in destructor: {e}", self.metadata)
|
f"Error in destructor: {e}", self.metadata)
|
||||||
|
|
||||||
def validate_connection(self) -> tuple[bool, dict[str, Any]]:
|
async def validate_connection(self) -> tuple[bool, dict[str, Any]]:
|
||||||
"""
|
"""
|
||||||
Validates the connection to the OPC server.
|
Validates the connection to the OPC server.
|
||||||
If the connection is not established, it attempts to reconnect.
|
If the connection is not established, it attempts to reconnect.
|
||||||
@@ -179,13 +180,13 @@ class OpcRepository():
|
|||||||
If the client is connected, it returns True.
|
If the client is connected, it returns True.
|
||||||
"""
|
"""
|
||||||
if self.client is None:
|
if self.client is None:
|
||||||
return self.connect()
|
return await self.connect()
|
||||||
|
|
||||||
if self.error_count > 5:
|
if self.error_count > 5:
|
||||||
self.logger.custom_warning(
|
self.logger.custom_warning(
|
||||||
f"OPC server {self.id} will be disconnected due to multiple errors", self.metadata)
|
f"OPC server {self.id} will be disconnected due to multiple errors", self.metadata)
|
||||||
try:
|
try:
|
||||||
self.disconnect()
|
await self.disconnect()
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
trace = traceback.format_exc()
|
trace = traceback.format_exc()
|
||||||
self.logger.custom_error(
|
self.logger.custom_error(
|
||||||
@@ -193,21 +194,22 @@ class OpcRepository():
|
|||||||
self.logger.custom_error(trace, self.metadata)
|
self.logger.custom_error(trace, self.metadata)
|
||||||
self.logger.custom_info(
|
self.logger.custom_info(
|
||||||
f"Attempting to reconnect to OPC server {self.id}...", self.metadata)
|
f"Attempting to reconnect to OPC server {self.id}...", self.metadata)
|
||||||
return self.connect()
|
return await self.connect()
|
||||||
|
|
||||||
if hasattr(self.client, 'aio_obj') and self.client.aio_obj.uaclient.protocol is None or \
|
|
||||||
(hasattr(self.client.aio_obj.uaclient, 'protocol') and
|
|
||||||
self.client.aio_obj.uaclient.protocol.state == "closed"):
|
|
||||||
|
|
||||||
|
# Check if client is connected using asyncua's connection state
|
||||||
|
try:
|
||||||
|
# Try to get a simple node to test connection
|
||||||
|
await self.client.get_node("ns=0;i=2253") # Server node
|
||||||
|
except Exception:
|
||||||
self.logger.custom_error(
|
self.logger.custom_error(
|
||||||
f"OPC server {self.id} is not connected", self.metadata)
|
f"OPC server {self.id} is not connected", self.metadata)
|
||||||
|
|
||||||
if self.last_reconnection_time is None or (datetime.now() - self.last_reconnection_time).total_seconds(
|
if self.last_reconnection_time is None or (datetime.now() - self.last_reconnection_time).total_seconds(
|
||||||
) > self.reconnection_interval:
|
) > self.reconnection_interval:
|
||||||
self.disconnect()
|
await self.disconnect()
|
||||||
self.logger.custom_info(
|
self.logger.custom_info(
|
||||||
f"Trying to reconnect to OPC server {self.id}...", self.metadata)
|
f"Trying to reconnect to OPC server {self.id}...", self.metadata)
|
||||||
return self.connect()
|
return await self.connect()
|
||||||
|
|
||||||
return False, {
|
return False, {
|
||||||
"notification_id": f"OPC_CONNECTION_AWAITING_RECONNECTION_WINDOW_{self.id}",
|
"notification_id": f"OPC_CONNECTION_AWAITING_RECONNECTION_WINDOW_{self.id}",
|
||||||
@@ -218,7 +220,7 @@ class OpcRepository():
|
|||||||
|
|
||||||
return True, {}
|
return True, {}
|
||||||
|
|
||||||
def write_data(self, node: str, value: Any, data_type: str,
|
async def write_data(self, node: str, value: Any, data_type: str,
|
||||||
logger: Logger, metadata: dict[str, Any]) -> tuple[bool, dict[str, Any]]:
|
logger: Logger, metadata: dict[str, Any]) -> tuple[bool, dict[str, Any]]:
|
||||||
"""
|
"""
|
||||||
Writes data to the OPC server.
|
Writes data to the OPC server.
|
||||||
@@ -231,7 +233,7 @@ class OpcRepository():
|
|||||||
If the client is connected, it returns True.
|
If the client is connected, it returns True.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
is_connected, error = self.validate_connection()
|
is_connected, error = await self.validate_connection()
|
||||||
|
|
||||||
if not is_connected:
|
if not is_connected:
|
||||||
return False, error
|
return False, error
|
||||||
@@ -239,7 +241,7 @@ class OpcRepository():
|
|||||||
start_time = time.time()
|
start_time = time.time()
|
||||||
|
|
||||||
try:
|
try:
|
||||||
node = self.client.get_node(node)
|
node_obj = self.client.get_node(node)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
trace = traceback.format_exc()
|
trace = traceback.format_exc()
|
||||||
logger.custom_error(trace, metadata.get('schedule_name', 'N/A'))
|
logger.custom_error(trace, metadata.get('schedule_name', 'N/A'))
|
||||||
@@ -278,7 +280,7 @@ class OpcRepository():
|
|||||||
)
|
)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
node.write_value(ua_data)
|
await node_obj.write_value(ua_data)
|
||||||
|
|
||||||
metrics.PREDICTION_OPC_WRITING_COUNT.labels(
|
metrics.PREDICTION_OPC_WRITING_COUNT.labels(
|
||||||
pod_id=self.pod_id,
|
pod_id=self.pod_id,
|
||||||
|
|||||||
@@ -64,6 +64,9 @@ async def main():
|
|||||||
notification_handler=notification_handler
|
notification_handler=notification_handler
|
||||||
)
|
)
|
||||||
|
|
||||||
|
logger.custom_info('Initializing OPC...', metadata)
|
||||||
|
await activities.init_opc()
|
||||||
|
|
||||||
logger.custom_info(
|
logger.custom_info(
|
||||||
f'Starting SDK Metrics Server on port {SDK_METRICS_PORT}...', metadata)
|
f'Starting SDK Metrics Server on port {SDK_METRICS_PORT}...', metadata)
|
||||||
|
|
||||||
@@ -74,7 +77,7 @@ async def main():
|
|||||||
)
|
)
|
||||||
)
|
)
|
||||||
|
|
||||||
logger.custom_info('Starting Temporal Client...', metadata)
|
logger.custom_info(f'Starting Temporal Client at {host}...', metadata)
|
||||||
|
|
||||||
temporal_client = await client.Client.connect(
|
temporal_client = await client.Client.connect(
|
||||||
target_host=host,
|
target_host=host,
|
||||||
@@ -151,7 +154,7 @@ async def main():
|
|||||||
if notification_handler:
|
if notification_handler:
|
||||||
notification_handler.shutdown()
|
notification_handler.shutdown()
|
||||||
if activities:
|
if activities:
|
||||||
activities.shutdown()
|
await activities.shutdown()
|
||||||
# Exit with a non-zero status code to indicate failure to Kubernetes
|
# Exit with a non-zero status code to indicate failure to Kubernetes
|
||||||
metrics.APP_UP.labels(pod_id=POD_ID).set(0) # Mark app as DOWN
|
metrics.APP_UP.labels(pod_id=POD_ID).set(0) # Mark app as DOWN
|
||||||
sys.exit(1)
|
sys.exit(1)
|
||||||
|
|||||||
18
run_local.sh
Executable file
18
run_local.sh
Executable file
@@ -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 laborious.worker.worker
|
||||||
Reference in New Issue
Block a user