Merge pull request #33 from Aignosi/feature/SIENTIAPDE-1478
SIENTIAPDE-1478: Orchestration Platform Enhancements, Worker Refactoring, and Documentation Updates
This commit is contained in:
131
README.md
131
README.md
@@ -60,15 +60,15 @@ The SIENTIA DataOps Orchestrator uses a Temporal-based workflow architecture wit
|
|||||||
### Architecture Principles
|
### Architecture Principles
|
||||||
|
|
||||||
#### 1. **Separation of Concerns**
|
#### 1. **Separation of Concerns**
|
||||||
- **Worker Layer**: Manages Temporal workers, task queues, and application lifecycle
|
- **Worker Layer**: Manages Temporal workers, task queues, Prometheus metrics, and application lifecycle
|
||||||
- **Workflow Layer**: Orchestrates business logic and process coordination
|
- **Workflow Layer**: Orchestrates business logic and process coordination with subworkflows for modularity
|
||||||
- **Activity Layer**: Implements specific operations and external system interactions
|
- **Activity Layer**: Implements specific operations and external system interactions with error handling
|
||||||
- **Data Layer**: Handles data persistence, caching, and external service connections
|
- **Data Layer**: Handles data persistence (MongoDB, PostgreSQL), caching (Redis), and external service connections
|
||||||
|
|
||||||
#### 2. **Task Queue Isolation**
|
#### 2. **Task Queue Isolation**
|
||||||
- **Orchestrator Queue**: Pipeline and resource management workflows
|
- **Orchestrator Queue**: Pipeline and resource management workflows (orchestrator workflow)
|
||||||
- **Alerts Queue**: Real-time error notification workflows
|
- **Alerts Queue**: Real-time error notification workflows (alerts workflow, load_notification_package, process_notifications subworkflows)
|
||||||
- **Reports Queue**: Scheduled reporting and summary workflows
|
- **Reports Queue**: Scheduled reporting and summary workflows (reports workflow, load_notification_package, process_notifications subworkflows)
|
||||||
|
|
||||||
#### 3. **Fault Tolerance & Resilience**
|
#### 3. **Fault Tolerance & Resilience**
|
||||||
- **Automatic Retry Policies**: Configurable retry strategies for transient failures
|
- **Automatic Retry Policies**: Configurable retry strategies for transient failures
|
||||||
@@ -97,12 +97,11 @@ The **Orchestrator** workflow is the main coordination workflow that manages pip
|
|||||||
- **Infrastructure Management**: Creates, updates, and deletes workflow schedules
|
- **Infrastructure Management**: Creates, updates, and deletes workflow schedules
|
||||||
|
|
||||||
**Execution Flow**:
|
**Execution Flow**:
|
||||||
1. **Configuration Loading**: Retrieves pipeline and OPC server configurations from MongoDB
|
1. **Parallel Data Loading**: Concurrently loads pipelines, OPC servers, orchestrated schedules, OPC slots, and active ingestors
|
||||||
2. **Resource Assessment**: Loads current OPC slots and active ingestors from Redis
|
2. **Parallel Processing**: Formats orchestrated schedules and processes new schedules and slots
|
||||||
3. **Schedule Processing**: Formats configurations for different workflow types
|
3. **Parallel Config Creation**: Creates schedule and slot action configurations, normalizes schedules, and creates collections with TTL indexes
|
||||||
4. **Deployment Operations**: Creates, updates, or deletes Temporal schedules
|
4. **Parallel Operations**: Executes slot deletions, slot updates, schedule deletions, schedule creations, and schedule updates concurrently
|
||||||
5. **Resource Updates**: Updates OPC slots and MongoDB timestamps
|
5. **Parallel Reports & Timestamps**: Generates orchestration reports and updates MongoDB timestamps for created/updated/deleted pipelines
|
||||||
6. **Reporting**: Generates comprehensive orchestration reports
|
|
||||||
|
|
||||||
**Input Parameters**:
|
**Input Parameters**:
|
||||||
```json
|
```json
|
||||||
@@ -188,12 +187,13 @@ The **Alerts** workflow processes and sends real-time error notifications to con
|
|||||||
- **Persistent Monitoring**: Tracks and escalates persistent issues
|
- **Persistent Monitoring**: Tracks and escalates persistent issues
|
||||||
|
|
||||||
**Execution Flow**:
|
**Execution Flow**:
|
||||||
1. **Notification Loading**: Retrieves ERROR-level notifications from MongoDB
|
1. **Notification Loading**: Subworkflow loads ERROR-level notifications from MongoDB with timestamp filtering
|
||||||
2. **Timestamp Filtering**: Applies incremental processing using Redis timestamps
|
2. **Configuration Loading**: Loads active receiver group configurations from MongoDB
|
||||||
3. **Group Filtering**: Filters notifications by user group configurations
|
3. **Alert Filtering**: Applies TTL-based filtering to identify core_alerts (new) and persistent_alerts (ongoing issues)
|
||||||
4. **TTL Processing**: Checks notification cache to prevent duplicate alerts
|
4. **Group Filtering**: Filters notifications by user group content policies and ignore lists
|
||||||
5. **Email Generation**: Creates HTML email content for each group
|
5. **Email Processing**: Subworkflow generates HTML emails and sends to receiver groups
|
||||||
6. **Delivery & Logging**: Sends emails and logs results to PostgreSQL
|
6. **Cache Storage**: Stores sent notification cache in Redis with configurable TTL to prevent duplicates
|
||||||
|
7. **Audit Logging**: Logs delivery results to PostgreSQL for monitoring
|
||||||
|
|
||||||
**Input Parameters**:
|
**Input Parameters**:
|
||||||
```json
|
```json
|
||||||
@@ -231,12 +231,11 @@ The **Reports** workflow generates and sends scheduled comprehensive reports to
|
|||||||
- **Audit Trail**: Complete logging of report delivery
|
- **Audit Trail**: Complete logging of report delivery
|
||||||
|
|
||||||
**Execution Flow**:
|
**Execution Flow**:
|
||||||
1. **Data Collection**: Loads all notifications from MongoDB (any level)
|
1. **Notification Loading**: Subworkflow loads all notifications from MongoDB (any level) with timestamp filtering
|
||||||
2. **Timestamp Processing**: Uses incremental loading with Redis timestamps
|
2. **Configuration Loading**: Loads active receiver group configurations from MongoDB
|
||||||
3. **Group Processing**: Applies user group filtering for report customization
|
3. **Report Filtering**: Filters notifications by user group content policies (must include "reports") and ignore lists
|
||||||
4. **Report Generation**: Creates HTML reports with comprehensive summaries
|
4. **Email Processing**: Subworkflow generates HTML reports organized by notification level and model, then sends to receiver groups
|
||||||
5. **Distribution**: Sends reports to configured recipients
|
5. **Audit Logging**: Logs delivery results to PostgreSQL for monitoring and tracking
|
||||||
6. **Audit Logging**: Records delivery status in PostgreSQL
|
|
||||||
|
|
||||||
**Input Parameters**:
|
**Input Parameters**:
|
||||||
```json
|
```json
|
||||||
@@ -288,10 +287,10 @@ flowchart LR
|
|||||||
- `sending_configs` (list[dict]): Active receiver group configurations
|
- `sending_configs` (list[dict]): Active receiver group configurations
|
||||||
|
|
||||||
**Key Activities**:
|
**Key Activities**:
|
||||||
- `get_last_data_timestamp`: Retrieves last processed timestamp from Redis
|
- `get_last_data_timestamp`: Retrieves last processed timestamp from Redis for incremental processing
|
||||||
- `find_documents_in_mongodb`: Loads receiver group configurations
|
- `find_documents_in_mongodb`: Loads active receiver group configurations from MongoDB
|
||||||
- `load_latest_data`: Loads notifications with timestamp filtering
|
- `load_latest_data`: Loads notifications with timestamp filtering from notification_queue collection
|
||||||
- `put_last_data_timestamp`: Updates last processed timestamp
|
- `put_last_data_timestamp`: Updates last processed timestamp in Redis with 5-hour TTL
|
||||||
|
|
||||||
**Architecture**:
|
**Architecture**:
|
||||||
|
|
||||||
@@ -343,10 +342,10 @@ flowchart LR
|
|||||||
- `log_report` (list[dict]): Detailed delivery status for each notification
|
- `log_report` (list[dict]): Detailed delivery status for each notification
|
||||||
|
|
||||||
**Key Activities**:
|
**Key Activities**:
|
||||||
- `build_email_html`: Generates HTML content using Jinja2 templates
|
- `build_email_html`: Generates HTML content using Jinja2 templates organized by notification level and model
|
||||||
- `send_email`: Delivers emails to receiver groups with error handling
|
- `send_email`: Delivers emails to receiver groups with attachment support and automatic SMTP reconnection
|
||||||
- `format_log_report`: Formats delivery results for database storage
|
- `format_log_report`: Formats delivery results into DataFrame structure for database storage, aggregating by notification ID and trigger
|
||||||
- `export_data_to_postgres`: Stores audit logs in PostgreSQL
|
- `export_data_to_postgres`: Stores audit logs in PostgreSQL with timestamp conversion
|
||||||
|
|
||||||
**Architecture**:
|
**Architecture**:
|
||||||
|
|
||||||
@@ -373,10 +372,11 @@ The orchestrator includes an advanced notification filtering system that prevent
|
|||||||
- **Ignore Lists**: Supports notification exclusion per group
|
- **Ignore Lists**: Supports notification exclusion per group
|
||||||
|
|
||||||
### **Report Filtering (`filter_notification_reports`)**
|
### **Report Filtering (`filter_notification_reports`)**
|
||||||
- **Purpose**: Filters all notification levels for comprehensive reports
|
- **Purpose**: Filters all notification levels for comprehensive scheduled reports
|
||||||
- **Comprehensive Coverage**: Includes ERROR, WARNING, INFO, and DEBUG levels
|
- **Comprehensive Coverage**: Includes ERROR, WARNING, INFO, and DEBUG levels
|
||||||
- **Group Customization**: Applies different content policies per receiver group
|
- **Group Customization**: Applies different content policies per receiver group with ignore list support
|
||||||
- **Scheduled Processing**: Designed for regular report generation
|
- **Scheduled Processing**: Designed for regular report generation without TTL-based duplicate prevention
|
||||||
|
- **Duplicate Prevention**: Prevents duplicate notifications within the same report using trigger and notification ID keys
|
||||||
|
|
||||||
### **Notification Caching (`store_notification_cache`)**
|
### **Notification Caching (`store_notification_cache`)**
|
||||||
- **Purpose**: Manages Redis-based notification cache for TTL enforcement
|
- **Purpose**: Manages Redis-based notification cache for TTL enforcement
|
||||||
@@ -389,32 +389,40 @@ The orchestrator includes an advanced notification filtering system that prevent
|
|||||||
#### **Worker (`orchestrator/worker/worker.py`)**
|
#### **Worker (`orchestrator/worker/worker.py`)**
|
||||||
- **Purpose**: Main application orchestrator managing Temporal workers and task queues
|
- **Purpose**: Main application orchestrator managing Temporal workers and task queues
|
||||||
- **Responsibilities**:
|
- **Responsibilities**:
|
||||||
- Temporal client initialization and connection management
|
- Temporal client initialization and connection management with SDK metrics
|
||||||
- Worker lifecycle management and graceful shutdown
|
- Worker lifecycle management and graceful shutdown
|
||||||
- Task queue configuration (orchestrator, alerts, reports)
|
- Task queue configuration (orchestrator, alerts, reports) with dedicated workers
|
||||||
- Prometheus metrics server initialization
|
- Prometheus metrics server initialization on HTTP_METRICS_PORT
|
||||||
|
- Temporal SDK metrics server initialization on HTTP_SDK_METRICS_PORT
|
||||||
- Notification handler setup and configuration
|
- Notification handler setup and configuration
|
||||||
- **Key Features**:
|
- **Key Features**:
|
||||||
- Multi-queue worker management with automatic scaling
|
- Multi-queue worker management with three dedicated workers (orchestrator, alerts, reports)
|
||||||
- Health check endpoints for Kubernetes liveness/readiness probes
|
- Application health metrics (app_up gauge) for Kubernetes liveness/readiness probes
|
||||||
- Graceful shutdown with cleanup procedures
|
- Graceful shutdown with cleanup procedures for all connections
|
||||||
- Comprehensive error handling and metrics collection
|
- Comprehensive error handling and metrics collection
|
||||||
|
- Parallel worker execution using asyncio.gather
|
||||||
|
|
||||||
#### **Activities (`orchestrator/activities/`)**
|
#### **Activities (`orchestrator/activities/`)**
|
||||||
- **Activities**: Main activity orchestrator combining all operations
|
- **Activities**: Main activity orchestrator combining all operations (TemporalManager, SlotManager, Formatters, MongoDB, Email, Postgres)
|
||||||
- **TemporalManager**: Temporal schedule CRUD operations across namespaces
|
- **TemporalManager**: Temporal schedule CRUD operations across scouter and laborious namespaces with search attributes
|
||||||
- **SlotManager**: Redis-based OPC slot and cache management with notification filtering
|
- **SlotManager**: Redis-based OPC slot and cache management with notification filtering and TTL-based duplicate prevention
|
||||||
- **MongoDB**: Document operations, aggregations, and TTL management
|
- **MongoDB**: Document operations, aggregations, timestamp management, and TTL index creation
|
||||||
- **Email**: SMTP operations with HTML generation and attachment support
|
- **Email**: SMTP operations with HTML generation, attachment support, and automatic reconnection handling
|
||||||
- **Formatters**: Configuration processing, slot distribution algorithms, and notification filtering for reports
|
- **Formatters**: Configuration processing, slot distribution algorithms, notification filtering for reports, and schedule/slot orchestration reporting
|
||||||
- **Postgres**: PostgreSQL operations for audit logging and data export (via sientia-dataops-library)
|
- **Postgres**: PostgreSQL operations for audit logging and data export (via sientia-dataops-library)
|
||||||
|
|
||||||
#### **Utilities (`orchestrator/utils/`)**
|
#### **Utilities (`orchestrator/utils/`)**
|
||||||
- **Connectors Configuration**: Database and service configuration management
|
- **Connectors Configuration**: Database and service configuration management from environment variables
|
||||||
- **Email Builder**: HTML email template generation and formatting using Jinja2 templates
|
- **Email Builder**: HTML email template generation and formatting using Jinja2 templates with support for alerts and reports
|
||||||
- **Orchestrator Functions**: Pipeline configuration transformation utilities for scouter, predictions_batch, minimal_retrain, drift, and simple_metrics workflows
|
- **Orchestrator Functions**: Pipeline configuration transformation utilities supporting:
|
||||||
- **Converters**: Data type conversion and validation utilities including frequency parsing
|
- `scouter`: OPC UA data collection workflows
|
||||||
- **Templates**: HTML email templates for alerts and reports
|
- `pi_web_api_scouter`: PI Web API data collection workflows
|
||||||
|
- `predictions_batch`: ML prediction workflows with OPC write-back and multi-stage filtering
|
||||||
|
- `minimal_retrain`: Model retraining workflows with SQL queries
|
||||||
|
- `drift`: Data drift detection workflows
|
||||||
|
- `simple_metrics`: Model performance metrics computation workflows
|
||||||
|
- **Converters**: Data type conversion and validation utilities including frequency parsing for Temporal schedules
|
||||||
|
- **Templates**: HTML email templates for alerts and reports (email_template.html, general_template.html)
|
||||||
|
|
||||||
## 📋 Prerequisites
|
## 📋 Prerequisites
|
||||||
|
|
||||||
@@ -621,18 +629,19 @@ Temporal input configuration sample:
|
|||||||
The Orchestrator system exposes comprehensive Prometheus metrics:
|
The Orchestrator system exposes comprehensive Prometheus metrics:
|
||||||
|
|
||||||
### Application Metrics
|
### Application Metrics
|
||||||
- `app_up`: Application health status (1=healthy, 0=unhealthy)
|
- `app_up`: Application health status gauge (1=healthy, 0=unhealthy) labeled by pod_id
|
||||||
- `email_sent_count`: Email delivery operation count by group
|
- `email_sent_count`: Email delivery counter labeled by pod_id, model_name, pipeline_name, and email_group
|
||||||
|
|
||||||
### Workflow Metrics
|
### Workflow Metrics
|
||||||
- Schedule creation, update, and deletion success rates
|
- Schedule creation, update, and deletion success rates (via notification reports)
|
||||||
- Notification processing times and error rates
|
- Notification processing times and error rates (via PostgreSQL audit logs)
|
||||||
- Resource allocation and slot management metrics
|
- Resource allocation and slot management metrics (via notification reports)
|
||||||
|
- Temporal SDK metrics exposed on HTTP_SDK_METRICS_PORT (default: 9091)
|
||||||
|
|
||||||
### Database Metrics
|
### Database Metrics
|
||||||
- MongoDB query performance and connection health
|
- MongoDB query performance and connection health (via sientia-dataops-library)
|
||||||
- Redis operation counts and response times
|
- Redis operation counts and response times (via sientia-dataops-library)
|
||||||
- PostgreSQL export operations and audit log metrics
|
- PostgreSQL export operations and audit log metrics (via sientia-dataops-library)
|
||||||
|
|
||||||
## 🧪 Testing
|
## 🧪 Testing
|
||||||
|
|
||||||
|
|||||||
File diff suppressed because one or more lines are too long
@@ -58,10 +58,18 @@ schedule_types: dict[str, ScheduleType] = {
|
|||||||
'namespace': 'laborious',
|
'namespace': 'laborious',
|
||||||
'function': predictions_batch,
|
'function': predictions_batch,
|
||||||
},
|
},
|
||||||
|
'xgboost_predictions_batch': {
|
||||||
|
'namespace': 'laborious',
|
||||||
|
'function': predictions_batch,
|
||||||
|
},
|
||||||
'minimal_retrain': {
|
'minimal_retrain': {
|
||||||
'namespace': 'laborious',
|
'namespace': 'laborious',
|
||||||
'function': minimal_retrain,
|
'function': minimal_retrain,
|
||||||
},
|
},
|
||||||
|
'xgboost_minimal_retrain': {
|
||||||
|
'namespace': 'laborious',
|
||||||
|
'function': minimal_retrain,
|
||||||
|
},
|
||||||
'drift': {
|
'drift': {
|
||||||
'namespace': 'laborious',
|
'namespace': 'laborious',
|
||||||
'function': drift,
|
'function': drift,
|
||||||
|
|||||||
@@ -256,17 +256,23 @@ def pi_web_api_scouter(config: dict[str, Any]):
|
|||||||
|
|
||||||
def overlap_filter_config(base_filter_config: dict[str, Any], config: list[dict[str, Any]]):
|
def overlap_filter_config(base_filter_config: dict[str, Any], config: list[dict[str, Any]]):
|
||||||
"""
|
"""
|
||||||
Overlap filter configuration with base filter config.
|
Merge filter configurations with base filter configuration.
|
||||||
|
|
||||||
|
Extends the base filter configuration dictionary by adding or overwriting
|
||||||
|
filters from the provided configuration list. Used in predictions_batch
|
||||||
|
workflows to combine default filters with user-defined custom filters.
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
base_filter_config (dict[str, Any]): Base filter configuration to extend.
|
base_filter_config (dict[str, Any]): Base filter configuration dictionary to extend.
|
||||||
config (list[dict[str, Any]]): List of filter configurations to add, each containing:
|
Each filter entry contains 'policy' and optionally 'config' keys.
|
||||||
- filter_name (str): Name of the filter
|
config (list[dict[str, Any]]): List of filter configurations to merge, each containing:
|
||||||
- policy (str): Filter policy
|
- filter_name (str): Name of the filter to add or update
|
||||||
- config (dict, optional): Additional filter configuration
|
- policy (str): Filter policy (e.g., 'STOP', 'CONTINUE', 'REPEAT')
|
||||||
|
- config (dict, optional): Additional filter-specific configuration
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
dict[str, Any]: Extended filter configuration with new filters added.
|
dict[str, Any]: Extended filter configuration dictionary with merged filters.
|
||||||
|
Filters from config list overwrite or add to base_filter_config entries.
|
||||||
"""
|
"""
|
||||||
for fil in config:
|
for fil in config:
|
||||||
base_filter_config[fil['filter_name']] = {
|
base_filter_config[fil['filter_name']] = {
|
||||||
@@ -330,11 +336,12 @@ def predictions_batch(config: dict[str, Any]):
|
|||||||
- model_retention_minutes (int, optional): Data retention time in minutes (default: 60)
|
- model_retention_minutes (int, optional): Data retention time in minutes (default: 60)
|
||||||
- save_transform (bool, optional): Save transformed data to database (default: True)
|
- save_transform (bool, optional): Save transformed data to database (default: True)
|
||||||
- predictions_storage_policy (str, optional): Prediction storage policy (default: 'lts:1')
|
- predictions_storage_policy (str, optional): Prediction storage policy (default: 'lts:1')
|
||||||
|
- pi_web_api_output_config (dict, optional): PI Web API output configuration for write-back (default: {})
|
||||||
- Additional fields from common_config
|
- Additional fields from common_config
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
dict[str, Any]: Complete predictions batch configuration with OPC output mappings,
|
dict[str, Any]: Complete predictions batch configuration with OPC output mappings,
|
||||||
multi-stage filters, SQL query, and retention policies
|
PI Web API output configuration, multi-stage filters, SQL query, and retention policies
|
||||||
"""
|
"""
|
||||||
tags: dict[str, Any] = {}
|
tags: dict[str, Any] = {}
|
||||||
for tag in config.get('write_tags', []):
|
for tag in config.get('write_tags', []):
|
||||||
@@ -367,6 +374,7 @@ def predictions_batch(config: dict[str, Any]):
|
|||||||
'transform_table_name': 'transformed_data',
|
'transform_table_name': 'transformed_data',
|
||||||
'retention_time': config.get('model_retention_minutes', 60) * 60,
|
'retention_time': config.get('model_retention_minutes', 60) * 60,
|
||||||
'opc_output_config': tags,
|
'opc_output_config': tags,
|
||||||
|
'pi_web_api_output_config': config.get('pi_web_api_output_config', {}),
|
||||||
'input_filters': overlap_filter_config(
|
'input_filters': overlap_filter_config(
|
||||||
{'EMPTY_DATA': {'policy': 'STOP', 'config': {}}}, config.get('input_filters', [])
|
{'EMPTY_DATA': {'policy': 'STOP', 'config': {}}}, config.get('input_filters', [])
|
||||||
),
|
),
|
||||||
@@ -388,17 +396,24 @@ def predictions_batch(config: dict[str, Any]):
|
|||||||
|
|
||||||
def gather_read_tags(pipelines: list[dict[str, Any]]) -> dict[str, Any]:
|
def gather_read_tags(pipelines: list[dict[str, Any]]) -> dict[str, Any]:
|
||||||
"""
|
"""
|
||||||
Gather all read tags from input pipelines.
|
Gather all read tags from scouter pipeline configurations.
|
||||||
|
|
||||||
Collects all read tags from scouter pipelines and organizes them by
|
Collects all read tags from scouter-type pipelines and organizes them by
|
||||||
server_id and tag_address, tracking which topics each tag is associated with.
|
server_id and tag_address, tracking which Kafka topics each tag is associated with.
|
||||||
|
This function is used during slot configuration to aggregate tags across multiple
|
||||||
|
scouter pipelines for efficient OPC server slot allocation.
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
pipelines (list[dict[str, Any]]): The pipeline configurations to process
|
pipelines (list[dict[str, Any]]): List of pipeline configurations to process.
|
||||||
|
Only pipelines with workflow_type 'scouter' are processed. Each scouter
|
||||||
|
pipeline should contain a 'read_tags' list with tag configurations.
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
dict[str, Any]: Dictionary of read tags keyed by "server_id:tag_address",
|
dict[str, Any]: Dictionary of read tags keyed by "server_id:tag_address",
|
||||||
each containing tag configuration and associated topics
|
where each entry contains:
|
||||||
|
- All original tag configuration fields
|
||||||
|
- topics (list[str]): List of Kafka topic names associated with this tag
|
||||||
|
(format: 'raw_{schedule_name}')
|
||||||
"""
|
"""
|
||||||
|
|
||||||
tags = {}
|
tags = {}
|
||||||
@@ -424,24 +439,31 @@ def build_tag_config(
|
|||||||
"""
|
"""
|
||||||
Build tag configuration for a specific slot and OPC server.
|
Build tag configuration for a specific slot and OPC server.
|
||||||
|
|
||||||
Organizes tags by OPC server and calculates the minimum subscription period
|
Organizes tags by OPC server name and calculates the minimum subscription period
|
||||||
based on tag frequencies. Validates that all server IDs exist in the OPC
|
based on tag frequencies. Validates that all server IDs exist in the OPC
|
||||||
servers configuration.
|
servers configuration. The subscription period is set to half of the minimum
|
||||||
|
tag frequency to ensure efficient data collection.
|
||||||
|
|
||||||
Args:
|
Args:
|
||||||
tags (list[dict[str, Any]]): List of tag configurations containing:
|
tags (list[dict[str, Any]]): List of tag configurations containing:
|
||||||
- server_id (str): ID of the OPC server
|
- server_id (str): ID of the OPC server
|
||||||
- tag_address (str): Address of the tag
|
- tag_address (str): Address/path of the OPC tag
|
||||||
- frequency (int): Tag read frequency in milliseconds
|
- frequency (int): Tag read frequency in milliseconds
|
||||||
opc_servers (dict[str, Any]): Dictionary of OPC server configurations
|
- Additional tag-specific configuration fields
|
||||||
|
opc_servers (dict[str, Any]): Dictionary of OPC server configurations keyed by server_id.
|
||||||
|
Each server configuration should contain:
|
||||||
|
- server_name (str): Human-readable server name
|
||||||
|
- url (str): OPC server URL
|
||||||
|
- uri (str): OPC server URI
|
||||||
|
- cert_path (str, optional): Certificate file path
|
||||||
|
- private_key_path (str, optional): Private key file path
|
||||||
|
- server_cert_path (str, optional): Server certificate file path
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
tuple[dict[str, Any], list]: A tuple containing:
|
tuple[dict[str, Any], list]: A tuple containing:
|
||||||
- Slot configuration dictionary organized by server name
|
- Slot configuration dictionary organized by server_name, where each server
|
||||||
- List of server IDs that were not found in opc_servers
|
contains connection details, tags dictionary, and subscription_period_ms
|
||||||
|
- List of server IDs (str) that were not found in opc_servers configuration
|
||||||
Raises:
|
|
||||||
ValueError: If the specified server_id is not found in opc_servers
|
|
||||||
"""
|
"""
|
||||||
|
|
||||||
slot_config = {}
|
slot_config = {}
|
||||||
|
|||||||
72
orchestrator/worker/prepare_worker.py
Normal file
72
orchestrator/worker/prepare_worker.py
Normal file
@@ -0,0 +1,72 @@
|
|||||||
|
import os
|
||||||
|
import re
|
||||||
|
from collections.abc import Sequence
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
from sientia_do.observability.logger import Logger
|
||||||
|
from temporalio.client import Client
|
||||||
|
from temporalio.worker import PollerBehaviorAutoscaling, Worker
|
||||||
|
|
||||||
|
parameters = [
|
||||||
|
('MAX_CONCURRENT_WORKFLOW_TASKS', '200'),
|
||||||
|
('MAX_CONCURRENT_ACTIVITIES', '200'),
|
||||||
|
('MAX_CONCURRENT_LOCAL_ACTIVITIES', '200'),
|
||||||
|
('MAX_CACHED_WORKFLOWS', '200'),
|
||||||
|
('WORKFLOW_POLLER_BEHAVIOUR_MINIMUM', '10'),
|
||||||
|
('WORKFLOW_POLLER_BEHAVIOUR_INITIAL', '100'),
|
||||||
|
('WORKFLOW_POLLER_BEHAVIOUR_MAXIMUM', '200'),
|
||||||
|
('ACTIVITY_POLLER_BEHAVIOUR_MINIMUM', '10'),
|
||||||
|
('ACTIVITY_POLLER_BEHAVIOUR_INITIAL', '100'),
|
||||||
|
('ACTIVITY_POLLER_BEHAVIOUR_MAXIMUM', '200'),
|
||||||
|
]
|
||||||
|
|
||||||
|
|
||||||
|
def camel_to_snake(text: str) -> str:
|
||||||
|
"""Convert camelCase or PascalCase to snake_case."""
|
||||||
|
text = re.sub('(.)([A-Z][a-z]+)', r'\1_\2', text)
|
||||||
|
text = re.sub('([a-z0-9])([A-Z])', r'\1_\2', text)
|
||||||
|
return text.lower()
|
||||||
|
|
||||||
|
|
||||||
|
def prepare_worker(
|
||||||
|
main_workflow: type,
|
||||||
|
other_workflows: Sequence[type],
|
||||||
|
activities: Sequence[Any],
|
||||||
|
temporal_client: Client,
|
||||||
|
logger: Logger,
|
||||||
|
) -> Worker:
|
||||||
|
main_workflow_name = main_workflow.__name__.upper()
|
||||||
|
|
||||||
|
queue_name = f'{camel_to_snake(main_workflow.__name__)}-queue'
|
||||||
|
|
||||||
|
local_workflow_parameters = {}
|
||||||
|
|
||||||
|
for parameter in parameters:
|
||||||
|
local_workflow_parameters[parameter[0]] = int(
|
||||||
|
os.getenv(main_workflow_name + '_' + parameter[0], parameter[1])
|
||||||
|
)
|
||||||
|
|
||||||
|
logger.info(f'Preparing worker for {main_workflow_name} with queue {queue_name}')
|
||||||
|
|
||||||
|
return Worker(
|
||||||
|
temporal_client,
|
||||||
|
task_queue=queue_name,
|
||||||
|
workflows=[main_workflow, *other_workflows],
|
||||||
|
activities=[*activities],
|
||||||
|
max_concurrent_workflow_tasks=local_workflow_parameters['MAX_CONCURRENT_WORKFLOW_TASKS'],
|
||||||
|
max_concurrent_activities=local_workflow_parameters['MAX_CONCURRENT_ACTIVITIES'],
|
||||||
|
max_concurrent_local_activities=local_workflow_parameters[
|
||||||
|
'MAX_CONCURRENT_LOCAL_ACTIVITIES'
|
||||||
|
],
|
||||||
|
max_cached_workflows=local_workflow_parameters['MAX_CACHED_WORKFLOWS'],
|
||||||
|
workflow_task_poller_behavior=PollerBehaviorAutoscaling(
|
||||||
|
minimum=local_workflow_parameters['WORKFLOW_POLLER_BEHAVIOUR_MINIMUM'],
|
||||||
|
initial=local_workflow_parameters['WORKFLOW_POLLER_BEHAVIOUR_INITIAL'],
|
||||||
|
maximum=local_workflow_parameters['WORKFLOW_POLLER_BEHAVIOUR_MAXIMUM'],
|
||||||
|
),
|
||||||
|
activity_task_poller_behavior=PollerBehaviorAutoscaling(
|
||||||
|
minimum=local_workflow_parameters['ACTIVITY_POLLER_BEHAVIOUR_MINIMUM'],
|
||||||
|
initial=local_workflow_parameters['ACTIVITY_POLLER_BEHAVIOUR_INITIAL'],
|
||||||
|
maximum=local_workflow_parameters['ACTIVITY_POLLER_BEHAVIOUR_MAXIMUM'],
|
||||||
|
),
|
||||||
|
)
|
||||||
@@ -1,6 +1,5 @@
|
|||||||
from temporalio import client, workflow
|
from temporalio import client, workflow
|
||||||
from temporalio.runtime import PrometheusConfig, Runtime, TelemetryConfig
|
from temporalio.runtime import PrometheusConfig, Runtime, TelemetryConfig
|
||||||
from temporalio.worker import Worker
|
|
||||||
|
|
||||||
with workflow.unsafe.imports_passed_through():
|
with workflow.unsafe.imports_passed_through():
|
||||||
import asyncio
|
import asyncio
|
||||||
@@ -17,10 +16,10 @@ with workflow.unsafe.imports_passed_through():
|
|||||||
build_email_config,
|
build_email_config,
|
||||||
build_mongodb_config,
|
build_mongodb_config,
|
||||||
build_postgres_config,
|
build_postgres_config,
|
||||||
# build_couchbase_config,
|
|
||||||
build_redis_config,
|
build_redis_config,
|
||||||
build_temporal_config,
|
build_temporal_config,
|
||||||
)
|
)
|
||||||
|
from orchestrator.worker.prepare_worker import prepare_worker
|
||||||
from orchestrator.workflows.alerts import Alerts
|
from orchestrator.workflows.alerts import Alerts
|
||||||
from orchestrator.workflows.orchestrator import Orchestrator
|
from orchestrator.workflows.orchestrator import Orchestrator
|
||||||
from orchestrator.workflows.reports import Reports
|
from orchestrator.workflows.reports import Reports
|
||||||
@@ -104,10 +103,10 @@ async def main():
|
|||||||
logger.custom_info('Starting Workers...', metadata=metadata)
|
logger.custom_info('Starting Workers...', metadata=metadata)
|
||||||
|
|
||||||
workers = [
|
workers = [
|
||||||
Worker(
|
prepare_worker(
|
||||||
temporal_client,
|
temporal_client=temporal_client,
|
||||||
task_queue='orchestrator-queue',
|
main_workflow=Orchestrator,
|
||||||
workflows=[Orchestrator],
|
other_workflows=[],
|
||||||
activities=[
|
activities=[
|
||||||
# Redis
|
# Redis
|
||||||
activities.load_active_ingestors,
|
activities.load_active_ingestors,
|
||||||
@@ -137,11 +136,12 @@ async def main():
|
|||||||
activities.report_slot_orchestration,
|
activities.report_slot_orchestration,
|
||||||
activities.format_schedule_config,
|
activities.format_schedule_config,
|
||||||
],
|
],
|
||||||
|
logger=logger,
|
||||||
),
|
),
|
||||||
Worker(
|
prepare_worker(
|
||||||
temporal_client,
|
temporal_client=temporal_client,
|
||||||
task_queue='alerts-queue',
|
main_workflow=Alerts,
|
||||||
workflows=[Alerts, LoadNotificationPackage, ProcessNotifications],
|
other_workflows=[LoadNotificationPackage, ProcessNotifications],
|
||||||
activities=[
|
activities=[
|
||||||
# Load notifications
|
# Load notifications
|
||||||
activities.get_last_data_timestamp,
|
activities.get_last_data_timestamp,
|
||||||
@@ -158,11 +158,12 @@ async def main():
|
|||||||
# Store notification cache
|
# Store notification cache
|
||||||
activities.store_notification_cache,
|
activities.store_notification_cache,
|
||||||
],
|
],
|
||||||
|
logger=logger,
|
||||||
),
|
),
|
||||||
Worker(
|
prepare_worker(
|
||||||
temporal_client,
|
temporal_client=temporal_client,
|
||||||
task_queue='reports-queue',
|
main_workflow=Reports,
|
||||||
workflows=[Reports, LoadNotificationPackage, ProcessNotifications],
|
other_workflows=[LoadNotificationPackage, ProcessNotifications],
|
||||||
activities=[
|
activities=[
|
||||||
# Load notifications
|
# Load notifications
|
||||||
activities.get_last_data_timestamp,
|
activities.get_last_data_timestamp,
|
||||||
@@ -177,6 +178,7 @@ async def main():
|
|||||||
activities.format_log_report,
|
activities.format_log_report,
|
||||||
activities.export_data_to_postgres,
|
activities.export_data_to_postgres,
|
||||||
],
|
],
|
||||||
|
logger=logger,
|
||||||
),
|
),
|
||||||
]
|
]
|
||||||
|
|
||||||
|
|||||||
@@ -4,5 +4,5 @@ sqlalchemy
|
|||||||
redis
|
redis
|
||||||
pymongo
|
pymongo
|
||||||
jinja2
|
jinja2
|
||||||
git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.7.1
|
git+ssh://git@github.com/Aignosi/sientia-dataops-library.git@1.8.2
|
||||||
prometheus-client
|
prometheus-client
|
||||||
|
|||||||
@@ -8,4 +8,4 @@ sonar.python.coverage.reportPaths=coverage.xml
|
|||||||
sonar.python.xunit.reportPath=pytest.xml
|
sonar.python.xunit.reportPath=pytest.xml
|
||||||
sonar.python.version=3.11
|
sonar.python.version=3.11
|
||||||
sonar.projectVersion=1.0.0
|
sonar.projectVersion=1.0.0
|
||||||
sonar.coverage.exclusions=orchestrator/worker/worker.py
|
sonar.coverage.exclusions=orchestrator/worker/*
|
||||||
|
|||||||
@@ -207,6 +207,15 @@ def test_predictions_batch(mock_process_path_priority, mock_overlap_filter_confi
|
|||||||
{'server_id': 'test_server_id', 'type': 'prediction', 'addr': 'test_addr'},
|
{'server_id': 'test_server_id', 'type': 'prediction', 'addr': 'test_addr'},
|
||||||
{'server_id': 'test_server_id', 'type': 'confidence', 'addr': 'test_addr'},
|
{'server_id': 'test_server_id', 'type': 'confidence', 'addr': 'test_addr'},
|
||||||
],
|
],
|
||||||
|
'pi_web_api_output_config': {
|
||||||
|
'endpoint': 'test_endpoint',
|
||||||
|
'prediction_tags': {
|
||||||
|
'tag_1': 'webid_1',
|
||||||
|
},
|
||||||
|
'confidence_tags': {
|
||||||
|
'tag_2': 'webid_2',
|
||||||
|
},
|
||||||
|
},
|
||||||
'input_filters': [{'filter_name': 'test_filter_name', 'policy': 'test_policy'}],
|
'input_filters': [{'filter_name': 'test_filter_name', 'policy': 'test_policy'}],
|
||||||
'mlflow_transform_filters': [{'filter_name': 'test_filter_name', 'policy': 'test_policy'}],
|
'mlflow_transform_filters': [{'filter_name': 'test_filter_name', 'policy': 'test_policy'}],
|
||||||
'mlflow_predict_filters': [{'filter_name': 'test_filter_name', 'policy': 'test_policy'}],
|
'mlflow_predict_filters': [{'filter_name': 'test_filter_name', 'policy': 'test_policy'}],
|
||||||
@@ -257,6 +266,15 @@ def test_predictions_batch(mock_process_path_priority, mock_overlap_filter_confi
|
|||||||
'confidence_tags': {'test_addr': {'data_type': 'float'}},
|
'confidence_tags': {'test_addr': {'data_type': 'float'}},
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
|
'pi_web_api_output_config': {
|
||||||
|
'endpoint': 'test_endpoint',
|
||||||
|
'prediction_tags': {
|
||||||
|
'tag_1': 'webid_1',
|
||||||
|
},
|
||||||
|
'confidence_tags': {
|
||||||
|
'tag_2': 'webid_2',
|
||||||
|
},
|
||||||
|
},
|
||||||
'input_filters': {'test_filter_name': {'policy': 'test_policy', 'config': {}}},
|
'input_filters': {'test_filter_name': {'policy': 'test_policy', 'config': {}}},
|
||||||
'mlflow_transform_filters': {'test_filter_name': {'policy': 'test_policy', 'config': {}}},
|
'mlflow_transform_filters': {'test_filter_name': {'policy': 'test_policy', 'config': {}}},
|
||||||
'mlflow_predict_filters': {'test_filter_name': {'policy': 'test_policy', 'config': {}}},
|
'mlflow_predict_filters': {'test_filter_name': {'policy': 'test_policy', 'config': {}}},
|
||||||
|
|||||||
@@ -164,7 +164,7 @@ env:
|
|||||||
- name: GITHUB_REPO_URL
|
- name: GITHUB_REPO_URL
|
||||||
value: "git@github.com:Aignosi/sientia-dataops-orchestrator_temporal.git"
|
value: "git@github.com:Aignosi/sientia-dataops-orchestrator_temporal.git"
|
||||||
- name: GITHUB_BRANCH
|
- name: GITHUB_BRANCH
|
||||||
value: "fix/SIENTIAPDE-1461"
|
value: "feature/SIENTIAPDE-1478"
|
||||||
- name: PYTHON_APP
|
- name: PYTHON_APP
|
||||||
value: "orchestrator.worker.worker"
|
value: "orchestrator.worker.worker"
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user