From 93e50883c1f1a828fa11a385102d2c8eb911fe1e Mon Sep 17 00:00:00 2001 From: vitor-aignosi Date: Mon, 1 Sep 2025 09:19:46 -0300 Subject: [PATCH] SIENTIAPDE-1182 Update README.md to enhance system architecture documentation - Removed outdated system overview diagram and replaced it with updated architecture diagrams for workflows, including Data Loading, Processing, and Notifications. - Added detailed architecture sections for the Alerts and Reports workflows, including input parameters and flowcharts to illustrate processes. - Improved clarity and organization of the README to better reflect the current system design and functionality. --- README.md | 217 ++++++++++++++++++++++++++++++++---------------------- 1 file changed, 127 insertions(+), 90 deletions(-) diff --git a/README.md b/README.md index 2d3a177..8665909 100644 --- a/README.md +++ b/README.md @@ -22,96 +22,6 @@ A high-performance, scalable workflow orchestration system built on Temporal.io The SIENTIA DataOps Orchestrator uses a Temporal-based workflow architecture with clear separation of concerns and robust error handling. The architecture is designed for high availability, scalability, and operational excellence in data pipeline environments. -### System Overview - -``` -┌───────────────────────────────────────────────────────────────────────────────┐ -│ Temporal Cluster │ -│ ┌─────────────────┐ ┌──────────────────┐ ┌─────────────────────────┐ │ -│ │ Main Worker │ │ Temporal Client │ │ Task Queues │ │ -│ │ │◄──►│ │◄──►│ │ │ -│ │ - Metrics Server│ │ - Namespace Mgmt │ │ - orchestrator-queue │ │ -│ │ - Notifications │ │ - Runtime Config │ │ - alerts-queue │ │ -│ │ - Lifecycle │ │ - Connection │ │ - reports-queue │ │ -│ │ - Health Checks │ │ - Security │ │ - Auto-scaling │ │ -│ └─────────────────┘ └──────────────────┘ └─────────────────────────┘ │ -└───────────────────────────────────────────────────────────────────────────────┘ - │ - ▼ - ┌────────────────────────────────────────────────────────┐ - │ Workflow Layer │ - │ ┌─────────────────┐ ┌─────────────────────────────┐ │ - │ │ Orchestrator │ │ Subworkflows │ │ - │ │ │ │ │ │ - │ │ - Pipeline Mgmt │ │ - LoadNotificationPackage │ │ - │ │ - Resource Mgmt │ │ - ProcessNotifications │ │ - │ │ - Schedule Mgmt │ │ - Error Handling │ │ - │ └─────────────────┘ └─────────────────────────────┘ │ - │ ┌─────────────────┐ ┌─────────────────────────────┐ │ - │ │ Alerts │ │ Reports │ │ - │ │ │ │ │ │ - │ │ - Error Alerts │ │ - Scheduled Reports │ │ - │ │ - TTL Filtering │ │ - Summary Generation │ │ - │ │ - Group Filtering│ │ - Comprehensive Logging │ │ - │ └─────────────────┘ └─────────────────────────────┘ │ - └────────────────────────────────────────────────────────┘ - │ - ▼ - ┌────────────────────────────────────────────────────────┐ - │ Activity Layer │ - │ ┌─────────────────┐ ┌─────────────────────────────┐ │ - │ │ Temporal Manager│ │ MongoDB Operations │ │ - │ │ │ │ │ │ - │ │ - Schedule CRUD │ │ - Document Queries │ │ - │ │ - Multi-namespace│ │ - Aggregation Pipelines │ │ - │ │ - Normalization │ │ - Timestamp Management │ │ - │ │ - Error Handling│ │ - TTL Collections │ │ - │ └─────────────────┘ └─────────────────────────────┘ │ - │ ┌─────────────────┐ ┌─────────────────────────────┐ │ - │ │ Slot Manager │ │ Email Services │ │ - │ │ │ │ │ │ - │ │ - OPC Slots │ │ - HTML Generation │ │ - │ │ - Active Ingest │ │ - SMTP Management │ │ - │ │ - Redis Ops │ │ - Attachment Handling │ │ - │ │ - Cache Mgmt │ │ - Auto Reconnection │ │ - │ └─────────────────┘ └─────────────────────────────┘ │ - └────────────────────────────────────────────────────────┘ - │ - ▼ - ┌────────────────────────────────────────────────────────┐ - │ Data Services │ - │ ┌─────────────────┐ ┌─────────────────────────────┐ │ - │ │ Redis │ │ MongoDB │ │ - │ │ │ │ │ │ - │ │ - OPC Slots │ │ - Pipeline Configs │ │ - │ │ - Timestamps │ │ - Notification Queue │ │ - │ │ - Notification │ │ - Receiver Groups │ │ - │ │ - Cache TTL │ │ - Orchestrated Schedules │ │ - │ └─────────────────┘ └─────────────────────────────┘ │ - │ ┌─────────────────┐ ┌─────────────────────────────┐ │ - │ │ PostgreSQL │ │ SMTP Services │ │ - │ │ │ │ │ │ - │ │ - Log Reports │ │ - Email Delivery │ │ - │ │ - Audit Trails │ │ - Group Management │ │ - │ │ - Metrics Data │ │ - Attachment Support │ │ - │ │ - Data Export │ │ - Security & Auth │ │ - │ └─────────────────┘ └─────────────────────────────┘ │ - └────────────────────────────────────────────────────────┘ - │ - ▼ - ┌─────────────────────────────────────────────────────────┐ - │ External Systems │ - │ ┌─────────────────┐ ┌─────────────────────────────┐ │ - │ │ Prometheus │ │ Kubernetes │ │ - │ │ │ │ │ │ - │ │ - App Metrics │ │ - Container Orchestration │ │ - │ │ - Email Metrics │ │ - Health Checks │ │ - │ │ - SDK Metrics │ │ - Auto-scaling │ │ - │ │ - Alerting │ │ - Resource Management │ │ - │ └─────────────────┘ └─────────────────────────────┘ │ - └─────────────────────────────────────────────────────────┘ -``` - ### Architecture Principles #### 1. **Separation of Concerns** @@ -177,6 +87,61 @@ The **Orchestrator** workflow is the main coordination workflow that manages pip } ``` +#### Architecture + +```mermaid +flowchart LR + subgraph "Parallel Data Loading" + A[1. aggregate_documents_in_mongodb] + B[2. find_documents_in_mongodb
OPC Servers] + C[3. find_documents_in_mongodb
Orchestrated Schedules] + D[4. load_opc_slots] + E[5. load_active_ingestors] + end + + subgraph "Parallel Processing" + F[6. format_schedule_config] + G[7. process_schedules] + H[8. process_slots] + end + + subgraph "Parallel Config Creation" + I[9. create_schedule_config] + J[10. create_slot_config] + K[11. normalize_schedules] + L[12. create_collection_with_ttl_index] + end + + subgraph "Parallel Operations" + M[13. delete_slots] + N[14. update_slots] + O[15. delete_schedules] + P[16. create_schedules] + Q[17. update_schedules] + end + + subgraph "Parallel Reports & Timestamps" + R[18. report_schedule_orchestration] + S[19. report_slot_orchestration] + T[20. update_pipelines_timestamps] + U[21. create_pipelines_timestamps] + V[22. delete_pipelines_timestamps] + end + + A --> F + F --> I + I --> M + M --> R + + A -.-> MongoDB1[(MongoDB)] + D -.-> Redis1[(Redis)] + K -.-> Temporal[(Temporal)] + R -.-> Reports[Reports] +``` + +Note: Green blocks (🟩) indicate parallel processing operations that run concurrently for improved performance. + + #### 2. Alerts Workflow (`alerts.py`) The **Alerts** workflow processes and sends real-time error notifications to configured user groups with intelligent filtering and duplicate prevention. @@ -204,6 +169,22 @@ The **Alerts** workflow processes and sends real-time error notifications to con } ``` +#### Architecture + +```mermaid +flowchart LR + A[1. load_notification_package🔃] --> B[2. filter_notification_alerts] --> C[3. process_notifications🔃] + + A -.-> MongoDB[(MongoDB)] + A -.-> Redis[(Redis)] + B -.-> Filters[Report Filters] + C -.-> Email[(Email)] + C -.-> PostgreSQL[(PostgreSQL)] + + style A fill:#000,color:#fff + style C fill:#000,color:#fff +``` + #### 3. Reports Workflow (`reports.py`) The **Reports** workflow generates and sends scheduled comprehensive reports to configured user groups. @@ -222,6 +203,29 @@ The **Reports** workflow generates and sends scheduled comprehensive reports to 5. **Distribution**: Sends reports to configured recipients 6. **Audit Logging**: Records delivery status in PostgreSQL +**Input Parameters**: +```json +{ + "schedule_name": "reports", +} +``` + +#### Architecture + +```mermaid +flowchart LR + A[1. load_notification_package🔃] --> B[2. filter_notification_reports] --> C[3. process_notifications🔃] + + A -.-> MongoDB[(MongoDB)] + A -.-> Redis[(Redis)] + B -.-> Filters[Report Filters] + C -.-> Email[(Email)] + C -.-> PostgreSQL[(PostgreSQL)] + + style A fill:#000,color:#fff + style C fill:#000,color:#fff +``` + ### Subworkflows #### 1. Load Notification Package (`load_notification_package.py`) @@ -248,6 +252,26 @@ The **Reports** workflow generates and sends scheduled comprehensive reports to - `notification_package` (list[dict]): Retrieved notifications - `sending_configs` (list[dict]): Active receiver group configurations +**Architecture**: + +```mermaid +flowchart LR + subgraph "Parallel Loading" + A[1. get_last_data_timestamp] + B[2. find_documents_in_mongodb
Receiver Groups] + end + + C[3. load_latest_data] --> D[4. put_last_data_timestamp] + + A --> C + B --> D + + A -.-> Redis1[(Redis)] + B -.-> MongoDB1[(MongoDB)] + C -.-> MongoDB2[(MongoDB)] + D -.-> Redis2[(Redis)] +``` + #### 2. Process Notifications (`process_notifications.py`) **Purpose**: Handles email generation, delivery, and audit logging for notification workflows. @@ -276,6 +300,19 @@ The **Reports** workflow generates and sends scheduled comprehensive reports to **Returns**: - `log_report` (list[dict]): Detailed delivery status for each notification +**Architecture**: + +```mermaid +flowchart LR + A[1. build_email_html] --> B[2. send_email] --> C[3. format_log_report] --> D[4. export_data_to_postgres] + + A -.-> HTML[HTML Generator] + B -.-> SMTP[(SMTP)] + C -.-> Formatter[Log Formatter] + D -.-> PostgreSQL[(PostgreSQL)] + +``` + ### Key Components #### **Worker (`orchestrator/worker/worker.py`)**