Merge pull request #26 from Aignosi/SIENTIAPDE-1312-melhorias-e-correcoes-nas-pipelines-de-dados

SIENTIAPDE-1312: Enhancements to Tag Management, Scheduling, and Error Reporting
This commit is contained in:
Matheus Demoner
2025-10-21 17:25:24 -03:00
committed by GitHub
10 changed files with 213 additions and 286 deletions

View File

@@ -1,9 +1,5 @@
from smtplib import SMTPServerDisconnected
from temporalio import activity, workflow
from orchestrator import metrics
with workflow.unsafe.imports_passed_through():
import smtplib
import traceback
@@ -11,12 +7,14 @@ with workflow.unsafe.imports_passed_through():
from email.mime.base import MIMEBase
from email.mime.multipart import MIMEMultipart
from email.mime.text import MIMEText
from smtplib import SMTPServerDisconnected
from typing import Any
from sientia_do.notifications.handlers import NotificationHandler
from sientia_do.observability.logger import Logger
from sientia_do.temporal.activities.base import BaseActivity
from orchestrator import metrics
from orchestrator.utils.email_builder import EmailBuilder

View File

@@ -1,9 +1,8 @@
from collections.abc import Hashable
from temporalio import activity, workflow
with workflow.unsafe.imports_passed_through():
import json
from collections.abc import Hashable
from logging import Logger
from math import ceil
from typing import Any
@@ -163,36 +162,31 @@ class Formatters(BaseActivity):
last_index = 0
for i in range(1, number_of_slots):
slot_config[f'{i}'] = {}
for tag in tags[last_index : last_index + tags_per_slot]:
try:
slot_config = build_tag_config(tag, slot_config.copy(), opc_servers, i)
except ValueError as e:
self.send_notification(
metadata=metadata,
notification_id='ORCHESTRATOR_BUILD_TAG_CONFIG_ERROR',
message=str(e),
block='orchestrator',
level=NotificationLevel.ERROR,
)
slot_tags = tags[last_index : last_index + tags_per_slot]
slot_config[f'{i}'], notifications = build_tag_config(slot_tags, opc_servers)
last_index += tags_per_slot
slot_config[f'{number_of_slots}'] = {}
for tag in tags[last_index:]:
try:
slot_config = build_tag_config(
tag, slot_config.copy(), opc_servers, number_of_slots
)
except ValueError as e:
if notifications:
self.send_notification(
metadata=metadata,
notification_id='ORCHESTRATOR_BUILD_TAG_CONFIG_ERROR',
message=str(e),
notification_id='ORCHESTRATOR_SERVER_NOT_FOUND_DURING_SLOT_CONFIGURATION',
message=f'Servers {", ".join(notifications)} not found in opc_servers',
block='orchestrator',
level=NotificationLevel.ERROR,
)
last_index += tags_per_slot
slot_tags = tags[last_index:]
slot_config[f'{number_of_slots}'], notifications = build_tag_config(slot_tags, opc_servers)
if notifications:
self.send_notification(
metadata=metadata,
notification_id='ORCHESTRATOR_SERVER_NOT_FOUND_DURING_SLOT_CONFIGURATION',
message=f'Servers {", ".join(notifications)} not found in opc_servers',
block='orchestrator',
level=NotificationLevel.ERROR,
)
self.info('Processed slots', metadata=metadata)
self.debug(json.dumps(slot_config, indent=4, sort_keys=True), metadata=metadata)

View File

@@ -1,10 +1,8 @@
from datetime import UTC
from temporalio import activity, workflow
with workflow.unsafe.imports_passed_through():
import traceback
from datetime import datetime
from datetime import UTC, datetime
from logging import Logger
from typing import Any

View File

@@ -129,7 +129,7 @@ class TemporalManager(BaseActivity):
trace = traceback.format_exc()
self.send_notification(
metadata=metadata,
notification_id='TEMPORAL_NORMALIZE_SCHEDULES_ERROR',
notification_id='SCHEDULER_NORMALIZE_SCHEDULES_ERROR',
message=f'Failed to normalize schedules: {e}',
block='normalize_schedules',
level=NotificationLevel.ERROR,
@@ -209,7 +209,10 @@ class TemporalManager(BaseActivity):
ScheduleIntervalSpec(
every=timedelta(
seconds=parse_frequency(schedule.get('frequency', '1m'))
)
),
offset=timedelta(
seconds=parse_frequency(schedule.get('offset', '0m'))
),
)
]
),

View File

@@ -22,6 +22,7 @@ def common_config(config: dict[str, Any]):
'workflow_type': config['workflow_type'],
'schedule_name': config['schedule_name'],
'frequency': config.get('frequency', '1m'),
'offset': config.get('offset', '0m'),
'max_retry_policy': config.get('max_retry_policy', 1),
'model_id': config['model_id'],
'model_name': model['name'],
@@ -96,6 +97,7 @@ def scouter(config: dict[str, Any]):
'retention_time': config.get('tag_retention_minutes', 60) * 60,
'model_tags': tags,
'debug_data_package': config.get('debug_data_package', False),
'fill_missing_tags': config.get('fill_missing_tags', False),
}
@@ -238,13 +240,13 @@ def gather_read_tags(pipelines: list[dict[str, Any]]) -> dict[str, Any]:
def build_tag_config(
tag: dict[str, Any], slot_config: dict[str, Any], opc_servers: dict[str, Any], i: int
):
tags: list[dict[str, Any]], opc_servers: dict[str, Any]
) -> tuple[dict[str, Any], list]:
"""
Build tag configuration for a specific slot and OPC server.
Args:
tag (dict[str, Any]): Tag configuration containing:
tags (list[dict[str, Any]]): List of tag configurations containing:
- server_id (str): ID of the OPC server
- tag_address (str): Address of the tag
slot_config (dict[str, Any]): Current slot configuration to update.
@@ -257,26 +259,39 @@ def build_tag_config(
Raises:
ValueError: If the specified server_id is not found in opc_servers.
"""
server_id = tag['server_id']
if server_id not in opc_servers:
raise ValueError(f'Server {server_id} not found in opc_servers')
slot_config = {}
notifications = []
server_name = opc_servers[server_id]['server_name']
if server_name not in slot_config[f'{i}']:
slot_config[f'{i}'][server_name] = {
'server_id': server_id,
'name': server_name,
'url': opc_servers[server_id]['url'],
'server_uri': opc_servers[server_id]['uri'],
'cert_path': opc_servers[server_id].get('cert_path', None),
'private_key_path': opc_servers[server_id].get('private_key_path', None),
'server_cert_path': opc_servers[server_id].get('server_cert_path', None),
'tags': {},
for tag in tags:
server_id = tag['server_id']
if server_id not in opc_servers:
notifications.append(server_id)
continue
server_name = opc_servers[server_id]['server_name']
if server_name not in slot_config:
slot_config[server_name] = {
'server_id': server_id,
'name': server_name,
'url': opc_servers[server_id]['url'],
'server_uri': opc_servers[server_id]['uri'],
'cert_path': opc_servers[server_id].get('cert_path', None),
'private_key_path': opc_servers[server_id].get('private_key_path', None),
'server_cert_path': opc_servers[server_id].get('server_cert_path', None),
'tags': {},
}
slot_config[server_name]['tags'][tag['tag_address']] = {
**tag,
}
slot_config[f'{i}'][server_name]['tags'][tag['tag_address']] = {
**tag,
}
for server_name in slot_config:
frequencies = [int(x['frequency']) for x in slot_config[server_name]['tags'].values()]
return slot_config
min_frequency = min(frequencies) if frequencies else 1000
slot_config[server_name]['subscription_period_ms'] = min_frequency / 2
return slot_config, notifications

View File

@@ -5,5 +5,5 @@ redis
couchbase
pymongo
jinja2
git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.6
git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.4.7
prometheus-client

View File

@@ -1,12 +1,11 @@
import json
from unittest.mock import ANY, MagicMock, call, patch
from unittest.mock import MagicMock, call, patch
from pandas import DataFrame
from pytest import fixture, mark
from sientia_do.notifications.models import NotificationLevel
from orchestrator.activities.formatters import Formatters
from orchestrator.utils.orchestrator_functions import build_tag_config
@fixture
@@ -102,22 +101,25 @@ async def test_process_schedules(
'server_name': 'test_server_name',
'tag_address': 'test_tag_address',
'topics': ['raw_test_schedule'],
'frequency': 1000,
},
'2:test_tag_address2': {
'server_id': '2',
'server_name': 'test_server_name2',
'tag_address': 'test_tag_address2',
'topics': ['raw_test_schedule2'],
'frequency': 1000,
},
'2:test_tag_address3': {
'server_id': '2',
'server_name': 'test_server_name2',
'tag_address': 'test_tag_address3',
'topics': ['raw_test_schedule2'],
'frequency': 1000,
},
},
)
@patch('orchestrator.activities.formatters.build_tag_config', side_effect=build_tag_config)
@patch('orchestrator.activities.formatters.build_tag_config')
async def test_process_slots(mock_build_tag_config, mock_gather_read_tags, formatters):
input_data = {
'opc_servers': [
@@ -126,222 +128,75 @@ async def test_process_slots(mock_build_tag_config, mock_gather_read_tags, forma
],
'active_ingestors': ['test_active_ingestor1', 'test_active_ingestor2'],
'pipelines': 'test_gather_read_tags',
**metadata,
}
expected_opc_servers = {
'1': {
'id': '1',
'server_name': 'test_server_name',
'url': 'test_url',
'uri': 'test_uri',
},
'2': {
'id': '2',
'server_name': 'test_server_name2',
'url': 'test_url2',
'uri': 'test_uri2',
},
}
slot_mock = {
'test_server_name': {
'server_id': '1',
'name': 'test_server_name',
'url': 'test_url',
'server_uri': 'test_uri',
'cert_path': None,
'private_key_path': None,
'server_cert_path': None,
}
}
mock_build_tag_config.return_value = (slot_mock, ['2'])
result = await formatters.process_slots(input_data)
tags = list(mock_gather_read_tags.return_value.values())
mock_gather_read_tags.assert_called_once_with(input_data['pipelines'])
mock_build_tag_config.assert_has_calls(
[
call(
{
'server_id': '1',
'server_name': 'test_server_name',
'tag_address': 'test_tag_address',
'topics': ['raw_test_schedule'],
},
ANY,
{
'1': {
'id': '1',
'server_name': 'test_server_name',
'url': 'test_url',
'uri': 'test_uri',
},
'2': {
'id': '2',
'server_name': 'test_server_name2',
'url': 'test_url2',
'uri': 'test_uri2',
},
},
1,
)
call(tags[:2], expected_opc_servers),
call(tags[2:], expected_opc_servers),
]
)
mock_build_tag_config.assert_has_calls(
[
call(
{
'server_id': '2',
'server_name': 'test_server_name2',
'tag_address': 'test_tag_address2',
'topics': ['raw_test_schedule2'],
},
ANY,
{
'1': {
'id': '1',
'server_name': 'test_server_name',
'url': 'test_url',
'uri': 'test_uri',
},
'2': {
'id': '2',
'server_name': 'test_server_name2',
'url': 'test_url2',
'uri': 'test_uri2',
},
},
1,
)
]
)
mock_build_tag_config.assert_has_calls(
[
call(
{
'server_id': '2',
'server_name': 'test_server_name2',
'tag_address': 'test_tag_address3',
'topics': ['raw_test_schedule2'],
},
ANY,
{
'1': {
'id': '1',
'server_name': 'test_server_name',
'url': 'test_url',
'uri': 'test_uri',
},
'2': {
'id': '2',
'server_name': 'test_server_name2',
'url': 'test_url2',
'uri': 'test_uri2',
},
},
2,
)
]
)
assert result == {
'1': {
'test_server_name': {
'server_id': '1',
'name': 'test_server_name',
'url': 'test_url',
'server_uri': 'test_uri',
'cert_path': None,
'private_key_path': None,
'server_cert_path': None,
'tags': {
'test_tag_address': {
'server_id': '1',
'server_name': 'test_server_name',
'tag_address': 'test_tag_address',
'topics': ['raw_test_schedule'],
}
},
},
'test_server_name2': {
'server_id': '2',
'name': 'test_server_name2',
'url': 'test_url2',
'server_uri': 'test_uri2',
'cert_path': None,
'private_key_path': None,
'server_cert_path': None,
'tags': {
'test_tag_address2': {
'server_id': '2',
'server_name': 'test_server_name2',
'tag_address': 'test_tag_address2',
'topics': ['raw_test_schedule2'],
}
},
},
},
'2': {
'test_server_name2': {
'server_id': '2',
'name': 'test_server_name2',
'url': 'test_url2',
'server_uri': 'test_uri2',
'cert_path': None,
'private_key_path': None,
'server_cert_path': None,
'tags': {
'test_tag_address3': {
'server_id': '2',
'server_name': 'test_server_name2',
'tag_address': 'test_tag_address3',
'topics': ['raw_test_schedule2'],
}
},
}
},
}
@mark.asyncio
@patch(
'orchestrator.activities.formatters.gather_read_tags',
return_value={
'1:test_tag_address': {
'server_id': '1',
'server_name': 'test_server_name',
'tag_address': 'test_tag_address',
'topics': ['raw_test_schedule'],
},
'2:test_tag_address2': {
'server_id': '2',
'server_name': 'test_server_name2',
'tag_address': 'test_tag_address2',
'topics': ['raw_test_schedule2'],
},
},
)
@patch('orchestrator.activities.formatters.build_tag_config', side_effect=ValueError('test_error'))
async def test_process_slots_exception(mock_build_tag_config, mock_gather_read_tags, formatters):
input_data = {
**metadata,
'opc_servers': [
{
'id': '1',
'server_name': 'test_server_name',
'url': 'test_url',
'uri': 'test_uri',
'cert_path': 'test_cert_path',
'private_key_path': 'test_private_key_path',
'server_cert_path': 'test_server_cert_path',
},
{
'id': '2',
'server_name': 'test_server_name2',
'url': 'test_url2',
'uri': 'test_uri2',
'cert_path': 'test_cert_path2',
'private_key_path': 'test_private_key_path2',
'server_cert_path': 'test_server_cert_path2',
},
],
'active_ingestors': ['test_active_ingestor1', 'test_active_ingestor2'],
'pipelines': 'test_gather_read_tags',
}
await formatters.process_slots(input_data)
formatters.send_notification.assert_has_calls(
[
call(
metadata=metadata['metadata'],
notification_id='ORCHESTRATOR_BUILD_TAG_CONFIG_ERROR',
message='test_error',
notification_id='ORCHESTRATOR_SERVER_NOT_FOUND_DURING_SLOT_CONFIGURATION',
message='Servers 2 not found in opc_servers',
block='orchestrator',
level=NotificationLevel.ERROR,
),
call(
metadata=metadata['metadata'],
notification_id='ORCHESTRATOR_BUILD_TAG_CONFIG_ERROR',
message='test_error',
notification_id='ORCHESTRATOR_SERVER_NOT_FOUND_DURING_SLOT_CONFIGURATION',
message='Servers 2 not found in opc_servers',
block='orchestrator',
level=NotificationLevel.ERROR,
),
]
)
assert result == {
'1': slot_mock,
'2': slot_mock,
}
@mark.asyncio
async def test_format_schedule_config(formatters):

View File

@@ -102,7 +102,7 @@ async def test_normalize_schedules_error(temporal_manager):
assert str(e) == 'Test exception'
temporal_manager.send_notification.assert_called_once_with(
metadata=metadata['metadata'],
notification_id='TEMPORAL_NORMALIZE_SCHEDULES_ERROR',
notification_id='SCHEDULER_NORMALIZE_SCHEDULES_ERROR',
message='Failed to normalize schedules: Test exception',
block='normalize_schedules',
level=NotificationLevel.ERROR,
@@ -148,6 +148,7 @@ async def test_create_schedule(
'model_name': 'test-model-name',
'workflow_type': 'test-workflow',
'frequency': '10y',
'offset': '5m',
'data': {'test': 'test'},
'execution_timeout_seconds': 400,
'task_timeout_seconds': 400,
@@ -159,6 +160,7 @@ async def test_create_schedule(
'model_name': 'test-model-name',
'workflow_type': 'test-workflow',
'frequency': '2m',
'offset': '1h',
'data': {'test': 'test'},
'execution_timeout_seconds': 500,
'task_timeout_seconds': 500,
@@ -239,10 +241,15 @@ async def test_create_schedule(
)
mock_schedule_interval_spec.assert_has_calls(
[call(every=timedelta(seconds=60)), call(every=timedelta(seconds=120))]
[
call(every=timedelta(seconds=60), offset=timedelta(seconds=0)),
call(every=timedelta(seconds=120), offset=timedelta(seconds=3600)),
]
)
mock_parse_frequency.assert_has_calls([call('1m'), call('10y'), call('2m')])
mock_parse_frequency.assert_has_calls(
[call('1m'), call('0m'), call('10y'), call('2m'), call('1h')]
)
mock_typed_search_attributes.assert_has_calls(
[

View File

@@ -24,6 +24,7 @@ def test_common_config():
'workflow_type': 'scouter',
'schedule_name': 'test_schedule',
'frequency': '1m',
'offset': '0m',
'max_retry_policy': 1,
'model_id': 'test_model_id',
'model_name': 'test_model_name',
@@ -48,6 +49,7 @@ def test_minimal_retrain():
'workflow_type': 'minimal_retrain',
'schedule_name': 'test_schedule',
'frequency': '1m',
'offset': '0m',
'max_retry_policy': 1,
'model_id': 'test_model_id',
'model_name': 'test_model_name',
@@ -81,6 +83,7 @@ def test_scouter():
'workflow_type': 'scouter',
'schedule_name': 'test_schedule',
'frequency': '1m',
'offset': '0m',
'max_retry_policy': 1,
'model_id': 'test_model_id',
'model_name': 'test_model_name',
@@ -95,6 +98,7 @@ def test_scouter():
'debug_data_package': False,
'execution_timeout_seconds': 300,
'task_timeout_seconds': 300,
'fill_missing_tags': False,
}
assert result == expected
@@ -163,6 +167,7 @@ def test_predictions_batch(mock_process_path_priority, mock_overlap_filter_confi
'workflow_type': 'predictions_batch',
'schedule_name': 'test_schedule',
'frequency': '1m',
'offset': '0m',
'max_retry_policy': 1,
'model_id': 'test_model_id',
'model_name': 'test_model_name',
@@ -245,39 +250,91 @@ def test_gather_read_tags():
def test_build_tag_config():
tag = {'server_id': '1', 'server_name': 'test_server_name', 'tag_address': 'test_tag_address'}
opc_servers = {'1': {'server_name': 'test_server_name', 'url': 'test_url', 'uri': 'test_uri'}}
slot_config = {'1': {}}
i = 1
result = build_tag_config(tag, slot_config, opc_servers, i)
expected = {
tags = [
{
'server_id': '1',
'server_name': 'test_server_name',
'tag_address': 'test_tag_address',
'frequency': 1000,
},
{
'server_id': '2',
'server_name': 'test_server_name2',
'tag_address': 'test_tag_address2',
'frequency': 1000,
},
{
'server_id': '2',
'server_name': 'test_server_name2',
'tag_address': 'test_tag_address3',
'frequency': 300,
},
{
'server_id': '3',
'server_name': 'test_server_name3',
'tag_address': 'test_tag_address4',
'frequency': 200,
},
]
opc_servers = {
'1': {
'test_server_name': {
'server_id': '1',
'name': 'test_server_name',
'url': 'test_url',
'server_uri': 'test_uri',
'cert_path': None,
'private_key_path': None,
'server_cert_path': None,
'tags': {
'test_tag_address': {
'server_id': '1',
'server_name': 'test_server_name',
'tag_address': 'test_tag_address',
}
},
}
}
'server_name': 'test_server_name',
'url': 'test_url',
'uri': 'test_uri',
},
'2': {
'server_name': 'test_server_name2',
'url': 'test_url2',
'uri': 'test_uri2',
},
}
assert result == expected
result = build_tag_config(tags, opc_servers)
def test_build_tag_config_no_server_id():
tag = {'server_id': '1', 'tag_address': 'test_tag_address'}
opc_servers = {}
expected = {
'test_server_name': {
'server_id': '1',
'name': 'test_server_name',
'url': 'test_url',
'server_uri': 'test_uri',
'cert_path': None,
'private_key_path': None,
'server_cert_path': None,
'subscription_period_ms': 500,
'tags': {
'test_tag_address': {
'server_id': '1',
'server_name': 'test_server_name',
'tag_address': 'test_tag_address',
'frequency': 1000,
},
},
},
'test_server_name2': {
'server_id': '2',
'name': 'test_server_name2',
'url': 'test_url2',
'server_uri': 'test_uri2',
'cert_path': None,
'private_key_path': None,
'server_cert_path': None,
'subscription_period_ms': 150,
'tags': {
'test_tag_address2': {
'server_id': '2',
'server_name': 'test_server_name2',
'tag_address': 'test_tag_address2',
'frequency': 1000,
},
'test_tag_address3': {
'server_id': '2',
'server_name': 'test_server_name2',
'tag_address': 'test_tag_address3',
'frequency': 300,
},
},
},
}
try:
build_tag_config(tag, {}, opc_servers, 1)
except ValueError as e:
assert str(e) == 'Server 1 not found in opc_servers'
assert result == (expected, ['3'])

View File

@@ -11,7 +11,7 @@ image:
# This sets the pull policy for images.
pullPolicy: Always
# Overrides the image tag whose default is the chart appVersion.
tag: "0.4.6"
tag: "0.4.9"
# This is for the secrets for pulling an image from a private repository more information can be found here: https://kubernetes.io/docs/tasks/configure-pod-container/pull-image-private-registry/
imagePullSecrets:
@@ -151,7 +151,7 @@ env:
- name: GITHUB_REPO_URL
value: "git@github.com:Aignosi/sientia-dataops-orchestrator_temporal.git"
- name: GITHUB_BRANCH
value: "SIENTIAPDE-1231-ajustar-o-retreino-do-courier-no-laborious"
value: "SIENTIAPDE-1312-melhorias-e-correcoes-nas-pipelines-de-dados"
- name: PYTHON_APP
value: "orchestrator.worker.worker"