diff --git a/orchestrator/worker/worker.py b/orchestrator/worker/worker.py index 383c1be..8aa019f 100644 --- a/orchestrator/worker/worker.py +++ b/orchestrator/worker/worker.py @@ -41,7 +41,6 @@ async def main(): activities = Activities( temporal_client=temporal_client, - # couchbase_config=build_couchbase_config(), redis_config=build_redis_config(), mongodb_config=build_mongodb_config(), logger=logger, diff --git a/test.ipynb b/test.ipynb index ecf11d1..948e129 100644 --- a/test.ipynb +++ b/test.ipynb @@ -120,7 +120,7 @@ "import os\n", "from unittest.mock import MagicMock\n", "\n", - "host = \"localhost:7233\"\n", + "host = \"localhost:44795\"\n", "logger = MagicMock(info=MagicMock(side_effect=print), debug=MagicMock(side_effect=print))\n", "\n", "temporal_client = await client.Client.connect(\n", @@ -136,17 +136,17 @@ }, { "cell_type": "code", - "execution_count": 6, + "execution_count": 4, "id": "bb750ae6", "metadata": {}, "outputs": [ { "data": { "text/plain": [ - "" + "" ] }, - "execution_count": 6, + "execution_count": 4, "metadata": {}, "output_type": "execute_result" } @@ -163,29 +163,70 @@ ")\n", "from temporalio.common import TypedSearchAttributes, SearchAttributeKey, SearchAttributePair\n", "\n", + "# temporal operator search-attribute create --namespace default --name model_id --type Text\n", + "# temporal operator search-attribute create --namespace default --name orchestrated --type Text\n", + "# temporal operator search-attribute create --namespace default --name model_name --type Text\n", "\n", - "customer_id_key = SearchAttributeKey.for_keyword(\"orchestrated\")\n", - "search_attributes = TypedSearchAttributes([\n", - " SearchAttributePair(customer_id_key, \"true\")\n", - "])\n", "await temporal_client.create_schedule(\n", - " \"meu-schedule-id5\",\n", + " \"orchestrator\",\n", " Schedule(\n", " action=ScheduleActionStartWorkflow(\n", - " 'scouter-test2',\n", + " 'orchestrator',\n", " {\n", - " 'args': {\n", - " 'arg1': 'value1'\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", + " \"models.active\": True\n", + " }\n", + " },\n", + " {\n", + " \"$match\": {\n", + " \"pipelines.active\": True\n", + " }\n", + " },\n", + " {\n", + " \"$addFields\": {\n", + " \"models\": {\n", + " \"$arrayElemAt\": [\n", + " \"$model_docs\",\n", + " 0\n", + " ]\n", + " }\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=\"workflow-id-unico\",\n", - " task_queue=\"nome-da-task-queue\",\n", + " id=\"orchestrator\",\n", + " task_queue=\"orchestrator-queue\",\n", + " execution_timeout=timedelta(minutes=600)\n", " ),\n", " spec=ScheduleSpec(\n", - " intervals=[ScheduleIntervalSpec(every=timedelta(minutes=10))]\n", + " intervals=[ScheduleIntervalSpec(every=timedelta(minutes=60))]\n", " )\n", - " ),\n", - " search_attributes=search_attributes,\n", + " )\n", ")\n" ] },