From 5050e0652b451cb6266ee89605e6ef27d85c4f3d Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Fri, 22 Aug 2025 08:11:42 -0300 Subject: [PATCH 1/7] SIENTIAPDE-1193 Update requirements and values for SIENTIAPDE-1193 - Bump sientia-dataops-library version from 1.4.1 to 1.4.2 in requirements.txt. - Change GITHUB_BRANCH in values.yaml to reflect the new task SIENTIAPDE-1193 regarding datetime writing in temporal. --- .../workflows/sub_workflows/format_and_export_prediction.py | 5 +++++ requirements.txt | 2 +- values.yaml | 2 +- 3 files changed, 7 insertions(+), 2 deletions(-) diff --git a/laborious/workflows/sub_workflows/format_and_export_prediction.py b/laborious/workflows/sub_workflows/format_and_export_prediction.py index a5c3afd..a3ed8e8 100644 --- a/laborious/workflows/sub_workflows/format_and_export_prediction.py +++ b/laborious/workflows/sub_workflows/format_and_export_prediction.py @@ -5,6 +5,7 @@ with workflow.unsafe.imports_passed_through(): from typing import Any from datetime import timedelta from sientia_do.temporal.policies import retry_policy + from sientia_do.temporal.constants import DATETIME_FORMAT @workflow.defn(name="format_and_export_prediction") @@ -91,6 +92,10 @@ class FormatAndExportPrediction(): 'schema': input_data['schema'], 'table_name': input_data['table_name'], 'data': prediction, + 'timestamp_conversion': { + 'column': 'timestamp', + 'format': DATETIME_FORMAT + } }, retry_policy=retry_policy, start_to_close_timeout=timedelta(seconds=60) diff --git a/requirements.txt b/requirements.txt index fdc9371..38bd3ea 100644 --- a/requirements.txt +++ b/requirements.txt @@ -3,6 +3,6 @@ psycopg2-binary sqlalchemy asyncua redis -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.1 +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.2 git+ssh://git@github.com/Aignosi/sientia-mlops-library.git@0.38.5 prometheus-client diff --git a/values.yaml b/values.yaml index 25180fb..797023c 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-1199-revisar-e-testar-observabilidade" + value: "SIENTIAPDE-1193-conferir-como-a-escrita-de-datetime-ocorre-no-temporal" - name: PYTHON_APP value: "laborious.worker.worker" From 2301dd63c7f6b151b7062e6b51cf68f52395ed25 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Fri, 22 Aug 2025 11:33:33 -0300 Subject: [PATCH 2/7] SIENTIAPDE-1193 Update requirements and refactor Gates activity for improved datetime handling - Update sientia-dataops-library reference in requirements.txt to version 1.4.3. - Refactor datetime handling in Gates activity to use a consistent format with timezone support, replacing direct datetime calls with a centralized now() function and DATETIME_FORMAT_WITH_TZ constant. --- laborious/activities/gates.py | 11 +++++++---- requirements.txt | 3 ++- 2 files changed, 9 insertions(+), 5 deletions(-) diff --git a/laborious/activities/gates.py b/laborious/activities/gates.py index 32d1006..caa9eaa 100644 --- a/laborious/activities/gates.py +++ b/laborious/activities/gates.py @@ -7,6 +7,7 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.notifications.models import NotificationLevel from sientia_do.temporal.activities.base import BaseActivity from sientia_do.observability.logger import Logger + from sientia_do.temporal.constants import DATETIME_FORMAT_WITH_TZ, now from laborious.utils.filters.mlflow_filters import nan_values_filter, api_error_filter from typing import Any from laborious.utils.filters.conditional_filters import ( @@ -14,7 +15,6 @@ with workflow.unsafe.imports_passed_through(): filter_specific_variables_null_values ) from pandas import DataFrame - from datetime import datetime from laborious import metrics input_filter_functions = { @@ -319,12 +319,15 @@ class Gates(BaseActivity): data = DataFrame(input_data['data']) if data.empty: - return datetime.now().strftime('%Y-%m-%d %H:%M:%S') + return now().strftime(DATETIME_FORMAT_WITH_TZ) + + max_timestamp = max( + data['timestamp'].values.tolist()).strftime(DATETIME_FORMAT_WITH_TZ) self.info( - f"Last timestamp: {max(data['timestamp'].values.tolist())}", metadata) + f"Last timestamp: {max_timestamp}", metadata) - return max(data['timestamp'].values.tolist()) + return max_timestamp @activity.defn(name="write_metrics") async def write_metrics(self, input_data: dict[str, Any]): diff --git a/requirements.txt b/requirements.txt index 38bd3ea..300e1f0 100644 --- a/requirements.txt +++ b/requirements.txt @@ -3,6 +3,7 @@ psycopg2-binary sqlalchemy asyncua redis -git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.2 +# git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.3 +/home/grezewave/Documents/projects/sientia/sientia-dataops-library git+ssh://git@github.com/Aignosi/sientia-mlops-library.git@0.38.5 prometheus-client From ae6d28df5410937e31c6f66792249d04923282f8 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Fri, 22 Aug 2025 14:17:54 -0300 Subject: [PATCH 3/7] SIENTIAPDE-1193 Update Gates activity and format_and_export_prediction workflow to use DATETIME_FORMAT_WITH_TZ for consistent timestamp handling - Modified Gates activity to correctly handle the maximum timestamp without formatting it prematurely. - Updated format_and_export_prediction workflow to utilize DATETIME_FORMAT_WITH_TZ for timestamp conversion. - Enhanced tests to ensure timestamp conversion is applied consistently across workflows. --- laborious/activities/gates.py | 2 +- .../sub_workflows/format_and_export_prediction.py | 4 ++-- .../test_format_and_export_prediction.py | 13 +++++++++++-- 3 files changed, 14 insertions(+), 5 deletions(-) diff --git a/laborious/activities/gates.py b/laborious/activities/gates.py index caa9eaa..c3d119f 100644 --- a/laborious/activities/gates.py +++ b/laborious/activities/gates.py @@ -322,7 +322,7 @@ class Gates(BaseActivity): return now().strftime(DATETIME_FORMAT_WITH_TZ) max_timestamp = max( - data['timestamp'].values.tolist()).strftime(DATETIME_FORMAT_WITH_TZ) + data['timestamp'].values.tolist()) self.info( f"Last timestamp: {max_timestamp}", metadata) diff --git a/laborious/workflows/sub_workflows/format_and_export_prediction.py b/laborious/workflows/sub_workflows/format_and_export_prediction.py index a3ed8e8..ee7c424 100644 --- a/laborious/workflows/sub_workflows/format_and_export_prediction.py +++ b/laborious/workflows/sub_workflows/format_and_export_prediction.py @@ -5,7 +5,7 @@ with workflow.unsafe.imports_passed_through(): from typing import Any from datetime import timedelta from sientia_do.temporal.policies import retry_policy - from sientia_do.temporal.constants import DATETIME_FORMAT + from sientia_do.temporal.constants import DATETIME_FORMAT_WITH_TZ @workflow.defn(name="format_and_export_prediction") @@ -94,7 +94,7 @@ class FormatAndExportPrediction(): 'data': prediction, 'timestamp_conversion': { 'column': 'timestamp', - 'format': DATETIME_FORMAT + 'format': DATETIME_FORMAT_WITH_TZ } }, retry_policy=retry_policy, 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 648a38f..54909a6 100644 --- a/tests/laborious/workflows/subworkflows/test_format_and_export_prediction.py +++ b/tests/laborious/workflows/subworkflows/test_format_and_export_prediction.py @@ -3,6 +3,7 @@ from pytest import mark, fixture from laborious.activities.activities import Activities from laborious.workflows.sub_workflows.format_and_export_prediction import FormatAndExportPrediction +from sientia_do.temporal.constants import DATETIME_FORMAT_WITH_TZ @fixture @@ -73,7 +74,11 @@ async def test_run_none_path_flag(workflow_mock, format_and_export_prediction): 'schema': input_data['schema'], 'table_name': input_data['table_name'], 'data': workflow_mock.execute_activity_method.return_value, - **metadata + **metadata, + 'timestamp_conversion': { + 'column': 'timestamp', + 'format': DATETIME_FORMAT_WITH_TZ + } }, retry_policy=ANY, start_to_close_timeout=ANY @@ -138,7 +143,11 @@ async def test_run_default_path_flag(workflow_mock, format_and_export_prediction 'schema': input_data['schema'], 'table_name': input_data['table_name'], 'data': workflow_mock.execute_activity_method.return_value, - **metadata + **metadata, + 'timestamp_conversion': { + 'column': 'timestamp', + 'format': DATETIME_FORMAT_WITH_TZ + } }, retry_policy=ANY, start_to_close_timeout=ANY From aba077fcc9406a4f5ce453962c1b2c24877a731c Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Fri, 22 Aug 2025 16:12:44 -0300 Subject: [PATCH 4/7] SIENTIAPDE-1193 Update image tag in values.yaml to 0.4.3 and fix timestamp formatting in Gates activity --- laborious/activities/gates.py | 2 +- values.yaml | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/laborious/activities/gates.py b/laborious/activities/gates.py index c3d119f..caa9eaa 100644 --- a/laborious/activities/gates.py +++ b/laborious/activities/gates.py @@ -322,7 +322,7 @@ class Gates(BaseActivity): return now().strftime(DATETIME_FORMAT_WITH_TZ) max_timestamp = max( - data['timestamp'].values.tolist()) + data['timestamp'].values.tolist()).strftime(DATETIME_FORMAT_WITH_TZ) self.info( f"Last timestamp: {max_timestamp}", metadata) diff --git a/values.yaml b/values.yaml index 797023c..a74e456 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.4.2" + tag: "0.4.3" # 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 4b9c4ecca03cfd7bdc67aa40996c53a160c5939c Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Fri, 22 Aug 2025 16:14:48 -0300 Subject: [PATCH 5/7] SIENTIAPDE-1193 Enhance logging in Gates activity by adding debug output for input data --- laborious/activities/gates.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/laborious/activities/gates.py b/laborious/activities/gates.py index caa9eaa..919a2b2 100644 --- a/laborious/activities/gates.py +++ b/laborious/activities/gates.py @@ -318,11 +318,13 @@ class Gates(BaseActivity): data = DataFrame(input_data['data']) + self.debug(f"Input data: {data.to_string()}", metadata) + if data.empty: return now().strftime(DATETIME_FORMAT_WITH_TZ) max_timestamp = max( - data['timestamp'].values.tolist()).strftime(DATETIME_FORMAT_WITH_TZ) + data['timestamp'].values.tolist()) self.info( f"Last timestamp: {max_timestamp}", metadata) From 9cdbc0afda850d6e24e9b744b1b5b5d4fb46995b Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 25 Aug 2025 08:27:58 -0300 Subject: [PATCH 6/7] SIENTIAPDE-1193 Update dependencies and enhance datetime handling - Bump sientia-dataops-library version in requirements.txt to 1.4.4. - Update image tag in values.yaml to 0.4.4. - Add 'datetime_columns' to metadata in MinimalRetrain and PredictionsBatch workflows for improved data handling. --- laborious/workflows/minimal_retrain.py | 1 + laborious/workflows/predictions_batch.py | 1 + requirements.txt | 3 +-- values.yaml | 2 +- 4 files changed, 4 insertions(+), 3 deletions(-) diff --git a/laborious/workflows/minimal_retrain.py b/laborious/workflows/minimal_retrain.py index b3da70a..0baf647 100644 --- a/laborious/workflows/minimal_retrain.py +++ b/laborious/workflows/minimal_retrain.py @@ -52,6 +52,7 @@ class MinimalRetrain(): { **metadata, 'query': input_data['query'], + 'datetime_columns': input_data.get('datetime_columns', []) }, retry_policy=retry_policy, start_to_close_timeout=timedelta(seconds=60) diff --git a/laborious/workflows/predictions_batch.py b/laborious/workflows/predictions_batch.py index a4b14f0..263e893 100644 --- a/laborious/workflows/predictions_batch.py +++ b/laborious/workflows/predictions_batch.py @@ -54,6 +54,7 @@ class PredictionsBatch(): { **metadata, 'query': input_data['query'], + 'datetime_columns': input_data.get('datetime_columns', []) }, retry_policy=retry_policy, start_to_close_timeout=timedelta(seconds=60) diff --git a/requirements.txt b/requirements.txt index 300e1f0..5a8d019 100644 --- a/requirements.txt +++ b/requirements.txt @@ -3,7 +3,6 @@ psycopg2-binary sqlalchemy asyncua redis -# git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.3 -/home/grezewave/Documents/projects/sientia/sientia-dataops-library +git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.4 git+ssh://git@github.com/Aignosi/sientia-mlops-library.git@0.38.5 prometheus-client diff --git a/values.yaml b/values.yaml index a74e456..0a59346 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.4.3" + tag: "0.4.4" # 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 8f8859aaa61a4fb128fe6e7a25d99f770c341eda Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 25 Aug 2025 09:11:52 -0300 Subject: [PATCH 7/7] SIENTIAPDE-1193 Enhance test cases by adding 'datetime_columns' to input data for MinimalRetrain and PredictionsBatch workflows. This improves metadata handling in asynchronous tests. --- coverage.sh | 1 + tests/laborious/workflows/test_minimal_retrain.py | 1 + tests/laborious/workflows/test_predictions_batch.py | 4 +++- 3 files changed, 5 insertions(+), 1 deletion(-) create mode 100755 coverage.sh diff --git a/coverage.sh b/coverage.sh new file mode 100755 index 0000000..5867692 --- /dev/null +++ b/coverage.sh @@ -0,0 +1 @@ +pytest --cov=laborious --cov-report=html && xdg-open htmlcov/index.html \ No newline at end of file diff --git a/tests/laborious/workflows/test_minimal_retrain.py b/tests/laborious/workflows/test_minimal_retrain.py index bcb2288..b3b03b5 100644 --- a/tests/laborious/workflows/test_minimal_retrain.py +++ b/tests/laborious/workflows/test_minimal_retrain.py @@ -48,6 +48,7 @@ async def test_run(workflow_mock: AsyncMock, minimal_retrain: MinimalRetrain): { **metadata, "query": input_data["query"], + 'datetime_columns': input_data.get('datetime_columns', []) }, retry_policy=ANY, start_to_close_timeout=ANY diff --git a/tests/laborious/workflows/test_predictions_batch.py b/tests/laborious/workflows/test_predictions_batch.py index 3c8fc82..e9b9bb6 100644 --- a/tests/laborious/workflows/test_predictions_batch.py +++ b/tests/laborious/workflows/test_predictions_batch.py @@ -32,7 +32,8 @@ async def test_run(workflow_mock: AsyncMock, predictions_batch: PredictionsBatch 'query': 'SELECT * FROM test', 'schema': 'test_schema', 'table_name': 'test_table', - 'opc_output_config': 'test_opc_output_config' + 'opc_output_config': 'test_opc_output_config', + 'datetime_columns': ['timestamp', 'created_at'] } await predictions_batch.run(input_data) @@ -43,6 +44,7 @@ async def test_run(workflow_mock: AsyncMock, predictions_batch: PredictionsBatch { **metadata, 'query': input_data['query'], + 'datetime_columns': input_data.get('datetime_columns', []) }, retry_policy=ANY, start_to_close_timeout=ANY