From f046d079b57e9519922ad493249d47a22fb62857 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 20 Oct 2025 15:55:51 -0300 Subject: [PATCH 1/6] SIENTIAPDE-1312 Update dependencies and enhance OpcRepository functionality - Updated sientia-dataops-library version to 1.4.7 in requirements files. - Modified GITHUB_BRANCH in values.yaml for improved pipeline management. - Refactored OpcRepository class to inherit from BaseActivity, adding enhanced logging and error handling during disconnection. - Implemented a disconnection fallback mechanism to ensure graceful handling of OPC server disconnections. --- laborious/activities/storage.py | 7 +-- laborious/utils/repository/opc_repository.py | 61 +++++++++++++++++--- requirements-light.txt | 2 +- requirements.txt | 2 +- values.yaml | 2 +- 5 files changed, 57 insertions(+), 17 deletions(-) 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" From c2315a94553619a0860e7235b049e5994c4f29a5 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 20 Oct 2025 16:26:01 -0300 Subject: [PATCH 2/6] SIENTIAPDE-1312 Enhance OpcRepository client initialization and error handling - Updated the Client instantiation in OpcRepository to include a timeout and watchdog interval for improved connection management. - Added a disconnection call in the exception handling block to ensure proper resource cleanup during connection failures. --- laborious/utils/repository/opc_repository.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/laborious/utils/repository/opc_repository.py b/laborious/utils/repository/opc_repository.py index 619ff62..d53cf1b 100644 --- a/laborious/utils/repository/opc_repository.py +++ b/laborious/utils/repository/opc_repository.py @@ -130,7 +130,7 @@ class OpcRepository(BaseActivity): Exception: If the connection to the OPC server fails. """ - self.client = Client(self.url) # type: ignore[attr-defined] + self.client = Client(self.url, timeout=10, watchdog_intervall=3600000) # type: ignore[attr-defined] self.client.name = self.pod_id self.client.application_name = self.pod_id @@ -170,6 +170,8 @@ class OpcRepository(BaseActivity): await self.client.connect() return True, {} except Exception as e: + self.disconnect() + trace = traceback.format_exc() self.logger.custom_error(trace, self.metadata) From 461a885eeb85d3ec2f23d0759dddfd1012e5b021 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 20 Oct 2025 16:48:21 -0300 Subject: [PATCH 3/6] SIENTIAPDE-1312 Refactor OPC handling by removing pod_id from initialization and updating logging format - Removed pod_id parameter from OPC class and repository initialization to streamline connection management. - Updated logging statements for improved readability during disconnection attempts and error handling. --- laborious/activities/opc.py | 1 - laborious/activities/storage.py | 2 +- laborious/utils/repository/opc_repository.py | 26 +++++----- tests/laborious/activities/test_opc.py | 2 - .../utils/repository/test_opc_repository.py | 51 ++++++++++++++++--- 5 files changed, 60 insertions(+), 22 deletions(-) diff --git a/laborious/activities/opc.py b/laborious/activities/opc.py index b4604d8..1951973 100644 --- a/laborious/activities/opc.py +++ b/laborious/activities/opc.py @@ -85,7 +85,6 @@ class OPC(BaseActivity): server_cert_path=server['server_cert_path'], notification_handler=self.notification_handler, reconnection_interval=server['reconnection_interval'], - pod_id=self.pod_id, ) is_connected, error_data = await self.opc_repository[opc_id].connect() if not is_connected: diff --git a/laborious/activities/storage.py b/laborious/activities/storage.py index 5439420..c04b6a4 100644 --- a/laborious/activities/storage.py +++ b/laborious/activities/storage.py @@ -10,7 +10,7 @@ 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, DATETIME_FORMAT_FILENAME + from sientia_do.temporal.constants import DATETIME_FORMAT_FILENAME, now from laborious.utils.repository.minio_repository import MinioRepository diff --git a/laborious/utils/repository/opc_repository.py b/laborious/utils/repository/opc_repository.py index d53cf1b..8ab8563 100644 --- a/laborious/utils/repository/opc_repository.py +++ b/laborious/utils/repository/opc_repository.py @@ -45,14 +45,13 @@ class OpcRepository(BaseActivity): 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 + server_cert_path: str | None = None, ): self.url = url self.id = opc_id @@ -66,7 +65,6 @@ class OpcRepository(BaseActivity): self.last_reconnection_time: None | datetime = None self.notification_handler = notification_handler self.client: None | Client = None - self.pod_id = pod_id BaseActivity.__init__(self, logger, notification_handler, set_error_counter=True) @@ -192,16 +190,20 @@ class OpcRepository(BaseActivity): error_stack = [] for i in range(5): try: - self.logger.info(f'Disconnecting from OPC UA server, attempt {i+1} of 5') + 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(), - }) + 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 @@ -221,13 +223,13 @@ class OpcRepository(BaseActivity): 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}', + message='Failed to disconnect from OPC server in 5 attempts.', 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.logger.warning(f'Disconnected from OPC server {self.id} successfully') self.client = None diff --git a/tests/laborious/activities/test_opc.py b/tests/laborious/activities/test_opc.py index 815df27..a9bbe67 100644 --- a/tests/laborious/activities/test_opc.py +++ b/tests/laborious/activities/test_opc.py @@ -105,7 +105,6 @@ async def test_init_opc(mock_send_notification, mock_opc_repository): server_cert_path='', notification_handler=mock_notification_handler, reconnection_interval=60, - pod_id='localhost', ), ] ) @@ -121,7 +120,6 @@ async def test_init_opc(mock_send_notification, mock_opc_repository): server_cert_path='', notification_handler=mock_notification_handler, reconnection_interval=60, - pod_id='localhost', ) ] ) diff --git a/tests/laborious/utils/repository/test_opc_repository.py b/tests/laborious/utils/repository/test_opc_repository.py index 6ec680d..2a47f13 100644 --- a/tests/laborious/utils/repository/test_opc_repository.py +++ b/tests/laborious/utils/repository/test_opc_repository.py @@ -1,3 +1,4 @@ +import json from datetime import datetime from unittest.mock import ANY, AsyncMock, MagicMock, Mock, call, patch @@ -15,7 +16,7 @@ def mock_logger(): @pytest.fixture def opc_repository(mock_logger): - return OpcRepository( + repository = OpcRepository( opc_id='test_repo', url='opc.tcp://localhost:4840', logger=mock_logger, @@ -26,6 +27,8 @@ def opc_repository(mock_logger): private_key_path='/path/to/key.pem', server_cert_path='/path/to/server_cert.pem', ) + repository.send_notification = MagicMock() + return repository @pytest.fixture @@ -162,11 +165,37 @@ async def test_try_connect_no_client(opc_repository): @pytest.mark.asyncio -async def test_disconnect(opc_repository, mock_client): +async def test_disconnection_fallback_success(opc_repository, mock_client): opc_repository.client = mock_client - await opc_repository.disconnect() + mock_client.disconnect.return_value = True + result = await opc_repository.disconnection_fallback() mock_client.disconnect.assert_called_once() + assert result == [] + + +@pytest.mark.asyncio +async def test_disconnection_fallback_fail(opc_repository, mock_client): + opc_repository.client = mock_client + mock_client.disconnect.side_effect = Exception('Test error') + result = await opc_repository.disconnection_fallback() + assert result == [ + {'attempt': 1, 'error': 'Test error', 'traceback': ANY}, + {'attempt': 2, 'error': 'Test error', 'traceback': ANY}, + {'attempt': 3, 'error': 'Test error', 'traceback': ANY}, + {'attempt': 4, 'error': 'Test error', 'traceback': ANY}, + {'attempt': 5, 'error': 'Test error', 'traceback': ANY}, + ] + assert mock_client.disconnect.call_count == 5 + + +@pytest.mark.asyncio +async def test_disconnect(opc_repository, mock_client): + opc_repository.client = mock_client + opc_repository.disconnection_fallback = AsyncMock(return_value=[]) + await opc_repository.disconnect() + + opc_repository.disconnection_fallback.assert_called_once() assert opc_repository.client is None @@ -179,11 +208,21 @@ async def test_disconnect_no_client(opc_repository): @pytest.mark.asyncio async def test_disconnect_error(opc_repository, mock_client): opc_repository.client = mock_client - mock_client.disconnect.side_effect = Exception('Test error') + opc_repository.disconnection_fallback = AsyncMock( + return_value=[{'attempt': 1, 'error': 'Test error', 'traceback': 'text'}] + ) await opc_repository.disconnect() - opc_repository.logger.custom_error.assert_called_once_with( - 'Failed to disconnect from OPC server: Test error', ANY + opc_repository.disconnection_fallback.assert_called_once() + opc_repository.send_notification.assert_called_once_with( + metadata=opc_repository.metadata, + notification_id=f'OPC_DISCONNECTION_ERROR_{opc_repository.id}', + message='Failed to disconnect from OPC server in 5 attempts.', + block='opc_repository', + level=NotificationLevel.ERROR, + attachment_content=json.dumps( + [{'attempt': 1, 'error': 'Test error', 'traceback': 'text'}], indent=4 + ), ) assert opc_repository.client is None From 4c1e9008abcc9ec28a75bf66f6380fe6223bde87 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 21 Oct 2025 14:32:35 -0300 Subject: [PATCH 4/6] SIENTIAPDE-1312 Update tests notebook execution count and bump image tag in values.yaml - Incremented the execution count in the tests notebook for accurate tracking. - Updated the image tag in values.yaml from "0.0.3" to "1.0.1" for versioning consistency. - Enhanced error handling in model_repository.py to ensure proper retrieval of created experiments. --- laborious/utils/repository/model_repository.py | 5 ++++- values.yaml | 2 +- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/laborious/utils/repository/model_repository.py b/laborious/utils/repository/model_repository.py index 3361d9a..20ac9f3 100644 --- a/laborious/utils/repository/model_repository.py +++ b/laborious/utils/repository/model_repository.py @@ -175,7 +175,10 @@ class MLFlowRepository: if experiment is None: if create_if_not_exists: - experiment = mlflow.create_experiment(experiment_name) + experiment_id = mlflow.create_experiment(experiment_name) + experiment = mlflow.get_experiment(experiment_id) + if experiment is None: + raise ValueError(f'Experiment {experiment_name} not found after creation, unknown reason') else: raise ValueError(f'Experiment {experiment_name} not found') diff --git a/values.yaml b/values.yaml index f100a7e..e4b3165 100644 --- a/values.yaml +++ b/values.yaml @@ -11,7 +11,7 @@ image: # This sets the pull policy for images. pullPolicy: Always # Overrides the image tag whose default is the chart appVersion. - tag: "0.0.3" + tag: "1.0.1" 0# This is for the secrets for pulling an image from a private repository more information can be found here: https://kubernetes.io/docs/tasks/configure-pod-container/pull-image-private-registry/ imagePullSecrets: From 197a25ebf7855acad2b0de230c178f649e3abc54 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 21 Oct 2025 14:40:52 -0300 Subject: [PATCH 5/6] SIENTIAPDE-1312 Ensure model_temp_path directory exists before saving retrain data in MLFlowRepository --- laborious/utils/repository/model_repository.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/laborious/utils/repository/model_repository.py b/laborious/utils/repository/model_repository.py index 20ac9f3..1ab72c1 100644 --- a/laborious/utils/repository/model_repository.py +++ b/laborious/utils/repository/model_repository.py @@ -768,6 +768,8 @@ class MLFlowRepository: data_path = f'{model_temp_path}/retrain_data.csv' + makedirs(model_temp_path, exist_ok=True) + data.to_csv(data_path, index=True) self.logger.custom_info( From c5fb921654700635bf924aa175b43829f0d7e18c Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Tue, 21 Oct 2025 15:58:28 -0300 Subject: [PATCH 6/6] SIENTIAPDE-1312 Refactor error handling in MLFlowRepository and update tests - Improved error message formatting in MLFlowRepository for better readability. - Updated test assertions to ensure correct calls to MLflow methods during experiment retrieval and creation. --- laborious/utils/repository/model_repository.py | 4 +++- tests/laborious/utils/repository/test_model_repository.py | 5 ++++- 2 files changed, 7 insertions(+), 2 deletions(-) diff --git a/laborious/utils/repository/model_repository.py b/laborious/utils/repository/model_repository.py index 1ab72c1..7400502 100644 --- a/laborious/utils/repository/model_repository.py +++ b/laborious/utils/repository/model_repository.py @@ -178,7 +178,9 @@ class MLFlowRepository: experiment_id = mlflow.create_experiment(experiment_name) experiment = mlflow.get_experiment(experiment_id) if experiment is None: - raise ValueError(f'Experiment {experiment_name} not found after creation, unknown reason') + raise ValueError( + f'Experiment {experiment_name} not found after creation, unknown reason' + ) else: raise ValueError(f'Experiment {experiment_name} not found') diff --git a/tests/laborious/utils/repository/test_model_repository.py b/tests/laborious/utils/repository/test_model_repository.py index a437653..c21608d 100644 --- a/tests/laborious/utils/repository/test_model_repository.py +++ b/tests/laborious/utils/repository/test_model_repository.py @@ -157,10 +157,13 @@ def test_get_experiment_none_create(mlflow, mlflow_repository): mlflow.get_experiment_by_name.return_value = None - mlflow.create_experiment.return_value = experiment + mlflow.get_experiment.return_value = experiment output = mlflow_repository.get_experiment('test', create_if_not_exists=True) + mlflow.create_experiment.assert_called_once_with('test') + mlflow.get_experiment.assert_called_once_with(mlflow.create_experiment.return_value) + assert output == experiment