From 1bddde17f4f6d916adea622f9e74716728d92ce5 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Fri, 9 Jan 2026 09:00:21 -0300 Subject: [PATCH] SIENTIAPDE-1478 Refactor Activities and API Integration for PI Web API - Reintroduced the API import in the Activities class for proper integration. - Cleaned up whitespace and formatting in the API class and related tests for improved readability. - Updated test cases to ensure consistent formatting in error messages and configuration structures for PI Web API. - Enhanced connectors_config.py with additional whitespace for better organization. --- laborious/activities/activities.py | 4 +-- laborious/activities/api.py | 10 +++---- laborious/utils/connectors_config.py | 2 ++ laborious/worker/worker.py | 13 ++++----- .../format_and_export_prediction.py | 3 +-- tests/laborious/activities/test_api.py | 27 +++++++++---------- .../test_format_and_export_prediction.py | 24 +++++++++++++---- .../subworkflows/test_prediction_process.py | 18 ++++++++++--- 8 files changed, 63 insertions(+), 38 deletions(-) diff --git a/laborious/activities/activities.py b/laborious/activities/activities.py index c279756..bd8553d 100644 --- a/laborious/activities/activities.py +++ b/laborious/activities/activities.py @@ -7,12 +7,12 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.observability.logger import Logger + from laborious.activities.api import API from laborious.activities.gates import Gates from laborious.activities.mlflow import MLFlow from laborious.activities.model_metrics import ModelMetrics from laborious.activities.opc import OPC from laborious.activities.storage import Storage - from laborious.activities.api import API class Activities(Storage, MLFlow, Gates, OPC, ModelMetrics, API): @@ -147,4 +147,4 @@ class Activities(Storage, MLFlow, Gates, OPC, ModelMetrics, API): Gates.close(self) await OPC.close(self) ModelMetrics.close(self) - API.close(self) \ No newline at end of file + API.close(self) diff --git a/laborious/activities/api.py b/laborious/activities/api.py index df186d5..50be4f0 100644 --- a/laborious/activities/api.py +++ b/laborious/activities/api.py @@ -5,18 +5,17 @@ with workflow.unsafe.imports_passed_through(): from typing import Any from pandas import DataFrame - from sientia_do.notifications.handlers import NotificationHandler from sientia_do.notifications.models import NotificationLevel from sientia_do.observability.logger import Logger from sientia_do.observability.metrics_controller import MetricsController from sientia_do.observability.sientia_monitoring import SientiaMonitoring - from sientia_do.temporal.constants import DATETIME_FORMAT_WITH_TZ from sientia_do.repository.pi_web_api_client import PIWebAPIClient PI_WEB_API_PREDICTION_ERROR_CONFIDENCE = 13 + class API(SientiaMonitoring): """ PI Web API operations for writing data to PI Web API. @@ -78,7 +77,7 @@ class API(SientiaMonitoring): metadata = input_data['metadata'] data = DataFrame(input_data['data']) pi_web_api_output_config = input_data['pi_web_api_output_config'] - + self.info('Writing data to PI Web API...', metadata) endpoint = pi_web_api_output_config['endpoint'] @@ -92,7 +91,6 @@ class API(SientiaMonitoring): confidence_value = data.head(1)['prediction_confidence'].values[0] try: - await self.pi_web_api_client.write_value( web_ids=prediction_tags, value={ @@ -138,5 +136,5 @@ class API(SientiaMonitoring): level=NotificationLevel.ERROR, attachment_content=trace, ) - - return data.to_dict() \ No newline at end of file + + return data.to_dict() diff --git a/laborious/utils/connectors_config.py b/laborious/utils/connectors_config.py index 4476c09..c151377 100644 --- a/laborious/utils/connectors_config.py +++ b/laborious/utils/connectors_config.py @@ -2,6 +2,7 @@ import json from os import getenv from typing import Any + def build_mlflow_config() -> dict[str, Any]: """ Build MLFlow server configuration from environment variables. @@ -66,6 +67,7 @@ def build_opc_config() -> dict[str, Any]: } } + def build_minio_config() -> dict[str, Any]: """ Build MinIO (S3-compatible) configuration from environment variables. diff --git a/laborious/worker/worker.py b/laborious/worker/worker.py index db85a78..60990cb 100644 --- a/laborious/worker/worker.py +++ b/laborious/worker/worker.py @@ -43,7 +43,6 @@ Poller Configuration: - POLLER_INITIAL: Initial number of pollers (default: 2) """ - from temporalio import client, workflow from temporalio.runtime import PrometheusConfig, Runtime, TelemetryConfig from temporalio.worker import ( @@ -55,10 +54,16 @@ from temporalio.worker import ( with workflow.unsafe.imports_passed_through(): import asyncio + import os import sys from datetime import timedelta from prometheus_client import start_http_server + from sientia_do.connectors_config import ( + build_api_config, + build_mongodb_config, + build_postgres_config, + ) from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.observability.logger import get_logger @@ -69,10 +74,6 @@ with workflow.unsafe.imports_passed_through(): build_mlflow_config, build_opc_config, ) - from sientia_do.connectors_config import ( - build_postgres_config, - build_mongodb_config, - ) from laborious.workflows.drift import Drift from laborious.workflows.minimal_retrain import MinimalRetrain from laborious.workflows.predictions_batch import PredictionsBatch @@ -81,7 +82,6 @@ with workflow.unsafe.imports_passed_through(): FormatAndExportPrediction, ) from laborious.workflows.sub_workflows.prediction_process import PredictionProcess - import os POD_ID = os.getenv('POD_ID') SDK_METRICS_PORT = int(os.getenv('HTTP_SDK_METRICS_PORT', '9091')) @@ -183,6 +183,7 @@ async def main(): mlflow_config=build_mlflow_config(), minio_config=build_minio_config(), opc_config=build_opc_config(), + pi_web_api_config=build_api_config(), logger=logger, notification_handler=notification_handler, ) diff --git a/laborious/workflows/sub_workflows/format_and_export_prediction.py b/laborious/workflows/sub_workflows/format_and_export_prediction.py index 1aca3c3..bde9ff3 100644 --- a/laborious/workflows/sub_workflows/format_and_export_prediction.py +++ b/laborious/workflows/sub_workflows/format_and_export_prediction.py @@ -154,7 +154,7 @@ class FormatAndExportPrediction: ) write_transformed_handler = None - + opc_metrics = {} # write to pi web api @@ -197,7 +197,6 @@ class FormatAndExportPrediction: start_to_close_timeout=timedelta(seconds=180), ) - if write_transformed_handler is not None: await write_transformed_handler diff --git a/tests/laborious/activities/test_api.py b/tests/laborious/activities/test_api.py index 38339ea..eaf23f2 100644 --- a/tests/laborious/activities/test_api.py +++ b/tests/laborious/activities/test_api.py @@ -1,7 +1,6 @@ from unittest.mock import ANY, AsyncMock, MagicMock, call, patch import pytest_asyncio -from pandas import DataFrame from pytest import fixture, mark from sientia_do.notifications.models import NotificationLevel @@ -21,7 +20,7 @@ def _create_mock_dataframe(to_dict_return=None): """Helper function to create a mocked DataFrame for testing.""" mock_df = MagicMock() mock_head = MagicMock() - + def get_column_values(key): if key == 'prediction': return MagicMock(values=[0.75]) @@ -29,10 +28,10 @@ def _create_mock_dataframe(to_dict_return=None): return MagicMock(values=[0.95]) else: return MagicMock(values=['2024-01-01T00:00:00+00:00']) - + mock_head.__getitem__.side_effect = get_column_values mock_df.head.return_value = mock_head - + if to_dict_return is None: to_dict_return = { 'prediction': [0.75], @@ -40,7 +39,7 @@ def _create_mock_dataframe(to_dict_return=None): 'timestamp': ['2024-01-01T00:00:00+00:00'], } mock_df.to_dict.return_value = to_dict_return - + return mock_df @@ -147,11 +146,13 @@ async def test_write_pi_web_api_data_success(mock_dataframe, api, base_input_dat @mark.asyncio @patch('laborious.activities.api.DataFrame') async def test_write_pi_web_api_data_prediction_error(mock_dataframe, api, base_input_data): - mock_dataframe.return_value = _create_mock_dataframe({ - 'prediction': [0.75], - 'prediction_confidence': [PI_WEB_API_PREDICTION_ERROR_CONFIDENCE], - 'timestamp': ['2024-01-01T00:00:00+00:00'], - }) + mock_dataframe.return_value = _create_mock_dataframe( + { + 'prediction': [0.75], + 'prediction_confidence': [PI_WEB_API_PREDICTION_ERROR_CONFIDENCE], + 'timestamp': ['2024-01-01T00:00:00+00:00'], + } + ) api.pi_web_api_client.write_value.side_effect = Exception('Prediction write failed') @@ -160,7 +161,7 @@ async def test_write_pi_web_api_data_prediction_error(mock_dataframe, api, base_ api.send_notification_async.assert_called_once_with( metadata=metadata['metadata'], notification_id='WRITE_PI_WEB_API_PREDICTION_ERROR', - message='Error writing prediction data to PI Web API: Prediction write failed\n Tags: {\'tag1\': \'web_id_1\'}', + message="Error writing prediction data to PI Web API: Prediction write failed\n Tags: {'tag1': 'web_id_1'}", block='write_pi_web_api_data', level=NotificationLevel.ERROR, attachment_content=ANY, @@ -185,7 +186,7 @@ async def test_write_pi_web_api_data_confidence_error(mock_dataframe, api, base_ api.send_notification_async.assert_called_once_with( metadata=metadata['metadata'], notification_id='WRITE_PI_WEB_API_CONFIDENCE_ERROR', - message='Error writing confidence data to PI Web API: Confidence write failed\n Tags: {\'tag2\': \'web_id_2\'}', + message="Error writing confidence data to PI Web API: Confidence write failed\n Tags: {'tag2': 'web_id_2'}", block='write_pi_web_api_data', level=NotificationLevel.ERROR, attachment_content=ANY, @@ -199,8 +200,6 @@ async def test_write_pi_web_api_data_confidence_error(mock_dataframe, api, base_ assert api.pi_web_api_client.write_value.call_count == 2 - - @mark.asyncio @patch('laborious.activities.api.DataFrame') async def test_write_pi_web_api_data_empty_tags(mock_dataframe, api, base_input_data): diff --git a/tests/laborious/workflows/subworkflows/test_format_and_export_prediction.py b/tests/laborious/workflows/subworkflows/test_format_and_export_prediction.py index ffac5c6..dd95ee9 100644 --- a/tests/laborious/workflows/subworkflows/test_format_and_export_prediction.py +++ b/tests/laborious/workflows/subworkflows/test_format_and_export_prediction.py @@ -392,7 +392,11 @@ async def test_run_none_path_flag_with_pi_web_api(workflow_mock, format_and_expo 'prediction_confidence': 0, 'schema': 'test_schema', 'table_name': 'test_table', - 'pi_web_api_output_config': {'endpoint': 'https://test-pi-server.com', 'prediction_tags': {}, 'confidence_tags': {}}, + 'pi_web_api_output_config': { + 'endpoint': 'https://test-pi-server.com', + 'prediction_tags': {}, + 'confidence_tags': {}, + }, 'prediction_store_policy': 'erl:1', } @@ -483,7 +487,9 @@ async def test_run_none_path_flag_with_pi_web_api(workflow_mock, format_and_expo 'laborious.workflows.sub_workflows.format_and_export_prediction.workflow', new_callable=AsyncMock, ) -async def test_run_none_path_flag_with_pi_web_api_and_opc(workflow_mock, format_and_export_prediction): +async def test_run_none_path_flag_with_pi_web_api_and_opc( + workflow_mock, format_and_export_prediction +): input_data = { 'metadata': metadata, 'path_flag': None, @@ -494,7 +500,11 @@ async def test_run_none_path_flag_with_pi_web_api_and_opc(workflow_mock, format_ 'schema': 'test_schema', 'table_name': 'test_table', 'opc_output_config': {'test': 'config'}, - 'pi_web_api_output_config': {'endpoint': 'https://test-pi-server.com', 'prediction_tags': {}, 'confidence_tags': {}}, + 'pi_web_api_output_config': { + 'endpoint': 'https://test-pi-server.com', + 'prediction_tags': {}, + 'confidence_tags': {}, + }, 'prediction_store_policy': 'erl:1', } @@ -532,7 +542,7 @@ async def test_run_none_path_flag_with_pi_web_api_and_opc(workflow_mock, format_ }, retry_policy=ANY, start_to_close_timeout=ANY, - ) + ), ] ) @@ -590,7 +600,11 @@ async def test_run_default_path_flag_with_pi_web_api(workflow_mock, format_and_e 'prediction_confidence': 0, 'schema': 'test_schema', 'table_name': 'test_table', - 'pi_web_api_output_config': {'endpoint': 'https://test-pi-server.com', 'prediction_tags': {}, 'confidence_tags': {}}, + 'pi_web_api_output_config': { + 'endpoint': 'https://test-pi-server.com', + 'prediction_tags': {}, + 'confidence_tags': {}, + }, 'comment': 'test_comment', } diff --git a/tests/laborious/workflows/subworkflows/test_prediction_process.py b/tests/laborious/workflows/subworkflows/test_prediction_process.py index a31ab2a..22df274 100644 --- a/tests/laborious/workflows/subworkflows/test_prediction_process.py +++ b/tests/laborious/workflows/subworkflows/test_prediction_process.py @@ -731,7 +731,11 @@ async def test_path_flag_handler_continue(workflow_mock, prediction_process): 'model_name': model_name, 'model_config': model_config, 'opc_output_config': {'test': 'config'}, - 'pi_web_api_output_config': {'endpoint': 'https://test-pi-server.com', 'prediction_tags': {}, 'confidence_tags': {}}, + 'pi_web_api_output_config': { + 'endpoint': 'https://test-pi-server.com', + 'prediction_tags': {}, + 'confidence_tags': {}, + }, 'prediction_store_policy': prediction_store_policy, }, confidence, @@ -758,7 +762,11 @@ async def test_path_flag_handler_continue(workflow_mock, prediction_process): 'transform_table_name': 'test_transform_table', 'comment': 'Prediction Process', 'opc_output_config': {'test': 'config'}, - 'pi_web_api_output_config': {'endpoint': 'https://test-pi-server.com', 'prediction_tags': {}, 'confidence_tags': {}}, + 'pi_web_api_output_config': { + 'endpoint': 'https://test-pi-server.com', + 'prediction_tags': {}, + 'confidence_tags': {}, + }, 'prediction_store_policy': prediction_store_policy, }, ) @@ -792,7 +800,11 @@ async def test_path_flag_handler_unknown(workflow_mock, prediction_process): 'model_name': model_name, 'model_config': model_config, 'opc_output_config': {'test': 'config'}, - 'pi_web_api_output_config': {'endpoint': 'https://test-pi-server.com', 'prediction_tags': {}, 'confidence_tags': {}}, + 'pi_web_api_output_config': { + 'endpoint': 'https://test-pi-server.com', + 'prediction_tags': {}, + 'confidence_tags': {}, + }, 'prediction_store_policy': prediction_store_policy, }, confidence,