From e17df007ab6913b0dff8fe3c242a686acd368716 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Wed, 27 Aug 2025 09:12:24 -0300 Subject: [PATCH 1/6] SIENTIAPDE-1205 Update orchestration notebook execution count, modify image tag in values.yaml, and standardize datetime handling in MongoDB activities - Changed execution count in init_orchestration.ipynb from 2 to 1. - Updated image tag in values.yaml from 0.4.4 to 0.4.5. - Replaced DATETIME_FORMAT_MS_WITH_TZ with DATETIME_FORMAT_MS in mongo_db.py for consistent timestamp handling. --- init_orchestration.ipynb | 105 ++-------------------------- orchestrator/activities/mongo_db.py | 6 +- values.yaml | 4 +- 3 files changed, 10 insertions(+), 105 deletions(-) diff --git a/init_orchestration.ipynb b/init_orchestration.ipynb index 14f66c0..7bfbd9f 100644 --- a/init_orchestration.ipynb +++ b/init_orchestration.ipynb @@ -2,7 +2,7 @@ "cells": [ { "cell_type": "code", - "execution_count": 2, + "execution_count": 1, "id": "3f8b77a4", "metadata": {}, "outputs": [], @@ -21,101 +21,6 @@ ")" ] }, - { - "cell_type": "code", - "execution_count": 3, - "id": "6ee4b5a7", - "metadata": {}, - "outputs": [ - { - "data": { - "text/plain": [ - "" - ] - }, - "execution_count": 3, - "metadata": {}, - "output_type": "execute_result" - } - ], - "source": [ - "from datetime import timedelta\n", - "from temporalio.client import (\n", - " Client,\n", - " Schedule,\n", - " ScheduleActionStartWorkflow,\n", - " ScheduleIntervalSpec,\n", - " ScheduleSpec,\n", - ")\n", - "from temporalio.common import TypedSearchAttributes, SearchAttributeKey, SearchAttributePair\n", - "\n", - "# temporal operator search-attribute create --namespace scouter --name model_id --type Text && temporal operator search-attribute create --namespace scouter --name orchestrated --type Text && temporal operator search-attribute create --namespace scouter --name model_name --type Text && temporal operator search-attribute create --namespace laborious --name model_id --type Text && temporal operator search-attribute create --namespace laborious --name orchestrated --type Text && temporal operator search-attribute create --namespace laborious --name model_name --type Text\n", - "\n", - "\n", - "\n", - "await temporal_client.create_schedule(\n", - " \"orchestrator\",\n", - " Schedule(\n", - " action=ScheduleActionStartWorkflow(\n", - " 'orchestrator',\n", - " {\n", - " \"schedule_name\": \"orchestrator-test\",\n", - " \"pipelines_query\": {\n", - " \"collection\": \"pipelines\",\n", - " \"aggregation\": [\n", - " {\n", - " \"$lookup\": {\n", - " \"from\": \"models\",\n", - " \"localField\": \"model_id\",\n", - " \"foreignField\": \"id\",\n", - " \"as\": \"model_docs\"\n", - " }\n", - " },\n", - " {\n", - " \"$match\": {\n", - " \"active\": True\n", - " }\n", - " },\n", - " {\n", - " \"$addFields\": {\n", - " \"models\": {\n", - " \"$arrayElemAt\": [\n", - " \"$model_docs\",\n", - " 0\n", - " ]\n", - " }\n", - " }\n", - " },\n", - " {\n", - " \"$match\": {\n", - " \"models.active\": True\n", - " }\n", - " },\n", - " {\n", - " \"$project\": {\n", - " \"model_docs\": 0\n", - " }\n", - " }\n", - " ]\n", - " },\n", - " \"opc_servers_query\": {\n", - " \"collection\": \"opc-servers\",\n", - " \"filters\": {\n", - "\n", - " }\n", - " }\n", - " },\n", - " id=\"orchestrator\",\n", - " task_queue=\"orchestrator-queue\",\n", - " execution_timeout=timedelta(minutes=600)\n", - " ),\n", - " spec=ScheduleSpec(\n", - " intervals=[ScheduleIntervalSpec(every=timedelta(minutes=60))]\n", - " )\n", - " )\n", - ")" - ] - }, { "cell_type": "code", "execution_count": 4, @@ -131,13 +36,13 @@ "\u001b[31mRPCError\u001b[39m Traceback (most recent call last)", "\u001b[36mFile \u001b[39m\u001b[32m~/Documents/projects/sientia/sientia-dataops-orchestrator_temporal/venv/lib/python3.11/site-packages/temporalio/service.py:1243\u001b[39m, in \u001b[36m_BridgeServiceClient._rpc_call\u001b[39m\u001b[34m(self, rpc, req, resp_type, service, retry, metadata, timeout)\u001b[39m\n\u001b[32m 1242\u001b[39m client = \u001b[38;5;28;01mawait\u001b[39;00m \u001b[38;5;28mself\u001b[39m._connected_client()\n\u001b[32m-> \u001b[39m\u001b[32m1243\u001b[39m resp = \u001b[38;5;28;01mawait\u001b[39;00m client.call(\n\u001b[32m 1244\u001b[39m service=service,\n\u001b[32m 1245\u001b[39m rpc=rpc,\n\u001b[32m 1246\u001b[39m req=req,\n\u001b[32m 1247\u001b[39m resp_type=resp_type,\n\u001b[32m 1248\u001b[39m retry=retry,\n\u001b[32m 1249\u001b[39m metadata=metadata,\n\u001b[32m 1250\u001b[39m timeout=timeout,\n\u001b[32m 1251\u001b[39m )\n\u001b[32m 1252\u001b[39m \u001b[38;5;28;01mif\u001b[39;00m LOG_PROTOS:\n", "\u001b[36mFile \u001b[39m\u001b[32m~/Documents/projects/sientia/sientia-dataops-orchestrator_temporal/venv/lib/python3.11/site-packages/temporalio/bridge/client.py:151\u001b[39m, in \u001b[36mClient.call\u001b[39m\u001b[34m(self, service, rpc, req, resp_type, retry, metadata, timeout)\u001b[39m\n\u001b[32m 150\u001b[39m resp = resp_type()\n\u001b[32m--> \u001b[39m\u001b[32m151\u001b[39m resp.ParseFromString(\u001b[38;5;28;01mawait\u001b[39;00m resp_fut)\n\u001b[32m 152\u001b[39m \u001b[38;5;28;01mreturn\u001b[39;00m resp\n", - "\u001b[31mRPCError\u001b[39m: (6, 'Workflow execution is already running. WorkflowId: temporal-sys-scheduler:orchestrator, RunId: 01989edb-863c-7585-8653-945d00384e2f.', b'\\x08\\x06\\x12\\x84\\x01Workflow execution is already running. WorkflowId: temporal-sys-scheduler:orchestrator, RunId: 01989edb-863c-7585-8653-945d00384e2f.\\x1a\\xa7\\x01\\nWtype.googleapis.com/temporal.api.errordetails.v1.WorkflowExecutionAlreadyStartedFailure\\x12L\\n$d248fcc1-8071-40fc-b5f3-e70b8cb1dbb5\\x12$01989edb-863c-7585-8653-945d00384e2f')", + "\u001b[31mRPCError\u001b[39m: (6, 'Workflow execution is already running. WorkflowId: temporal-sys-scheduler:orchestrator, RunId: 0198eb44-22d8-7316-9b99-dc99c4b50e08.', b'\\x08\\x06\\x12\\x84\\x01Workflow execution is already running. WorkflowId: temporal-sys-scheduler:orchestrator, RunId: 0198eb44-22d8-7316-9b99-dc99c4b50e08.\\x1a\\xa7\\x01\\nWtype.googleapis.com/temporal.api.errordetails.v1.WorkflowExecutionAlreadyStartedFailure\\x12L\\n$fa9c61bb-4176-4474-9d45-5f4bde5fa2c5\\x12$0198eb44-22d8-7316-9b99-dc99c4b50e08')", "\nDuring handling of the above exception, another exception occurred:\n", "\u001b[31mRPCError\u001b[39m Traceback (most recent call last)", "\u001b[36mFile \u001b[39m\u001b[32m~/Documents/projects/sientia/sientia-dataops-orchestrator_temporal/venv/lib/python3.11/site-packages/temporalio/client.py:6430\u001b[39m, in \u001b[36m_ClientImpl.create_schedule\u001b[39m\u001b[34m(self, input)\u001b[39m\n\u001b[32m 6427\u001b[39m temporalio.converter.encode_search_attributes(\n\u001b[32m 6428\u001b[39m \u001b[38;5;28minput\u001b[39m.search_attributes, request.search_attributes\n\u001b[32m 6429\u001b[39m )\n\u001b[32m-> \u001b[39m\u001b[32m6430\u001b[39m \u001b[38;5;28;01mawait\u001b[39;00m \u001b[38;5;28mself\u001b[39m._client.workflow_service.create_schedule(\n\u001b[32m 6431\u001b[39m request,\n\u001b[32m 6432\u001b[39m retry=\u001b[38;5;28;01mTrue\u001b[39;00m,\n\u001b[32m 6433\u001b[39m metadata=\u001b[38;5;28minput\u001b[39m.rpc_metadata,\n\u001b[32m 6434\u001b[39m timeout=\u001b[38;5;28minput\u001b[39m.rpc_timeout,\n\u001b[32m 6435\u001b[39m )\n\u001b[32m 6436\u001b[39m \u001b[38;5;28;01mexcept\u001b[39;00m RPCError \u001b[38;5;28;01mas\u001b[39;00m err:\n", "\u001b[36mFile \u001b[39m\u001b[32m~/Documents/projects/sientia/sientia-dataops-orchestrator_temporal/venv/lib/python3.11/site-packages/temporalio/service.py:1170\u001b[39m, in \u001b[36mServiceCall.__call__\u001b[39m\u001b[34m(self, req, retry, metadata, timeout)\u001b[39m\n\u001b[32m 1155\u001b[39m \u001b[38;5;250m\u001b[39m\u001b[33;03m\"\"\"Invoke underlying client with the given request.\u001b[39;00m\n\u001b[32m 1156\u001b[39m \n\u001b[32m 1157\u001b[39m \u001b[33;03mArgs:\u001b[39;00m\n\u001b[32m (...)\u001b[39m\u001b[32m 1168\u001b[39m \u001b[33;03m RPCError: Any RPC error that occurs during the call.\u001b[39;00m\n\u001b[32m 1169\u001b[39m \u001b[33;03m\"\"\"\u001b[39;00m\n\u001b[32m-> \u001b[39m\u001b[32m1170\u001b[39m \u001b[38;5;28;01mreturn\u001b[39;00m \u001b[38;5;28;01mawait\u001b[39;00m \u001b[38;5;28mself\u001b[39m.service_client._rpc_call(\n\u001b[32m 1171\u001b[39m \u001b[38;5;28mself\u001b[39m.name,\n\u001b[32m 1172\u001b[39m req,\n\u001b[32m 1173\u001b[39m \u001b[38;5;28mself\u001b[39m.resp_type,\n\u001b[32m 1174\u001b[39m service=\u001b[38;5;28mself\u001b[39m.service,\n\u001b[32m 1175\u001b[39m retry=retry,\n\u001b[32m 1176\u001b[39m metadata=metadata,\n\u001b[32m 1177\u001b[39m timeout=timeout,\n\u001b[32m 1178\u001b[39m )\n", "\u001b[36mFile \u001b[39m\u001b[32m~/Documents/projects/sientia/sientia-dataops-orchestrator_temporal/venv/lib/python3.11/site-packages/temporalio/service.py:1258\u001b[39m, in \u001b[36m_BridgeServiceClient._rpc_call\u001b[39m\u001b[34m(self, rpc, req, resp_type, service, retry, metadata, timeout)\u001b[39m\n\u001b[32m 1257\u001b[39m status, message, details = err.args\n\u001b[32m-> \u001b[39m\u001b[32m1258\u001b[39m \u001b[38;5;28;01mraise\u001b[39;00m RPCError(message, RPCStatusCode(status), details)\n", - "\u001b[31mRPCError\u001b[39m: Workflow execution is already running. WorkflowId: temporal-sys-scheduler:orchestrator, RunId: 01989edb-863c-7585-8653-945d00384e2f.", + "\u001b[31mRPCError\u001b[39m: Workflow execution is already running. WorkflowId: temporal-sys-scheduler:orchestrator, RunId: 0198eb44-22d8-7316-9b99-dc99c4b50e08.", "\nDuring handling of the above exception, another exception occurred:\n", "\u001b[31mScheduleAlreadyRunningError\u001b[39m Traceback (most recent call last)", "\u001b[36mCell\u001b[39m\u001b[36m \u001b[39m\u001b[32mIn[4]\u001b[39m\u001b[32m, line 13\u001b[39m\n\u001b[32m 2\u001b[39m \u001b[38;5;28;01mfrom\u001b[39;00m\u001b[38;5;250m \u001b[39m\u001b[34;01mtemporalio\u001b[39;00m\u001b[34;01m.\u001b[39;00m\u001b[34;01mclient\u001b[39;00m\u001b[38;5;250m \u001b[39m\u001b[38;5;28;01mimport\u001b[39;00m (\n\u001b[32m 3\u001b[39m Schedule,\n\u001b[32m 4\u001b[39m ScheduleActionStartWorkflow,\n\u001b[32m 5\u001b[39m ScheduleIntervalSpec,\n\u001b[32m 6\u001b[39m ScheduleSpec,\n\u001b[32m 7\u001b[39m )\n\u001b[32m 9\u001b[39m \u001b[38;5;66;03m# temporal operator search-attribute create --namespace scouter --name model_id --type Text && temporal operator search-attribute create --namespace scouter --name orchestrated --type Text && temporal operator search-attribute create --namespace scouter --name model_name --type Text && temporal operator search-attribute create --namespace laborious --name model_id --type Text && temporal operator search-attribute create --namespace laborious --name orchestrated --type Text && temporal operator search-attribute create --namespace laborious --name model_name --type Text\u001b[39;00m\n\u001b[32m---> \u001b[39m\u001b[32m13\u001b[39m \u001b[38;5;28;01mawait\u001b[39;00m temporal_client.create_schedule(\n\u001b[32m 14\u001b[39m \u001b[33m\"\u001b[39m\u001b[33morchestrator\u001b[39m\u001b[33m\"\u001b[39m,\n\u001b[32m 15\u001b[39m Schedule(\n\u001b[32m 16\u001b[39m action=ScheduleActionStartWorkflow(\n\u001b[32m 17\u001b[39m \u001b[33m'\u001b[39m\u001b[33morchestrator\u001b[39m\u001b[33m'\u001b[39m,\n\u001b[32m 18\u001b[39m {\n\u001b[32m 19\u001b[39m \u001b[33m\"\u001b[39m\u001b[33mschedule_name\u001b[39m\u001b[33m\"\u001b[39m: \u001b[33m\"\u001b[39m\u001b[33morchestrator-test\u001b[39m\u001b[33m\"\u001b[39m,\n\u001b[32m 20\u001b[39m \u001b[33m\"\u001b[39m\u001b[33mpipelines_query\u001b[39m\u001b[33m\"\u001b[39m: {\n\u001b[32m 21\u001b[39m \u001b[33m\"\u001b[39m\u001b[33mcollection\u001b[39m\u001b[33m\"\u001b[39m: \u001b[33m\"\u001b[39m\u001b[33mpipelines\u001b[39m\u001b[33m\"\u001b[39m,\n\u001b[32m 22\u001b[39m \u001b[33m\"\u001b[39m\u001b[33maggregation\u001b[39m\u001b[33m\"\u001b[39m: [\n\u001b[32m 23\u001b[39m {\n\u001b[32m 24\u001b[39m \u001b[33m\"\u001b[39m\u001b[33m$lookup\u001b[39m\u001b[33m\"\u001b[39m: {\n\u001b[32m 25\u001b[39m \u001b[33m\"\u001b[39m\u001b[33mfrom\u001b[39m\u001b[33m\"\u001b[39m: \u001b[33m\"\u001b[39m\u001b[33mmodels\u001b[39m\u001b[33m\"\u001b[39m,\n\u001b[32m 26\u001b[39m \u001b[33m\"\u001b[39m\u001b[33mlocalField\u001b[39m\u001b[33m\"\u001b[39m: \u001b[33m\"\u001b[39m\u001b[33mmodel_id\u001b[39m\u001b[33m\"\u001b[39m,\n\u001b[32m 27\u001b[39m \u001b[33m\"\u001b[39m\u001b[33mforeignField\u001b[39m\u001b[33m\"\u001b[39m: \u001b[33m\"\u001b[39m\u001b[33mid\u001b[39m\u001b[33m\"\u001b[39m,\n\u001b[32m 28\u001b[39m \u001b[33m\"\u001b[39m\u001b[33mas\u001b[39m\u001b[33m\"\u001b[39m: \u001b[33m\"\u001b[39m\u001b[33mmodel_docs\u001b[39m\u001b[33m\"\u001b[39m\n\u001b[32m 29\u001b[39m }\n\u001b[32m 30\u001b[39m },\n\u001b[32m 31\u001b[39m {\n\u001b[32m 32\u001b[39m \u001b[33m\"\u001b[39m\u001b[33m$match\u001b[39m\u001b[33m\"\u001b[39m: {\n\u001b[32m 33\u001b[39m \u001b[33m\"\u001b[39m\u001b[33mactive\u001b[39m\u001b[33m\"\u001b[39m: \u001b[38;5;28;01mTrue\u001b[39;00m\n\u001b[32m 34\u001b[39m }\n\u001b[32m 35\u001b[39m },\n\u001b[32m 36\u001b[39m {\n\u001b[32m 37\u001b[39m \u001b[33m\"\u001b[39m\u001b[33m$addFields\u001b[39m\u001b[33m\"\u001b[39m: {\n\u001b[32m 38\u001b[39m \u001b[33m\"\u001b[39m\u001b[33mmodels\u001b[39m\u001b[33m\"\u001b[39m: {\n\u001b[32m 39\u001b[39m \u001b[33m\"\u001b[39m\u001b[33m$arrayElemAt\u001b[39m\u001b[33m\"\u001b[39m: [\n\u001b[32m 40\u001b[39m \u001b[33m\"\u001b[39m\u001b[33m$model_docs\u001b[39m\u001b[33m\"\u001b[39m,\n\u001b[32m 41\u001b[39m \u001b[32m0\u001b[39m\n\u001b[32m 42\u001b[39m ]\n\u001b[32m 43\u001b[39m }\n\u001b[32m 44\u001b[39m }\n\u001b[32m 45\u001b[39m },\n\u001b[32m 46\u001b[39m {\n\u001b[32m 47\u001b[39m \u001b[33m\"\u001b[39m\u001b[33m$match\u001b[39m\u001b[33m\"\u001b[39m: {\n\u001b[32m 48\u001b[39m \u001b[33m\"\u001b[39m\u001b[33mmodels.active\u001b[39m\u001b[33m\"\u001b[39m: \u001b[38;5;28;01mTrue\u001b[39;00m\n\u001b[32m 49\u001b[39m }\n\u001b[32m 50\u001b[39m },\n\u001b[32m 51\u001b[39m {\n\u001b[32m 52\u001b[39m \u001b[33m\"\u001b[39m\u001b[33m$project\u001b[39m\u001b[33m\"\u001b[39m: {\n\u001b[32m 53\u001b[39m \u001b[33m\"\u001b[39m\u001b[33mmodel_docs\u001b[39m\u001b[33m\"\u001b[39m: \u001b[32m0\u001b[39m\n\u001b[32m 54\u001b[39m }\n\u001b[32m 55\u001b[39m }\n\u001b[32m 56\u001b[39m ]\n\u001b[32m 57\u001b[39m },\n\u001b[32m 58\u001b[39m \u001b[33m\"\u001b[39m\u001b[33mopc_servers_query\u001b[39m\u001b[33m\"\u001b[39m: {\n\u001b[32m 59\u001b[39m \u001b[33m\"\u001b[39m\u001b[33mcollection\u001b[39m\u001b[33m\"\u001b[39m: \u001b[33m\"\u001b[39m\u001b[33mopc-servers\u001b[39m\u001b[33m\"\u001b[39m,\n\u001b[32m 60\u001b[39m \u001b[33m\"\u001b[39m\u001b[33mfilters\u001b[39m\u001b[33m\"\u001b[39m: {\n\u001b[32m 61\u001b[39m \n\u001b[32m 62\u001b[39m }\n\u001b[32m 63\u001b[39m }\n\u001b[32m 64\u001b[39m },\n\u001b[32m 65\u001b[39m \u001b[38;5;28mid\u001b[39m=\u001b[33m\"\u001b[39m\u001b[33morchestrator\u001b[39m\u001b[33m\"\u001b[39m,\n\u001b[32m 66\u001b[39m task_queue=\u001b[33m\"\u001b[39m\u001b[33morchestrator-queue\u001b[39m\u001b[33m\"\u001b[39m,\n\u001b[32m 67\u001b[39m execution_timeout=timedelta(minutes=\u001b[32m600\u001b[39m)\n\u001b[32m 68\u001b[39m ),\n\u001b[32m 69\u001b[39m spec=ScheduleSpec(\n\u001b[32m 70\u001b[39m intervals=[ScheduleIntervalSpec(every=timedelta(minutes=\u001b[32m60\u001b[39m))]\n\u001b[32m 71\u001b[39m )\n\u001b[32m 72\u001b[39m )\n\u001b[32m 73\u001b[39m )\n", @@ -232,7 +137,7 @@ { "data": { "text/plain": [ - "" + "" ] }, "execution_count": 5, @@ -271,7 +176,7 @@ { "data": { "text/plain": [ - "" + "" ] }, "execution_count": 6, diff --git a/orchestrator/activities/mongo_db.py b/orchestrator/activities/mongo_db.py index e53f51a..592dd0d 100644 --- a/orchestrator/activities/mongo_db.py +++ b/orchestrator/activities/mongo_db.py @@ -9,7 +9,7 @@ with workflow.unsafe.imports_passed_through(): 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.temporal.constants import DATETIME_FORMAT_MS_WITH_TZ, now + from sientia_do.temporal.constants import DATETIME_FORMAT_MS, now from datetime import datetime @@ -442,7 +442,7 @@ class MongoDB(BaseActivity): data_filter = { **base_data_filter, "timestamp": { - "$gt": datetime.strptime(last_data_timestamp, DATETIME_FORMAT_MS_WITH_TZ) + "$gt": datetime.strptime(last_data_timestamp, DATETIME_FORMAT_MS) } } @@ -460,7 +460,7 @@ class MongoDB(BaseActivity): for item in data: item['timestamp'] = item['timestamp'].strftime( - DATETIME_FORMAT_MS_WITH_TZ) + DATETIME_FORMAT_MS) self.info( f"Loaded {len(data)} documents from MongoDB", diff --git a/values.yaml b/values.yaml index 300ea80..08c34a9 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.4" + tag: "0.4.5" # 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: @@ -151,7 +151,7 @@ env: - name: GITHUB_REPO_URL value: "git@github.com:Aignosi/sientia-dataops-orchestrator_temporal.git" - name: GITHUB_BRANCH - value: "SIENTIAPDE-1193-conferir-como-a-escrita-de-datetime-ocorre-no-temporal" + value: "main" - name: PYTHON_APP value: "orchestrator.worker.worker" From 93915841d6923f7569ed50e92e38163516837cf0 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Wed, 27 Aug 2025 09:25:17 -0300 Subject: [PATCH 2/6] SIENTIAPDE-1205 SIENTIAPDE-1205: Update GITHUB_BRANCH and standardize datetime handling in MongoDB activities - Changed GITHUB_BRANCH in values.yaml from "main" to "SIENTIAPDE-1205-alterar-opc-para-assincrono". - Updated datetime parsing in mongo_db.py to use DATETIME_FORMAT_MS_WITH_TZ for consistent timestamp handling. --- orchestrator/activities/mongo_db.py | 6 +++--- values.yaml | 2 +- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/orchestrator/activities/mongo_db.py b/orchestrator/activities/mongo_db.py index 592dd0d..e53f51a 100644 --- a/orchestrator/activities/mongo_db.py +++ b/orchestrator/activities/mongo_db.py @@ -9,7 +9,7 @@ with workflow.unsafe.imports_passed_through(): 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.temporal.constants import DATETIME_FORMAT_MS, now + from sientia_do.temporal.constants import DATETIME_FORMAT_MS_WITH_TZ, now from datetime import datetime @@ -442,7 +442,7 @@ class MongoDB(BaseActivity): data_filter = { **base_data_filter, "timestamp": { - "$gt": datetime.strptime(last_data_timestamp, DATETIME_FORMAT_MS) + "$gt": datetime.strptime(last_data_timestamp, DATETIME_FORMAT_MS_WITH_TZ) } } @@ -460,7 +460,7 @@ class MongoDB(BaseActivity): for item in data: item['timestamp'] = item['timestamp'].strftime( - DATETIME_FORMAT_MS) + DATETIME_FORMAT_MS_WITH_TZ) self.info( f"Loaded {len(data)} documents from MongoDB", diff --git a/values.yaml b/values.yaml index 08c34a9..442dd0d 100644 --- a/values.yaml +++ b/values.yaml @@ -151,7 +151,7 @@ env: - name: GITHUB_REPO_URL value: "git@github.com:Aignosi/sientia-dataops-orchestrator_temporal.git" - name: GITHUB_BRANCH - value: "main" + value: "SIENTIAPDE-1205-alterar-opc-para-assincrono" - name: PYTHON_APP value: "orchestrator.worker.worker" From 2a3e5b8bc4aaeb9be8ece8dcb2c566b37aa923c1 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Wed, 27 Aug 2025 09:32:42 -0300 Subject: [PATCH 3/6] SIENTIAPDE-1205 SIENTIAPDE-1205: Refactor timestamp handling in MongoDB activities - Updated the import statement in mongo_db.py to include timezone support. - Modified timestamp formatting to ensure UTC timezone is applied before string conversion, enhancing consistency in datetime handling. --- orchestrator/activities/mongo_db.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/orchestrator/activities/mongo_db.py b/orchestrator/activities/mongo_db.py index e53f51a..d456d7e 100644 --- a/orchestrator/activities/mongo_db.py +++ b/orchestrator/activities/mongo_db.py @@ -10,7 +10,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.temporal.constants import DATETIME_FORMAT_MS_WITH_TZ, now - from datetime import datetime + from datetime import datetime, timezone def clear_mongo_id(docs: list) -> list: @@ -459,7 +459,7 @@ class MongoDB(BaseActivity): ) for item in data: - item['timestamp'] = item['timestamp'].strftime( + item['timestamp'] = item['timestamp'].replace(tzinfo=timezone.utc).strftime( DATETIME_FORMAT_MS_WITH_TZ) self.info( From 90e9e571b2b75f6c346030bf428741356d028eb0 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Wed, 27 Aug 2025 09:35:55 -0300 Subject: [PATCH 4/6] SIENTIAPDE-1205 SIENTIAPDE-1205: Simplify timestamp formatting in MongoDB activities - Removed timezone conversion from timestamp handling in mongo_db.py, streamlining the formatting process while maintaining consistency with DATETIME_FORMAT_MS_WITH_TZ. --- orchestrator/activities/mongo_db.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/orchestrator/activities/mongo_db.py b/orchestrator/activities/mongo_db.py index d456d7e..e38a8eb 100644 --- a/orchestrator/activities/mongo_db.py +++ b/orchestrator/activities/mongo_db.py @@ -459,7 +459,7 @@ class MongoDB(BaseActivity): ) for item in data: - item['timestamp'] = item['timestamp'].replace(tzinfo=timezone.utc).strftime( + item['timestamp'] = item['timestamp'].strftime( DATETIME_FORMAT_MS_WITH_TZ) self.info( From fbf006bb99cd1a9dfe5a9a003a068870dae0ce2d Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Wed, 27 Aug 2025 09:37:45 -0300 Subject: [PATCH 5/6] SIENTIAPDE-1205 SIENTIAPDE-1205: Update timestamp handling in MongoDB activities - Modified timestamp formatting in mongo_db.py to apply UTC timezone before string conversion, ensuring consistency with DATETIME_FORMAT_MS_WITH_TZ. --- orchestrator/activities/mongo_db.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/orchestrator/activities/mongo_db.py b/orchestrator/activities/mongo_db.py index e38a8eb..d456d7e 100644 --- a/orchestrator/activities/mongo_db.py +++ b/orchestrator/activities/mongo_db.py @@ -459,7 +459,7 @@ class MongoDB(BaseActivity): ) for item in data: - item['timestamp'] = item['timestamp'].strftime( + item['timestamp'] = item['timestamp'].replace(tzinfo=timezone.utc).strftime( DATETIME_FORMAT_MS_WITH_TZ) self.info( From 8f49745b0ae1f9c1637b23aac0eb821be4f54d0c Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Thu, 28 Aug 2025 09:00:39 -0300 Subject: [PATCH 6/6] SIENTIAPDE-1205 SIENTIAPDE-1205: Update timestamp formatting in MongoDB tests - Modified timestamp handling in test_mongo_db.py to utilize DATETIME_FORMAT_MS_WITH_TZ for consistent formatting across test cases. - Ensured all timestamp strings include UTC timezone information for improved accuracy in datetime comparisons. --- run_coverage.sh | 11 +++++++++++ run_local.sh | 18 ++++++++++++++++++ tests/orchestrator/activities/test_mongo_db.py | 9 +++++---- 3 files changed, 34 insertions(+), 4 deletions(-) create mode 100755 run_coverage.sh create mode 100755 run_local.sh diff --git a/run_coverage.sh b/run_coverage.sh new file mode 100755 index 0000000..dcd534e --- /dev/null +++ b/run_coverage.sh @@ -0,0 +1,11 @@ +#!/bin/bash + +# Exit on any error +set -e + +echo "Activating virtual environment..." +source ./venv/bin/activate + +pytest --cov=orchestrator --cov-report=html + +xdg-open htmlcov/index.html \ No newline at end of file diff --git a/run_local.sh b/run_local.sh new file mode 100755 index 0000000..9aaa688 --- /dev/null +++ b/run_local.sh @@ -0,0 +1,18 @@ +#!/bin/bash + +# Exit on any error +set -e + +echo "Activating virtual environment..." +source ./venv/bin/activate + +echo "Loading environment variables from .env..." +if [ -f .env ]; then + export $(cat .env | grep -v '^#' | xargs) + echo "Environment variables loaded from .env" +else + echo "Warning: .env file not found. Continuing without environment variables." +fi + +echo "Starting orchestrator application..." +python -m orchestrator.app diff --git a/tests/orchestrator/activities/test_mongo_db.py b/tests/orchestrator/activities/test_mongo_db.py index 2ae912d..7dedc60 100644 --- a/tests/orchestrator/activities/test_mongo_db.py +++ b/tests/orchestrator/activities/test_mongo_db.py @@ -5,6 +5,7 @@ from pytest import fixture, mark from orchestrator.activities.mongo_db import clear_mongo_id from orchestrator.activities.mongo_db import MongoDB from sientia_do.notifications.models import NotificationLevel +from sientia_do.temporal.constants import DATETIME_FORMAT_MS_WITH_TZ def test_clear_mongo_id(): @@ -532,7 +533,7 @@ async def test_load_latest_data_none_last_data_timestamp(mongo_db): 'name': 'test1', 'value': 1, 'timestamp': 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) } ] @@ -558,7 +559,7 @@ async def test_load_latest_data_none_last_data_timestamp(mongo_db): assert result == [{ 'name': 'test1', 'value': 1, - 'timestamp': '2023-01-01 12:00:00.000000' + 'timestamp': '2023-01-01 12:00:00.000000+0000' }] @@ -573,7 +574,7 @@ async def test_load_latest_data_not_none_last_data_timestamp(mongo_db): 'name': 'test1', 'value': 1, 'timestamp': datetime.strptime( - '2023-01-01 12:00:00.000000+0000', '%Y-%m-%d %H:%M:%S.%f%z') + '2023-01-01 12:00:00.000000+0000', DATETIME_FORMAT_MS_WITH_TZ) } ] @@ -594,7 +595,7 @@ async def test_load_latest_data_not_none_last_data_timestamp(mongo_db): 'level': 'ERROR', 'timestamp': { '$gt': datetime.strptime( - '2023-01-01 12:00:00.000000+0000', '%Y-%m-%d %H:%M:%S.%f%z') + '2023-01-01 12:00:00.000000+0000', DATETIME_FORMAT_MS_WITH_TZ) } }, {"_id": 0}