Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions backend/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,8 @@ soroban-sdk = "21.0"
# OpenAPI documentation
utoipa = { version = "4.0", features = ["axum_extras"] }
utoipa-swagger-ui = { version = "6.0", features = ["axum"] }
# Streaming and event processing
reqwest = { version = "0.11", features = ["json"] }

[dev-dependencies]
tower-test = "0.4"
Expand Down
50 changes: 50 additions & 0 deletions backend/migrations/20260626000000_add_saga_management.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
-- Add saga management tables for distributed transaction orchestration

-- Create enum for saga states
CREATE TYPE saga_state AS ENUM ('Pending', 'InProgress', 'Compensating', 'Completed', 'Failed', 'Aborted');

-- Saga instances table
CREATE TABLE saga_instances (
id UUID PRIMARY KEY DEFAULT uuid_generate_v4(),
saga_type TEXT NOT NULL,
state saga_state NOT NULL DEFAULT 'Pending',
current_step TEXT,
completed_steps TEXT[] DEFAULT '{}',
failed_step TEXT,
context JSONB NOT NULL DEFAULT '{}',
error_message TEXT,
started_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT NOW(),
updated_at TIMESTAMP WITH TIME ZONE NOT NULL DEFAULT NOW(),
completed_at TIMESTAMP WITH TIME ZONE,
metadata JSONB DEFAULT '{}'
);

-- Create indexes for saga queries
CREATE INDEX idx_saga_instances_state ON saga_instances(state);
CREATE INDEX idx_saga_instances_type ON saga_instances(saga_type);
CREATE INDEX idx_saga_instances_started_at ON saga_instances(started_at);
CREATE INDEX idx_saga_instances_updated_at ON saga_instances(updated_at);

-- Create trigger for updated_at
CREATE OR REPLACE FUNCTION update_saga_updated_at()
RETURNS TRIGGER AS $$
BEGIN
NEW.updated_at = NOW();
RETURN NEW;
END;
$$ LANGUAGE plpgsql;

CREATE TRIGGER trigger_update_saga_updated_at
BEFORE UPDATE ON saga_instances
FOR EACH ROW
EXECUTE FUNCTION update_saga_updated_at();

-- Add comments
COMMENT ON TABLE saga_instances IS 'Stores saga instances for distributed transaction orchestration';
COMMENT ON COLUMN saga_instances.saga_type IS 'Type of saga (e.g., product_registration)';
COMMENT ON COLUMN saga_instances.state IS 'Current state of the saga';
COMMENT ON COLUMN saga_instances.current_step IS 'ID of the currently executing step';
COMMENT ON COLUMN saga_instances.completed_steps IS 'List of completed step IDs';
COMMENT ON COLUMN saga_instances.failed_step IS 'ID of the step that failed (if any)';
COMMENT ON COLUMN saga_instances.context IS 'Execution context and data';
COMMENT ON COLUMN saga_instances.error_message IS 'Error message if saga failed';
67 changes: 66 additions & 1 deletion backend/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,10 @@ use services::{
AnalyticsService, ApiKeyService, AuditService, BatchService, CarbonService,
CollaborationService, EventService, FinancialService, IoTService, PhysicsModelService, PredictiveRoutingService,
ProductService, QualityService, RecallService, RegulatoryService, SupplierService,
SyncService, UserService,
SyncService, UserService, MercuryIndexer, MercuryConfig, RuleEngine, SagaManager,
RedisWorkerPool, WorkerConfig, TrackingEventProcessor, get_default_rules,
get_product_registration_saga, NoopAction, EventProcessingHandler, RuleEvaluationHandler,
NotificationHandler, AlertHandler, WebhookHandler,
};
use utils::CronService;

Expand Down Expand Up @@ -63,6 +66,10 @@ pub struct AppState {
pub redis_client: redis::Client,
pub config: Config,
pub monitoring_system: MonitoringSystem,
pub mercury_indexer: Arc<MercuryIndexer>,
pub rule_engine: Arc<RuleEngine>,
pub saga_manager: Arc<SagaManager>,
pub worker_pool: Arc<RedisWorkerPool>,
}

impl AppState {
Expand Down Expand Up @@ -106,6 +113,36 @@ impl AppState {
Arc::new(PredictiveRoutingService::new(db.pool().clone()));
let physics_model_service = Arc::new(PhysicsModelService::new(db.pool().clone()));

// Initialize Mercury streaming indexer
let mercury_config = MercuryConfig::default();
let (mercury_indexer, _event_rx) =
MercuryIndexer::new(mercury_config, db.pool().clone(), redis_client.clone());
mercury_indexer.add_processor(Arc::new(TrackingEventProcessor::new(db.pool().clone())));
let mercury_indexer = Arc::new(mercury_indexer);

// Initialize rule engine with default rules
let mut rule_engine = RuleEngine::new();
for rule in get_default_rules() {
rule_engine.add_rule(rule);
}
rule_engine.register_handler("alert".to_string(), Arc::new(AlertHandler::new()));
rule_engine.register_handler("webhook".to_string(), Arc::new(WebhookHandler::new()));
let rule_engine = Arc::new(rule_engine);

// Initialize saga manager
let mut saga_manager = SagaManager::new(db.pool().clone(), redis_client.clone());
saga_manager.register_saga(get_product_registration_saga());
saga_manager.register_action("noop".to_string(), Arc::new(NoopAction));
let saga_manager = Arc::new(saga_manager);

// Initialize Redis worker pool
let worker_config = WorkerConfig::default();
let mut worker_pool = RedisWorkerPool::new(worker_config, redis_client.clone());
worker_pool.register_handler(Arc::new(EventProcessingHandler::new()));
worker_pool.register_handler(Arc::new(RuleEvaluationHandler::new()));
worker_pool.register_handler(Arc::new(NotificationHandler::new()));
let worker_pool = Arc::new(worker_pool);

// Initialize comprehensive monitoring system
let monitoring_system = MonitoringSystem::new();

Expand All @@ -132,6 +169,10 @@ impl AppState {
redis_client,
config,
monitoring_system,
mercury_indexer,
rule_engine,
saga_manager,
worker_pool,
})
}
}
Expand Down Expand Up @@ -159,6 +200,30 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
CronService::new(app_state.db.pool().clone(), app_state.redis_client.clone());
cron_service.start_scheduler().await;

// Start Mercury streaming indexer
let mercury_indexer = app_state.mercury_indexer.clone();
tokio::spawn(async move {
if let Err(e) = mercury_indexer.start().await {
tracing::error!("Mercury indexer failed: {}", e);
}
});

// Recover any in-progress sagas
let saga_manager = app_state.saga_manager.clone();
tokio::spawn(async move {
if let Err(e) = saga_manager.recover_sagas().await {
tracing::error!("Saga recovery failed: {}", e);
}
});

// Start Redis worker pool
let worker_pool = app_state.worker_pool.clone();
tokio::spawn(async move {
if let Err(e) = worker_pool.start().await {
tracing::error!("Worker pool failed: {}", e);
}
});

// Build router with security middleware
let app = Router::new()
.merge(crate::routes::health_routes())
Expand Down
11 changes: 11 additions & 0 deletions backend/src/services.rs
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,17 @@ pub use regulatory_service::RegulatoryService;
pub mod predictive_routing_service;
pub use predictive_routing_service::PredictiveRoutingService;

pub mod mercury_indexer;
pub use mercury_indexer::{MercuryIndexer, MercuryConfig, EventProcessor, TrackingEventProcessor};

pub mod rule_engine;
pub use rule_engine::{RuleEngine, Rule, RuleContext, ActionHandler, AlertHandler, WebhookHandler, get_default_rules};

pub mod saga_manager;
pub use saga_manager::{SagaManager, SagaInstance, SagaState, SagaAction, get_product_registration_saga, NoopAction};

pub mod redis_workers;
pub use redis_workers::{RedisWorkerPool, WorkerTask, WorkerConfig, TaskHandler, EventProcessingHandler, RuleEvaluationHandler, NotificationHandler};
pub mod physics_model_service;
pub use physics_model_service::PhysicsModelService;

Expand Down
Loading
Loading