diff --git a/orchestrator/activities/formatters.py b/orchestrator/activities/formatters.py index 4d186e1..e51012d 100644 --- a/orchestrator/activities/formatters.py +++ b/orchestrator/activities/formatters.py @@ -163,36 +163,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) diff --git a/orchestrator/utils/orchestrator_functions.py b/orchestrator/utils/orchestrator_functions.py index b59bfac..02371a1 100644 --- a/orchestrator/utils/orchestrator_functions.py +++ b/orchestrator/utils/orchestrator_functions.py @@ -240,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. @@ -259,32 +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()] - frequencies = [int(x['frequency']) for x in slot_config[f'{i}'][server_name]['tags'].values()] + min_frequency = min(frequencies) if frequencies else 1000 - min_frequency = min(frequencies) if frequencies else 1000 + slot_config[server_name]['subscription_period_ms'] = min_frequency / 2 - slot_config[f"{i}"][server_name]['subscription_period_ms'] = min_frequency/2 - - return slot_config + return slot_config, notifications diff --git a/tests/orchestrator/activities/test_formatters.py b/tests/orchestrator/activities/test_formatters.py index f111123..234a0e2 100644 --- a/tests/orchestrator/activities/test_formatters.py +++ b/tests/orchestrator/activities/test_formatters.py @@ -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): diff --git a/tests/orchestrator/activities/test_temporal_manager.py b/tests/orchestrator/activities/test_temporal_manager.py index d8af937..0d34bcf 100644 --- a/tests/orchestrator/activities/test_temporal_manager.py +++ b/tests/orchestrator/activities/test_temporal_manager.py @@ -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, diff --git a/tests/orchestrator/utils/test_orchestrator_functions.py b/tests/orchestrator/utils/test_orchestrator_functions.py index 5a1f64f..f74a936 100644 --- a/tests/orchestrator/utils/test_orchestrator_functions.py +++ b/tests/orchestrator/utils/test_orchestrator_functions.py @@ -98,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 @@ -249,41 +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', 'frequency': 1000} - opc_servers = {'1': {'server_name': 'test_server_name', 'url': 'test_url', 'uri': 'test_uri', 'subscription_period_ms': 1000}} - slot_config: dict = {'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, - 'subscription_period_ms': 500, - 'tags': { - 'test_tag_address': { - 'server_id': '1', - 'server_name': 'test_server_name', - 'tag_address': 'test_tag_address', - 'frequency': 1000, - } - }, - } - } + '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: dict = {} + 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'])