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.
This commit is contained in:
vitor-aignosi
2026-01-09 09:00:21 -03:00
parent 892823df11
commit 1bddde17f4
8 changed files with 63 additions and 38 deletions

View File

@@ -7,12 +7,12 @@ with workflow.unsafe.imports_passed_through():
from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler
from sientia_do.observability.logger import Logger from sientia_do.observability.logger import Logger
from laborious.activities.api import API
from laborious.activities.gates import Gates from laborious.activities.gates import Gates
from laborious.activities.mlflow import MLFlow from laborious.activities.mlflow import MLFlow
from laborious.activities.model_metrics import ModelMetrics from laborious.activities.model_metrics import ModelMetrics
from laborious.activities.opc import OPC from laborious.activities.opc import OPC
from laborious.activities.storage import Storage from laborious.activities.storage import Storage
from laborious.activities.api import API
class Activities(Storage, MLFlow, Gates, OPC, ModelMetrics, API): class Activities(Storage, MLFlow, Gates, OPC, ModelMetrics, API):
@@ -147,4 +147,4 @@ class Activities(Storage, MLFlow, Gates, OPC, ModelMetrics, API):
Gates.close(self) Gates.close(self)
await OPC.close(self) await OPC.close(self)
ModelMetrics.close(self) ModelMetrics.close(self)
API.close(self) API.close(self)

View File

@@ -5,18 +5,17 @@ with workflow.unsafe.imports_passed_through():
from typing import Any from typing import Any
from pandas import DataFrame from pandas import DataFrame
from sientia_do.notifications.handlers import NotificationHandler from sientia_do.notifications.handlers import NotificationHandler
from sientia_do.notifications.models import NotificationLevel from sientia_do.notifications.models import NotificationLevel
from sientia_do.observability.logger import Logger from sientia_do.observability.logger import Logger
from sientia_do.observability.metrics_controller import MetricsController from sientia_do.observability.metrics_controller import MetricsController
from sientia_do.observability.sientia_monitoring import SientiaMonitoring 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 from sientia_do.repository.pi_web_api_client import PIWebAPIClient
PI_WEB_API_PREDICTION_ERROR_CONFIDENCE = 13 PI_WEB_API_PREDICTION_ERROR_CONFIDENCE = 13
class API(SientiaMonitoring): class API(SientiaMonitoring):
""" """
PI Web API operations for writing data to PI Web API. PI Web API operations for writing data to PI Web API.
@@ -78,7 +77,7 @@ class API(SientiaMonitoring):
metadata = input_data['metadata'] metadata = input_data['metadata']
data = DataFrame(input_data['data']) data = DataFrame(input_data['data'])
pi_web_api_output_config = input_data['pi_web_api_output_config'] pi_web_api_output_config = input_data['pi_web_api_output_config']
self.info('Writing data to PI Web API...', metadata) self.info('Writing data to PI Web API...', metadata)
endpoint = pi_web_api_output_config['endpoint'] endpoint = pi_web_api_output_config['endpoint']
@@ -92,7 +91,6 @@ class API(SientiaMonitoring):
confidence_value = data.head(1)['prediction_confidence'].values[0] confidence_value = data.head(1)['prediction_confidence'].values[0]
try: try:
await self.pi_web_api_client.write_value( await self.pi_web_api_client.write_value(
web_ids=prediction_tags, web_ids=prediction_tags,
value={ value={
@@ -138,5 +136,5 @@ class API(SientiaMonitoring):
level=NotificationLevel.ERROR, level=NotificationLevel.ERROR,
attachment_content=trace, attachment_content=trace,
) )
return data.to_dict() return data.to_dict()

View File

@@ -2,6 +2,7 @@ import json
from os import getenv from os import getenv
from typing import Any from typing import Any
def build_mlflow_config() -> dict[str, Any]: def build_mlflow_config() -> dict[str, Any]:
""" """
Build MLFlow server configuration from environment variables. 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]: def build_minio_config() -> dict[str, Any]:
""" """
Build MinIO (S3-compatible) configuration from environment variables. Build MinIO (S3-compatible) configuration from environment variables.

View File

@@ -43,7 +43,6 @@ Poller Configuration:
- POLLER_INITIAL: Initial number of pollers (default: 2) - POLLER_INITIAL: Initial number of pollers (default: 2)
""" """
from temporalio import client, workflow from temporalio import client, workflow
from temporalio.runtime import PrometheusConfig, Runtime, TelemetryConfig from temporalio.runtime import PrometheusConfig, Runtime, TelemetryConfig
from temporalio.worker import ( from temporalio.worker import (
@@ -55,10 +54,16 @@ from temporalio.worker import (
with workflow.unsafe.imports_passed_through(): with workflow.unsafe.imports_passed_through():
import asyncio import asyncio
import os
import sys import sys
from datetime import timedelta from datetime import timedelta
from prometheus_client import start_http_server 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.notifications.handlers import CoreNotificationHandler as NotificationHandler
from sientia_do.observability.logger import get_logger from sientia_do.observability.logger import get_logger
@@ -69,10 +74,6 @@ with workflow.unsafe.imports_passed_through():
build_mlflow_config, build_mlflow_config,
build_opc_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.drift import Drift
from laborious.workflows.minimal_retrain import MinimalRetrain from laborious.workflows.minimal_retrain import MinimalRetrain
from laborious.workflows.predictions_batch import PredictionsBatch from laborious.workflows.predictions_batch import PredictionsBatch
@@ -81,7 +82,6 @@ with workflow.unsafe.imports_passed_through():
FormatAndExportPrediction, FormatAndExportPrediction,
) )
from laborious.workflows.sub_workflows.prediction_process import PredictionProcess from laborious.workflows.sub_workflows.prediction_process import PredictionProcess
import os
POD_ID = os.getenv('POD_ID') POD_ID = os.getenv('POD_ID')
SDK_METRICS_PORT = int(os.getenv('HTTP_SDK_METRICS_PORT', '9091')) SDK_METRICS_PORT = int(os.getenv('HTTP_SDK_METRICS_PORT', '9091'))
@@ -183,6 +183,7 @@ async def main():
mlflow_config=build_mlflow_config(), mlflow_config=build_mlflow_config(),
minio_config=build_minio_config(), minio_config=build_minio_config(),
opc_config=build_opc_config(), opc_config=build_opc_config(),
pi_web_api_config=build_api_config(),
logger=logger, logger=logger,
notification_handler=notification_handler, notification_handler=notification_handler,
) )

View File

@@ -154,7 +154,7 @@ class FormatAndExportPrediction:
) )
write_transformed_handler = None write_transformed_handler = None
opc_metrics = {} opc_metrics = {}
# write to pi web api # write to pi web api
@@ -197,7 +197,6 @@ class FormatAndExportPrediction:
start_to_close_timeout=timedelta(seconds=180), start_to_close_timeout=timedelta(seconds=180),
) )
if write_transformed_handler is not None: if write_transformed_handler is not None:
await write_transformed_handler await write_transformed_handler

View File

@@ -1,7 +1,6 @@
from unittest.mock import ANY, AsyncMock, MagicMock, call, patch from unittest.mock import ANY, AsyncMock, MagicMock, call, patch
import pytest_asyncio import pytest_asyncio
from pandas import DataFrame
from pytest import fixture, mark from pytest import fixture, mark
from sientia_do.notifications.models import NotificationLevel 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.""" """Helper function to create a mocked DataFrame for testing."""
mock_df = MagicMock() mock_df = MagicMock()
mock_head = MagicMock() mock_head = MagicMock()
def get_column_values(key): def get_column_values(key):
if key == 'prediction': if key == 'prediction':
return MagicMock(values=[0.75]) return MagicMock(values=[0.75])
@@ -29,10 +28,10 @@ def _create_mock_dataframe(to_dict_return=None):
return MagicMock(values=[0.95]) return MagicMock(values=[0.95])
else: else:
return MagicMock(values=['2024-01-01T00:00:00+00:00']) return MagicMock(values=['2024-01-01T00:00:00+00:00'])
mock_head.__getitem__.side_effect = get_column_values mock_head.__getitem__.side_effect = get_column_values
mock_df.head.return_value = mock_head mock_df.head.return_value = mock_head
if to_dict_return is None: if to_dict_return is None:
to_dict_return = { to_dict_return = {
'prediction': [0.75], 'prediction': [0.75],
@@ -40,7 +39,7 @@ def _create_mock_dataframe(to_dict_return=None):
'timestamp': ['2024-01-01T00:00:00+00:00'], 'timestamp': ['2024-01-01T00:00:00+00:00'],
} }
mock_df.to_dict.return_value = to_dict_return mock_df.to_dict.return_value = to_dict_return
return mock_df return mock_df
@@ -147,11 +146,13 @@ async def test_write_pi_web_api_data_success(mock_dataframe, api, base_input_dat
@mark.asyncio @mark.asyncio
@patch('laborious.activities.api.DataFrame') @patch('laborious.activities.api.DataFrame')
async def test_write_pi_web_api_data_prediction_error(mock_dataframe, api, base_input_data): async def test_write_pi_web_api_data_prediction_error(mock_dataframe, api, base_input_data):
mock_dataframe.return_value = _create_mock_dataframe({ mock_dataframe.return_value = _create_mock_dataframe(
'prediction': [0.75], {
'prediction_confidence': [PI_WEB_API_PREDICTION_ERROR_CONFIDENCE], 'prediction': [0.75],
'timestamp': ['2024-01-01T00:00:00+00:00'], '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') 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( api.send_notification_async.assert_called_once_with(
metadata=metadata['metadata'], metadata=metadata['metadata'],
notification_id='WRITE_PI_WEB_API_PREDICTION_ERROR', 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', block='write_pi_web_api_data',
level=NotificationLevel.ERROR, level=NotificationLevel.ERROR,
attachment_content=ANY, 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( api.send_notification_async.assert_called_once_with(
metadata=metadata['metadata'], metadata=metadata['metadata'],
notification_id='WRITE_PI_WEB_API_CONFIDENCE_ERROR', 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', block='write_pi_web_api_data',
level=NotificationLevel.ERROR, level=NotificationLevel.ERROR,
attachment_content=ANY, 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 assert api.pi_web_api_client.write_value.call_count == 2
@mark.asyncio @mark.asyncio
@patch('laborious.activities.api.DataFrame') @patch('laborious.activities.api.DataFrame')
async def test_write_pi_web_api_data_empty_tags(mock_dataframe, api, base_input_data): async def test_write_pi_web_api_data_empty_tags(mock_dataframe, api, base_input_data):

View File

@@ -392,7 +392,11 @@ async def test_run_none_path_flag_with_pi_web_api(workflow_mock, format_and_expo
'prediction_confidence': 0, 'prediction_confidence': 0,
'schema': 'test_schema', 'schema': 'test_schema',
'table_name': 'test_table', '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', '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', 'laborious.workflows.sub_workflows.format_and_export_prediction.workflow',
new_callable=AsyncMock, 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 = { input_data = {
'metadata': metadata, 'metadata': metadata,
'path_flag': None, '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', 'schema': 'test_schema',
'table_name': 'test_table', 'table_name': 'test_table',
'opc_output_config': {'test': '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': 'erl:1', '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, retry_policy=ANY,
start_to_close_timeout=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, 'prediction_confidence': 0,
'schema': 'test_schema', 'schema': 'test_schema',
'table_name': 'test_table', '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', 'comment': 'test_comment',
} }

View File

@@ -731,7 +731,11 @@ async def test_path_flag_handler_continue(workflow_mock, prediction_process):
'model_name': model_name, 'model_name': model_name,
'model_config': model_config, 'model_config': model_config,
'opc_output_config': {'test': '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, 'prediction_store_policy': prediction_store_policy,
}, },
confidence, confidence,
@@ -758,7 +762,11 @@ async def test_path_flag_handler_continue(workflow_mock, prediction_process):
'transform_table_name': 'test_transform_table', 'transform_table_name': 'test_transform_table',
'comment': 'Prediction Process', 'comment': 'Prediction Process',
'opc_output_config': {'test': '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, '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_name': model_name,
'model_config': model_config, 'model_config': model_config,
'opc_output_config': {'test': '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, 'prediction_store_policy': prediction_store_policy,
}, },
confidence, confidence,