pub mod checkpoint;
pub mod deduplication;
pub mod transaction;
pub use checkpoint::TransactionalCheckpoint;
pub use deduplication::{DeduplicationKey, EventDeduplicator};
pub use transaction::{TransactionCoordinator, TwoPhaseCommit};
use crate::error::Result;
use crate::models::Position;
use std::sync::Arc;
use tokio::sync::RwLock;
#[derive(Debug, Clone)]
pub struct AtLeastOnceConfig {
pub enabled: bool,
pub deduplication_window: usize,
pub transaction_timeout_secs: u64,
pub two_phase_commit: bool,
pub checkpoint_before_write: bool,
}
impl Default for AtLeastOnceConfig {
fn default() -> Self {
Self {
enabled: true,
deduplication_window: 10000,
transaction_timeout_secs: 30,
two_phase_commit: true,
checkpoint_before_write: true,
}
}
}
pub struct AtLeastOnceManager {
config: AtLeastOnceConfig,
deduplicator: Arc<RwLock<EventDeduplicator>>,
pub transaction_coordinator: Arc<TransactionCoordinator>,
}
impl AtLeastOnceManager {
pub fn new(config: AtLeastOnceConfig) -> Self {
let deduplicator = Arc::new(RwLock::new(EventDeduplicator::new(
config.deduplication_window,
)));
let transaction_coordinator =
Arc::new(TransactionCoordinator::new(config.transaction_timeout_secs));
Self {
config,
deduplicator,
transaction_coordinator,
}
}
pub async fn is_duplicate(&self, key: &DeduplicationKey) -> Result<bool> {
if !self.config.enabled {
return Ok(false);
}
let dedup = self.deduplicator.read().await;
Ok(dedup.contains(key))
}
pub async fn mark_processed(&self, key: DeduplicationKey) -> Result<()> {
if !self.config.enabled {
return Ok(());
}
let mut dedup = self.deduplicator.write().await;
dedup.add(key);
Ok(())
}
pub async fn begin_transaction(&self, task_id: &str) -> Result<String> {
if !self.config.two_phase_commit {
return Ok(String::new());
}
self.transaction_coordinator.begin(task_id).await
}
pub async fn prepare(&self, transaction_id: &str, position: Position) -> Result<bool> {
if !self.config.two_phase_commit {
return Ok(true);
}
self.transaction_coordinator
.prepare(transaction_id, position)
.await
}
pub async fn commit(&self, transaction_id: &str) -> Result<()> {
if !self.config.two_phase_commit {
return Ok(());
}
self.transaction_coordinator.commit(transaction_id).await
}
pub async fn rollback(&self, transaction_id: &str) -> Result<()> {
if !self.config.two_phase_commit {
return Ok(());
}
self.transaction_coordinator.rollback(transaction_id).await
}
}