Enhance logging in OPC activity by adding debug statements for output configuration and server details to improve traceability.
160 lines
6.4 KiB
Python
160 lines
6.4 KiB
Python
from temporalio import activity, workflow
|
|
|
|
|
|
with workflow.unsafe.imports_passed_through():
|
|
from sientia_do.notifications.handlers import NotificationHandler
|
|
from sientia_do.notifications.models import NotificationLevel
|
|
from sientia_do.temporal.activities.base import BaseActivity
|
|
from sientia_do.temporal.utils.logger import Logger
|
|
from laborious.utils.repository.opc_repository import OpcRepository
|
|
from typing import Any
|
|
import traceback
|
|
from pandas import DataFrame
|
|
|
|
OPC_WRITTING_ERROR_CONFIDENCE = 12
|
|
|
|
|
|
class OPC(BaseActivity):
|
|
def __init__(self, opc_servers: dict[str, dict[str, Any]],
|
|
logger: Logger, notification_handler: NotificationHandler):
|
|
|
|
self.logger = logger
|
|
self.notification_handler = notification_handler
|
|
self.opc_servers = opc_servers
|
|
|
|
self.opc_repository = {}
|
|
self.logger.debug(f"OPC servers: {opc_servers}")
|
|
for id, server in opc_servers.items():
|
|
self.opc_repository[id] = OpcRepository(
|
|
id=server['id'],
|
|
url=server['url'],
|
|
logger=self.logger,
|
|
server_uri=server['server_uri'],
|
|
cert_path=server['cert_path'],
|
|
private_key_path=server['private_key_path'],
|
|
server_cert_path=server['server_cert_path'],
|
|
notification_handler=self.notification_handler,
|
|
reconnection_interval=server['reconnection_interval'],
|
|
)
|
|
self.opc_repository[id].connect()
|
|
|
|
BaseActivity.__init__(self, logger, notification_handler)
|
|
|
|
def write_data(self, server_id: str, tag: str, data: Any,
|
|
data_type: str, tag_type: str) -> bool:
|
|
"""
|
|
Write data to OPC server.
|
|
|
|
Args:
|
|
- server_id (str): The id of the OPC server.
|
|
- tag (str): The tag to write to.
|
|
- data (Any): The data to write.
|
|
- data_type (str): The data type.
|
|
- tag_type (str): The tag type.
|
|
|
|
Returns:
|
|
- bool: True if the data was written successfully, False otherwise.
|
|
"""
|
|
|
|
try:
|
|
return self.opc_repository[server_id].write_data(
|
|
tag, data, data_type)
|
|
self.logger.debug(f"Wrote {tag_type} to {tag}")
|
|
except Exception as e:
|
|
trace = traceback.format_exc()
|
|
self.notification_handler.build_and_send_notification(
|
|
notification_id=f"WRITE_OPC_{tag_type.upper()}_ERROR",
|
|
message=f"Error writing data to OPC server: {e}",
|
|
block="write_opc_data",
|
|
level=NotificationLevel.ERROR,
|
|
attachment_content=trace
|
|
)
|
|
raise e
|
|
|
|
@activity.defn(name='write_opc_data')
|
|
async def write_opc_data(self, input_data: dict[str, Any]) -> dict[Any, Any]:
|
|
"""
|
|
Write prediction and confidence data to OPC servers. The two writing
|
|
operations are optional and independent of each other.
|
|
|
|
Args:
|
|
- input_data(dict[str, Any]): The input data. Contains the following keys:
|
|
- data(dict[str, Any]): The dataframe that contains the data to write
|
|
to the OPC servers.
|
|
- opc_output_config(dict[str, Any]): The OPC writing configuration.
|
|
The keys are the OPC server names and the values contain:
|
|
- prediction_tags(dict[str, Any]): The tags to write to the OPC servers.
|
|
- confidence_tags(dict[str, Any]): The tags to write to the OPC servers.
|
|
|
|
Returns:
|
|
- dict[Any, Any]: The data that was written to the OPC servers.
|
|
|
|
"""
|
|
metadata = input_data['metadata']
|
|
self.debug("Writing data to OPC servers...", metadata)
|
|
data = DataFrame(input_data['data'])
|
|
opc_output_config = input_data['opc_output_config']
|
|
self.debug(data, metadata)
|
|
|
|
success = True
|
|
|
|
self.logger.debug(f"OPC output config: {opc_output_config}")
|
|
self.logger.debug(f"OPC servers: {self.opc_repository}")
|
|
for server_id, config in opc_output_config.items():
|
|
if self.opc_repository.get(server_id) is None:
|
|
self.error(f"OPC server {server_id} not found", metadata)
|
|
continue
|
|
|
|
if 'prediction_tags' in config:
|
|
for tag, tag_config in config['prediction_tags'].items():
|
|
success = success and self.write_data(
|
|
server_id=server_id,
|
|
tag=tag,
|
|
data=data.head(1)['prediction'].values[0],
|
|
data_type=tag_config['data_type'],
|
|
tag_type='prediction'
|
|
)
|
|
|
|
if 'confidence_tags' in config:
|
|
for tag, tag_config in config['confidence_tags'].items():
|
|
self.write_data(
|
|
server_id=server_id,
|
|
tag=tag,
|
|
data=data.head(1)['prediction_confidence'].values[0],
|
|
data_type=tag_config['data_type'],
|
|
tag_type='confidence'
|
|
)
|
|
|
|
return self.process_confidence(data, success, metadata)
|
|
|
|
def process_confidence(self, data: DataFrame, success: bool, metadata: dict[str, Any]) -> dict[Any, Any]:
|
|
"""
|
|
Processes the confidence of OPC server write operations and updates the DataFrame accordingly.
|
|
|
|
If the write operation was not successful, sets the 'prediction_confidence' column in the DataFrame
|
|
to a predefined error confidence value and logs a debug message. Otherwise, logs a success message.
|
|
|
|
Args:
|
|
data (DataFrame): The DataFrame containing the data to be processed.
|
|
success (bool): Indicates whether the data was successfully written to the OPC servers.
|
|
|
|
Returns:
|
|
dict[Any, Any]: The processed data as a dictionary.
|
|
"""
|
|
|
|
if not success:
|
|
data['prediction_confidence'] = OPC_WRITTING_ERROR_CONFIDENCE
|
|
self.debug(
|
|
"Some data could not be written to OPC servers, setting confidence to "
|
|
f"{OPC_WRITTING_ERROR_CONFIDENCE}."
|
|
)
|
|
|
|
else:
|
|
self.debug("Data written to OPC servers successfully.", metadata)
|
|
|
|
return data.to_dict()
|
|
|
|
def shutdown(self):
|
|
for opc in self.opc_repository.values():
|
|
opc.disconnect()
|