Enhance data handling and export processes in Laborious workflows - Updated `gates.py` to improve data quality validation, filtering, and formatting operations, including enhanced metrics recording. - Refined `mlflow.py` to better manage model transformations and reference data retrieval from MLflow Model Registry. - Enhanced `format_and_export_prediction.py` to support separate export of transformed data, improving flexibility in data handling. - Added comprehensive test coverage for new functionalities, including transformed data formatting and retrain report generation. - Improved documentation in `README.md` to reflect changes in activities and workflows, ensuring clarity on data processing and export paths.
946 lines
33 KiB
Markdown
946 lines
33 KiB
Markdown
# Sientia DataOps Laborious
|
|
|
|
A comprehensive, Temporal-based ML orchestration system for industrial data processing and model inference. Laborious delivers enterprise-grade batch prediction, model management, optional real-time export (OPC), and automated retraining with strong data quality validation and observability.
|
|
|
|
## 📑 Table of Contents
|
|
|
|
- [Features](#features)
|
|
- [Core Functionality](#core-functionality)
|
|
- [Advanced Capabilities](#advanced-capabilities)
|
|
- [Development & Quality Assurance](#development--quality-assurance)
|
|
- [Architecture](#architecture)
|
|
- [Architecture Principles](#architecture-principles)
|
|
- [Key Components](#key-components)
|
|
- [Data Flow Architecture](#data-flow-architecture)
|
|
- [Security Architecture](#security-architecture)
|
|
- [Workflows](#workflows)
|
|
- [Predictions Batch Workflow](#1-predictions-batch-workflow-predictions_batchpy)
|
|
- [Prediction Process Workflow](#2-prediction-process-workflow-prediction_processpy)
|
|
- [Format and Export Prediction Workflow](#3-format-and-export-prediction-workflow-format_and_export_predictionpy)
|
|
- [Minimal Retrain Workflow](#4-minimal-retrain-workflow-minimal_retrainpy)
|
|
- [Installation & Setup](#installation--setup)
|
|
- [Prerequisites](#prerequisites)
|
|
- [Environment Setup](#environment-setup)
|
|
- [Temporal Namespace Setup](#temporal-namespace-setup)
|
|
- [Local Development Setup](#local-development-setup)
|
|
- [How to Run](#how-to-run)
|
|
- [Running the Laborious Application](#running-the-laborious-application)
|
|
- [Running Tests and Coverage](#running-tests-and-coverage)
|
|
- [Manual Test Execution](#manual-test-execution)
|
|
- [Manual Application Execution](#manual-application-execution)
|
|
- [Code Quality & Validation](#code-quality--validation)
|
|
- [Overview](#overview)
|
|
- [Validation Tools](#validation-tools)
|
|
- [Tools Installation](#tools-installation)
|
|
- [Complete Validation](#complete-validation)
|
|
- [Automatic Fixes](#automatic-fixes)
|
|
- [Configuration](#configuration)
|
|
- [CI/CD Integration](#cicd-integration)
|
|
- [Best Practices](#best-practices)
|
|
- [Testing](#testing)
|
|
- [Test Structure](#test-structure)
|
|
- [Test Execution](#test-execution)
|
|
- [Monitoring and Metrics](#monitoring-and-metrics)
|
|
- [Application Health Metrics](#application-health-metrics)
|
|
- [Prediction Operation Metrics](#prediction-operation-metrics)
|
|
- [OPC Export Metrics](#opc-export-metrics)
|
|
- [Data Quality Metrics](#data-quality-metrics)
|
|
- [Configuration](#configuration-1)
|
|
- [Environment Variables](#environment-variables)
|
|
- [OPC Configuration](#opc-configuration)
|
|
- [Workflow Configuration](#workflow-configuration)
|
|
- [Development](#development)
|
|
- [Project Structure](#project-structure)
|
|
- [Adding New Features](#adding-new-features)
|
|
- [Troubleshooting](#troubleshooting)
|
|
- [Common Issues](#common-issues)
|
|
- [Debug Mode](#debug-mode)
|
|
- [Performance Tuning](#performance-tuning)
|
|
- [Key Parameters](#key-parameters)
|
|
- [Scaling Considerations](#scaling-considerations)
|
|
- [Contributing](#contributing)
|
|
- [Code Quality Standards](#code-quality-standards)
|
|
- [License](#license)
|
|
- [Support](#support)
|
|
|
|
## Features
|
|
|
|
### Core Functionality
|
|
- **Batch Prediction Processing**: High-throughput ML inference using MLFlow models
|
|
- **Temporal Workflow Orchestration**: Robust workflow management with retries and fault tolerance
|
|
- **Data Quality Gates**: Configurable filtering for input data and MLFlow API responses
|
|
- **Multi-Model Support**: Flexible model management with retention and versioning
|
|
- **Optional Real-time Export**: PostgreSQL persistence and OPC server integration for industrial systems
|
|
- **Comprehensive Monitoring**: Prometheus metrics and structured logging for observability
|
|
|
|
### Advanced Capabilities
|
|
- **Incremental Data Processing**: Timestamp-based loading to avoid reprocessing
|
|
- **Configurable Data Retention**: Model retention policies with automatic cleanup
|
|
- **Notification System**: Integrated alerting via MongoDB
|
|
- **Scalable Architecture**: Kubernetes-ready with horizontal scaling
|
|
- **Model Retraining**: Automated retraining workflows with production model updates
|
|
|
|
### Development & Quality Assurance
|
|
- **Code Quality Tools**: Ruff (lint/format), mypy (types), Bandit (security)
|
|
- **Automated Validation**: `validate.sh` and CI quality gates
|
|
- **Comprehensive Testing**: pytest with async support and high coverage
|
|
- **Type Safety**: Static type checking with mypy
|
|
- **Coverage Visualization**: Coverage Gutters integration
|
|
|
|
## Architecture
|
|
|
|
Laborious uses a Temporal-based architecture with strong separation of concerns and defensive error handling for production ML.
|
|
|
|
### Architecture Principles
|
|
|
|
#### 1. **Separation of Concerns**
|
|
- **Worker Layer**: Temporal workers, task queues, lifecycle
|
|
- **Workflow Layer**: Business orchestration and coordination
|
|
- **Activity Layer**: External system interactions and isolated operations
|
|
- **Data Layer**: Persistence, caching, connectors
|
|
|
|
#### 2. **Fault Tolerance & Resilience**
|
|
- **Automatic Retry Policies** for transient failures
|
|
- **Graceful Degradation** and circuit breaking for dependencies
|
|
- **Detailed Error Handling** with notifications
|
|
|
|
#### 3. **Scalability & Performance**
|
|
- **Horizontal Scaling**: Multiple worker instances for load distribution
|
|
- **Task Queue Isolation**: Separate queues for different workflow types
|
|
- **Connection Pooling**: Optimized database and external service connections
|
|
- **Asynchronous Processing**: Non-blocking operations for improved throughput
|
|
|
|
#### 4. **Observability & Monitoring**
|
|
- **Prometheus Metrics**: Comprehensive system and business metrics
|
|
- **Structured Logging**: Consistent log format with correlation IDs
|
|
- **Health Checks**: Endpoint health monitoring and alerting
|
|
- **Performance Tracing**: Request flow tracking and bottleneck identification
|
|
|
|
### Key Components
|
|
|
|
#### **Worker (`laborious/worker/worker.py`)**
|
|
- Temporal client setup, worker lifecycle, task queues
|
|
- Metrics server initialization, notification handler setup
|
|
- Graceful shutdown and autoscaling-friendly behavior
|
|
|
|
#### **Workflows (`laborious/workflows/`)**
|
|
- `predictions_batch.py`: Batch prediction entry point
|
|
- `sub_workflows/prediction_process.py`: Core prediction pipeline
|
|
- `sub_workflows/format_and_export_prediction.py`: Formatting and export
|
|
- `minimal_retrain.py`: Automated model retraining and production update
|
|
|
|
#### **Activities (`laborious/activities/`)**
|
|
- `gates.py`: Data quality validation, filtering, and data formatting operations
|
|
- Input/response/content gates for quality validation
|
|
- Prediction and transformed data formatting
|
|
- Retrain report formatting and metrics recording
|
|
- `mlflow.py`: Transform, predict, and model management operations
|
|
- MLFlow model transformation and prediction
|
|
- Model retraining and production updates
|
|
- Reference data retrieval from MLflow Model Registry
|
|
- `opc.py`: OPC UA export to industrial systems (optional)
|
|
- `activities.py`: Aggregates activity interfaces
|
|
|
|
#### **Data Services (`laborious/utils/`)**
|
|
- `connectors_config.py`: Env-driven configuration builders
|
|
- `repository/model_repository.py`: MLFlow operations and retraining
|
|
- `repository/opc_repository.py`: OPC communication and writes
|
|
- `filters/conditional_filters.py` and `filters/mlflow_filters.py`
|
|
|
|
### Data Flow Architecture
|
|
|
|
#### **1. Batch Prediction Pipeline**
|
|
```
|
|
Input Data (PostgreSQL) → Data Quality Gates → MLFlow Transform →
|
|
MLFlow Prediction → Response Validation → Format & Export
|
|
├─→ Predictions → PostgreSQL [+ OPC]
|
|
└─→ Transformed Data → PostgreSQL (optional)
|
|
```
|
|
|
|
#### **2. Model Retraining Pipeline**
|
|
```
|
|
Training Data → Model Retraining → Quality Validation →
|
|
Production Update → Notification & Monitoring
|
|
```
|
|
|
|
### Security Architecture
|
|
|
|
#### **Authentication & Authorization**
|
|
- **MLFlow API Authentication**: Username/password
|
|
- **Database Security**: Encrypted connections and credential management
|
|
- **OPC Certificates** (if enabled): Client/server certs
|
|
- **Kubernetes Secrets**: Secure secret storage
|
|
|
|
#### **Network Security**
|
|
- TLS/SSL, network policies, service mesh, firewalls, VPN
|
|
|
|
#### **Data Security**
|
|
- At-rest/in-transit encryption, RBAC, audit logging, lifecycle management
|
|
|
|
## Workflows
|
|
|
|
### 1. Predictions Batch Workflow (`predictions_batch.py`)
|
|
|
|
The **PredictionsBatch** workflow is the main entry point for batch prediction pipelines. It orchestrates the complete prediction process and implements a robust data loading and processing pattern.
|
|
|
|
#### Purpose
|
|
- **Batch Prediction Orchestration**: Coordinates data loading and prediction processing
|
|
- **Data Preparation**: Loads data using custom SQL queries with configurable schemas
|
|
- **Workflow Delegation**: Delegates actual prediction processing to the PredictionProcess workflow
|
|
- **Configuration Management**: Handles model configuration, filters, and retention policies
|
|
|
|
#### Execution Flow
|
|
1. **Data Loading**: Executes custom SQL query to load data from PostgreSQL
|
|
2. **Input Preparation**: Prepares prediction input with metadata and configuration
|
|
3. **Workflow Delegation**: Spawns PredictionProcess child workflow for actual processing
|
|
4. **Error Handling**: Implements comprehensive error handling with retry policies
|
|
|
|
#### Key Features
|
|
- **Custom Query Support**: Flexible SQL-based data loading
|
|
- **Schema Configuration**: Configurable data schema definitions
|
|
- **Automatic Retry**: Implements Temporal retry policies for fault tolerance
|
|
- **Timeout Management**: 60-second timeout for all activities
|
|
- **Comprehensive Error Handling**: Detailed error reporting and notification integration
|
|
|
|
#### Input Parameters
|
|
```json
|
|
{
|
|
"schedule_name": "hourly_predictions",
|
|
"model_name": "temperature_prediction_model",
|
|
"model_id": "temp_pred_001",
|
|
"query": "SELECT * FROM sensor_data WHERE timestamp > NOW() - INTERVAL '1 hour'",
|
|
"schema": {
|
|
"timestamp": "datetime",
|
|
"temperature": "float",
|
|
"humidity": "float"
|
|
},
|
|
"table_name": "predictions",
|
|
"input_filters": {
|
|
"EMPTY_DATA": {"POLICY": "STOP"}
|
|
},
|
|
"mlflow_transform_filters": {
|
|
"API_ERROR": {"POLICY": "STOP"}
|
|
},
|
|
"mlflow_predict_filters": {
|
|
"API_ERROR": {"POLICY": "STOP"}
|
|
},
|
|
"model_retention": 60,
|
|
"path_priority": ["STOP", "CONTINUE", "REPEAT"],
|
|
"opc_output_config": {
|
|
"server_id": "opc_server_1",
|
|
"tags": ["prediction_output"]
|
|
}
|
|
}
|
|
```
|
|
|
|
#### Architecture Diagram
|
|
```mermaid
|
|
flowchart LR
|
|
A[1. load_custom_query] --> B[2. prediction_process 🔃]
|
|
|
|
A -.-> Database[(Database)]
|
|
```
|
|
|
|
### 2. Prediction Process Workflow (`prediction_process.py`)
|
|
|
|
The **PredictionProcess** workflow implements the core prediction pipeline for ML model inference. It handles data quality validation, MLFlow model interactions, and prediction processing.
|
|
|
|
#### Purpose
|
|
- **Data Quality Validation**: Applies configurable filters for data integrity
|
|
- **MLFlow Integration**: Manages model transformation and prediction requests
|
|
- **Response Validation**: Filters MLFlow API responses for quality assurance
|
|
- **Prediction Export**: Delegates prediction formatting and export operations
|
|
|
|
#### Execution Flow
|
|
1. **Timestamp Retrieval**: Gets the last processed timestamp for incremental processing
|
|
2. **Input Data Gate**: Applies configured filters for data quality validation
|
|
3. **Path Decision**: Determines processing path based on filter results
|
|
4. **MLFlow Transform**: Requests data transformation using MLFlow models
|
|
5. **Response Validation**: Filters transform responses for quality assurance
|
|
6. **MLFlow Prediction**: Executes prediction using transformed data
|
|
7. **Content Validation**: Filters prediction responses for final quality check
|
|
8. **Export Delegation**: Delegates to FormatAndExportPrediction workflow
|
|
|
|
#### Key Features
|
|
- **Configurable Quality Gates**: Multiple filter types with policy-based configuration
|
|
- **Flexible Path Handling**: Configurable decision paths (STOP, CONTINUE, REPEAT)
|
|
- **MLFlow Integration**: Comprehensive model management and inference
|
|
- **Incremental Processing**: Timestamp-based data processing optimization
|
|
- **Comprehensive Monitoring**: Detailed metrics and error reporting
|
|
|
|
#### Input Parameters
|
|
```json
|
|
{
|
|
"metadata": {
|
|
"schedule_name": "hourly_predictions",
|
|
"model_name": "temperature_prediction_model",
|
|
"model_id": "temp_pred_001",
|
|
"workflow_name": "predictions_batch"
|
|
},
|
|
"data": {...},
|
|
"schema": {...},
|
|
"table_name": "predictions",
|
|
"model_id": "temp_pred_001",
|
|
"model_name": "temperature_prediction_model",
|
|
"input_filters": {
|
|
"EMPTY_DATA": {"POLICY": "STOP"},
|
|
"SPECIFIC_VARIABLES_NULL_VALUES": {
|
|
"POLICY": "STOP",
|
|
"config": {"variables": ["temperature", "humidity"]}
|
|
}
|
|
},
|
|
"mlflow_transform_filters": {
|
|
"API_ERROR": {"POLICY": "STOP"}
|
|
},
|
|
"mlflow_predict_filters": {
|
|
"API_ERROR": {"POLICY": "STOP"},
|
|
"NAN_VALUES": {"POLICY": "STOP"}
|
|
},
|
|
"model_retention": 60,
|
|
"path_priority": ["STOP", "CONTINUE", "REPEAT"],
|
|
"opc_output_config": {...}
|
|
}
|
|
```
|
|
|
|
#### Architecture Diagram
|
|
```mermaid
|
|
flowchart LR
|
|
A[1. get_last_timestamp] --> B[2. input_gate] --> C[3. request_transform] --> D[4. mlflow_response_gate] --> E[5. mlflow_content_gate] --> F[6. request_predict] --> G[7. mlflow_response_gate] --> H[8. format_and_export_prediction🔃]
|
|
|
|
A -.-> Redis[(Redis)]
|
|
C -.-> MLFlow[MLFlow]
|
|
F -.-> MLFlow[MLFlow]
|
|
G -.-> Filters[MLFlow Filters]
|
|
```
|
|
|
|
### 3. Format and Export Prediction Workflow (`format_and_export_prediction.py`)
|
|
|
|
The **FormatAndExportPrediction** workflow handles prediction data formatting and export operations to multiple destinations.
|
|
|
|
#### Purpose
|
|
- **Data Formatting**: Formats prediction data for different output destinations
|
|
- **PostgreSQL Export**: Persists predictions to database with metrics
|
|
- **OPC Integration**: Writes predictions to OPC servers for real-time access
|
|
- **Metrics Recording**: Tracks export operations and performance metrics
|
|
|
|
#### Execution Flow
|
|
1. **Path Decision**: Determines formatting path based on configuration
|
|
2. **Data Formatting**: Formats prediction data for specific output requirements
|
|
3. **Transformed Data Processing**: Optionally formats and exports transformed data separately
|
|
4. **PostgreSQL Export**: Writes formatted predictions to database
|
|
5. **OPC Export**: Writes predictions to OPC servers
|
|
6. **Metrics Recording**: Records export performance and success metrics
|
|
|
|
#### Key Features
|
|
- **Flexible Formatting**: Configurable output formats for different destinations
|
|
- **Multi-Destination Export**: PostgreSQL and OPC server integration
|
|
- **Transformed Data Export**: Optional separate export of MLFlow transformed data
|
|
- **Performance Monitoring**: Comprehensive metrics for export operations
|
|
- **Error Handling**: Robust error handling with notification integration
|
|
|
|
#### Architecture Diagram
|
|
```mermaid
|
|
flowchart LR
|
|
A[1. format_prediction/format_default_prediction] --> B[2. format_transformed_data] --> C[3. write_opc_data] --> D[4. export_data_to_postgres] --> E[5. write_metrics]
|
|
|
|
A -.-> Format[Data Formatting]
|
|
B -.-> Transform[Transformed Data]
|
|
C -.-> OPC[OPC Servers]
|
|
D -.-> PostgreSQL[(PostgreSQL)]
|
|
E -.-> Prometheus[Prometheus]
|
|
```
|
|
|
|
#### Transformed Data Export
|
|
When `transformed_data` is provided in the input, the workflow will:
|
|
- Format the transformed data using `format_transformed_data` activity
|
|
- Export it to a separate table (`transform_table_name`) asynchronously
|
|
- Wait for both prediction and transformed data exports to complete
|
|
- This enables separate tracking of model transformations for analysis and debugging
|
|
|
|
### 4. Minimal Retrain Workflow (`minimal_retrain.py`)
|
|
|
|
The **MinimalRetrain** workflow handles automated model retraining and production model updates.
|
|
|
|
#### Purpose
|
|
- **Model Retraining**: Automates ML model retraining processes
|
|
- **Production Updates**: Manages production model version updates
|
|
- **Data Export**: Exports training data for model development
|
|
- **Quality Assurance**: Ensures model quality before production deployment
|
|
|
|
#### Execution Flow
|
|
1. **Data Loading**: Loads training data using custom queries
|
|
2. **Model Retraining**: Executes model retraining process
|
|
3. **Quality Validation**: Validates retrained model performance
|
|
4. **Production Update**: Updates production model if quality criteria met
|
|
5. **Data Export**: Exports training data for analysis
|
|
|
|
#### Architecture Diagram
|
|
```mermaid
|
|
flowchart LR
|
|
A[1. load_custom_query] --> B[2. retrain_model] --> C[3. update_production_model] --> D[4. export_data_to_postgres]
|
|
|
|
A -.-> Database[(Database)]
|
|
B -.-> MLFlow[MLFlow]
|
|
C -.-> MLFlow[MLFlow]
|
|
D -.-> PostgreSQL[(PostgreSQL)]
|
|
```
|
|
|
|
## 📋 Prerequisites
|
|
|
|
- Python 3.11+
|
|
- Temporal server/cluster
|
|
- PostgreSQL database
|
|
- MLFlow server
|
|
- MinIO object storage (for MLFlow artifacts)
|
|
- MongoDB server (for notifications)
|
|
- OPC server(s) if using OPC export
|
|
|
|
**Note**: External dependencies must be available either through:
|
|
- Kubernetes cluster deployment
|
|
- Docker Compose setup
|
|
- Cloud-managed services
|
|
- Local installations
|
|
|
|
## 🚀 Installation
|
|
|
|
### Local Development Setup
|
|
|
|
1. **Clone the repository**
|
|
```bash
|
|
git clone <repository-url>
|
|
cd sientia-dataops-laborious
|
|
```
|
|
|
|
2. **Create virtual environment**
|
|
```bash
|
|
python3.11 -m venv venv
|
|
source ./venv/bin/activate
|
|
```
|
|
|
|
3. **Install dependencies**
|
|
|
|
1. **Install github cli**
|
|
```bash
|
|
sudo apt update
|
|
sudo apt install gh -y
|
|
```
|
|
|
|
2. **Authenticate with github**
|
|
```bash
|
|
gh auth login
|
|
```
|
|
|
|
3. **Run the install_dependencies.sh script**
|
|
```bash
|
|
chmod +x install_dependencies.sh
|
|
./install_dependencies.sh
|
|
```
|
|
|
|
4. **Create environment configuration file**
|
|
```bash
|
|
cp .env.example .env
|
|
# Edit .env with your connection details
|
|
```
|
|
|
|
5. **Configure external dependencies**
|
|
|
|
You'll need to set up port forwarding or connections to external services. For example:
|
|
|
|
```bash
|
|
# Port forwarding from Kubernetes cluster
|
|
kubectl port-forward svc/postgresql 5432:5432
|
|
kubectl port-forward svc/mlflow 5000:5000
|
|
kubectl port-forward svc/mongodb 27017:27017
|
|
|
|
# Or connect to external services
|
|
# Ensure services are accessible on localhost with appropriate ports
|
|
```
|
|
|
|
## 📦 How to Run
|
|
|
|
### Running the Laborious Application
|
|
|
|
Use the provided script to run the application locally:
|
|
|
|
```bash
|
|
# Make script executable (first time only)
|
|
chmod +x run_local.sh
|
|
|
|
# Run the application
|
|
./run_local.sh
|
|
```
|
|
|
|
The script will:
|
|
- Activate the virtual environment
|
|
- Load environment variables from `.env`
|
|
- Start the laborious worker application
|
|
|
|
### Running Tests and Coverage
|
|
|
|
Use the provided script to run tests with coverage:
|
|
|
|
```bash
|
|
# Make script executable (first time only)
|
|
chmod +x run_coverage.sh
|
|
|
|
# Run tests with coverage
|
|
./run_coverage.sh
|
|
```
|
|
|
|
The script will:
|
|
- Activate the virtual environment
|
|
- Run pytest with coverage reporting
|
|
- Generate HTML coverage report
|
|
- Open the coverage report in your browser
|
|
|
|
### Manual Test Execution
|
|
|
|
You can also run tests manually:
|
|
|
|
```bash
|
|
# Activate virtual environment
|
|
source ./venv/bin/activate
|
|
|
|
# Run all tests
|
|
pytest
|
|
|
|
# Run with coverage
|
|
pytest --cov=laborious --cov-report=html
|
|
|
|
# Run specific test categories
|
|
pytest tests/activities/
|
|
pytest tests/workflow/
|
|
```
|
|
|
|
### Manual Application Execution
|
|
|
|
For manual execution without scripts:
|
|
|
|
```bash
|
|
# Activate virtual environment
|
|
source ./venv/bin/activate
|
|
|
|
# Load environment variables (if using .env file)
|
|
if [ -f .env ]; then
|
|
export $(cat .env | grep -v '^#' | xargs)
|
|
fi
|
|
|
|
# Start the laborious worker
|
|
python -m laborious.worker.worker
|
|
```
|
|
|
|
## Code Quality & Validation
|
|
|
|
### Overview
|
|
|
|
Since Python is not compiled, we validate quality, security, and correctness before execution.
|
|
|
|
### Validation Tools
|
|
|
|
- Ruff: Linting and formatting
|
|
- mypy: Static type checking
|
|
- Bandit: Security analysis
|
|
- pytest: Unit/integration testing with coverage
|
|
|
|
### Tools Installation
|
|
|
|
```bash
|
|
pip install -r requirements-dev.txt
|
|
```
|
|
|
|
### Complete Validation
|
|
|
|
Option 1 (recommended):
|
|
```bash
|
|
./validate.sh
|
|
```
|
|
The script runs, in order:
|
|
1. Format check (Ruff)
|
|
2. Linting (Ruff)
|
|
3. Type checking (mypy)
|
|
4. Security analysis (Bandit)
|
|
5. Tests with coverage (pytest)
|
|
|
|
Option 2 (individual commands):
|
|
```bash
|
|
ruff format --check laborious/ tests/
|
|
ruff check laborious/ tests/
|
|
mypy laborious/
|
|
bandit -r laborious/ -ll
|
|
pytest tests/ --cov=laborious --cov-report=term-missing
|
|
```
|
|
|
|
### Automatic Fixes
|
|
|
|
```bash
|
|
ruff format laborious/ tests/
|
|
ruff check --fix laborious/ tests/
|
|
```
|
|
|
|
### Configuration
|
|
|
|
All settings reside in `pyproject.toml` (Ruff, mypy, pytest, Bandit).
|
|
|
|
### CI/CD Integration
|
|
|
|
The workflow at `.github/workflows/quality-gate.yml` executes validations on each push/PR.
|
|
|
|
### Best Practices
|
|
|
|
- Run `./validate.sh` before committing
|
|
- Use `ruff check --watch` for continuous feedback
|
|
- Add type hints and tests for new code
|
|
|
|
## 🧪 Testing
|
|
|
|
### Test Structure
|
|
```
|
|
tests/
|
|
├── activities/ # Activity implementation tests
|
|
│ ├── test_gates.py # Data quality gates and formatting tests
|
|
│ ├── test_mlflow.py # MLFlow operations and reference data tests
|
|
│ └── ... # Other activity tests
|
|
├── workflows/ # Workflow orchestration tests
|
|
│ └── subworkflows/ # Sub-workflow tests
|
|
│ └── test_format_and_export_prediction.py # Export workflow tests
|
|
├── utils/ # Utility function tests
|
|
└── integration/ # End-to-end workflow tests
|
|
```
|
|
|
|
### Test Coverage
|
|
The test suite provides comprehensive coverage for:
|
|
- **Data Quality Gates**: Input, response, and content validation filters
|
|
- **Data Formatting**: Prediction, transformed data, and retrain report formatting
|
|
- **MLFlow Operations**: Transform, predict, retrain, and reference data retrieval
|
|
- **Workflow Orchestration**: Complete workflow execution paths and error handling
|
|
- **Metrics Recording**: Performance monitoring and OPC export metrics
|
|
|
|
### Test Execution
|
|
```bash
|
|
# Install test dependencies
|
|
pip install pytest pytest-cov pytest-asyncio
|
|
|
|
# Run tests with coverage
|
|
pytest --cov=laborious --cov-report=html
|
|
|
|
# Run specific test modules
|
|
pytest tests/activities/test_gates.py
|
|
pytest tests/workflow/test_predictions_batch.py
|
|
```
|
|
|
|
## 📊 Monitoring and Metrics
|
|
|
|
The Laborious system exposes comprehensive Prometheus metrics for operational visibility and performance monitoring:
|
|
|
|
### Application Health Metrics
|
|
- `app_up`: Application health status (1=healthy, 0=unhealthy)
|
|
- Labels: `pod_id`
|
|
|
|
### Prediction Operation Metrics
|
|
- `laborious_predictions_written_count`: Counter for successful prediction exports
|
|
- Labels: `pod_id`, `model_name`, `workflow_name`
|
|
- `laborious_prediction_confidence_monitor`: Gauge for current prediction confidence levels
|
|
- Labels: `pod_id`, `model_name`, `workflow_name`
|
|
- `laborious_prediction_response_time_monitor`: Histogram for prediction response times
|
|
- Labels: `pod_id`, `model_name`, `workflow_name`
|
|
- Buckets: [0.01, 0.05, 0.1, 0.2, 0.5, 1.0, 2.0, 5.0, 10.0]
|
|
|
|
### OPC Export Metrics
|
|
- `laborious_prediction_opc_writing_count`: Counter for OPC server write operations
|
|
- Labels: `pod_id`, `model_name`, `workflow_name`, `opc_server_id`
|
|
- `laborious_prediction_opc_writing_response_time_monitor`: Histogram for OPC write response times
|
|
- Labels: `pod_id`, `model_name`, `workflow_name`, `opc_server_id`
|
|
- Buckets: [0.01, 0.05, 0.1, 0.2, 0.5, 1.0, 2.0, 5.0, 10.0]
|
|
|
|
### Data Quality Metrics
|
|
- Filter pass/fail rates through notification system
|
|
- MLFlow API response validation metrics
|
|
- Data quality gate performance tracking
|
|
|
|
## ⚙️ Configuration
|
|
|
|
### Environment Variables
|
|
|
|
| Variable | Description | Default | Required |
|
|
|----------|-------------|---------|----------|
|
|
| `TEMPORAL_HOST` | Temporal server address | `localhost:7233` | Yes |
|
|
| `TEMPORAL_NAMESPACE` | Temporal namespace | `laborious` | No |
|
|
| `POSTGRES_HOST` | PostgreSQL hostname | `localhost` | Yes |
|
|
| `POSTGRES_PORT` | PostgreSQL port | `5432` | Yes |
|
|
| `POSTGRES_USER` | PostgreSQL username | `sientia` | Yes |
|
|
| `POSTGRES_PASSWORD` | PostgreSQL password | `sientia` | Yes |
|
|
| `POSTGRES_DBNAME` | PostgreSQL database | `sientia` | Yes |
|
|
| `POSTGRES_MIN_CONNECTIONS` | Minimum PostgreSQL connections | `5` | No |
|
|
| `POSTGRES_MAX_CONNECTIONS` | Maximum PostgreSQL connections | `20` | No |
|
|
| `MLFLOW_HOST` | MLFlow server hostname | `http://localhost` | Yes |
|
|
| `MLFLOW_PORT` | MLFlow server port | `5080` | Yes |
|
|
| `MLFLOW_USERNAME` | MLFlow username | `aignosi` | Yes |
|
|
| `MLFLOW_PASSWORD` | MLFlow password | `aignosi` | Yes |
|
|
| `OPC_CONFIG` | OPC server configuration (JSON) | `{}` | No |
|
|
| `OPC_ID` | OPC server identifier | `1` | No |
|
|
| `OPC_URL` | OPC server URL | `opc.tcp://localhost:4840` | No |
|
|
| `OPC_SERVER_URI` | OPC server URI | `opc.tcp://localhost:4840` | No |
|
|
| `OPC_CERT_PATH` | OPC client certificate path | `None` | No |
|
|
| `OPC_PRIVATE_KEY_PATH` | OPC private key path | `None` | No |
|
|
| `OPC_SERVER_CERT_PATH` | OPC server certificate path | `None` | No |
|
|
| `OPC_RECONNECTION_INTERVAL` | OPC reconnection interval (ms) | `120` | No |
|
|
| `MONGODB_URL` | MongoDB connection URI | `localhost:27018` | Yes |
|
|
| `MONGODB_USERNAME` | MongoDB username | `root` | Yes |
|
|
| `MONGODB_PASSWORD` | MongoDB password | `wKZDbMNU1c` | Yes |
|
|
| `MONGODB_DATABASE_NAME` | MongoDB database name | `sientia` | Yes |
|
|
| `MONGODB_TTL_INDEX_HOURS` | MongoDB TTL index hours | `1` | No |
|
|
| `LOG_LEVEL` | Application log level | `INFO` | No |
|
|
| `PROJECT_NAME` | Project name for metrics | `laborious` | No |
|
|
| `HTTP_METRICS_PORT` | Prometheus metrics port | `9090` | No |
|
|
| `HTTP_SDK_METRICS_PORT` | Temporal SDK metrics port | `9091` | No |
|
|
| `POD_ID` | Kubernetes pod identifier | `None` | No |
|
|
|
|
|
|
|
|
|
|
### OPC Configuration
|
|
|
|
For multiple OPC servers, use the `OPC_CONFIG` environment variable:
|
|
|
|
```json
|
|
{
|
|
"opc_server_1": {
|
|
"url": "opc.tcp://server1:4840",
|
|
"name": "Server1",
|
|
"server_uri": "urn:server1:opcua",
|
|
"cert_path": "/path/to/cert.pem",
|
|
"private_key_path": "/path/to/key.pem",
|
|
"server_cert_path": "/path/to/server_cert.pem",
|
|
"reconnection_interval": 5000
|
|
},
|
|
"opc_server_2": {
|
|
"url": "opc.tcp://server2:4840",
|
|
"name": "Server2",
|
|
"server_uri": "urn:server2:opcua",
|
|
"cert_path": "/path/to/cert.pem",
|
|
"private_key_path": "/path/to/key.pem",
|
|
"server_cert_path": "/path/to/server_cert.pem",
|
|
"reconnection_interval": 5000
|
|
}
|
|
}
|
|
```
|
|
|
|
For single OPC server, use individual environment variables:
|
|
- `OPC_URL`
|
|
- `OPC_NAME`
|
|
- `OPC_SERVER_URI`
|
|
- `OPC_CERT_PATH`
|
|
- `OPC_PRIVATE_KEY_PATH`
|
|
- `OPC_SERVER_CERT_PATH`
|
|
- `OPC_RECONNECTION_INTERVAL`
|
|
|
|
### Workflow Configuration
|
|
|
|
MongoDB pipeline configuration:
|
|
|
|
#### Predictions Batch Workflow configuration sample
|
|
|
|
This is the configuration for the Predictions Batch Workflow, to be inserted into the MongoDB pipeline collection.
|
|
|
|
```json
|
|
{
|
|
"schedule_name": "laborious-orchestrated-pipeline",
|
|
"model_id": "1",
|
|
"workflow_type": "predictions_batch",
|
|
"frequency": "30s",
|
|
"max_retry_policy": 1,
|
|
"query": "select * from sientia_data.laborious_data where model_id = 1 and \"timestamp\" > NOW() - INTERVAL '5 minutes' order by \"timestamp\" desc limit 30;",
|
|
"write_tags": [
|
|
{
|
|
"server_id": "server1",
|
|
"type": "prediction",
|
|
"addr": "ns=2;i=5",
|
|
"data_type": "double"
|
|
},
|
|
{
|
|
"server_id": "server1",
|
|
"type": "confidence",
|
|
"addr": "ns=2;i=6",
|
|
"data_type": "double"
|
|
}
|
|
],
|
|
"input_filters": {
|
|
"EMPTY_DATA": {"POLICY": "STOP"},
|
|
"SPECIFIC_VARIABLES_NULL_VALUES": {
|
|
"POLICY": "CONTINUE",
|
|
"config": {"variables": ["Counter"]}
|
|
}
|
|
},
|
|
"mlflow_transform_filters": {
|
|
"API_ERROR": {"POLICY": "REPEAT"},
|
|
"NAN_VALUES": {"POLICY": "STOP"}
|
|
},
|
|
"mlflow_predict_filters": {
|
|
"API_ERROR": {"POLICY": "CONTINUE"}
|
|
},
|
|
"path_priority": ["STOP", "CONTINUE", "REPEAT"],
|
|
"active": true,
|
|
"updated_at": {
|
|
"$date": "2025-09-16T10:00:00.000Z"
|
|
},
|
|
"datetime_columns": ["timestamp", "created_at"],
|
|
"predictions_storage_policy": "lts:1"
|
|
}
|
|
```
|
|
|
|
This is the configuration created by the Orchestrator in Temporal.
|
|
|
|
```json
|
|
{
|
|
"datetime_columns":["timestamp","created_at"],
|
|
"frequency":"15m",
|
|
"input_filters":{"EMPTY_DATA":{"config":{},"policy":"STOP"}},
|
|
"max_retry_policy":1,
|
|
"mlflow_predict_filters":{"API_ERROR":{"config":{},"policy":"CONTINUE"}},
|
|
"mlflow_transform_filters":{
|
|
"API_ERROR":{"config":{},"policy":"CONTINUE"},
|
|
"EMPTY_DATA":{"config":{},"policy":"STOP"}
|
|
},
|
|
"model_config":{
|
|
"is_compressed":true,
|
|
"predict_flavor":"pyfunc",
|
|
"retention_minutes":60,
|
|
"retention_target":"artifact",
|
|
"transform_function_keyword":"transform"
|
|
},
|
|
"model_id":"352",
|
|
"model_name":"courier",
|
|
"opc_output_config":{},
|
|
"path_priority":["STOP","CONTINUE","REPEAT"],
|
|
"predictions_storage_policy":"lts:1",
|
|
"query":"select * from sientia_data.laborious_data where model_id = 352 order by \"timestamp\" desc limit 300;",
|
|
"retention_time":3600,
|
|
"schedule_name":"laborious-courier",
|
|
"schema":"sientia_data",
|
|
"table_name":"predictions",
|
|
"updated_at":"2025-09-12 19:35:01.600000+0000",
|
|
"workflow_type":"predictions_batch"
|
|
}
|
|
```
|
|
|
|
## 🔧 Development
|
|
|
|
### Project Structure
|
|
```
|
|
laborious/
|
|
├── activities/ # Temporal activity implementations
|
|
│ ├── activities.py # Main activities orchestrator
|
|
│ ├── gates.py # Data quality gates and filtering
|
|
│ ├── mlflow.py # MLFlow model operations
|
|
│ └── opc.py # OPC server operations
|
|
├── workflows/ # Temporal workflow definitions
|
|
│ ├── predictions_batch.py # Main batch prediction workflow
|
|
│ ├── minimal_retrain.py # Model retraining workflow
|
|
│ └── sub_workflows/ # Sub-workflow implementations
|
|
│ ├── prediction_process.py # Core prediction workflow
|
|
│ └── format_and_export_prediction.py # Export workflow
|
|
├── worker/ # Worker implementation
|
|
│ └── worker.py # Main worker orchestrator
|
|
├── utils/ # Utility functions
|
|
│ ├── connectors_config.py # Database configuration
|
|
│ ├── filters/ # Data quality filters
|
|
│ │ ├── conditional_filters.py # Conditional data filters
|
|
│ │ └── mlflow_filters.py # MLFlow response filters
|
|
│ └── repository/ # Data access layer
|
|
│ ├── model_repository.py # MLFlow model operations
|
|
│ └── opc_repository.py # OPC server operations
|
|
├── metrics.py # Prometheus metrics definitions
|
|
└── __init__.py
|
|
```
|
|
|
|
### Adding New Features
|
|
|
|
1. **Follow Temporal patterns** for new workflows and activities
|
|
2. **Add comprehensive docstrings** for all public methods
|
|
3. **Include Prometheus metrics** for monitoring
|
|
4. **Add unit tests** for new functionality
|
|
5. **Update this README** with new features and configuration
|
|
|
|
## 🐛 Troubleshooting
|
|
|
|
### Common Issues
|
|
|
|
1. **Temporal Connection Failures**
|
|
- Verify Temporal server is running and accessible
|
|
- Check namespace configuration and permissions
|
|
- Review server logs for connection issues
|
|
|
|
2. **MLFlow Connection Issues**
|
|
- Verify MLFlow server is running and accessible
|
|
- Check authentication credentials and permissions
|
|
- Ensure model names and versions exist
|
|
|
|
3. **Database Connection Issues**
|
|
- Verify PostgreSQL service is running
|
|
- Check connection credentials and network access
|
|
- Ensure proper connection pool configuration
|
|
|
|
4. **OPC Connection Failures**
|
|
- Verify OPC server is accessible
|
|
- Check certificate and key file paths
|
|
- Review OPC server logs for connection issues
|
|
|
|
5. **Workflow Execution Failures**
|
|
- Review activity error logs and notifications
|
|
- Check data quality filter configurations
|
|
- Verify input data format and required fields
|
|
|
|
### Debug Mode
|
|
|
|
Enable debug logging by setting the log level:
|
|
```bash
|
|
export LOG_LEVEL=DEBUG
|
|
```
|
|
|
|
## ⚡ Performance Tuning
|
|
|
|
### Key Parameters
|
|
|
|
- **Worker Concurrency**: Adjust `max_concurrent_workflow_tasks` and `max_concurrent_activities`
|
|
- **Connection Pools**: Optimize database connection pool sizes
|
|
- **Model Retention**: Configure MLFlow model retention based on requirements
|
|
- **Batch Sizes**: Adjust data processing batch sizes for optimal throughput
|
|
|
|
### Scaling Considerations
|
|
|
|
- **Horizontal Scaling**: Deploy multiple worker instances
|
|
- **Task Queue Distribution**: Use multiple task queues for different workflow types
|
|
- **Database Performance**: Optimize indexes and connection pooling
|
|
- **MLFlow Performance**: Configure appropriate model serving resources
|
|
|
|
## 🤝 Contributing
|
|
|
|
1. Fork the repository
|
|
2. Create a feature branch
|
|
3. Make your changes with comprehensive testing
|
|
4. Update documentation and docstrings
|
|
5. Submit a pull request
|
|
|
|
### Code Quality Standards
|
|
|
|
- Follow PEP 8 style guidelines
|
|
- Include comprehensive docstrings for all public methods
|
|
- Maintain test coverage above 80%
|
|
- Use type hints where appropriate
|
|
- Follow Temporal.io best practices
|
|
|
|
## 📄 License
|
|
|
|
This project is licensed under the terms specified in the LICENSE file.
|
|
|
|
## 🆘 Support
|
|
|
|
For support and questions:
|
|
- Check the troubleshooting section above
|
|
- Review the metrics and logs for error patterns
|
|
- Open an issue in the project repository
|
|
- Contact the development team
|
|
|
|
---
|
|
|
|
**Note**: The Laborious system is designed for production use in industrial ML environments. Ensure proper security configuration and network isolation for production deployments. |