diff --git a/init_orchestration.ipynb b/init_orchestration.ipynb new file mode 100644 index 0000000..14f66c0 --- /dev/null +++ b/init_orchestration.ipynb @@ -0,0 +1,332 @@ +{ + "cells": [ + { + "cell_type": "code", + "execution_count": 2, + "id": "3f8b77a4", + "metadata": {}, + "outputs": [], + "source": [ + "from temporalio import client\n", + "from orchestrator.activities.temporal_manager import TemporalManager\n", + "import os\n", + "from unittest.mock import MagicMock\n", + "\n", + "host = \"localhost:7233\"\n", + "logger = MagicMock(info=MagicMock(side_effect=print), debug=MagicMock(side_effect=print))\n", + "\n", + "temporal_client = await client.Client.connect(\n", + " target_host=host,\n", + " namespace=os.getenv('TEMPORAL_NAMESPACE', 'default')\n", + ")" + ] + }, + { + "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, + "id": "d9d2a242", + "metadata": {}, + "outputs": [ + { + "ename": "ScheduleAlreadyRunningError", + "evalue": "Schedule already running", + "output_type": "error", + "traceback": [ + "\u001b[31m---------------------------------------------------------------------------\u001b[39m", + "\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')", + "\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.", + "\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", + "\u001b[36mFile \u001b[39m\u001b[32m~/Documents/projects/sientia/sientia-dataops-orchestrator_temporal/venv/lib/python3.11/site-packages/temporalio/client.py:1308\u001b[39m, in \u001b[36mClient.create_schedule\u001b[39m\u001b[34m(self, id, schedule, trigger_immediately, backfill, memo, search_attributes, static_summary, static_details, rpc_metadata, rpc_timeout)\u001b[39m\n\u001b[32m 1275\u001b[39m \u001b[38;5;250m\u001b[39m\u001b[33;03m\"\"\"Create a schedule and return its handle.\u001b[39;00m\n\u001b[32m 1276\u001b[39m \n\u001b[32m 1277\u001b[39m \u001b[33;03mArgs:\u001b[39;00m\n\u001b[32m (...)\u001b[39m\u001b[32m 1305\u001b[39m \u001b[33;03m running.\u001b[39;00m\n\u001b[32m 1306\u001b[39m \u001b[33;03m\"\"\"\u001b[39;00m\n\u001b[32m 1307\u001b[39m temporalio.common._warn_on_deprecated_search_attributes(search_attributes)\n\u001b[32m-> \u001b[39m\u001b[32m1308\u001b[39m \u001b[38;5;28;01mreturn\u001b[39;00m \u001b[38;5;28;01mawait\u001b[39;00m \u001b[38;5;28mself\u001b[39m._impl.create_schedule(\n\u001b[32m 1309\u001b[39m CreateScheduleInput(\n\u001b[32m 1310\u001b[39m \u001b[38;5;28mid\u001b[39m=\u001b[38;5;28mid\u001b[39m,\n\u001b[32m 1311\u001b[39m schedule=schedule,\n\u001b[32m 1312\u001b[39m trigger_immediately=trigger_immediately,\n\u001b[32m 1313\u001b[39m backfill=backfill,\n\u001b[32m 1314\u001b[39m memo=memo,\n\u001b[32m 1315\u001b[39m search_attributes=search_attributes,\n\u001b[32m 1316\u001b[39m rpc_metadata=rpc_metadata,\n\u001b[32m 1317\u001b[39m rpc_timeout=rpc_timeout,\n\u001b[32m 1318\u001b[39m )\n\u001b[32m 1319\u001b[39m )\n", + "\u001b[36mFile \u001b[39m\u001b[32m~/Documents/projects/sientia/sientia-dataops-orchestrator_temporal/venv/lib/python3.11/site-packages/temporalio/client.py:6445\u001b[39m, in \u001b[36m_ClientImpl.create_schedule\u001b[39m\u001b[34m(self, input)\u001b[39m\n\u001b[32m 6437\u001b[39m already_started = (\n\u001b[32m 6438\u001b[39m err.status == RPCStatusCode.ALREADY_EXISTS\n\u001b[32m 6439\u001b[39m \u001b[38;5;129;01mand\u001b[39;00m err.grpc_status.details\n\u001b[32m (...)\u001b[39m\u001b[32m 6442\u001b[39m )\n\u001b[32m 6443\u001b[39m )\n\u001b[32m 6444\u001b[39m \u001b[38;5;28;01mif\u001b[39;00m already_started:\n\u001b[32m-> \u001b[39m\u001b[32m6445\u001b[39m \u001b[38;5;28;01mraise\u001b[39;00m ScheduleAlreadyRunningError()\n\u001b[32m 6446\u001b[39m \u001b[38;5;28;01mraise\u001b[39;00m\n\u001b[32m 6447\u001b[39m \u001b[38;5;28;01mreturn\u001b[39;00m ScheduleHandle(\u001b[38;5;28mself\u001b[39m._client, \u001b[38;5;28minput\u001b[39m.id)\n", + "\u001b[31mScheduleAlreadyRunningError\u001b[39m: Schedule already running" + ] + } + ], + "source": [ + "from datetime import timedelta\n", + "from temporalio.client import (\n", + " Schedule,\n", + " ScheduleActionStartWorkflow,\n", + " ScheduleIntervalSpec,\n", + " ScheduleSpec,\n", + ")\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": 5, + "id": "6ea0f616", + "metadata": {}, + "outputs": [ + { + "data": { + "text/plain": [ + "" + ] + }, + "execution_count": 5, + "metadata": {}, + "output_type": "execute_result" + } + ], + "source": [ + "await temporal_client.create_schedule(\n", + " \"alerts\",\n", + " Schedule(\n", + " action=ScheduleActionStartWorkflow(\n", + " 'alerts',\n", + " {\n", + " \"schedule_name\": \"alerts\",\n", + " \"notification_ttl\": 5*60,\n", + " \"sent_ttl\": 10*60\n", + " },\n", + " id=\"alerts\",\n", + " task_queue=\"alerts-queue\",\n", + " execution_timeout=timedelta(minutes=600)\n", + " ),\n", + " spec=ScheduleSpec(\n", + " intervals=[ScheduleIntervalSpec(every=timedelta(seconds=30))]\n", + " )\n", + " )\n", + ")" + ] + }, + { + "cell_type": "code", + "execution_count": 6, + "id": "05c46ec8", + "metadata": {}, + "outputs": [ + { + "data": { + "text/plain": [ + "" + ] + }, + "execution_count": 6, + "metadata": {}, + "output_type": "execute_result" + } + ], + "source": [ + "await temporal_client.create_schedule(\n", + " \"reports\",\n", + " Schedule(\n", + " action=ScheduleActionStartWorkflow(\n", + " 'reports',\n", + " {\n", + " \"schedule_name\": \"reports\",\n", + " },\n", + " id=\"reports\",\n", + " task_queue=\"reports-queue\",\n", + " execution_timeout=timedelta(minutes=600)\n", + " ),\n", + " spec=ScheduleSpec(\n", + " intervals=[ScheduleIntervalSpec(every=timedelta(minutes=20))]\n", + " )\n", + " )\n", + ")" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "76e3d8a7", + "metadata": {}, + "outputs": [], + "source": [] + } + ], + "metadata": { + "kernelspec": { + "display_name": "venv", + "language": "python", + "name": "python3" + }, + "language_info": { + "codemirror_mode": { + "name": "ipython", + "version": 3 + }, + "file_extension": ".py", + "mimetype": "text/x-python", + "name": "python", + "nbconvert_exporter": "python", + "pygments_lexer": "ipython3", + "version": "3.11.13" + } + }, + "nbformat": 4, + "nbformat_minor": 5 +} diff --git a/orchestrator/activities/mongo_db.py b/orchestrator/activities/mongo_db.py index 105264e..4f9ba78 100644 --- a/orchestrator/activities/mongo_db.py +++ b/orchestrator/activities/mongo_db.py @@ -1,4 +1,3 @@ -from pandas import DataFrame from temporalio import workflow, activity @@ -257,7 +256,8 @@ class MongoDB(BaseActivity): data_filter = argument if argument else {} try: - collection.insert_many(data_filter) + if data_filter: + collection.insert_many(data_filter) except Exception as e: trace = traceback.format_exc() self.send_notification(