From 2a5534cf64f6b4262fbc10e411c186fab3332146 Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Fri, 25 Jul 2025 16:03:39 -0300 Subject: [PATCH] SIENTIAPDE-1172 feat: enhance orchestrator activities with new email and Postgres integrations - Added Email and Postgres classes to the Activities class for improved functionality. - Introduced new methods in MongoDB and SlotManager for loading and managing data. - Updated requirements.txt to include jinja2. - Added new formatting activity for log reports in Formatters class. - Enhanced test coverage for MongoDB and SlotManager activities. --- email.html | 247 +++++++++++++++ notification_generator.ipynb | 287 ++++++++++++++++++ orchestrator/activities/activities.py | 26 +- orchestrator/activities/email.py | 149 +++++++++ orchestrator/activities/formatters.py | 45 +++ orchestrator/activities/mongo_db.py | 89 +++++- orchestrator/activities/slot_manager.py | 70 +++++ orchestrator/utils/email_builder.py | 75 +++++ .../utils/templates/email_template.html | 24 ++ .../utils/templates/general_template.html | 24 ++ orchestrator/workflows/alerts.py | 22 ++ orchestrator/workflows/orchestrator.py | 9 + orchestrator/workflows/reports.py | 19 ++ .../subworkflows/load_notification_package.py | 93 ++++++ .../subworkflows/process_notifications.py | 87 ++++++ requirements.txt | 1 + .../orchestrator/activities/test_mongo_db.py | 131 ++++++++ .../activities/test_slot_manager.py | 153 ++++++++++ .../test_load_notification_package.py | 115 +++++++ 19 files changed, 1660 insertions(+), 6 deletions(-) create mode 100644 email.html create mode 100644 notification_generator.ipynb create mode 100644 orchestrator/activities/email.py create mode 100644 orchestrator/utils/email_builder.py create mode 100644 orchestrator/utils/templates/email_template.html create mode 100644 orchestrator/utils/templates/general_template.html create mode 100644 orchestrator/workflows/alerts.py create mode 100644 orchestrator/workflows/reports.py create mode 100644 orchestrator/workflows/subworkflows/load_notification_package.py create mode 100644 orchestrator/workflows/subworkflows/process_notifications.py create mode 100644 tests/orchestrator/workflows/subworkflows/test_load_notification_package.py diff --git a/email.html b/email.html new file mode 100644 index 0000000..7f53648 --- /dev/null +++ b/email.html @@ -0,0 +1,247 @@ + + + + + + SIENTIA™ Report + + + +

SIENTIA™ Alerts

+ +

Errors detected:

+ +

Model: ut aliqua ut ipsum

+ + + + + + + + + + + + + + + + + + + +
Notification IDBlockTimestampMessage
TAG_node5:tag_87_LISTENNING_STOPPEDdolor_elit_aliqua2025-07-25 16:50:48.925421+00:00sit elit ipsum labore incididunt magna do et sit incididunt
+ +

Model: elit dolore

+ + + + + + + + + + + + + + + + + + + +
Notification IDBlockTimestampMessage
OPC_LISTENNING_BACK__server_10magna_magna_et2025-07-25 16:50:48.925478+00:00aliqua eiusmod ut eiusmod ut et ipsum
+ +

Model: lorem lorem ipsum amet

+ + + + + + + + + + + + + + + + + + + +
Notification IDBlockTimestampMessage
TAG_node8:tag_11_LISTENNING_STOPPEDdo_sit2025-07-25 16:50:48.925506+00:00labore do adipiscing adipiscing incididunt eiusmod
+ +

Warnings detected:

+ +

Model: amet eiusmod dolor ipsum

+ + + + + + + + + + + + + + + + + + + +
Notification IDBlockTimestampMessage
REPORT_PARTITION_MANAGERadipiscing_amet_ut_et2025-07-25 16:50:48.925368+00:00eiusmod et amet sit lorem incididunt
+ +

Model: sit et labore

+ + + + + + + + + + + + + + + + + + + +
Notification IDBlockTimestampMessage
SIT_CONSECTETUR_EIUSMODadipiscing_elit_ut2025-07-25 16:50:48.925392+00:00magna do eiusmod et labore elit
+ +

Model: consectetur labore do consectetur

+ + + + + + + + + + + + + + + + + + + +
Notification IDBlockTimestampMessage
SED_SIT_DOLOR_CONSECTETURdolor_dolore_ut2025-07-25 16:50:48.925435+00:00magna sed incididunt dolor tempor
+ +

Model: et sed ut et

+ + + + + + + + + + + + + + + + + + + +
Notification IDBlockTimestampMessage
UT_UT_LABOREsit_do2025-07-25 16:50:48.925492+00:00do sed sed elit do eiusmod
+ +

Infos detected:

+ +

Model: eiusmod eiusmod

+ + + + + + + + + + + + + + + + + + + +
Notification IDBlockTimestampMessage
OPC_CONNECTION_RETRY__server_2labore_sit_do2025-07-25 16:50:48.925407+00:00magna et ipsum amet dolore adipiscing adipiscing do
+ +

Model: adipiscing consectetur

+ + + + + + + + + + + + + + + + + + + +
Notification IDBlockTimestampMessage
REPORT_PARTITION_MANAGERincididunt_aliqua_ipsum_ut2025-07-25 16:50:48.925447+00:00labore magna adipiscing tempor elit et et magna elit lorem
+ +

Model: et amet

+ + + + + + + + + + + + + + + + + + + +
Notification IDBlockTimestampMessage
OPC_LISTENNING_BACK__server_2sed_elit_ut_do2025-07-25 16:50:48.925462+00:00dolor et adipiscing adipiscing dolore ipsum ipsum aliqua et elit
+ + + + + \ No newline at end of file diff --git a/notification_generator.ipynb b/notification_generator.ipynb new file mode 100644 index 0000000..197e89b --- /dev/null +++ b/notification_generator.ipynb @@ -0,0 +1,287 @@ +{ + "cells": [ + { + "cell_type": "code", + "execution_count": 2, + "id": "1786cbf0", + "metadata": {}, + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + "Generated 10 notifications\n" + ] + } + ], + "source": [ + "from unittest.mock import MagicMock\n", + "from orchestrator.activities.mongo_db import MongoDB\n", + "from orchestrator.utils.connectors_config import build_mongodb_config\n", + "import random\n", + "from datetime import datetime, timezone\n", + "import json\n", + "\n", + "word_list = [\n", + " 'lorem', 'ipsum', 'dolor', 'sit', 'amet', 'consectetur',\n", + " 'adipiscing', 'elit', 'sed', 'do', 'eiusmod', 'tempor',\n", + " 'incididunt', 'ut', 'labore', 'et', 'dolore', 'magna', 'aliqua'\n", + " ]\n", + "\n", + "config = build_mongodb_config()\n", + "\n", + "mongo_db = MongoDB(\n", + " connection_string=config['connection_string'],\n", + " database_name=config['database_name'],\n", + " ttl_index_seconds=config['ttl_index_seconds'],\n", + " logger=MagicMock(),\n", + " notification_handler=MagicMock()\n", + ")\n", + "\n", + "AMOUT = 10\n", + "\n", + "notification_ids = [\n", + " \"TAG_node:_LISTENNING_STOPPED\",\n", + " \"OPC___\",\n", + " \"REPORT_PARTITION_MANAGER\",\n", + " \"\"\n", + "]\n", + "\n", + "notifications = []\n", + "\n", + "for _ in range(AMOUT):\n", + " # Choose random notification ID and fill in placeholders\n", + " notification_id = random.choice(notification_ids)\n", + " \n", + " if \"\" in notification_id:\n", + " notification_id = notification_id.replace(\"\", str(random.randint(1,100)))\n", + " if \"\" in notification_id:\n", + " notification_id = notification_id.replace(\"\", f\"tag_{random.randint(1,100)}\")\n", + " if \"\" in notification_id:\n", + " events = [\"LISTENNING_STOPPED\", \"CONNECTION_RETRY\", \"LISTENNING_BACK\", \"SUBSCRIPTION\", \"QUEUE_RETRIEVE\"]\n", + " notification_id = notification_id.replace(\"\", random.choice(events))\n", + " if \"\" in notification_id:\n", + " notification_id = notification_id.replace(\"\", f\"server_{random.randint(1,10)}\")\n", + " if \"\" in notification_id:\n", + " notification_id = notification_id.replace(\"\", \"_\".join(random.choices(word_list, k=random.randint(3,5))).upper())\n", + "\n", + " notification = {\n", + " \"project\": \"sientia-laborious\",\n", + " \"pipeline\": \"_\".join(random.choices(word_list, k=random.randint(2,4))), \n", + " \"trigger\": \"_\".join(random.choices(word_list, k=random.randint(3,6))),\n", + " \"model_name\": \" \".join(random.choices(word_list, k=random.randint(2,4))),\n", + " \"model_id\": str(random.randint(1,100)),\n", + " \"block\": \"_\".join(random.choices(word_list, k=random.randint(2,4))),\n", + " \"level\": random.choice([\"INFO\", \"WARNING\", \"ERROR\"]),\n", + " \"message\": \" \".join(random.choices(word_list, k=random.randint(5,10))),\n", + " \"attachment_content\": \" \".join(random.choices(word_list, k=random.randint(20,30))),\n", + " \"notification_id\": notification_id,\n", + " \"timestamp\": datetime.now(timezone.utc)\n", + " }\n", + " \n", + " notifications.append(notification)\n", + "\n", + "print(f\"Generated {len(notifications)} notifications\")" + ] + }, + { + "cell_type": "code", + "execution_count": 3, + "id": "c1467b72", + "metadata": {}, + "outputs": [ + { + "name": "stdout", + "output_type": "stream", + "text": [ + "{'WARNING': {'section_name': 'Warnings detected:', 'models': {'amet eiusmod dolor ipsum': {'model_name': 'amet eiusmod dolor ipsum', 'events': [{'project': 'sientia-laborious', 'pipeline': 'ut_eiusmod_do', 'trigger': 'aliqua_magna_labore_magna_do_ut', 'model_name': 'amet eiusmod dolor ipsum', 'model_id': '26', 'block': 'adipiscing_amet_ut_et', 'level': 'WARNING', 'message': 'eiusmod et amet sit lorem incididunt', 'attachment_content': 'sed et ipsum aliqua et consectetur incididunt aliqua incididunt aliqua eiusmod adipiscing incididunt sed et ut ipsum aliqua elit ut', 'notification_id': 'REPORT_PARTITION_MANAGER', 'timestamp': datetime.datetime(2025, 7, 25, 16, 50, 48, 925368, tzinfo=datetime.timezone.utc)}]}, 'sit et labore': {'model_name': 'sit et labore', 'events': [{'project': 'sientia-laborious', 'pipeline': 'tempor_eiusmod_amet_amet', 'trigger': 'labore_elit_adipiscing_adipiscing', 'model_name': 'sit et labore', 'model_id': '8', 'block': 'adipiscing_elit_ut', 'level': 'WARNING', 'message': 'magna do eiusmod et labore elit', 'attachment_content': 'incididunt incididunt do adipiscing consectetur adipiscing ipsum incididunt et sed sed magna dolore ut ipsum amet dolor sed tempor lorem incididunt do labore elit aliqua magna eiusmod dolor', 'notification_id': 'SIT_CONSECTETUR_EIUSMOD', 'timestamp': datetime.datetime(2025, 7, 25, 16, 50, 48, 925392, tzinfo=datetime.timezone.utc)}]}, 'consectetur labore do consectetur': {'model_name': 'consectetur labore do consectetur', 'events': [{'project': 'sientia-laborious', 'pipeline': 'et_consectetur_ut', 'trigger': 'labore_aliqua_tempor_amet', 'model_name': 'consectetur labore do consectetur', 'model_id': '50', 'block': 'dolor_dolore_ut', 'level': 'WARNING', 'message': 'magna sed incididunt dolor tempor', 'attachment_content': 'lorem labore tempor amet adipiscing lorem dolor aliqua do et incididunt adipiscing labore aliqua dolor elit aliqua eiusmod sed aliqua aliqua consectetur incididunt sed ipsum eiusmod labore do', 'notification_id': 'SED_SIT_DOLOR_CONSECTETUR', 'timestamp': datetime.datetime(2025, 7, 25, 16, 50, 48, 925435, tzinfo=datetime.timezone.utc)}]}, 'et sed ut et': {'model_name': 'et sed ut et', 'events': [{'project': 'sientia-laborious', 'pipeline': 'eiusmod_consectetur_eiusmod_dolor', 'trigger': 'aliqua_sit_eiusmod', 'model_name': 'et sed ut et', 'model_id': '60', 'block': 'sit_do', 'level': 'WARNING', 'message': 'do sed sed elit do eiusmod', 'attachment_content': 'sit labore magna labore lorem eiusmod aliqua magna et eiusmod magna labore et amet adipiscing consectetur sit ut dolore tempor tempor amet ut', 'notification_id': 'UT_UT_LABORE', 'timestamp': datetime.datetime(2025, 7, 25, 16, 50, 48, 925492, tzinfo=datetime.timezone.utc)}]}}}, 'INFO': {'section_name': 'Infos detected:', 'models': {'eiusmod eiusmod': {'model_name': 'eiusmod eiusmod', 'events': [{'project': 'sientia-laborious', 'pipeline': 'do_et_et_dolor', 'trigger': 'adipiscing_dolor_dolore_adipiscing_tempor', 'model_name': 'eiusmod eiusmod', 'model_id': '70', 'block': 'labore_sit_do', 'level': 'INFO', 'message': 'magna et ipsum amet dolore adipiscing adipiscing do', 'attachment_content': 'aliqua amet lorem amet elit tempor amet aliqua do aliqua ut dolore sed eiusmod consectetur tempor dolor sit magna eiusmod consectetur ipsum', 'notification_id': 'OPC_CONNECTION_RETRY__server_2', 'timestamp': datetime.datetime(2025, 7, 25, 16, 50, 48, 925407, tzinfo=datetime.timezone.utc)}]}, 'adipiscing consectetur': {'model_name': 'adipiscing consectetur', 'events': [{'project': 'sientia-laborious', 'pipeline': 'sit_do_aliqua', 'trigger': 'aliqua_amet_magna_ipsum_lorem', 'model_name': 'adipiscing consectetur', 'model_id': '48', 'block': 'incididunt_aliqua_ipsum_ut', 'level': 'INFO', 'message': 'labore magna adipiscing tempor elit et et magna elit lorem', 'attachment_content': 'magna labore dolor ipsum dolor lorem sed magna eiusmod ipsum do aliqua magna labore et elit consectetur et adipiscing lorem incididunt lorem', 'notification_id': 'REPORT_PARTITION_MANAGER', 'timestamp': datetime.datetime(2025, 7, 25, 16, 50, 48, 925447, tzinfo=datetime.timezone.utc)}]}, 'et amet': {'model_name': 'et amet', 'events': [{'project': 'sientia-laborious', 'pipeline': 'consectetur_labore_magna', 'trigger': 'eiusmod_incididunt_consectetur_tempor', 'model_name': 'et amet', 'model_id': '31', 'block': 'sed_elit_ut_do', 'level': 'INFO', 'message': 'dolor et adipiscing adipiscing dolore ipsum ipsum aliqua et elit', 'attachment_content': 'tempor dolore ipsum ut ut et dolor eiusmod do magna magna magna et tempor labore labore dolor labore eiusmod do dolor incididunt lorem dolore ipsum ipsum aliqua sed incididunt', 'notification_id': 'OPC_LISTENNING_BACK__server_2', 'timestamp': datetime.datetime(2025, 7, 25, 16, 50, 48, 925462, tzinfo=datetime.timezone.utc)}]}}}, 'ERROR': {'section_name': 'Errors detected:', 'models': {'ut aliqua ut ipsum': {'model_name': 'ut aliqua ut ipsum', 'events': [{'project': 'sientia-laborious', 'pipeline': 'sit_amet', 'trigger': 'aliqua_sit_consectetur_lorem_consectetur', 'model_name': 'ut aliqua ut ipsum', 'model_id': '80', 'block': 'dolor_elit_aliqua', 'level': 'ERROR', 'message': 'sit elit ipsum labore incididunt magna do et sit incididunt', 'attachment_content': 'amet ut et ut sed ipsum aliqua elit incididunt sed dolore dolore ipsum sed elit sit adipiscing incididunt elit labore do magna et consectetur aliqua', 'notification_id': 'TAG_node5:tag_87_LISTENNING_STOPPED', 'timestamp': datetime.datetime(2025, 7, 25, 16, 50, 48, 925421, tzinfo=datetime.timezone.utc)}]}, 'elit dolore': {'model_name': 'elit dolore', 'events': [{'project': 'sientia-laborious', 'pipeline': 'sed_et_tempor_labore', 'trigger': 'aliqua_tempor_sit_lorem_amet_et', 'model_name': 'elit dolore', 'model_id': '20', 'block': 'magna_magna_et', 'level': 'ERROR', 'message': 'aliqua eiusmod ut eiusmod ut et ipsum', 'attachment_content': 'labore do dolore dolore adipiscing amet tempor eiusmod labore elit consectetur incididunt dolor et elit eiusmod ut ut aliqua incididunt', 'notification_id': 'OPC_LISTENNING_BACK__server_10', 'timestamp': datetime.datetime(2025, 7, 25, 16, 50, 48, 925478, tzinfo=datetime.timezone.utc)}]}, 'lorem lorem ipsum amet': {'model_name': 'lorem lorem ipsum amet', 'events': [{'project': 'sientia-laborious', 'pipeline': 'amet_eiusmod_eiusmod', 'trigger': 'magna_aliqua_dolore_incididunt_do_do', 'model_name': 'lorem lorem ipsum amet', 'model_id': '47', 'block': 'do_sit', 'level': 'ERROR', 'message': 'labore do adipiscing adipiscing incididunt eiusmod', 'attachment_content': 'sit amet dolor magna et sit adipiscing sed consectetur labore sit ut adipiscing magna amet dolor lorem dolore eiusmod sed dolor dolor sit consectetur sit aliqua dolore tempor', 'notification_id': 'TAG_node8:tag_11_LISTENNING_STOPPED', 'timestamp': datetime.datetime(2025, 7, 25, 16, 50, 48, 925506, tzinfo=datetime.timezone.utc)}]}}}}\n" + ] + } + ], + "source": [ + "from orchestrator.utils.email_builder import EmailBuilder\n", + "from unittest.mock import MagicMock\n", + "\n", + "\n", + "email_builder = EmailBuilder(logger=MagicMock())\n", + "\n", + "html = email_builder.build_email(notifications, \"Alerts\")\n", + "\n", + "with open('email.html', 'w') as file:\n", + " file.write(html)" + ] + }, + { + "cell_type": "code", + "execution_count": 5, + "id": "ffd7be77", + "metadata": {}, + "outputs": [ + { + "data": { + "text/plain": [ + "[{'project': 'sientia-laborious',\n", + " 'pipeline': 'lorem_sed_ut',\n", + " 'trigger': 'dolore_elit_lorem',\n", + " 'model_name': 'lorem tempor',\n", + " 'model_id': '71',\n", + " 'block': 'aliqua_aliqua_et',\n", + " 'level': 'ERROR',\n", + " 'message': 'tempor eiusmod dolor eiusmod sed amet labore et magna',\n", + " 'attachment_content': 'sed amet incididunt et amet lorem lorem ut adipiscing magna incididunt sit aliqua ut dolore dolore ipsum aliqua amet labore dolor amet et elit dolore dolor consectetur',\n", + " 'notification_id': 'OPC_LISTENNING_BACK__server_1',\n", + " 'timestamp': datetime.datetime(2025, 7, 25, 15, 44, 24, 147390, tzinfo=datetime.timezone.utc)},\n", + " {'project': 'sientia-laborious',\n", + " 'pipeline': 'et_do_eiusmod_adipiscing',\n", + " 'trigger': 'consectetur_magna_elit',\n", + " 'model_name': 'elit incididunt elit',\n", + " 'model_id': '85',\n", + " 'block': 'amet_elit_adipiscing',\n", + " 'level': 'WARNING',\n", + " 'message': 'do labore consectetur lorem lorem sit do aliqua magna magna',\n", + " 'attachment_content': 'aliqua tempor aliqua adipiscing sit elit ut magna sit ipsum magna tempor do incididunt sit magna sit consectetur amet elit lorem ipsum incididunt ipsum',\n", + " 'notification_id': 'ADIPISCING_DO_DO_DOLOR',\n", + " 'timestamp': datetime.datetime(2025, 7, 25, 15, 44, 24, 147411, tzinfo=datetime.timezone.utc)},\n", + " {'project': 'sientia-laborious',\n", + " 'pipeline': 'magna_tempor',\n", + " 'trigger': 'et_elit_labore_ut_eiusmod',\n", + " 'model_name': 'sit ipsum et eiusmod',\n", + " 'model_id': '14',\n", + " 'block': 'incididunt_do',\n", + " 'level': 'INFO',\n", + " 'message': 'sit elit eiusmod consectetur magna labore',\n", + " 'attachment_content': 'amet sit et do dolore eiusmod incididunt sed dolor dolor ipsum tempor lorem dolore sed amet ipsum ut lorem aliqua',\n", + " 'notification_id': 'IPSUM_ADIPISCING_IPSUM_MAGNA_LABORE',\n", + " 'timestamp': datetime.datetime(2025, 7, 25, 15, 44, 24, 147426, tzinfo=datetime.timezone.utc)},\n", + " {'project': 'sientia-laborious',\n", + " 'pipeline': 'adipiscing_adipiscing_adipiscing',\n", + " 'trigger': 'elit_dolor_sit_et',\n", + " 'model_name': 'sit adipiscing tempor sed',\n", + " 'model_id': '46',\n", + " 'block': 'tempor_tempor_sit',\n", + " 'level': 'WARNING',\n", + " 'message': 'tempor adipiscing tempor labore dolore tempor elit',\n", + " 'attachment_content': 'elit sit dolore consectetur ut ipsum ipsum aliqua incididunt elit eiusmod labore sed eiusmod eiusmod aliqua eiusmod eiusmod labore magna incididunt consectetur ipsum et aliqua sit adipiscing labore do',\n", + " 'notification_id': 'TAG_node57:tag_97_LISTENNING_STOPPED',\n", + " 'timestamp': datetime.datetime(2025, 7, 25, 15, 44, 24, 147440, tzinfo=datetime.timezone.utc)},\n", + " {'project': 'sientia-laborious',\n", + " 'pipeline': 'eiusmod_amet_aliqua',\n", + " 'trigger': 'ipsum_lorem_et',\n", + " 'model_name': 'consectetur adipiscing',\n", + " 'model_id': '72',\n", + " 'block': 'tempor_et_adipiscing_sed',\n", + " 'level': 'ERROR',\n", + " 'message': 'lorem ipsum do amet elit eiusmod elit lorem',\n", + " 'attachment_content': 'ipsum ut dolor eiusmod dolore elit consectetur ipsum dolor do sed dolor aliqua sit consectetur aliqua sit sit consectetur et incididunt adipiscing lorem do adipiscing',\n", + " 'notification_id': 'REPORT_PARTITION_MANAGER',\n", + " 'timestamp': datetime.datetime(2025, 7, 25, 15, 44, 24, 147453, tzinfo=datetime.timezone.utc)},\n", + " {'project': 'sientia-laborious',\n", + " 'pipeline': 'eiusmod_et_incididunt',\n", + " 'trigger': 'dolore_do_lorem_aliqua',\n", + " 'model_name': 'aliqua sed tempor ipsum',\n", + " 'model_id': '2',\n", + " 'block': 'et_elit_amet',\n", + " 'level': 'WARNING',\n", + " 'message': 'consectetur aliqua dolore amet eiusmod aliqua',\n", + " 'attachment_content': 'sit amet amet aliqua incididunt aliqua amet ipsum magna do adipiscing aliqua tempor magna dolor dolor labore et lorem consectetur dolore do ut adipiscing dolore',\n", + " 'notification_id': 'AMET_ELIT_ELIT_MAGNA_ALIQUA',\n", + " 'timestamp': datetime.datetime(2025, 7, 25, 15, 44, 24, 147482, tzinfo=datetime.timezone.utc)},\n", + " {'project': 'sientia-laborious',\n", + " 'pipeline': 'amet_dolore',\n", + " 'trigger': 'labore_magna_labore',\n", + " 'model_name': 'dolor ut incididunt adipiscing',\n", + " 'model_id': '3',\n", + " 'block': 'dolor_sit',\n", + " 'level': 'WARNING',\n", + " 'message': 'ut ipsum sed consectetur eiusmod',\n", + " 'attachment_content': 'aliqua dolor magna magna sed et magna sit aliqua labore aliqua sed ipsum aliqua amet dolore ipsum aliqua consectetur adipiscing eiusmod',\n", + " 'notification_id': 'REPORT_PARTITION_MANAGER',\n", + " 'timestamp': datetime.datetime(2025, 7, 25, 15, 44, 24, 147494, tzinfo=datetime.timezone.utc)},\n", + " {'project': 'sientia-laborious',\n", + " 'pipeline': 'dolor_ipsum',\n", + " 'trigger': 'do_sed_et_dolore_elit',\n", + " 'model_name': 'adipiscing aliqua labore magna',\n", + " 'model_id': '74',\n", + " 'block': 'incididunt_dolore',\n", + " 'level': 'INFO',\n", + " 'message': 'amet do amet sed sit',\n", + " 'attachment_content': 'et sed labore adipiscing labore dolore amet tempor incididunt aliqua dolor labore tempor adipiscing incididunt consectetur magna tempor tempor aliqua ut sit magna dolore ipsum',\n", + " 'notification_id': 'REPORT_PARTITION_MANAGER',\n", + " 'timestamp': datetime.datetime(2025, 7, 25, 15, 44, 24, 147510, tzinfo=datetime.timezone.utc)},\n", + " {'project': 'sientia-laborious',\n", + " 'pipeline': 'magna_sed_ut_sit',\n", + " 'trigger': 'ut_dolore_consectetur_labore_incididunt',\n", + " 'model_name': 'aliqua et consectetur',\n", + " 'model_id': '15',\n", + " 'block': 'ipsum_et',\n", + " 'level': 'ERROR',\n", + " 'message': 'ut dolore ipsum magna tempor sed magna labore amet magna',\n", + " 'attachment_content': 'ipsum adipiscing et amet elit dolor ut aliqua consectetur ut do elit amet tempor incididunt eiusmod eiusmod dolore ipsum sit dolore adipiscing aliqua elit adipiscing labore sed eiusmod ut',\n", + " 'notification_id': 'LABORE_INCIDIDUNT_ADIPISCING_DO_LABORE',\n", + " 'timestamp': datetime.datetime(2025, 7, 25, 15, 44, 24, 147525, tzinfo=datetime.timezone.utc)},\n", + " {'project': 'sientia-laborious',\n", + " 'pipeline': 'lorem_adipiscing_lorem_incididunt',\n", + " 'trigger': 'adipiscing_magna_dolor_consectetur_elit',\n", + " 'model_name': 'dolor sit lorem',\n", + " 'model_id': '51',\n", + " 'block': 'dolor_ipsum',\n", + " 'level': 'INFO',\n", + " 'message': 'labore aliqua dolor eiusmod incididunt adipiscing tempor',\n", + " 'attachment_content': 'lorem magna dolore consectetur eiusmod sit dolor eiusmod ut do dolor labore eiusmod et dolore sit tempor consectetur ipsum adipiscing eiusmod dolore consectetur labore sit do',\n", + " 'notification_id': 'CONSECTETUR_MAGNA_AMET_LOREM_TEMPOR',\n", + " 'timestamp': datetime.datetime(2025, 7, 25, 15, 44, 24, 147539, tzinfo=datetime.timezone.utc)}]" + ] + }, + "execution_count": 5, + "metadata": {}, + "output_type": "execute_result" + } + ], + "source": [ + "notifications" + ] + }, + { + "cell_type": "code", + "execution_count": 6, + "id": "55aed216", + "metadata": {}, + "outputs": [ + { + "data": { + "text/plain": [ + "InsertManyResult([ObjectId('68839c5572e2953521d60ada'), ObjectId('68839c5572e2953521d60adb'), ObjectId('68839c5572e2953521d60adc'), ObjectId('68839c5572e2953521d60add'), ObjectId('68839c5572e2953521d60ade'), ObjectId('68839c5572e2953521d60adf'), ObjectId('68839c5572e2953521d60ae0'), ObjectId('68839c5572e2953521d60ae1'), ObjectId('68839c5572e2953521d60ae2'), ObjectId('68839c5572e2953521d60ae3')], acknowledged=True)" + ] + }, + "execution_count": 6, + "metadata": {}, + "output_type": "execute_result" + } + ], + "source": [ + "mongo_db.database['notification_queue'].insert_many(notifications)" + ] + } + ], + "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/activities.py b/orchestrator/activities/activities.py index 6c769cc..10c1812 100644 --- a/orchestrator/activities/activities.py +++ b/orchestrator/activities/activities.py @@ -1,6 +1,7 @@ from temporalio import activity, workflow from temporalio.client import Client +from orchestrator.activities.email import Email from orchestrator.activities.mongo_db import MongoDB with workflow.unsafe.imports_passed_through(): @@ -10,17 +11,21 @@ with workflow.unsafe.imports_passed_through(): from orchestrator.activities.formatters import Formatters from typing import Any from logging import Logger + from sientia_do.temporal.activities.postgres import Postgres from sientia_do.notifications.handlers import CoreNotificationHandler as NotificationHandler class Activities( # Couchbase, - TemporalManager, SlotManager, Formatters, MongoDB): + TemporalManager, SlotManager, Formatters, MongoDB, Email, + Postgres): def __init__(self, temporal_config: dict[str, Any], # couchbase_config: dict[str, Any], redis_config: dict[str, Any], mongodb_config: dict[str, Any], + email_config: dict[str, Any], + postgres_config: dict[str, Any], logger: Logger, notification_handler: NotificationHandler): @@ -59,5 +64,24 @@ class Activities( # Couchbase, logger=logger, notification_handler=notification_handler) + Email.__init__(self, + sender_email=email_config['sender_email'], + sender_password=email_config['sender_password'], + smpt_server=email_config['smpt_server'], + port=email_config['port'], + logger=logger, + notification_handler=notification_handler) + + Postgres.__init__(self, + host=postgres_config['host'], + port=postgres_config['port'], + user=postgres_config['username'], + password=postgres_config['password'], + dbname=postgres_config['database_name'], + min_connections=postgres_config['min_connections'], + max_connections=postgres_config['max_connections'], + logger=logger, + notification_handler=notification_handler) + def shutdown(self): MongoDB.shutdown(self) diff --git a/orchestrator/activities/email.py b/orchestrator/activities/email.py new file mode 100644 index 0000000..3355ebb --- /dev/null +++ b/orchestrator/activities/email.py @@ -0,0 +1,149 @@ +from email import encoders +from email.mime.base import MIMEBase +import traceback +from temporalio import workflow, activity + +with workflow.unsafe.imports_passed_through(): + import smtplib + from typing import Any + from sientia_do.temporal.activities.base import BaseActivity + from sientia_do.temporal.utils.logger import Logger + from sientia_do.notifications.handlers import NotificationHandler + from orchestrator.utils.email_builder import EmailBuilder + from email.mime.multipart import MIMEMultipart + from email.mime.text import MIMEText + + +class Email(BaseActivity): + def __init__(self, sender_email: str, sender_password: str, + smpt_server: str, port: int, + logger: Logger, notification_handler: NotificationHandler): + + self.email_builder = EmailBuilder(logger=logger) + + self.sender_email = sender_email + self.sender_password = sender_password + self.port = port + self.logger = logger + + if self.sender_password: + self.server = smtplib.SMTP_SSL(smpt_server, port) + self.server.login(self.sender_email, self.sender_password) + else: + self.server = smtplib.SMTP(smpt_server, port) + + BaseActivity.__init__(self, + logger=logger, + notification_handler=notification_handler) + + @activity.defn(name="build_email_html") + async def build_email_html(self, input_data: dict[str, Any]) -> str: + """ + Builds the email html for each receiver group. + input_data: + - receiver_groups (dict): The receiver groups. + - mail_type (str): The mail type. + """ + metadata = input_data['metadata'] + receiver_groups = input_data['receiver_groups'] + mail_type = input_data['mail_type'] + + self.info(f"Building email html for {mail_type} mail type.", + metadata=metadata) + + for group_name, group_config in receiver_groups.items(): + + html = self.email_builder.build_email( + group_config['notifications'], mail_type) + + group_config['html'] = html + + self.info(f"Email html built for {mail_type} mail type.", + metadata=metadata) + + return receiver_groups + + def handle_attachments(self, attachments: list[dict], msg: MIMEMultipart) -> MIMEMultipart: + """ + Attaches a list of attachments to an email message. + Args: + attachments (List[Dict]): A list of dictionaries where each dictionary contains + the keys 'attachment_id', 'trigger', and 'timestamp' representing the attachment details. + msg (MIMEMultipart): The email message object to which the attachments will be added. + Returns: + MIMEMultipart: The email message object with the attachments added. + Raises: + Exception: If an attachment cannot be added, an error is logged. + """ + + for attachment in attachments: + att_name = attachment['filename'] + try: + # Create the attachment as a MIMEBase object + part = MIMEBase('application', 'octet-stream') + part.set_payload( + attachment['attachment_content'].encode('utf-8')) + encoders.encode_base64(part) + part.add_header( + 'Content-Disposition', + f'attachment; filename="{att_name}"' + ) + msg.attach(part) + except Exception as e: + self.logger.error( + f"Failed to attach content of {att_name}: {e}") + + return msg + + @activity.defn(name="send_email") + async def send_email(self, input_data: dict[str, Any]) -> dict[str, Any]: + """ + Sends an email to the receivers of each group. + input_data: + - receiver_groups (dict): The receiver groups. + - mail_type (str): The mail type. + """ + metadata = input_data['metadata'] + receiver_groups = input_data['receiver_groups'] + mail_type = input_data['mail_type'] + + self.info(f"Sending email for {mail_type} mail type.", + metadata=metadata) + + for group_name, group_config in receiver_groups.items(): + + receivers = ", ".join(group_config['members']) + + self.info(f"Sending email to {group_name}: {receivers}", + metadata=metadata) + + msg = MIMEMultipart() + msg.attach(MIMEText(group_config['html'], 'html')) + msg['From'] = self.sender_email + msg['To'] = receivers + msg['Subject'] = f"SIENTIA™ {mail_type}" + + msg = self.handle_attachments( + [notification['attachment_content'] + for notification in group_config['notifications'] + if notification['attachment_content']], + msg) + + try: + self.server.sendmail( + self.sender_email, receivers, msg.as_string()) + except Exception as e: + self.error(f"Failed to send email to {group_name}: {e}", + metadata=metadata) + traceback.print_exc() + group_config['status'] = 'failed' + else: + group_config['status'] = 'sent' + + self.info(f"Email sent to {group_name}.", + metadata=metadata) + + self.info(f"Email sent for {mail_type} mail type.", + metadata=metadata) + + return receiver_groups diff --git a/orchestrator/activities/formatters.py b/orchestrator/activities/formatters.py index 5078b12..9557b31 100644 --- a/orchestrator/activities/formatters.py +++ b/orchestrator/activities/formatters.py @@ -1,4 +1,5 @@ +from pandas import DataFrame from temporalio import activity, workflow from orchestrator.utils.orchestrator_functions import minimal_retrain @@ -488,3 +489,47 @@ class Formatters(BaseActivity): notification_id="REPORT_ORCHESTRATION_DELETED_SLOTS", attachment=deleted_slots ) + + @activity.defn(name="format_log_report") + async def format_log_report(self, input_data: dict[str, Any]) -> dict[str, Any]: + """ + Formats the receiver_groups status to a dataframe to be stored in the database. + input_data: + - receiver_groups (dict): The receiver groups. + """ + metadata = input_data["metadata"] + mail_type = input_data["mail_type"] + + self.info("Formatting log report...", metadata=metadata) + + receiver_groups = input_data['receiver_groups'] + + data = {} + + for group_name, group_config in receiver_groups.items(): + + notification_id = group_config['notification_id'] + trigger = group_config['trigger'] + + key = f"{notification_id}:{trigger}" + + if key not in data: + data[key] = { + 'status': group_config['status'], + 'timestamp': group_config['timestamp'], + 'groups': [group_name], + 'message': group_config['message'], + 'level': group_config['level'], + 'notification_id': notification_id, + 'block': group_config['block'], + 'schedule': trigger, + 'pipeline': group_config['pipeline'], + 'project': group_config['project'], + 'model_name': group_config['model_name'], + 'model_id': group_config['model_id'], + 'mail_type': mail_type + } + else: + data[key]['groups'].append(group_name) + + return DataFrame(list(data.values())).to_dict() diff --git a/orchestrator/activities/mongo_db.py b/orchestrator/activities/mongo_db.py index 515a9b4..105264e 100644 --- a/orchestrator/activities/mongo_db.py +++ b/orchestrator/activities/mongo_db.py @@ -1,3 +1,4 @@ +from pandas import DataFrame from temporalio import workflow, activity @@ -80,6 +81,15 @@ class MongoDB(BaseActivity): """ self.shutdown() + def find(self, collection_name: str, filters: dict[str, Any]) -> list[dict[str, Any]]: + collection = self.database[collection_name] + + documents = list(collection.find(filters, {"_id": 0})) + + documents = clear_mongo_id(documents) + + return documents + @activity.defn(name="find_documents_in_mongodb",) async def find_documents_in_mongodb(self, input_data: dict[str, Any]) -> list[dict[str, Any]]: """ @@ -106,11 +116,7 @@ class MongoDB(BaseActivity): f"Loading documents from collection '{collection_name}' with filters: {filters}", metadata=metadata) try: - collection = self.database[collection_name] - - documents = list(collection.find(filters, {"_id": 0})) - - documents = clear_mongo_id(documents) + documents = self.find(collection_name, filters) self.info( f"Loaded {len(documents)} documents from collection '{collection_name}'", metadata=metadata) @@ -352,3 +358,76 @@ class MongoDB(BaseActivity): ) self.error(trace, metadata=metadata) raise e + + @activity.defn(name="load_latest_data") + async def load_latest_data(self, input_data: dict[str, Any]) -> list[dict[str, Any]]: + """ + Loads the latest data from MongoDB. + input_data: + - metadata (dict): The metadata of the workflow. + - collection_name (str): The name of the collection to load data from. + - last_data_timestamp (str): The timestamp of the last data to load. + - base_data_filter (dict): The base data filter to apply to the query. + returns: + - data (list[dict]): The data loaded from MongoDB. + """ + metadata = input_data['metadata'] + collection_name = input_data['collection_name'] + last_data_timestamp = input_data['last_data_timestamp'] + base_data_filter = input_data['base_data_filter'] + + self.debug( + f"Loading data from MongoDB: {input_data}", + metadata=metadata + ) + + try: + + if last_data_timestamp is None: + data_filter = base_data_filter + else: + data_filter = { + **base_data_filter, + "timestamp": { + "$gt": datetime.strptime(last_data_timestamp, DEFAULT_DATE_FORMAT) + } + } + + self.debug( + f"Data filter: {data_filter}", + metadata=metadata + ) + + data = self.find(collection_name, data_filter) + + self.debug( + f"Collected: {data}", + metadata=metadata + ) + + for item in data: + item['timestamp'] = item['timestamp'].strftime( + DEFAULT_DATE_FORMAT) + + self.info( + f"Loaded {len(data)} documents from MongoDB", + metadata=metadata + ) + + self.debug( + f"Loaded data: {data}", + metadata=metadata + ) + + return data + except Exception as e: + trace = traceback.format_exc() + self.send_notification( + metadata=metadata, + notification_id="MONGO_LOAD_ERROR", + message=f"Error loading data from MongoDB: {e}", + block="load_latest_data", + level=NotificationLevel.ERROR, + attachment_content=trace + ) + raise e diff --git a/orchestrator/activities/slot_manager.py b/orchestrator/activities/slot_manager.py index 404c8d8..b5cf4b7 100644 --- a/orchestrator/activities/slot_manager.py +++ b/orchestrator/activities/slot_manager.py @@ -1,3 +1,4 @@ +from pandas import DataFrame from temporalio import activity, workflow with workflow.unsafe.imports_passed_through(): @@ -193,3 +194,72 @@ class SlotManager(Redis): f"Report: \n {json.dumps(report, indent=4, sort_keys=True)}", metadata=metadata) return report + + @activity.defn(name="get_last_data_timestamp") + async def get_last_data_timestamp(self, input_data: dict[str, Any]) -> str | None: + """ + Gets the last data timestamp from redis. + """ + metadata = input_data['metadata'] + key = "notification_last_timestamp" + + try: + data_hold = self.get(key) + except Exception as e: + self.send_notification( + metadata=metadata, + notification_id="REDIS_GET_ERROR", + message=f"Error getting last data timestamp: {e}", + block="get_last_data_timestamp", + level=NotificationLevel.ERROR, + attachment_content=traceback.format_exc() + ) + raise e + + self.debug( + f"Last collected timestamp: {data_hold}", + metadata=metadata + ) + + if not data_hold: + return None + + return data_hold + + @activity.defn(name="put_last_data_timestamp") + async def put_last_data_timestamp(self, input_data: dict[str, Any]): + """ + Puts the last data timestamp into redis. + """ + metadata = input_data['metadata'] + key = "notification_last_timestamp" + + data = DataFrame(input_data['data']) + + if data.empty: + self.warning("No data to insert", + metadata=metadata + ) + return None + + last_data_timestamp = data['timestamp'].max() + + self.debug( + f"Last collected timestamp to insert: {last_data_timestamp}", + metadata=metadata + ) + + try: + self.set(key, last_data_timestamp, ttl=60*60*5) + except Exception as e: + self.send_notification( + metadata=metadata, + notification_id="REDIS_SET_ERROR", + message=f"Error setting last data timestamp: {e}", + block="put_last_data_timestamp", + level=NotificationLevel.ERROR, + attachment_content=traceback.format_exc() + ) + raise e + + return last_data_timestamp diff --git a/orchestrator/utils/email_builder.py b/orchestrator/utils/email_builder.py new file mode 100644 index 0000000..8dc59e8 --- /dev/null +++ b/orchestrator/utils/email_builder.py @@ -0,0 +1,75 @@ +import json +from sientia_do.temporal.utils.logger import Logger +from sientia_do.notifications.models import NotificationLevel +from jinja2 import Template +import re + + +class EmailBuilder: + def __init__(self, logger: Logger): + self.logger = logger + + self.report_template_file = './orchestrator/utils/templates/email_template.html' + self.general_template_file = './orchestrator/utils/templates/general_template.html' + + self.tag_template_file = './orchestrator/utils/templates/opc_tag_report_template.html' + self.opc_template_file = './orchestrator/utils/templates/opc_connection_report_template.html' + + self.partition_manager_template_file = './orchestrator/utils/templates/partition_manager_report_template.html' + + with open(self.report_template_file, 'r') as file: + self.report_template = file.read() + with open(self.general_template_file, 'r') as file: + self.general_template = file.read() + + def replace_parameters(self, template: str, parameters: dict) -> str: + # Criar um template Jinja2 + template = Template(template) + + return template.render(parameters) + + def parameters(self, report_data: dict, general_events: dict, mail_type: str) -> dict: + return { + 'project_name': report_data[0]['project'], + 'mail_type': mail_type, + 'error_events': self.replace_parameters(self.general_template, + general_events['ERROR']) if general_events['ERROR']['models'] else '', + 'warning_events': self.replace_parameters(self.general_template, + general_events['WARNING']) if general_events['WARNING']['models'] else '', + 'info_events': self.replace_parameters(self.general_template, + general_events['INFO']) if general_events['INFO']['models'] else '', + } + + def build_email(self, report_data: list[dict], mail_type: str) -> str: + """ + Builds the email html. + """ + general_events = {} + + for report in report_data: + + level = report['level'] + model_name = report['model_name'] + + if level not in general_events: + general_events[level] = { + 'section_name': f'{level.capitalize()}s detected:', + 'models': {} + } + + if model_name not in general_events[level]['models']: + general_events[level]['models'][model_name] = { + 'model_name': model_name, + 'events': [] + } + + general_events[level]['models'][model_name]['events'].append( + report) + + print(general_events) + for _type, content in general_events.items(): + content['models'] = list(content['models'].values()) + + return self.replace_parameters(self.report_template, self.parameters( + report_data, general_events, mail_type + )) diff --git a/orchestrator/utils/templates/email_template.html b/orchestrator/utils/templates/email_template.html new file mode 100644 index 0000000..e269e53 --- /dev/null +++ b/orchestrator/utils/templates/email_template.html @@ -0,0 +1,24 @@ + + + + + + SIENTIA™ Report + + + +

SIENTIA™ {{ mail_type }}

+ + {{ error_events }} + {{ warning_events }} + {{ info_events }} + + {{ special_events }} + + diff --git a/orchestrator/utils/templates/general_template.html b/orchestrator/utils/templates/general_template.html new file mode 100644 index 0000000..7736178 --- /dev/null +++ b/orchestrator/utils/templates/general_template.html @@ -0,0 +1,24 @@ +

{{ section_name }}

+{% for model in models %} +

Model: {{ model.model_name }}

+ + + + + + + + + + + {% for event in model.events %} + + + + + + + {% endfor %} + +
Notification IDBlockTimestampMessage
{{ event.notification_id }}{{ event.block }}{{ event.timestamp }}{{ event.message }}
+{% endfor %} \ No newline at end of file diff --git a/orchestrator/workflows/alerts.py b/orchestrator/workflows/alerts.py new file mode 100644 index 0000000..eee7f1c --- /dev/null +++ b/orchestrator/workflows/alerts.py @@ -0,0 +1,22 @@ +from temporalio import workflow + +with workflow.unsafe.imports_passed_through(): + from orchestrator.activities.activities import Activities + from typing import Any + + +@workflow.defn(name="alerts") +class Alerts: + @workflow.run + async def run(self, input_data: dict[str, Any]): + # Call subworkflow "load_notification_package" passing the static filters + # (level = "ERROR" and timestamp > last timestamp) + + # Filter notification package by groups custom configs, levels and + # timestamp cached + + # Call subworkflow "process_notifications" passing the notification package + + # Store the notification_id sendings to avoid sending them again + + pass diff --git a/orchestrator/workflows/orchestrator.py b/orchestrator/workflows/orchestrator.py index f61587d..cb81d13 100644 --- a/orchestrator/workflows/orchestrator.py +++ b/orchestrator/workflows/orchestrator.py @@ -11,6 +11,15 @@ with workflow.unsafe.imports_passed_through(): class Orchestrator: @workflow.run async def run(self, input_data: dict[str, Any]): + """ + Orchestrates the pipeline and slot management. Gets configuration from MongoDB and Redis, + creates the configuration and deploys the schedules and slots in the Temporal server and + Redis server. + input_data: + - schedule_name (str): The name of the schedule. + - pipelines_query (dict): The query to get the pipelines. + - opc_servers_query (dict): The query to get the OPC servers. + """ input_data['workflow_name'] = 'orchestrator' diff --git a/orchestrator/workflows/reports.py b/orchestrator/workflows/reports.py new file mode 100644 index 0000000..002034f --- /dev/null +++ b/orchestrator/workflows/reports.py @@ -0,0 +1,19 @@ +from temporalio import workflow + +with workflow.unsafe.imports_passed_through(): + from orchestrator.activities.activities import Activities + from typing import Any + + +@workflow.defn(name="reports") +class Reports: + @workflow.run + async def run(self, input_data: dict[str, Any]): + # Call subworkflow "load_notification_package" passing the static filters + # (timestamp > last timestamp) + + # Filter notification package by groups custom configs + + # Call subworkflow "process_notifications" passing the notification package + + pass diff --git a/orchestrator/workflows/subworkflows/load_notification_package.py b/orchestrator/workflows/subworkflows/load_notification_package.py new file mode 100644 index 0000000..165a778 --- /dev/null +++ b/orchestrator/workflows/subworkflows/load_notification_package.py @@ -0,0 +1,93 @@ +from temporalio import workflow + +with workflow.unsafe.imports_passed_through(): + from orchestrator.activities.activities import Activities + from typing import Any + from sientia_do.temporal.utils.policies import retry_policy + from datetime import timedelta + + +@workflow.defn(name="load_notification_package") +class LoadNotificationPackage: + @workflow.run + async def run(self, input_data: dict[str, Any]): + """ + Loads the notification package from the MongoDB collection "notification_queue" + and the sending configs from the MongoDB collection "receiver_groups". + + input_data: + - metadata (dict): The metadata of the workflow. + + returns: + - last_timestamp (str): The last timestamp of the notification package. + - notification_package (list[dict]): The notification package. + - sending_configs (list[dict]): The sending configs. + """ + metadata = input_data['metadata'] + + # Load last timestamp from redis "notification_last_timestamp" + last_timestamp_handler = workflow.start_local_activity_method( + Activities.get_last_data_timestamp, + { + **metadata, + }, + start_to_close_timeout=timedelta(seconds=60), + retry_policy=retry_policy + ) + + # In parallel, load sending configs from collection "receiver_groups" + sending_configs_handler = workflow.start_local_activity_method( + Activities.find_documents_in_mongodb, + { + **metadata, + 'query': { + 'collection': 'receiver_groups', + 'filters': { + 'active': True + } + } + }, + start_to_close_timeout=timedelta(seconds=60), + retry_policy=retry_policy + ) + + last_timestamp = await last_timestamp_handler + + # Load notification package from collection "notification_queue", using a + # static filter + + notification_package = await workflow.start_local_activity_method( + Activities.load_latest_data, + { + **metadata, + 'collection_name': 'notification_queue', + 'last_data_timestamp': last_timestamp, + 'base_data_filter': { + 'level': 'ERROR' + } + }, + start_to_close_timeout=timedelta(seconds=60), + retry_policy=retry_policy + ) + + # Put last collected timestamp in redis "notification_last_timestamp" + + await workflow.start_activity_method( + Activities.put_last_data_timestamp, + { + **metadata, + 'data': notification_package + }, + start_to_close_timeout=timedelta(seconds=60), + retry_policy=retry_policy + ) + + # Return a dict with the following keys: + # - last_timestamp + # - notification_package + # - sending_configs + return { + 'last_timestamp': last_timestamp, + 'notification_package': notification_package, + 'sending_configs': await sending_configs_handler + } diff --git a/orchestrator/workflows/subworkflows/process_notifications.py b/orchestrator/workflows/subworkflows/process_notifications.py new file mode 100644 index 0000000..de71b27 --- /dev/null +++ b/orchestrator/workflows/subworkflows/process_notifications.py @@ -0,0 +1,87 @@ +from temporalio import workflow + +with workflow.unsafe.imports_passed_through(): + from orchestrator.activities.activities import Activities + from typing import Any + from datetime import timedelta + from sientia_do.temporal.utils.policies import retry_policy + + +@workflow.defn(name="process_notifications") +class ProcessNotifications: + @workflow.run + async def run(self, input_data: dict[str, Any]) -> dict[str, Any]: + """ + Processes the notifications. Builds the report html for each group and each model, + sends the report html to the receivers of each group, stores the sending log in the + postgres database "log_report", and returns the log report to the caller. + + input_data: + - metadata (dict): The metadata of the workflow. + - mail_type (str): The mail type. + - schema (str): The schema of the table. + - table_name (str): The name of the table. + - notification_package (list[dict]): The notification package. the format of each + notification package is: + { + 'group_name' (str) + 'group_members' (list[str]) + 'notifications' (dict) + { + 'model_name' (dict[str, list[dict]]) + } + } + returns: + - log_report (dict) + """ + + metadata = input_data["metadata"] + + # Use notification package to create the report html for each group and each model + data_to_sent = await workflow.execute_local_activity_method( + Activities.build_email_html, + { + **metadata, + "receiver_groups": input_data["notification_package"] + }, + schedule_to_close_timeout=timedelta(seconds=60), + retry_policy=retry_policy + ) + + # Send the report html to the receivers of each group + log_report = await workflow.execute_activity_method( + Activities.send_email, + { + **metadata, + "receiver_groups": data_to_sent, + "mail_type": input_data["mail_type"] + }, + schedule_to_close_timeout=timedelta(seconds=60), + retry_policy=retry_policy + ) + + # Format the log report to a dataframe to be stored in the database + log_report = await workflow.execute_local_activity_method( + Activities.format_log_report, + { + **metadata, + "receiver_groups": log_report, + "mail_type": input_data["mail_type"] + }, + ) + + # Store sending log in postgres database "log_report" + await workflow.execute_activity_method( + Activities.export_data_to_postgres, + { + **metadata, + "schema": input_data["schema"], + "table_name": input_data["table_name"], + "data": log_report + }, + schedule_to_close_timeout=timedelta(seconds=60), + retry_policy=retry_policy + ) + + # Return the log report to the caller + return log_report diff --git a/requirements.txt b/requirements.txt index 23d31f9..d06bf0e 100644 --- a/requirements.txt +++ b/requirements.txt @@ -4,4 +4,5 @@ sqlalchemy redis couchbase pymongo +jinja2 git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.3.4 diff --git a/tests/orchestrator/activities/test_mongo_db.py b/tests/orchestrator/activities/test_mongo_db.py index 707b866..4c4c956 100644 --- a/tests/orchestrator/activities/test_mongo_db.py +++ b/tests/orchestrator/activities/test_mongo_db.py @@ -1,4 +1,5 @@ from curses import meta +from datetime import datetime from unittest.mock import MagicMock, patch, ANY from pytest import fixture, mark from orchestrator.activities.mongo_db import clear_mongo_id @@ -518,3 +519,133 @@ async def test_create_collection_with_ttl_index_failure(mongo_db): ) else: assert False, "Expected an exception to be raised" + + +@mark.asyncio +async def test_load_latest_data_none_last_data_timestamp(mongo_db): + """Test load_latest_data""" + collection = MagicMock() + mongo_db.database.__getitem__.return_value = collection + + collection.find.return_value = [ + { + 'name': 'test1', + 'value': 1, + 'timestamp': datetime.strptime( + '2023-01-01 12:00:00.000000', '%Y-%m-%d %H:%M:%S.%f') + } + ] + + result = await mongo_db.load_latest_data({ + 'metadata': {'workflow_name': 'test_pipeline', 'schedule_name': 'test_schedule'}, + 'collection_name': 'test_collection', + 'last_data_timestamp': None, + 'base_data_filter': { + 'level': 'ERROR' + } + }) + + mongo_db.database.__getitem__.assert_called_once_with( + 'test_collection') + + collection.find.assert_called_once_with( + { + 'level': 'ERROR' + }, + {"_id": 0} + ) + + assert result == { + 'name': { + 0: 'test1' + }, + 'value': { + 0: 1 + }, + 'timestamp': { + 0: '2023-01-01 12:00:00.000000' + } + } + + +@mark.asyncio +async def test_load_latest_data_not_none_last_data_timestamp(mongo_db): + """Test load_latest_data""" + collection = MagicMock() + mongo_db.database.__getitem__.return_value = collection + + collection.find.return_value = [ + { + 'name': 'test1', + 'value': 1, + 'timestamp': datetime.strptime( + '2023-01-01 12:00:00.000000', '%Y-%m-%d %H:%M:%S.%f') + } + ] + + result = await mongo_db.load_latest_data({ + 'metadata': {'workflow_name': 'test_pipeline', 'schedule_name': 'test_schedule'}, + 'collection_name': 'test_collection', + 'last_data_timestamp': '2023-01-01 12:00:00.000000', + 'base_data_filter': { + 'level': 'ERROR' + } + }) + + mongo_db.database.__getitem__.assert_called_once_with( + 'test_collection') + + collection.find.assert_called_once_with( + { + 'level': 'ERROR', + 'timestamp': { + '$gt': datetime.strptime( + '2023-01-01 12:00:00.000000', '%Y-%m-%d %H:%M:%S.%f') + } + }, + {"_id": 0} + ) + + assert result == { + 'name': { + 0: 'test1' + }, + 'value': { + 0: 1 + }, + 'timestamp': { + 0: '2023-01-01 12:00:00.000000' + } + } + + +@mark.asyncio +async def test_load_latest_data_error(mongo_db): + """Test load_latest_data""" + collection = MagicMock() + mongo_db.send_notification = MagicMock() + mongo_db.database.__getitem__.return_value = collection + + collection.find.side_effect = Exception('test') + + try: + await mongo_db.load_latest_data({ + 'metadata': {'workflow_name': 'test_pipeline', 'schedule_name': 'test_schedule'}, + 'collection_name': 'test_collection', + 'last_data_timestamp': '2023-01-01 12:00:00.000000', + 'base_data_filter': { + 'level': 'ERROR' + } + }) + except Exception as e: + assert str(e) == 'test' + + mongo_db.send_notification.assert_called_once_with( + metadata={'workflow_name': 'test_pipeline', + 'schedule_name': 'test_schedule'}, + notification_id='MONGO_LOAD_ERROR', + message='Error loading data from MongoDB: test', + block='load_latest_data', + level=NotificationLevel.ERROR, + attachment_content=ANY + ) diff --git a/tests/orchestrator/activities/test_slot_manager.py b/tests/orchestrator/activities/test_slot_manager.py index 9be7686..ad2e9a5 100644 --- a/tests/orchestrator/activities/test_slot_manager.py +++ b/tests/orchestrator/activities/test_slot_manager.py @@ -1,4 +1,5 @@ from unittest.mock import MagicMock, patch, call, ANY +from pandas import DataFrame from pytest import mark, fixture from orchestrator.activities.slot_manager import SlotManager from sientia_do.notifications.models import NotificationLevel @@ -206,3 +207,155 @@ async def test_delete_slots(slot_manager): "message": "Test exception" } } + + +@mark.asyncio +async def test_get_last_data_timestamp_none(slot_manager): + """Test get_last_data_timestamp""" + test_data = { + **metadata, + 'workflow_name': 'test_pipeline', + 'schedule_name': 'test_schedule' + } + + slot_manager.get = MagicMock(return_value=None) + + result = await slot_manager.get_last_data_timestamp(test_data) + + assert result is None + + +@mark.asyncio +async def test_get_last_data_timestamp_not_none(slot_manager): + """Test get_last_data_timestamp""" + test_data = { + **metadata, + 'workflow_name': 'test_pipeline', + 'schedule_name': 'test_schedule' + } + + slot_manager.get = MagicMock(return_value='2023-01-01 12:00:00') + + result = await slot_manager.get_last_data_timestamp(test_data) + + slot_manager.get.assert_called_once_with( + 'notification_last_timestamp' + ) + + assert result == '2023-01-01 12:00:00' + + +@mark.asyncio +async def test_get_last_data_timestamp_error(slot_manager): + """Test get_last_data_timestamp error""" + test_data = { + **metadata, + 'workflow_name': 'test_pipeline', + 'schedule_name': 'test_schedule' + } + + slot_manager.send_notification = MagicMock() + slot_manager.get = MagicMock(side_effect=Exception('test')) + + try: + + await slot_manager.get_last_data_timestamp(test_data) + + except Exception as e: + assert str(e) == 'test' + + slot_manager.send_notification.assert_called_once_with( + metadata=metadata['metadata'], + notification_id="REDIS_GET_ERROR", + message="Error getting last data timestamp: test", + block="get_last_data_timestamp", + level=NotificationLevel.ERROR, + attachment_content=ANY + ) + + else: + assert False, "Expected exception" + + +@mark.asyncio +async def test_put_last_data_timestamp_empty_dataframe(slot_manager): + """Test put_last_data_timestamp with empty dataframe""" + test_data = { + **metadata, + 'workflow_name': 'test_pipeline', + 'schedule_name': 'test_schedule', + 'data': DataFrame(columns=['name', 'value', 'timestamp']).to_dict('records') + } + + slot_manager.set = MagicMock() + + result = await slot_manager.put_last_data_timestamp(test_data) + + assert result is None + + slot_manager.set.assert_not_called() + + +@mark.asyncio +async def test_put_last_data_timestamp_not_empty_dataframe(slot_manager): + """Test put_last_data_timestamp with not empty dataframe""" + + data = DataFrame({ + 'name': ['sensor1', 'sensor2'], + 'value': [25.5, 30.0], + 'timestamp': ['2023-01-01 12:00:00', '2023-01-01 12:00:01'] + }) + test_data = { + **metadata, + 'workflow_name': 'test_pipeline', + 'schedule_name': 'test_schedule', + 'data': data.to_dict('records') + } + + slot_manager.set = MagicMock() + + result = await slot_manager.put_last_data_timestamp(test_data) + + assert result == '2023-01-01 12:00:01' + + slot_manager.set.assert_called_once_with( + 'notification_last_timestamp', + '2023-01-01 12:00:01', + ttl=18000 + ) + + +@mark.asyncio +async def test_put_last_data_timestamp_error(slot_manager): + """Test put_last_data_timestamp error""" + test_data = { + **metadata, + 'workflow_name': 'test_pipeline', + 'schedule_name': 'test_schedule', + 'data': DataFrame({ + 'name': ['sensor1', 'sensor2'], + 'value': [25.5, 30.0], + 'timestamp': ['2023-01-01 12:00:00'] * 2 + }).to_dict('records') + } + + slot_manager.send_notification = MagicMock() + slot_manager.set = MagicMock(side_effect=Exception('test')) + + try: + await slot_manager.put_last_data_timestamp(test_data) + + except Exception as e: + assert str(e) == 'test' + + slot_manager.send_notification.assert_called_once_with( + metadata=metadata['metadata'], + notification_id="REDIS_SET_ERROR", + message="Error setting last data timestamp: test", + block="put_last_data_timestamp", + level=NotificationLevel.ERROR, + attachment_content=ANY + ) + + else: + assert False, "Expected exception" diff --git a/tests/orchestrator/workflows/subworkflows/test_load_notification_package.py b/tests/orchestrator/workflows/subworkflows/test_load_notification_package.py new file mode 100644 index 0000000..952cd3f --- /dev/null +++ b/tests/orchestrator/workflows/subworkflows/test_load_notification_package.py @@ -0,0 +1,115 @@ +from unittest.mock import AsyncMock, patch, ANY, call +from pytest import fixture, mark +from orchestrator.workflows.subworkflows.load_notification_package import LoadNotificationPackage +from orchestrator.activities.activities import Activities + + +@fixture +def load_notification_package(): + return LoadNotificationPackage() + + +metadata = { + 'metadata': { + 'schedule_name': 'test-schedule-name', + 'workflow_name': 'test-workflow', + 'model_name': '-', + 'model_id': '-', + } +} + + +@mark.asyncio +@patch("orchestrator.workflows.subworkflows.load_notification_package.workflow", new_callable=AsyncMock) +async def test_run(workflow_mock, load_notification_package): + input_data = { + 'metadata': metadata + } + + workflow_mock.start_local_activity_method.side_effect = [ + '2023-01-01 12:00:00', + [ + { + 'id': '1', + } + ], + [ + { + 'id_r': '1', + } + ] + ] + + output = await load_notification_package.run(input_data) + + assert output == { + 'last_timestamp': '2023-01-01 12:00:00', + 'notification_package': [ + { + 'id': '1', + } + ], + 'sending_configs': [ + { + 'id_r': '1', + } + ] + } + + workflow_mock.start_local_activity_method.assert_has_calls([ + call( + Activities.get_last_data_timestamp, + input_data['metadata'], + start_to_close_timeout=ANY, + retry_policy=ANY + ) + ]) + + workflow_mock.start_local_activity_method.assert_has_calls([ + call( + Activities.find_documents_in_mongodb, + { + **input_data['metadata'], + 'query': { + 'collection': 'receiver_groups', + 'filters': { + 'active': True + } + } + }, + start_to_close_timeout=ANY, + retry_policy=ANY + ) + ]) + + workflow_mock.start_local_activity_method.assert_has_calls([ + call( + Activities.load_latest_data, + { + **input_data['metadata'], + 'collection_name': 'notification_queue', + 'last_data_timestamp': '2023-01-01 12:00:00', + 'base_data_filter': { + 'level': 'ERROR' + } + }, + start_to_close_timeout=ANY, + retry_policy=ANY + ) + ]) + + workflow_mock.start_activity_method.assert_has_calls([ + call( + Activities.put_last_data_timestamp, + { + **input_data['metadata'], + 'data': [ + { + 'id': '1', + } + ] + }, + start_to_close_timeout=ANY, + retry_policy=ANY + ) + ])