diff --git a/coverage.sh b/coverage.sh new file mode 100755 index 0000000..fb1d880 --- /dev/null +++ b/coverage.sh @@ -0,0 +1 @@ +pytest --cov=sientia --cov-report=html && xdg-open htmlcov/index.html \ No newline at end of file diff --git a/requirements.txt b/requirements.txt index 91864f7..8c7edc9 100644 --- a/requirements.txt +++ b/requirements.txt @@ -5,7 +5,7 @@ asyncua redis aiokafka pymongo -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.3 git+ssh://git@github.com/Aignosi/sientia-mlops-library.git@0.38.5 pydruid[pandas] prometheus-client \ No newline at end of file diff --git a/scouter/activities/mongodb.py b/scouter/activities/mongodb.py index e36c563..a40cd83 100644 --- a/scouter/activities/mongodb.py +++ b/scouter/activities/mongodb.py @@ -3,13 +3,14 @@ from temporalio import workflow, activity with workflow.unsafe.imports_passed_through(): from typing import Any import traceback - from datetime import datetime + from datetime import datetime, timezone from pymongo import MongoClient from pandas import DataFrame from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler 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_MS_WITH_TZ def clear_mongo_id(docs: list) -> list: @@ -99,7 +100,7 @@ class MongoDB(BaseActivity): else: data_filter = { "inserted_at": { - "$gt": datetime.strptime(last_data_timestamp, "%Y-%m-%d %H:%M:%S.%f") + "$gt": datetime.strptime(last_data_timestamp, DATETIME_FORMAT_MS_WITH_TZ) } } @@ -119,8 +120,8 @@ class MongoDB(BaseActivity): ) for item in data: - item['inserted_at'] = item['inserted_at'].strftime( - "%Y-%m-%d %H:%M:%S.%f") + item['inserted_at'] = item['inserted_at'].replace( + tzinfo=timezone.utc).strftime(DATETIME_FORMAT_MS_WITH_TZ) self.info( f"Loaded {len(data)} documents from MongoDB", diff --git a/scouter/activities/redis.py b/scouter/activities/redis.py index fb3f556..7b62ee2 100644 --- a/scouter/activities/redis.py +++ b/scouter/activities/redis.py @@ -9,8 +9,8 @@ with workflow.unsafe.imports_passed_through(): from sientia_do.observability.logger import Logger from typing import Any from pandas import DataFrame - from datetime import datetime from scouter import metrics + from sientia_do.temporal.constants import DATETIME_FORMAT, now class Redis(RedisBase): @@ -27,7 +27,7 @@ class Redis(RedisBase): Gets the last data timestamp from redis. """ metadata = input_data['metadata'] - key = f"last_data_timestamp_{input_data['workflow_name']}_{input_data['schedule_name']}" + key = f"last_data_timestamp:{input_data['workflow_name']}:{input_data['schedule_name']}" self.info(f"Getting last data timestamp for {key}") @@ -60,7 +60,7 @@ class Redis(RedisBase): Puts the last data timestamp into redis. """ metadata = input_data['metadata'] - key = f"last_data_timestamp_{input_data['workflow_name']}_{input_data['schedule_name']}" + key = f"last_data_timestamp:{input_data['workflow_name']}:{input_data['schedule_name']}" self.info(f"Putting last data timestamp for {key}") @@ -149,6 +149,7 @@ class Redis(RedisBase): # Remove possibly removed tags tags = list(model_tags.keys()) + tags.append('timestamp') self.debug( f"Tags to keep: {tags}", metadata=metadata @@ -169,7 +170,7 @@ class Redis(RedisBase): (row['name'], value)) data_hold['timestamp'] = data['timestamp'].max() if not data.empty else \ - datetime.now().strftime("%Y-%m-%d %H:%M:%S") + data_hold['timestamp'] self.set(key, data_hold, ttl=retention_time) @@ -224,7 +225,7 @@ class Redis(RedisBase): data: The data used to collect the data. """ metadata = input_data['metadata'] - key = f"data_package_{input_data['workflow_name']}_{input_data['schedule_name']}_{datetime.now().strftime('%Y-%m-%d_%H-%M-%S')}" + key = f"data_package_{input_data['workflow_name']}_{input_data['schedule_name']}_{now().strftime(DATETIME_FORMAT)}" data = DataFrame(input_data['data']) held_data = DataFrame(input_data['held_data']) diff --git a/scouter/worker/worker.py b/scouter/worker/worker.py index f636d59..95af238 100644 --- a/scouter/worker/worker.py +++ b/scouter/worker/worker.py @@ -86,7 +86,7 @@ async def main(): temporal_client = await client.Client.connect( target_host=host, - namespace=os.getenv('TEMPORAL_NAMESPACE', 'laborious'), + namespace=os.getenv('TEMPORAL_NAMESPACE', 'scouter'), runtime=new_runtime ) diff --git a/scouter/workflow/sub_workflows/core_scouter.py b/scouter/workflow/sub_workflows/core_scouter.py index d6c83f4..89bf7ff 100644 --- a/scouter/workflow/sub_workflows/core_scouter.py +++ b/scouter/workflow/sub_workflows/core_scouter.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_WITH_TZ @workflow.defn(name="core_scouter") @@ -81,7 +82,11 @@ class CoreScouter: **metadata, 'schema': input_data['schema'], 'table_name': input_data['table_name'], - 'data': held_data + 'data': held_data, + 'timestamp_conversion': { + 'column': 'timestamp', + 'format': DATETIME_FORMAT_WITH_TZ + } }, retry_policy=retry_policy, start_to_close_timeout=timedelta(seconds=60) @@ -104,7 +109,7 @@ class CoreScouter: 'data': input_data['data'], 'held_data': held_data, 'workflow_name': input_data['workflow_name'], - 'schedule_name': input_data['schedule_name'] + 'schedule_name': input_data['schedule_name'], }, retry_policy=retry_policy, start_to_close_timeout=timedelta(seconds=60) diff --git a/tests/activities/test_mongo.py b/tests/activities/test_mongo.py index 24ae1c4..5f91e1c 100644 --- a/tests/activities/test_mongo.py +++ b/tests/activities/test_mongo.py @@ -2,6 +2,7 @@ from datetime import datetime from unittest.mock import ANY, MagicMock, call, patch from pytest import fixture, mark from sientia_do.notifications.models import NotificationLevel +from sientia_do.temporal.constants import DATETIME_FORMAT_MS_WITH_TZ from scouter.activities.mongodb import MongoDB, clear_mongo_id @@ -95,7 +96,7 @@ async def test_load_latest_data_none_last_data_timestamp(mongodb_activity): 'name': 'test1', 'value': 1, 'inserted_at': datetime.strptime( - '2023-01-01 12:00:00.000000', '%Y-%m-%d %H:%M:%S.%f') + '2023-01-01 12:00:00.000000+0000', DATETIME_FORMAT_MS_WITH_TZ) } ] @@ -121,7 +122,7 @@ async def test_load_latest_data_none_last_data_timestamp(mongodb_activity): 0: 1 }, 'inserted_at': { - 0: '2023-01-01 12:00:00.000000' + 0: '2023-01-01 12:00:00.000000+0000' } } @@ -137,14 +138,14 @@ async def test_load_latest_data_not_none_last_data_timestamp(mongodb_activity): 'name': 'test1', 'value': 1, 'inserted_at': datetime.strptime( - '2023-01-01 12:00:00.000000', '%Y-%m-%d %H:%M:%S.%f') + '2023-01-01 12:00:00.000000+0000', DATETIME_FORMAT_MS_WITH_TZ) } ] result = await mongodb_activity.load_latest_data({ 'metadata': {'workflow_name': 'test_pipeline', 'schedule_name': 'test_schedule'}, 'collection_name': 'test_collection', - 'last_data_timestamp': '2023-01-01 12:00:00.000000' + 'last_data_timestamp': '2023-01-01 12:00:00.000000+0000' }) mongodb_activity.database.__getitem__.assert_called_once_with( @@ -154,7 +155,7 @@ async def test_load_latest_data_not_none_last_data_timestamp(mongodb_activity): { 'inserted_at': { '$gt': datetime.strptime( - '2023-01-01 12:00:00.000000', '%Y-%m-%d %H:%M:%S.%f') + '2023-01-01 12:00:00.000000+0000', DATETIME_FORMAT_MS_WITH_TZ) } }, {"_id": 0} @@ -168,7 +169,7 @@ async def test_load_latest_data_not_none_last_data_timestamp(mongodb_activity): 0: 1 }, 'inserted_at': { - 0: '2023-01-01 12:00:00.000000' + 0: '2023-01-01 12:00:00.000000+0000' } } @@ -186,7 +187,7 @@ async def test_load_latest_data_error(mongodb_activity): await mongodb_activity.load_latest_data({ 'metadata': {'workflow_name': 'test_pipeline', 'schedule_name': 'test_schedule'}, 'collection_name': 'test_collection', - 'last_data_timestamp': '2023-01-01 12:00:00.000000' + 'last_data_timestamp': '2023-01-01 12:00:00.000000+0000' }) except Exception as e: assert str(e) == 'test' diff --git a/tests/activities/test_redis.py b/tests/activities/test_redis.py index f68e48f..f87d46f 100644 --- a/tests/activities/test_redis.py +++ b/tests/activities/test_redis.py @@ -83,7 +83,7 @@ async def test_get_last_data_timestamp_not_none(redis_activity): result = await redis_activity.get_last_data_timestamp(test_data) redis_activity.get.assert_called_once_with( - 'last_data_timestamp_test_pipeline_test_schedule' + 'last_data_timestamp:test_pipeline:test_schedule' ) assert result == '2023-01-01 12:00:00' @@ -163,7 +163,7 @@ async def test_put_last_data_timestamp_not_empty_dataframe(redis_activity): assert result == '2023-01-01 12:00:01' redis_activity.set.assert_called_once_with( - 'last_data_timestamp_test_pipeline_test_schedule', + 'last_data_timestamp:test_pipeline:test_schedule', '2023-01-01 12:00:01', ttl=18000 ) diff --git a/tests/workflow/sub_workflows/test_core_scouter.py b/tests/workflow/sub_workflows/test_core_scouter.py index 5e88c5b..16ecff6 100644 --- a/tests/workflow/sub_workflows/test_core_scouter.py +++ b/tests/workflow/sub_workflows/test_core_scouter.py @@ -2,6 +2,7 @@ from unittest.mock import AsyncMock, patch, call, ANY import pytest from scouter.workflow.sub_workflows.core_scouter import CoreScouter from scouter.activities.activities import Activities +from sientia_do.temporal.constants import DATETIME_FORMAT_WITH_TZ @pytest.fixture @@ -94,7 +95,12 @@ async def test_core_scouter_workflow_success(mock_workflow, core_scouter): **expected_metadata, 'schema': 'test_schema', 'table_name': 'test_table', - 'data': 'held_data'}, + 'data': 'held_data', + 'timestamp_conversion': { + 'column': 'timestamp', + 'format': DATETIME_FORMAT_WITH_TZ + } + }, retry_policy=ANY, start_to_close_timeout=ANY ) diff --git a/values.yaml b/values.yaml index 5d413f0..e4a8620 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.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: @@ -150,7 +150,7 @@ env: - name: GITHUB_REPO_URL value: "git@github.com:Aignosi/sientia-dataops-scouter_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: "scouter.worker.worker"