diff --git a/orchestrator/activities/temporal_manager.py b/orchestrator/activities/temporal_manager.py index 7fd520e..fc4c5a1 100644 --- a/orchestrator/activities/temporal_manager.py +++ b/orchestrator/activities/temporal_manager.py @@ -1,3 +1,4 @@ +from asyncio import sleep from temporalio import activity, workflow from temporalio.client import ( Client, Schedule, ScheduleActionStartWorkflow, ScheduleIntervalSpec, ScheduleSpec, ScheduleUpdate, ScheduleUpdateInput) @@ -52,11 +53,21 @@ class TemporalManager(BaseActivity): if search_attrs.get("orchestrated", ["false"]) == ["true"]: schedule_id = schedule.id + self.logger.debug("Schedule id: %s", schedule_id) + handle = self.temporal_client.get_schedule_handle(schedule_id) + self.logger.debug("Handle acquired") + self.schedule_handles[schedule_id] = handle - desc = await handle.describe() + self.logger.debug("Describing schedule...") + + desc = await handle.describe( + rpc_timeout=timedelta(seconds=60) + ) + + self.logger.debug("Parsing args...") for arg in desc.schedule.action.args: data = MessageToDict(arg)['data'] @@ -69,6 +80,8 @@ class TemporalManager(BaseActivity): 'data': json.loads(data), } + await sleep(0.1) + self.logger.info("Found %d orchestrated schedules", len(orchestrated_schedules)) 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/orchestrator/workflows/orchestrator.py b/orchestrator/workflows/orchestrator.py index b731eeb..7bece5f 100644 --- a/orchestrator/workflows/orchestrator.py +++ b/orchestrator/workflows/orchestrator.py @@ -14,27 +14,6 @@ class Orchestrator: input_data['workflow_name'] = 'orchestrator' - await workflow.execute_local_activity_method( - Activities.prepare_activity, - { - 'workflow_name': input_data['workflow_name'], - 'schedule_name': input_data['schedule_name'], - 'model_name': '-', - 'model_id': '-' - }, - retry_policy=retry_policy, - start_to_close_timeout=timedelta(seconds=60) - ) - - # pipeline_config_handler = workflow.execute_local_activity_method( - # Activities.load_query_from_couchbase, - # { - # 'query': input_data['pipelines_query'] - # }, - # retry_policy=retry_policy, - # start_to_close_timeout=timedelta(seconds=60) - # ) - pipeline_config_handler = workflow.execute_local_activity_method( Activities.aggregate_documents_in_mongodb, { @@ -44,15 +23,6 @@ class Orchestrator: start_to_close_timeout=timedelta(seconds=60) ) - # opc_servers_handler = workflow.execute_local_activity_method( - # Activities.load_query_from_couchbase, - # { - # 'query': input_data['opc_servers_query'] - # }, - # retry_policy=retry_policy, - # start_to_close_timeout=timedelta(seconds=60) - # ) - opc_servers_handler = workflow.execute_local_activity_method( Activities.find_documents_in_mongodb, { @@ -65,7 +35,7 @@ class Orchestrator: orchestrated_schedules_handler = workflow.execute_local_activity_method( Activities.load_schedule, retry_policy=retry_policy, - start_to_close_timeout=timedelta(seconds=60) + start_to_close_timeout=timedelta(seconds=600) ) current_slot_config_handler = workflow.execute_local_activity_method( diff --git a/test.ipynb b/test.ipynb index ecf11d1..7a97f4b 100644 --- a/test.ipynb +++ b/test.ipynb @@ -110,7 +110,7 @@ }, { "cell_type": "code", - "execution_count": 1, + "execution_count": 14, "id": "7d01f160", "metadata": {}, "outputs": [], @@ -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,32 +163,154 @@ ")\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", + " \"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=\"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" ] }, + { + "cell_type": "code", + "execution_count": null, + "id": "b840347d", + "metadata": {}, + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + "Error deleting schedule scouter-load-test-1: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-2: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-3: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-4: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-5: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-6: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-7: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-8: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-9: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-10: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-11: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-12: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-13: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-14: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-15: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-16: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-17: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-18: workflow execution already completed\n", + "Error deleting schedule scouter-load-test-19: workflow execution already completed\n" + ] + } + ], + "source": [ + "# Delete schedule by ID\n", + "schedule_id = \"scouter-load-test-num\"\n", + "\n", + "for i in range(0, 70):\n", + "\n", + " try:\n", + " handle = temporal_client.get_schedule_handle(\n", + " schedule_id.replace(\"num\", str(i)))\n", + " \n", + " await handle.delete()\n", + " except Exception as e:\n", + " print(f\"Error deleting schedule {schedule_id.replace('num', str(i))}: {e}\")\n", + "\n", + "# Delete schedule by ID\n", + "schedule_id = \"laborious-load-test-num\"\n", + "\n", + "for i in range(0, 70):\n", + "\n", + " try:\n", + " handle = temporal_client.get_schedule_handle(\n", + " schedule_id.replace(\"num\", str(i)))\n", + " \n", + " await handle.delete()\n", + " except Exception as e:\n", + " print(f\"Error deleting schedule {schedule_id.replace('num', str(i))}: {e}\")\n" + ] + }, + { + "cell_type": "code", + "execution_count": 7, + "id": "374b5b0e", + "metadata": {}, + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + "scouter-load-test-num\n" + ] + } + ], + "source": [ + "schedule_id = \"scouter-load-test-num\"\n", + "schedule_id.replace(\"-\", \"_\")\n", + "\n", + "print(schedule_id)" + ] + }, { "cell_type": "code", "execution_count": 2, @@ -258,6 +380,29 @@ " " ] }, + { + "cell_type": "code", + "execution_count": 17, + "id": "f81b3728", + "metadata": {}, + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + "{'_client': , 'id': 'orchestrator'}\n" + ] + } + ], + "source": [ + "async for schedule in await temporal_client.list_schedules():\n", + " id = schedule.id\n", + "\n", + " handle = temporal_client.get_schedule_handle(id)\n", + "\n", + " print(vars(handle))" + ] + }, { "cell_type": "code", "execution_count": 13, diff --git a/tests/orchestrator/activities/test_activities.py b/tests/orchestrator/activities/test_activities.py index 4ebbc18..557ec8f 100644 --- a/tests/orchestrator/activities/test_activities.py +++ b/tests/orchestrator/activities/test_activities.py @@ -8,21 +8,13 @@ from orchestrator.activities.slot_manager import SlotManager from orchestrator.activities.formatters import Formatters -@patch('orchestrator.activities.couchbase.Couchbase.__init__') @patch('orchestrator.activities.mongo_db.MongoDB.__init__') @patch('orchestrator.activities.temporal_manager.TemporalManager.__init__') @patch('orchestrator.activities.slot_manager.SlotManager.__init__') @patch('orchestrator.activities.formatters.Formatters.__init__') def test___init__(mock_formatters_init, mock_slot_manager_init, mock_temporal_manager_init, - mock_mongodb_init, - mock_couchbase_init): - - couchbase_config = { - 'connection_string': 'couchbase://localhost', - 'username': 'admin', - 'password': 'password' - } + mock_mongodb_init): mongo_db_config = { 'connection_string': 'mongodb://localhost:27017', @@ -42,7 +34,6 @@ def test___init__(mock_formatters_init, mock_slot_manager_init, activities = Activities( temporal_client=temporal_client, - # couchbase_config=couchbase_config, redis_config=redis_config, mongodb_config=mongo_db_config, logger=logger, @@ -50,7 +41,6 @@ def test___init__(mock_formatters_init, mock_slot_manager_init, ) assert isinstance(activities, Activities) - # assert isinstance(activities, Couchbase) assert isinstance(activities, MongoDB) assert isinstance(activities, TemporalManager) assert isinstance(activities, SlotManager) @@ -66,15 +56,6 @@ def test___init__(mock_formatters_init, mock_slot_manager_init, notification_handler=notification_handler ) - # mock_couchbase_init.assert_called_once_with( - # ANY, - # connection_string=couchbase_config['connection_string'], - # username=couchbase_config['username'], - # password=couchbase_config['password'], - # logger=logger, - # notification_handler=notification_handler - # ) - mock_mongodb_init.assert_called_once_with( ANY, connection_string='mongodb://localhost:27017', diff --git a/tests/orchestrator/workflows/test_orchestrator.py b/tests/orchestrator/workflows/test_orchestrator.py index 3a131a4..c4fd4c3 100644 --- a/tests/orchestrator/workflows/test_orchestrator.py +++ b/tests/orchestrator/workflows/test_orchestrator.py @@ -20,20 +20,6 @@ async def test_run(workflow_mock, orchestrator): await orchestrator.run(input_data) - workflow_mock.execute_local_activity_method.assert_has_calls([ - call( - Activities.prepare_activity, - { - "workflow_name": "orchestrator", - "schedule_name": "test-schedule-name", - "model_name": "-", - "model_id": "-" - }, - retry_policy=ANY, - start_to_close_timeout=ANY - ) - ]) - workflow_mock.execute_local_activity_method.assert_has_calls([ call( Activities.aggregate_documents_in_mongodb, diff --git a/values.yaml b/values.yaml index e346cfd..7bb0447 100644 --- a/values.yaml +++ b/values.yaml @@ -123,7 +123,7 @@ env: - name: GITHUB_REPO_URL value: "git@github.com:Aignosi/sientia-dataops-orchestrator_temporal.git" - name: GITHUB_BRANCH - value: "SIENTIAPDE-1110-criar-testes-e-2-e" + value: "SIENTIAPDE-1141-criar-testes-de-carga" - name: PYTHON_APP value: "orchestrator.worker.worker"