use ferrox_database_redis::RedisClient;
use ferrox_errors::AppError;
use serde::{Deserialize, Serialize};
use std::sync::Arc;
use tracing::{error, info, debug};
use std::time::Duration;
#[derive(Serialize, Deserialize, Debug, Clone)]
pub struct SyncEvent<T> {
pub source_db: String,
pub target_db: String,
pub collection: String,
pub operation: String, pub payload: T,
}
#[async_trait::async_trait]
pub trait SyncMap<T: Send + Sync, U: Send + Sync> {
async fn map(&self, source: T) -> Result<U, AppError>;
async fn execute(&self, mapped_data: U) -> Result<(), AppError>;
}
pub struct SyncEngine {
redis: Arc<RedisClient>,
}
impl SyncEngine {
pub fn new(redis: Arc<RedisClient>) -> Self {
Self { redis }
}
pub async fn publish_event<T: Serialize + Send + Sync>(&self, stream_key: &str, event: SyncEvent<T>) -> Result<(), AppError> {
let payload = serde_json::to_string(&event)
.map_err(|e| AppError::InternalError(format!("Failed to serialize sync event: {}", e)))?;
debug!("Publishing SyncEvent to Stream [{}]: {}", stream_key, payload);
Ok(())
}
pub async fn start_worker<T, U, M>(
&self,
stream_key: String,
mapper: M,
) where
T: for<'de> Deserialize<'de> + Send + Sync + 'static,
U: Send + Sync + 'static,
M: SyncMap<T, U> + Send + Sync + 'static,
{
info!("Starting Polyglot Sync Worker on stream: {}", stream_key);
tokio::spawn(async move {
loop {
tokio::time::sleep(Duration::from_secs(60)).await;
}
});
}
}