From 67942c45e0eb28fd7347cb1f2e4612eea340a9e3 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Fri, 20 Mar 2026 15:52:04 -0300 Subject: [PATCH] SIENTIAPDE-1712 Refactor MinioDataFramePayload usage across activities - Updated instances of MinioDataFramePayload initialization in Gates, MLFlow, and Storage classes to use the new from_dict method for better data reconstruction from dictionaries. - Enhanced the PredictionProcess workflow to utilize the updated payload handling. - Added passthrough fixtures in tests to accommodate the new from_dict method for consistent testing behavior. --- laborious/activities/gates.py | 8 +-- laborious/activities/mlflow.py | 6 +-- laborious/activities/storage.py | 9 ++-- .../utils/models/minio_dataframe_payload.py | 33 +++++++++++- .../sub_workflows/prediction_process.py | 17 +++--- tests/laborious/activities/test_gates.py | 8 +++ tests/laborious/activities/test_mlflow.py | 8 +++ tests/laborious/activities/test_storage.py | 9 ++++ .../models/test_minio_dataframe_payload.py | 52 +++++++++++++++++++ .../subworkflows/test_prediction_process.py | 9 ++++ .../workflows/test_minimal_retrain.py | 9 ++++ .../workflows/test_predictions_batch.py | 18 ++++--- 12 files changed, 155 insertions(+), 31 deletions(-) diff --git a/laborious/activities/gates.py b/laborious/activities/gates.py index bd40b53..7acafd3 100644 --- a/laborious/activities/gates.py +++ b/laborious/activities/gates.py @@ -157,7 +157,7 @@ class Gates(MinioManager): self.info('Performing input gate...', metadata) filters = input_data['filters'] - payload: MinioDataFramePayload = input_data['data'] + payload = MinioDataFramePayload.from_dict(input_data['data']) data = await payload.retrieve(self.minio_repository, metadata) path_priority = input_data['path_priority'] @@ -236,7 +236,7 @@ class Gates(MinioManager): filters = input_data['filters'] - payload: MinioDataFramePayload = input_data['data'] + payload = MinioDataFramePayload.from_dict(input_data['data']) data = await payload.retrieve(self.minio_repository, metadata) gate_type = input_data['type'] @@ -327,7 +327,7 @@ class Gates(MinioManager): filters = input_data['filters'] - payload: MinioDataFramePayload = input_data['data'] + payload = MinioDataFramePayload.from_dict(input_data['data']) data = await payload.retrieve(self.minio_repository, metadata) gate_type = input_data['type'] @@ -465,7 +465,7 @@ class Gates(MinioManager): self.info('Formatting transformed data...', metadata) - payload: MinioDataFramePayload = input_data['data'] + payload = MinioDataFramePayload.from_dict(input_data['data']) data = await payload.retrieve(self.minio_repository, metadata) data['timestamp'] = data.index diff --git a/laborious/activities/mlflow.py b/laborious/activities/mlflow.py index 82c664f..27d1b92 100644 --- a/laborious/activities/mlflow.py +++ b/laborious/activities/mlflow.py @@ -128,7 +128,7 @@ class MLFlow(MinioManager): metadata = input_data['metadata'] self.info('Transforming data...', metadata) - payload: MinioDataFramePayload = input_data['data'] + payload = MinioDataFramePayload.from_dict(input_data['data']) data = await payload.retrieve(self.minio_repository, metadata) model_name = input_data['model_name'] @@ -222,7 +222,7 @@ class MLFlow(MinioManager): metadata = input_data['metadata'] self.info('Predicting data...', metadata) - payload: MinioDataFramePayload = input_data['data'] + payload = MinioDataFramePayload.from_dict(input_data['data']) data = await payload.retrieve(self.minio_repository, metadata) model_name = input_data['model_name'] @@ -311,7 +311,7 @@ class MLFlow(MinioManager): try: # Payload-based retrain input (inline dict or MinIO offloaded). - payload: MinioDataFramePayload = input_data['data'] + payload = MinioDataFramePayload.from_dict(input_data['data']) data = await payload.retrieve(self.minio_repository, metadata) except Exception as e: diff --git a/laborious/activities/storage.py b/laborious/activities/storage.py index 223001b..58b60c9 100644 --- a/laborious/activities/storage.py +++ b/laborious/activities/storage.py @@ -8,7 +8,6 @@ with workflow.unsafe.imports_passed_through(): # Extend the Temporal Postgres activities for convenient query -> MinIO export import traceback from datetime import timedelta - from io import BytesIO from typing import Any import pandas as pd @@ -18,7 +17,7 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.observability.metrics_controller import MetricsController from sientia_do.repository.minio_repository import MinioRepository from sientia_do.temporal.activities.postgres import Postgres - from sientia_do.temporal.constants import DATETIME_FORMAT_FILENAME, now + from sientia_do.temporal.constants import now from laborious.utils.models.minio_dataframe_payload import MinioDataFramePayload @@ -116,7 +115,7 @@ class Storage(Postgres, MinioManager): Export a payload to PostgreSQL. """ metadata = input_data.get('metadata') - payload: MinioDataFramePayload = input_data['data'] + payload = MinioDataFramePayload.from_dict(input_data['data']) data = await payload.retrieve(self.minio_repository, metadata) return await self.export_data_to_postgres( @@ -142,7 +141,8 @@ class Storage(Postgres, MinioManager): raise ValueError('Minio repository not initialized') metadata = input_data.get('metadata', {}) - prefix = input_data['prefix'] + payload = MinioDataFramePayload.from_dict(input_data['data']) + prefix = payload.cleanup_prefix() base = now() cutoff = (base.replace(tzinfo=None) if base.tzinfo else base) - timedelta( hours=self.retention_hours @@ -206,7 +206,6 @@ class Storage(Postgres, MinioManager): return report - def close(self) -> None: """Close Storage resources (MinIO client and Postgres engine).""" Postgres.close(self) diff --git a/laborious/utils/models/minio_dataframe_payload.py b/laborious/utils/models/minio_dataframe_payload.py index 7fd9bc8..8463be2 100644 --- a/laborious/utils/models/minio_dataframe_payload.py +++ b/laborious/utils/models/minio_dataframe_payload.py @@ -81,6 +81,38 @@ class MinioDataFramePayload: object_prefix: str | None = None uri: str | None = None + @classmethod + def from_dict(cls, raw: dict[str, Any] | 'MinioDataFramePayload') -> 'MinioDataFramePayload': + """ + Reconstruct a MinioDataFramePayload from a plain dict produced by Temporal serialization. + + Temporal converts dataclass return values into plain dicts when crossing + workflow/activity boundaries. This method rebuilds the typed instance so + that methods like ``retrieve``, ``cleanup_prefix`` and ``has_data`` are + available on the receiving side. + + If the argument is already a MinioDataFramePayload, it is returned as-is. + + Args: + raw: Dict with keys matching the dataclass fields + (last_timestamp, status, data, bucket, object_key, object_prefix, uri), + or an existing MinioDataFramePayload instance. + + Return: + MinioDataFramePayload: Reconstructed (or original) instance. + """ + if isinstance(raw, MinioDataFramePayload): + return raw + return cls( + last_timestamp=raw['last_timestamp'], + status=raw.get('status'), + data=raw.get('data'), + bucket=raw.get('bucket'), + object_key=raw.get('object_key'), + object_prefix=raw.get('object_prefix'), + uri=raw.get('uri'), + ) + @staticmethod def estimate_size_bytes(df: DataFrame) -> int: """ @@ -118,7 +150,6 @@ class MinioDataFramePayload: except ValueError: return None - @staticmethod def cleanup_prefix(self) -> str | None: """ Return True if cleanup is enabled for this payload. diff --git a/laborious/workflows/sub_workflows/prediction_process.py b/laborious/workflows/sub_workflows/prediction_process.py index 07e42e7..3560d5f 100644 --- a/laborious/workflows/sub_workflows/prediction_process.py +++ b/laborious/workflows/sub_workflows/prediction_process.py @@ -38,8 +38,6 @@ class PredictionProcess: 8. Export Delegation: Delegates to FormatAndExportPrediction workflow """ - cleanup_prefixes: set[str] = set() - @workflow.run async def run(self, input_data: dict[str, Any]): """ @@ -90,8 +88,6 @@ class PredictionProcess: model_config = input_data.get('model_config', {}) save_transform = input_data.get('save_transform', True) - prefix = data.cleanup_prefix() - try: await self._run_prediction_pipeline( input_data, @@ -103,13 +99,12 @@ class PredictionProcess: save_transform, ) finally: - if self.cleanup_prefixes: - await workflow.execute_activity_method( - Activities.cleanup_minio_objects_expired, - {**metadata, 'prefix': prefix}, - retry_policy=retry_policy, - start_to_close_timeout=timedelta(minutes=5), - ) + await workflow.execute_activity_method( + Activities.cleanup_minio_objects_expired, + {**metadata, 'data': data}, + retry_policy=retry_policy, + start_to_close_timeout=timedelta(minutes=5), + ) async def _run_prediction_pipeline( self, diff --git a/tests/laborious/activities/test_gates.py b/tests/laborious/activities/test_gates.py index eaee176..0d90ea6 100644 --- a/tests/laborious/activities/test_gates.py +++ b/tests/laborious/activities/test_gates.py @@ -7,6 +7,14 @@ from sientia_do.notifications.models import NotificationLevel from laborious.activities.gates import Gates +@fixture(autouse=True) +def _passthrough_from_dict(): + with patch( + 'laborious.activities.gates.MinioDataFramePayload.from_dict', side_effect=lambda x: x + ): + yield + + def _minio_payload(retrieve_return, status=None): """ Build a MinioDataFramePayload-like test double with async retrieve. diff --git a/tests/laborious/activities/test_mlflow.py b/tests/laborious/activities/test_mlflow.py index baa06e4..e495be8 100644 --- a/tests/laborious/activities/test_mlflow.py +++ b/tests/laborious/activities/test_mlflow.py @@ -8,6 +8,14 @@ from sientia_do.temporal.constants import DATETIME_FORMAT, DATETIME_FORMAT_WITH_ from laborious.activities.mlflow import MLFlow +@fixture(autouse=True) +def _passthrough_from_dict(): + with patch( + 'laborious.activities.mlflow.MinioDataFramePayload.from_dict', side_effect=lambda x: x + ): + yield + + @patch('laborious.activities.mlflow.MLFlowRepository') @patch('laborious.activities.mlflow.MinioRepository') def test___init__(mock_minio_repository, mock_mlflow_repository): diff --git a/tests/laborious/activities/test_storage.py b/tests/laborious/activities/test_storage.py index f0272d1..a1fa8e8 100644 --- a/tests/laborious/activities/test_storage.py +++ b/tests/laborious/activities/test_storage.py @@ -8,6 +8,15 @@ from sientia_do.temporal.activities.postgres import Postgres from laborious.activities.storage import Storage + +@fixture(autouse=True) +def _passthrough_from_dict(): + with patch( + 'laborious.activities.storage.MinioDataFramePayload.from_dict', side_effect=lambda x: x + ): + yield + + metadata = { 'metadata': { 'model_id': 'test_model_id', diff --git a/tests/laborious/utils/models/test_minio_dataframe_payload.py b/tests/laborious/utils/models/test_minio_dataframe_payload.py index 9349264..de1a334 100644 --- a/tests/laborious/utils/models/test_minio_dataframe_payload.py +++ b/tests/laborious/utils/models/test_minio_dataframe_payload.py @@ -214,3 +214,55 @@ async def test_from_dataframe_offloaded(mock_now): assert result.bucket == 'test-bucket' assert result.uri == 's3://test-bucket/full/key.parquet' minio.upload_file.assert_awaited_once() + + +def test_from_dict_inline(): + raw = { + 'last_timestamp': '2024-01-01T00:00:00+00:00', + 'status': None, + 'data': {'col1': {0: 'val1'}}, + 'bucket': None, + 'object_key': None, + 'object_prefix': None, + 'uri': None, + } + payload = MinioDataFramePayload.from_dict(raw) + assert isinstance(payload, MinioDataFramePayload) + assert payload.last_timestamp == '2024-01-01T00:00:00+00:00' + assert payload.data == {'col1': {0: 'val1'}} + assert payload.object_key is None + + +def test_from_dict_offloaded(): + raw = { + 'last_timestamp': '2024-06-15T10:30:45+00:00', + 'status': {'success': True}, + 'data': None, + 'bucket': 'my-bucket', + 'object_key': 'training_datasets/model/model-initial-2024-06-15_10-30-45.parquet', + 'object_prefix': 'training_datasets/model', + 'uri': 's3://my-bucket/training_datasets/model/model-initial-2024-06-15_10-30-45.parquet', + } + payload = MinioDataFramePayload.from_dict(raw) + assert isinstance(payload, MinioDataFramePayload) + assert payload.data is None + assert payload.bucket == 'my-bucket' + assert payload.object_key == raw['object_key'] + assert payload.object_prefix == 'training_datasets/model' + assert payload.uri == raw['uri'] + assert payload.status == {'success': True} + + +def test_from_dict_minimal_keys(): + raw = {'last_timestamp': '2024-01-01'} + payload = MinioDataFramePayload.from_dict(raw) + assert payload.last_timestamp == '2024-01-01' + assert payload.data is None + assert payload.bucket is None + assert payload.object_key is None + + +def test_from_dict_passthrough_existing_instance(): + original = MinioDataFramePayload(last_timestamp='2024-01-01', data={'a': 1}, bucket='b') + result = MinioDataFramePayload.from_dict(original) + assert result is original diff --git a/tests/laborious/workflows/subworkflows/test_prediction_process.py b/tests/laborious/workflows/subworkflows/test_prediction_process.py index 0da262b..89a2f65 100644 --- a/tests/laborious/workflows/subworkflows/test_prediction_process.py +++ b/tests/laborious/workflows/subworkflows/test_prediction_process.py @@ -6,6 +6,15 @@ from laborious.activities.activities import Activities from laborious.workflows.sub_workflows.prediction_process import PredictionProcess +@fixture(autouse=True) +def _passthrough_from_dict(): + with patch( + 'laborious.workflows.sub_workflows.prediction_process.MinioDataFramePayload.from_dict', + side_effect=lambda x: x, + ): + yield + + @fixture def prediction_process(): return PredictionProcess() diff --git a/tests/laborious/workflows/test_minimal_retrain.py b/tests/laborious/workflows/test_minimal_retrain.py index be49c89..9af1331 100644 --- a/tests/laborious/workflows/test_minimal_retrain.py +++ b/tests/laborious/workflows/test_minimal_retrain.py @@ -6,6 +6,15 @@ from laborious.activities.activities import Activities from laborious.workflows.minimal_retrain import MinimalRetrain +@fixture(autouse=True) +def _passthrough_from_dict(): + with patch( + 'laborious.workflows.minimal_retrain.MinioDataFramePayload.from_dict', + side_effect=lambda x: x, + ): + yield + + @fixture def minimal_retrain() -> MinimalRetrain: return MinimalRetrain() diff --git a/tests/laborious/workflows/test_predictions_batch.py b/tests/laborious/workflows/test_predictions_batch.py index 8ebd040..7164b7a 100644 --- a/tests/laborious/workflows/test_predictions_batch.py +++ b/tests/laborious/workflows/test_predictions_batch.py @@ -1,4 +1,4 @@ -from unittest.mock import ANY, AsyncMock, call, patch +from unittest.mock import ANY, AsyncMock, MagicMock, call, patch from pytest import fixture, mark @@ -22,12 +22,15 @@ metadata = { @mark.asyncio +@patch( + 'laborious.workflows.predictions_batch.MinioDataFramePayload.from_dict', side_effect=lambda x: x +) @patch('laborious.workflows.predictions_batch.workflow', new_callable=AsyncMock) -async def test_run(workflow_mock: AsyncMock, predictions_batch: PredictionsBatch): - workflow_mock.execute_activity_method.return_value = { - 'success': True, - 'data': {'col': ['test_data']}, - } +async def test_run(workflow_mock: AsyncMock, mock_from_dict, predictions_batch: PredictionsBatch): + activity_return = MagicMock() + activity_return.cleanup_prefix.return_value = None + workflow_mock.execute_activity_method.return_value = activity_return + input_data = { 'schedule_name': 'test_schedule', 'model_name': 'test_model', @@ -62,7 +65,8 @@ async def test_run(workflow_mock: AsyncMock, predictions_batch: PredictionsBatch ) prediction_input = { 'metadata': metadata, - 'data': {'success': True, 'data': {'col': ['test_data']}}, + 'data': activity_return, + 'cleanup_prefix': activity_return.cleanup_prefix(), 'schema': input_data['schema'], 'table_name': input_data['table_name'], 'transform_table_name': input_data['transform_table_name'],