SIENTIAPDE-1445

SIENTIAPDE-1273 Add base_scouter and pi_web_api_scouter functions to orchestrator utilities for enhanced configuration management. Update tests to validate new functionalities and ensure proper integration with existing workflows.
This commit is contained in:
vitor-aignosi
2025-12-18 10:02:47 -03:00
parent 62e7a7dee9
commit 61d4b23e69
2 changed files with 252 additions and 11 deletions

View File

@@ -1,5 +1,7 @@
from typing import Any
from orchestrator.utils.converters import parse_frequency
def common_config(config: dict[str, Any]):
"""
@@ -110,6 +112,26 @@ def minimal_retrain(config: dict[str, Any]):
'datetime_columns': config.get('datetime_columns', []),
}
def base_scouter(config: dict[str, Any]):
"""
Base scouter configuration.
"""
filters = {}
for f in config.get('filters', []):
filters[f['filter_name']] = {'policy': f['policy']}
return {
**common_config(config),
'trigger_laborious': False,
'filters': filters,
'schema': 'sientia_data',
'table_name': 'laborious_data',
'retention_time': config.get('tag_retention_minutes', 60) * 60,
'debug_data_package': config.get('debug_data_package', False),
'fill_missing_tags': config.get('fill_missing_tags', False),
}
def scouter(config: dict[str, Any]):
"""
@@ -131,9 +153,7 @@ def scouter(config: dict[str, Any]):
Returns:
dict[str, Any]: Scouter configuration with topic, filters, tags, and retention settings.
"""
filters = {}
for f in config.get('filters', []):
filters[f['filter_name']] = {'policy': f['policy']}
tags = {}
for tag in config['read_tags']:
@@ -143,16 +163,42 @@ def scouter(config: dict[str, Any]):
}
return {
**common_config(config),
**base_scouter(config),
'topic': f'raw_{config["schedule_name"]}',
'trigger_laborious': False,
'filters': filters,
'schema': 'sientia_data',
'table_name': 'laborious_data',
'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),
}
def pi_web_api_scouter(config: dict[str, Any]):
"""
Build PI Web API scouter configuration from pipeline config.
"""
tags = {}
for tag in config['read_tags']:
tags[tag['tag_name']] = {
'webid': tag['webid'],
'aggr_func': tag.get('aggr_func', 'lts'),
'data_range': tag.get('data_range', [-100, 100]),
}
base_config = base_scouter(config)
pi_web_api_config = config['pi_web_api_config']
config_timeout = pi_web_api_config.get('api_timeout', None)
frequency = parse_frequency(base_config['frequency'])
if config_timeout is None or config_timeout > frequency:
config_timeout = frequency
return {
**base_config,
'model_tags': tags,
'pi_web_api_query': {
'endpoint': pi_web_api_config['endpoint'],
'period': pi_web_api_config.get('period', '*-1d'),
'max_count': pi_web_api_config.get('max_count', 1),
'api_timeout': config_timeout,
}
}