diff --git a/laborious/activities/storage.py b/laborious/activities/storage.py index 8aec756..5439420 100644 --- a/laborious/activities/storage.py +++ b/laborious/activities/storage.py @@ -10,14 +10,11 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.notifications.models import NotificationLevel from sientia_do.observability.logger import Logger from sientia_do.temporal.activities.postgres import Postgres - from sientia_do.temporal.constants import now + from sientia_do.temporal.constants import now, DATETIME_FORMAT_FILENAME from laborious.utils.repository.minio_repository import MinioRepository -DATETIME_FILENAME_FORMAT = '%Y-%m-%d_%H-%M-%S' - - class Storage(Postgres): """ Extensions for Postgres activities with a helper to export query results @@ -84,7 +81,7 @@ class Storage(Postgres): metadata = input_data.get('metadata', {}) object_prefix = input_data.get('object_prefix', 'datasets/retrain') - timestamp = now().strftime(DATETIME_FILENAME_FORMAT) + timestamp = now().strftime(DATETIME_FORMAT_FILENAME) object_name = f'{object_prefix}_{timestamp}.parquet' uri = f's3://{self.minio_repository.minio_bucket}/{object_name}' diff --git a/laborious/utils/repository/opc_repository.py b/laborious/utils/repository/opc_repository.py index be4edb5..619ff62 100644 --- a/laborious/utils/repository/opc_repository.py +++ b/laborious/utils/repository/opc_repository.py @@ -1,3 +1,5 @@ +import asyncio +import json import time import traceback from datetime import datetime @@ -10,6 +12,7 @@ from asyncua.ua import DataValue, DateTime, Variant, VariantType from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.notifications.models import NotificationLevel from sientia_do.observability.logger import Logger +from sientia_do.temporal.activities.base import BaseActivity from laborious import metrics @@ -37,19 +40,19 @@ data_type_map = { } -class OpcRepository: +class OpcRepository(BaseActivity): def __init__( self, opc_id: str, url: str, + pod_id: str, logger: Logger, notification_handler: NotificationHandler, reconnection_interval: int = 60, server_uri: str | None = None, cert_path: str | None = None, private_key_path: str | None = None, - server_cert_path: str | None = None, - pod_id: str | None = None, + server_cert_path: str | None = None ): self.url = url self.id = opc_id @@ -65,6 +68,8 @@ class OpcRepository: self.client: None | Client = None self.pod_id = pod_id + BaseActivity.__init__(self, logger, notification_handler, set_error_counter=True) + self.metadata = { 'model_name': '-', 'model_id': '-', @@ -125,7 +130,14 @@ class OpcRepository: Exception: If the connection to the OPC server fails. """ - self.client = Client(self.url) + self.client = Client(self.url) # type: ignore[attr-defined] + + self.client.name = self.pod_id + self.client.application_name = self.pod_id + pod_uri = self.pod_id.replace('-', ':') + self.client.application_uri = pod_uri + self.client.product_uri = pod_uri + if self.cert_path: await self.set_security() self.logger.custom_info(f'Starting connection to OPC server {self.id}...', self.metadata) @@ -169,6 +181,28 @@ class OpcRepository: 'attachment_content': trace, } + async def disconnection_fallback(self) -> list: + """ + Tries 5 times to disconnect from the OPC UA server, with a delay of 100ms x try. + """ + + assert self.client is not None + error_stack = [] + for i in range(5): + try: + self.logger.info(f'Disconnecting from OPC UA server, attempt {i+1} of 5') + await self.client.disconnect() + return [] + except Exception as e: + self.logger.error(f'Failed to disconnect from OPC UA server in attempt {i+1} of 5: {e}') + error_stack.append({ + 'attempt': i+1, + 'error': str(e), + 'traceback': traceback.format_exc(), + }) + await asyncio.sleep(0.1 * i) + return error_stack + async def disconnect(self): """ Gracefully disconnect from the OPC server. @@ -179,11 +213,20 @@ class OpcRepository: """ if self.client is None: return - try: - await self.client.disconnect() - self.logger.custom_info('Disconnected from OPC server', self.metadata) - except Exception as e: - self.logger.custom_error(f'Failed to disconnect from OPC server: {e}', self.metadata) + + errors = await self.disconnection_fallback() + if errors: + self.send_notification( + metadata=self.metadata, + notification_id=f'OPC_DISCONNECTION_ERROR_{self.id}', + message=f'Failed to disconnect from OPC server in 5 attempts: {errors}', + block='opc_repository', + level=NotificationLevel.ERROR, + attachment_content=json.dumps(errors, indent=4), + ) + else: + self.logger.warning(f'Disconnected from OPC server {self.id} successfully') + self.client = None async def validate_connection(self) -> tuple[bool, dict[str, Any]]: diff --git a/requirements-light.txt b/requirements-light.txt index c145097..278fb58 100644 --- a/requirements-light.txt +++ b/requirements-light.txt @@ -3,7 +3,7 @@ psycopg2-binary sqlalchemy asyncua redis -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.6 +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.7 prometheus-client botocore boto3 diff --git a/requirements.txt b/requirements.txt index 3ba54e8..7122a3e 100644 --- a/requirements.txt +++ b/requirements.txt @@ -3,7 +3,7 @@ psycopg2-binary sqlalchemy asyncua redis -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.6 +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.7 git+ssh://git@github.com/Aignosi/sientia-mlops-library.git@0.39.0 prometheus-client botocore diff --git a/values.yaml b/values.yaml index 1300546..f100a7e 100644 --- a/values.yaml +++ b/values.yaml @@ -151,7 +151,7 @@ env: - name: GITHUB_REPO_URL value: "git@github.com:Aignosi/sientia-dataops-laborious_temporal.git" - name: GITHUB_BRANCH - value: SIENTIAPDE-1231-ajustar-o-retreino-do-courier-no-laborious + value: "SIENTIAPDE-1312-melhorias-e-correcoes-nas-pipelines-de-dados" - name: PYTHON_APP value: "laborious.worker.worker"