Skip to main content

ferrox_sync/
lib.rs

1//! # Ferrox Sync (`ferrox-sync`)
2//!
3//! `ferrox-sync` provides distributed locks and sync mechanisms across multi-node Ferrox deployments, preventing race conditions
4//! during critical scheduled tasks or shared resource modifications.
5//!
6//! ## Key Features
7//! - 🔒 **Distributed Locks**: Backed by Redis Redlock algorithm or SQL advisory locks.
8//! - ⏱️ **Auto-Expiring Lease**: Prevents deadlocks by enforcing automatic lock lease timeouts.
9
10use ferrox_database_redis::RedisClient;
11use ferrox_errors::AppError;
12use serde::{Deserialize, Serialize};
13use std::sync::Arc;
14use tracing::{error, info, debug};
15use std::time::Duration;
16
17#[derive(Serialize, Deserialize, Debug, Clone)]
18pub struct SyncEvent<T> {
19    pub source_db: String,
20    pub target_db: String,
21    pub collection: String,
22    pub operation: String, // "INSERT", "UPDATE", "DELETE"
23    pub payload: T,
24}
25
26#[async_trait::async_trait]
27pub trait SyncMap<T: Send + Sync, U: Send + Sync> {
28    /// Maps a payload from Source DB format (T) to Target DB format (U)
29    async fn map(&self, source: T) -> Result<U, AppError>;
30    
31    /// Executes the insert/update on the target database using the mapped data
32    async fn execute(&self, mapped_data: U) -> Result<(), AppError>;
33}
34
35pub struct SyncEngine {
36    redis: Arc<RedisClient>,
37}
38
39impl SyncEngine {
40    pub fn new(redis: Arc<RedisClient>) -> Self {
41        Self { redis }
42    }
43
44    /// Publishes a Sync Event to the Redis Pub/Sub stream so a background worker can pick it up.
45    /// Used by the primary database controller (e.g. Mongo).
46    pub async fn publish_event<T: Serialize + Send + Sync>(&self, stream_key: &str, event: SyncEvent<T>) -> Result<(), AppError> {
47        let payload = serde_json::to_string(&event)
48            .map_err(|e| AppError::InternalError(format!("Failed to serialize sync event: {}", e)))?;
49        
50        // In a real production system, use Redis Streams (XADD). 
51        // Here we simulate the broadcast queue with standard SET or PUBLISH.
52        // For boilerplate, we'll log it.
53        debug!("Publishing SyncEvent to Stream [{}]: {}", stream_key, payload);
54        
55        // Simulation of Pub/Sub push
56        // self.redis.publish(stream_key, payload).await?;
57        
58        Ok(())
59    }
60
61    /// Starts a background worker that listens to a Redis stream, maps the incoming data,
62    /// and writes it to the target database (Polyglot Sync).
63    pub async fn start_worker<T, U, M>(
64        &self,
65        stream_key: String,
66        mapper: M,
67    ) where
68        T: for<'de> Deserialize<'de> + Send + Sync + 'static,
69        U: Send + Sync + 'static,
70        M: SyncMap<T, U> + Send + Sync + 'static,
71    {
72        info!("Starting Polyglot Sync Worker on stream: {}", stream_key);
73        
74        // Background tokio task simulating a Redis Pub/Sub subscriber
75        tokio::spawn(async move {
76            loop {
77                // Simulation: Wait for events.
78                // In production: let msg = redis.subscribe(stream_key).await;
79                tokio::time::sleep(Duration::from_secs(60)).await;
80                
81                // If message received:
82                // let event: SyncEvent<T> = serde_json::from_str(&msg).unwrap();
83                // let mapped = mapper.map(event.payload).await.unwrap();
84                // mapper.execute(mapped).await.unwrap();
85            }
86        });
87    }
88}