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}