Files
sientia-dataops-orchestrato…/test.ipynb
vitor-aignosi f43e962392 SIENTIAPDE-1141
fix: update execution counts and add schedule deletion logic in test.ipynb for improved error handling
2025-07-11 08:08:08 -03:00

554 lines
19 KiB
Plaintext

{
"cells": [
{
"cell_type": "code",
"execution_count": 4,
"id": "a287fa45",
"metadata": {},
"outputs": [
{
"name": "stdout",
"output_type": "stream",
"text": [
"Initializing Couchbase connection...\n",
"Awaiting Couchbase connection...\n",
"Couchbase connection ready\n"
]
}
],
"source": [
"from orchestrator.activities.couchbase import Couchbase\n",
"from unittest.mock import MagicMock\n",
"logger = MagicMock(info=MagicMock(side_effect=print), debug=MagicMock(side_effect=print))\n",
"couchbase = Couchbase(\n",
" connection_string=\"couchbase://localhost\",\n",
" username=\"sientia\",\n",
" password=\"sientia\",\n",
" logger=logger,\n",
" notification_handler=MagicMock()\n",
")\n"
]
},
{
"cell_type": "code",
"execution_count": 3,
"id": "d9a0f9c4",
"metadata": {},
"outputs": [
{
"name": "stdout",
"output_type": "stream",
"text": [
"Executing couchbase query: %s \n",
"SELECT\n",
" pipelines.*,\n",
" models as model\n",
"FROM\n",
" `pipelines`\n",
"JOIN\n",
" `models` ON KEYS pipelines.model_id;\n",
"\n",
"Fetched %d rows from couchbase 1\n",
"Rows: \n",
" %s [\n",
" {\n",
" \"filters\": [\n",
" {\n",
" \"filter_name\": \"OUT_OF_BOUNDS_FILTER\",\n",
" \"policy\": \"DISCARD\"\n",
" },\n",
" {\n",
" \"filter_name\": \"NULL_VALUES_FILTER\",\n",
" \"policy\": \"DISCARD\"\n",
" }\n",
" ],\n",
" \"frequency\": \"5s\",\n",
" \"max_retry_policy\": 1,\n",
" \"model\": {\n",
" \"name\": \"Demo Model-Demo2\"\n",
" },\n",
" \"model_id\": \"1\",\n",
" \"name\": \"scouter-opcua-pipeline\",\n",
" \"read_tags\": [\n",
" {\n",
" \"aggr_func\": \"avg\",\n",
" \"data_range\": [\n",
" -100,\n",
" 100\n",
" ],\n",
" \"tag_name\": \"Counter\"\n",
" }\n",
" ],\n",
" \"tag_retention_minutes\": 60,\n",
" \"workflow_type\": \"scouter\"\n",
" }\n",
"]\n",
"[{'filters': [{'filter_name': 'OUT_OF_BOUNDS_FILTER', 'policy': 'DISCARD'}, {'filter_name': 'NULL_VALUES_FILTER', 'policy': 'DISCARD'}], 'frequency': '5s', 'max_retry_policy': 1, 'model': {'name': 'Demo Model-Demo2'}, 'model_id': '1', 'name': 'scouter-opcua-pipeline', 'read_tags': [{'aggr_func': 'avg', 'data_range': [-100, 100], 'tag_name': 'Counter'}], 'tag_retention_minutes': 60, 'workflow_type': 'scouter'}]\n"
]
}
],
"source": [
"query = \"\"\"\n",
"SELECT\n",
" pipelines.*,\n",
" models as model\n",
"FROM\n",
" `pipelines`\n",
"JOIN\n",
" `models` ON KEYS pipelines.model_id;\n",
"\"\"\"\n",
"if __name__ == \"__main__\":\n",
"\n",
"\n",
" result = await couchbase.load_query_from_couchbase({\n",
" \"query\": query\n",
" })\n",
"\n",
" print(result)\n",
" "
]
},
{
"cell_type": "code",
"execution_count": 14,
"id": "7d01f160",
"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",
")\n",
"\n",
"manager = TemporalManager(temporal_client=temporal_client,\n",
" logger=logger,\n",
" notification_handler=MagicMock())\n",
" "
]
},
{
"cell_type": "code",
"execution_count": 4,
"id": "bb750ae6",
"metadata": {},
"outputs": [
{
"data": {
"text/plain": [
"<temporalio.client.ScheduleHandle at 0x7700871c6150>"
]
},
"execution_count": 4,
"metadata": {},
"output_type": "execute_result"
}
],
"source": [
"import asyncio\n",
"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 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",
"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",
")\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,
"id": "1bd82225",
"metadata": {},
"outputs": [
{
"name": "stdout",
"output_type": "stream",
"text": [
"Getting orchestrated schedules...\n"
]
}
],
"source": [
"schedules = await manager.load_schedule()"
]
},
{
"cell_type": "code",
"execution_count": 24,
"id": "1948670e",
"metadata": {},
"outputs": [],
"source": [
"from google.protobuf.json_format import MessageToDict\n",
"import base64\n",
"import json\n",
"for arg in schedules.schedule.action.args:\n",
" data = MessageToDict(arg)\n",
" data = base64.b64decode(data['data']).decode('utf-8')"
]
},
{
"cell_type": "code",
"execution_count": 12,
"id": "a4c777dd",
"metadata": {},
"outputs": [],
"source": [
"from google.protobuf.json_format import MessageToDict\n",
"import base64\n",
"import json\n",
"\n",
"schedules_config = {}\n",
"async for schedule in await temporal_client.list_schedules():\n",
" id = schedule.id\n",
"\n",
" handle = temporal_client.get_schedule_handle(id)\n",
"\n",
" desc = await handle.describe()\n",
"\n",
" for arg in desc.schedule.action.args:\n",
" data = MessageToDict(arg)['data']\n",
" data = base64.b64decode(data).decode('utf-8')\n",
"\n",
" frequency = desc.schedule.spec.intervals[0].every.seconds\n",
"\n",
" schedules_config[id] = {\n",
" 'frequency': frequency,\n",
" 'data': json.loads(data),\n",
" 'handle': handle\n",
" }\n",
" \n",
" \n",
" \n",
" "
]
},
{
"cell_type": "code",
"execution_count": 17,
"id": "f81b3728",
"metadata": {},
"outputs": [
{
"name": "stdout",
"output_type": "stream",
"text": [
"{'_client': <temporalio.client.Client object at 0x770086544f90>, '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,
"id": "988d1718",
"metadata": {},
"outputs": [],
"source": [
"schedules_config = {\n",
" \"scouter-opcua-orchestrated-pipeline\": schedules_config['scouter-opcua-orchestrated-pipeline'],\n",
"}"
]
},
{
"cell_type": "code",
"execution_count": 14,
"id": "c12f5e75",
"metadata": {},
"outputs": [
{
"name": "stdout",
"output_type": "stream",
"text": [
"scouter-opcua-orchestrated-pipeline {'frequency': 5, 'data': {'filters': {'NULL_VALUES_FILTER': {'policy': 'DISCARD'}, 'OUT_OF_BOUNDS_FILTER': {'policy': 'DISCARD'}}, 'frequency': '5s', 'max_retry_policy': 1, 'model_id': '1', 'model_name': 'Demo Model-Demo2', 'model_tags': {'Counter': {'aggr_func': 'avg', 'data_range': [-100, 100]}}, 'retention_time': 3600, 'schedule_name': 'scouter-opcua-orchestrated-pipeline', 'schema': 'sientia_data', 'table_name': 'laborious_data', 'topic': 'raw_scouter-opcua-orchestrated-pipeline', 'trigger_laborious': False, 'workflow_type': 'scouter'}, 'handle': <temporalio.client.ScheduleHandle object at 0x7dfab3b2ab50>}\n",
"scouter-opcua-orchestrated-pipeline\n",
"{'workflow': 'scouter', 'args': [metadata {\n",
" key: \"encoding\"\n",
" value: \"json/plain\"\n",
"}\n",
"data: \"{\\\"filters\\\":{\\\"NULL_VALUES_FILTER\\\":{\\\"policy\\\":\\\"DISCARD\\\"},\\\"OUT_OF_BOUNDS_FILTER\\\":{\\\"policy\\\":\\\"DISCARD\\\"}},\\\"frequency\\\":\\\"5s\\\",\\\"max_retry_policy\\\":1,\\\"model_id\\\":\\\"1\\\",\\\"model_name\\\":\\\"Demo Model-Demo2\\\",\\\"model_tags\\\":{\\\"Counter\\\":{\\\"aggr_func\\\":\\\"avg\\\",\\\"data_range\\\":[-100,100]}},\\\"retention_time\\\":3600,\\\"schedule_name\\\":\\\"scouter-opcua-orchestrated-pipeline\\\",\\\"schema\\\":\\\"sientia_data\\\",\\\"table_name\\\":\\\"laborious_data\\\",\\\"topic\\\":\\\"raw_scouter-opcua-orchestrated-pipeline\\\",\\\"trigger_laborious\\\":false,\\\"workflow_type\\\":\\\"scouter\\\"}\"\n",
"], 'id': 'scouter-opcua-orchestrated-pipeline', 'task_queue': 'scouter-queue', 'execution_timeout': datetime.timedelta(seconds=120), 'run_timeout': None, 'task_timeout': None, 'retry_policy': None, 'memo': None, 'typed_search_attributes': TypedSearchAttributes(search_attributes=[]), 'headers': None, 'untyped_search_attributes': {}, 'static_summary': None, 'static_details': None, 'priority': Priority(priority_key=None)}\n"
]
}
],
"source": [
"from temporalio.client import Client, ScheduleUpdateInput, ScheduleUpdate, ScheduleSpec\n",
"\n",
"\n",
"for id, schedule in schedules_config.items():\n",
" print(id, schedule)\n",
" \n",
" if schedule['frequency'] != 600:\n",
" print(id)\n",
"\n",
" handle = schedule['handle']\n",
" \n",
" async def update_schedule(input: ScheduleUpdateInput) -> ScheduleUpdate:\n",
" schedule_action = input.description.schedule.action\n",
"\n",
" print(schedule_action.__dict__)\n",
" \n",
" if hasattr(schedule_action, 'args'):\n",
" schedule_action.args = [{}]\n",
" \n",
" # Atualiza o intervalo de execução\n",
" input.description.schedule.spec.intervals = [\n",
" ScheduleIntervalSpec(every=timedelta(minutes=10))\n",
" ]\n",
" \n",
" return ScheduleUpdate(schedule=input.description.schedule)\n",
"\n",
" await handle.update(update_schedule)\n",
" \n"
]
},
{
"cell_type": "code",
"execution_count": 1,
"id": "87eb10c8",
"metadata": {},
"outputs": [],
"source": [
"from redis import Redis\n",
"\n",
"redis = Redis(host='localhost', port=6379)"
]
},
{
"cell_type": "code",
"execution_count": 6,
"id": "fe4ad9dc",
"metadata": {},
"outputs": [
{
"name": "stdout",
"output_type": "stream",
"text": [
"Chave: slot:opc_tags:2, Valor: {\n",
"\"slot2\": \"value\"\n",
"}\n",
"Chave: slot:opc_tags:1, Valor: {\n",
"\"slot1\": \"value\"\n",
"}\n"
]
}
],
"source": [
"matching_keys = redis.keys(\"slot:opc_tags:*\")\n",
"\n",
"if matching_keys:\n",
" decoded_keys = [key.decode('utf-8') for key in matching_keys]\n",
" values = redis.mget(decoded_keys)\n",
"\n",
" items = {}\n",
" for i, key in enumerate(decoded_keys):\n",
" value = values[i]\n",
" if value is not None:\n",
" try:\n",
" items[key] = value.decode('utf-8')\n",
" except (UnicodeDecodeError, AttributeError):\n",
" items[key] = value\n",
" else:\n",
" items[key] = None\n",
"\n",
" for key, value in items.items():\n",
" print(f\"Chave: {key}, Valor: {value}\")"
]
},
{
"cell_type": "code",
"execution_count": null,
"id": "87a9dc89",
"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
}