use super::BoxFuture;
use parking_lot::Mutex;
use std::collections::{HashMap, HashSet, VecDeque};
use std::sync::Arc;
use std::time::{Duration, Instant};
use super::{DedupStorage, ReconciliationStorage};
use crate::reliability::ReconciliationRequest;
#[derive(Debug, Default)]
pub struct InMemoryDedupStorage {
seen: Mutex<HashSet<String>>,
}
impl DedupStorage for InMemoryDedupStorage {
fn is_durable(&self) -> bool {
false
}
fn first_seen<'a>(
&'a self,
idempotency_key: &'a str,
) -> BoxFuture<'a, crate::core::Result<bool>> {
Box::pin(async move {
let mut seen = self.seen.lock();
if seen.contains(idempotency_key) {
return Ok(false);
}
seen.insert(idempotency_key.to_owned());
Ok(true)
})
}
}
#[derive(Debug)]
pub struct InMemoryReconciliationStorage {
seen: Mutex<HashSet<String>>,
queued: Mutex<Vec<ReconciliationRequest>>,
max_queued: usize,
}
impl Default for InMemoryReconciliationStorage {
fn default() -> Self {
Self {
seen: Mutex::new(HashSet::new()),
queued: Mutex::new(Vec::new()),
max_queued: usize::MAX,
}
}
}
impl InMemoryReconciliationStorage {
pub fn with_capacity(max_queued: usize) -> Self {
Self {
seen: Mutex::new(HashSet::new()),
queued: Mutex::new(Vec::new()),
max_queued,
}
}
}
impl ReconciliationStorage for InMemoryReconciliationStorage {
fn is_durable(&self) -> bool {
false
}
fn enqueue<'a>(
&'a self,
request: ReconciliationRequest,
) -> super::BoxFuture<'a, crate::core::Result<bool>> {
Box::pin(async move {
let mut seen = self.seen.lock();
if seen.contains(&request.idempotency_key) {
return Ok(false);
}
let mut queued = self.queued.lock();
if queued.len() >= self.max_queued {
return Err(crate::core::AsxError::new(
crate::core::ErrorCode::PolicyViolation,
format!(
"reconciliation queue at capacity ({} entries); resolve pending requests before enqueuing more",
self.max_queued
),
crate::core::ErrorContext::new("reconciliation_storage_enqueue"),
));
}
seen.insert(request.idempotency_key.clone());
queued.push(request);
Ok(true)
})
}
fn queued_requests(
&self,
) -> super::BoxFuture<'_, crate::core::Result<Vec<ReconciliationRequest>>> {
Box::pin(async move {
let queued = self.queued.lock();
Ok(queued.clone())
})
}
fn resolve<'a>(
&'a self,
idempotency_key: &'a str,
) -> super::BoxFuture<'a, crate::core::Result<bool>> {
Box::pin(async move {
let mut seen = self.seen.lock();
let mut queued = self.queued.lock();
let before = queued.len();
queued.retain(|r| r.idempotency_key != idempotency_key);
let removed = queued.len() < before;
if removed {
seen.remove(idempotency_key);
}
Ok(removed)
})
}
}
type DedupState = (VecDeque<Arc<str>>, HashSet<Arc<str>>);
#[derive(Debug)]
pub struct BoundedFifoDedupStorage {
state: Mutex<DedupState>,
capacity: usize,
}
impl BoundedFifoDedupStorage {
pub fn new(capacity: usize) -> Self {
assert!(capacity > 0, "BoundedFifoDedupStorage capacity must be > 0");
Self {
state: Mutex::new((
VecDeque::with_capacity(capacity),
HashSet::with_capacity(capacity),
)),
capacity,
}
}
pub fn capacity(&self) -> usize {
self.capacity
}
}
impl DedupStorage for BoundedFifoDedupStorage {
fn is_durable(&self) -> bool {
false
}
fn first_seen<'a>(
&'a self,
idempotency_key: &'a str,
) -> BoxFuture<'a, crate::core::Result<bool>> {
Box::pin(async move {
let mut state = self.state.lock();
let (queue, set) = &mut *state;
if set.contains(idempotency_key) {
return Ok(false);
}
if queue.len() >= self.capacity
&& let Some(oldest) = queue.pop_front()
{
set.remove(&*oldest);
}
let owned: Arc<str> = Arc::from(idempotency_key);
queue.push_back(Arc::clone(&owned));
set.insert(owned);
Ok(true)
})
}
}
#[derive(Debug)]
pub struct TtlDedupStorage {
inner: Mutex<TtlDedupInner>,
ttl: Duration,
sweep_interval: Duration,
}
#[derive(Debug)]
struct TtlDedupInner {
entries: HashMap<String, Instant>,
last_sweep: Instant,
}
impl TtlDedupStorage {
pub fn new(ttl: Duration) -> Self {
let sweep_interval = ttl / 10;
let sweep_interval = if sweep_interval > Duration::from_secs(60) {
Duration::from_secs(60)
} else {
sweep_interval
};
Self::with_sweep_interval(ttl, sweep_interval)
}
pub fn with_sweep_interval(ttl: Duration, sweep_interval: Duration) -> Self {
Self {
inner: Mutex::new(TtlDedupInner {
entries: HashMap::new(),
last_sweep: Instant::now(),
}),
ttl,
sweep_interval,
}
}
}
impl Default for TtlDedupStorage {
fn default() -> Self {
Self::new(Duration::from_secs(86_400))
}
}
impl DedupStorage for TtlDedupStorage {
fn is_durable(&self) -> bool {
false
}
fn first_seen<'a>(
&'a self,
idempotency_key: &'a str,
) -> BoxFuture<'a, crate::core::Result<bool>> {
Box::pin(async move {
let mut inner = self.inner.lock();
let now = Instant::now();
if inner
.entries
.get(idempotency_key)
.is_some_and(|exp| now < *exp)
{
return Ok(false);
}
if now.duration_since(inner.last_sweep) >= self.sweep_interval {
inner.entries.retain(|_, exp| now < *exp);
inner.last_sweep = now;
}
inner
.entries
.insert(idempotency_key.to_string(), now + self.ttl);
Ok(true)
})
}
}
#[cfg(feature = "testing")]
#[derive(Debug, Default)]
pub struct DurableInMemoryDedupBackend(TtlDedupStorage);
#[cfg(feature = "testing")]
impl DurableInMemoryDedupBackend {
pub fn new(ttl: std::time::Duration) -> Self {
Self(TtlDedupStorage::new(ttl))
}
}
#[cfg(feature = "testing")]
impl DedupStorage for DurableInMemoryDedupBackend {
fn is_durable(&self) -> bool {
true
}
fn first_seen<'a>(
&'a self,
idempotency_key: &'a str,
) -> BoxFuture<'a, crate::core::Result<bool>> {
self.0.first_seen(idempotency_key)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::reliability::DeliveryOutcome;
use crate::storage::drive_dedup_future;
use std::sync::Arc;
use std::thread;
#[test]
fn in_memory_dedup_accepts_first_and_rejects_duplicate() {
let storage = InMemoryDedupStorage::default();
let key = "dedup:test:msg-1";
assert!(drive_dedup_future(storage.first_seen(key)).unwrap());
assert!(!drive_dedup_future(storage.first_seen(key)).unwrap());
}
#[test]
fn in_memory_dedup_is_correct_under_parallel_load() {
let storage = Arc::new(InMemoryDedupStorage::default());
let key = "dedup:test:parallel-msg";
let mut handles = Vec::new();
for _ in 0..16 {
let storage = Arc::clone(&storage);
handles.push(thread::spawn(move || {
drive_dedup_future(storage.first_seen(key)).unwrap()
}));
}
let accepted = handles
.into_iter()
.map(|h| h.join().expect("thread join"))
.filter(|accepted| *accepted)
.count();
assert_eq!(accepted, 1);
}
#[test]
fn in_memory_reconciliation_deduplicates_by_key() {
let storage = InMemoryReconciliationStorage::default();
let first = ReconciliationRequest::for_outcome("m1", "p1", DeliveryOutcome::Indeterminate)
.expect("request");
let duplicate = first.clone();
assert!(drive_dedup_future(storage.enqueue(first)).unwrap());
assert!(!drive_dedup_future(storage.enqueue(duplicate)).unwrap());
assert_eq!(
drive_dedup_future(storage.queued_requests()).unwrap().len(),
1
);
}
#[test]
fn in_memory_reconciliation_resolve_removes_from_queue() {
let storage = InMemoryReconciliationStorage::default();
let req =
ReconciliationRequest::for_outcome("msg-resolve", "p1", DeliveryOutcome::Indeterminate)
.expect("request");
let key = req.idempotency_key.clone();
assert!(drive_dedup_future(storage.enqueue(req)).unwrap());
assert_eq!(
drive_dedup_future(storage.queued_requests()).unwrap().len(),
1
);
assert!(
drive_dedup_future(storage.resolve(&key)).unwrap(),
"resolve should return true when found"
);
assert_eq!(
drive_dedup_future(storage.queued_requests()).unwrap().len(),
0
);
assert!(!drive_dedup_future(storage.resolve(&key)).unwrap());
}
#[test]
fn ttl_dedup_rejects_duplicates_just_before_expiry() {
let storage =
TtlDedupStorage::with_sweep_interval(Duration::from_secs(2), Duration::from_secs(60));
let key = "dedup:ttl:edge-before-expiry";
assert!(drive_dedup_future(storage.first_seen(key)).unwrap());
thread::sleep(Duration::from_millis(1900));
assert!(
!drive_dedup_future(storage.first_seen(key)).unwrap(),
"duplicate must be rejected inside TTL window"
);
}
#[test]
fn ttl_dedup_accepts_replay_after_expiry_with_buffer() {
let storage =
TtlDedupStorage::with_sweep_interval(Duration::from_secs(2), Duration::from_secs(60));
let key = "dedup:ttl:edge-after-expiry";
assert!(drive_dedup_future(storage.first_seen(key)).unwrap());
thread::sleep(Duration::from_millis(2300));
assert!(
drive_dedup_future(storage.first_seen(key)).unwrap(),
"key must be accepted after TTL expiry"
);
}
#[test]
fn in_memory_reconciliation_bounded_capacity_rejects_overflow() {
let storage = InMemoryReconciliationStorage::with_capacity(2);
let r1 =
ReconciliationRequest::for_outcome("msg-cap-1", "p1", DeliveryOutcome::Indeterminate)
.unwrap();
let r2 =
ReconciliationRequest::for_outcome("msg-cap-2", "p1", DeliveryOutcome::Indeterminate)
.unwrap();
let r3 =
ReconciliationRequest::for_outcome("msg-cap-3", "p1", DeliveryOutcome::Indeterminate)
.unwrap();
assert!(drive_dedup_future(storage.enqueue(r1)).unwrap());
assert!(drive_dedup_future(storage.enqueue(r2)).unwrap());
assert!(
drive_dedup_future(storage.enqueue(r3)).is_err(),
"should fail at capacity"
);
}
#[test]
fn in_memory_reconciliation_resolve_makes_room_after_capacity_hit() {
let storage = InMemoryReconciliationStorage::with_capacity(1);
let r1 =
ReconciliationRequest::for_outcome("msg-room-1", "p1", DeliveryOutcome::Indeterminate)
.unwrap();
let r2 =
ReconciliationRequest::for_outcome("msg-room-2", "p1", DeliveryOutcome::Indeterminate)
.unwrap();
let key1 = r1.idempotency_key.clone();
assert!(drive_dedup_future(storage.enqueue(r1)).unwrap());
assert!(
drive_dedup_future(storage.enqueue(r2)).is_err(),
"should fail at capacity"
);
assert!(drive_dedup_future(storage.resolve(&key1)).unwrap());
let r3 =
ReconciliationRequest::for_outcome("msg-room-3", "p1", DeliveryOutcome::Indeterminate)
.unwrap();
assert!(
drive_dedup_future(storage.enqueue(r3)).unwrap(),
"new request should be accepted after resolve"
);
}
#[test]
fn in_memory_reconciliation_resolve_allows_reenqueue_of_same_key() {
let storage = InMemoryReconciliationStorage::default();
let r1 =
ReconciliationRequest::for_outcome("msg-retry", "p1", DeliveryOutcome::Indeterminate)
.unwrap();
let key = r1.idempotency_key.clone();
assert!(
drive_dedup_future(storage.enqueue(r1)).unwrap(),
"initial enqueue"
);
assert!(
!drive_dedup_future(
storage.enqueue(
ReconciliationRequest::for_outcome(
"msg-retry",
"p1",
DeliveryOutcome::Indeterminate
)
.unwrap()
)
)
.unwrap(),
"duplicate enqueue before resolve must be rejected"
);
assert!(
drive_dedup_future(storage.resolve(&key)).unwrap(),
"resolve returns true"
);
let r2 =
ReconciliationRequest::for_outcome("msg-retry", "p1", DeliveryOutcome::Indeterminate)
.unwrap();
assert!(
drive_dedup_future(storage.enqueue(r2)).unwrap(),
"re-enqueue after resolve must succeed (retry path)"
);
}
#[test]
fn bounded_fifo_accepts_first_and_rejects_duplicate() {
let storage = BoundedFifoDedupStorage::new(10);
assert!(drive_dedup_future(storage.first_seen("msg-a")).unwrap());
assert!(!drive_dedup_future(storage.first_seen("msg-a")).unwrap());
assert!(drive_dedup_future(storage.first_seen("msg-b")).unwrap());
}
#[test]
fn bounded_fifo_evicts_oldest_at_capacity() {
let storage = BoundedFifoDedupStorage::new(3);
assert!(drive_dedup_future(storage.first_seen("k1")).unwrap());
assert!(drive_dedup_future(storage.first_seen("k2")).unwrap());
assert!(drive_dedup_future(storage.first_seen("k3")).unwrap());
assert!(
drive_dedup_future(storage.first_seen("k4")).unwrap(),
"k1 evicted, k4 should insert"
);
assert!(
drive_dedup_future(storage.first_seen("k1")).unwrap(),
"k1 was evicted, should be accepted again"
);
assert!(
!drive_dedup_future(storage.first_seen("k4")).unwrap(),
"k4 is still in window"
);
}
#[test]
fn bounded_fifo_capacity_one_always_evicts_previous() {
let storage = BoundedFifoDedupStorage::new(1);
assert!(drive_dedup_future(storage.first_seen("only-key")).unwrap());
assert!(
!drive_dedup_future(storage.first_seen("only-key")).unwrap(),
"duplicate in window"
);
assert!(
drive_dedup_future(storage.first_seen("new-key")).unwrap(),
"evicts old, inserts new"
);
assert!(
drive_dedup_future(storage.first_seen("only-key")).unwrap(),
"evicted, accepted again"
);
}
#[test]
fn bounded_fifo_is_correct_under_parallel_load() {
let storage = Arc::new(BoundedFifoDedupStorage::new(128));
let key = "dedup:fifo:parallel-msg";
let mut handles = Vec::new();
for _ in 0..16 {
let storage = Arc::clone(&storage);
handles.push(thread::spawn(move || {
drive_dedup_future(storage.first_seen(key)).unwrap()
}));
}
let accepted = handles
.into_iter()
.map(|h| h.join().expect("thread join"))
.filter(|a| *a)
.count();
assert_eq!(accepted, 1);
}
#[test]
fn queued_requests_returns_all_items() {
let storage = InMemoryReconciliationStorage::default();
let r1 =
ReconciliationRequest::for_outcome("m-each-1", "p1", DeliveryOutcome::Indeterminate)
.unwrap();
let r2 =
ReconciliationRequest::for_outcome("m-each-2", "p1", DeliveryOutcome::Indeterminate)
.unwrap();
drive_dedup_future(storage.enqueue(r1)).unwrap();
drive_dedup_future(storage.enqueue(r2)).unwrap();
let mut collected: Vec<String> = drive_dedup_future(storage.queued_requests())
.unwrap()
.into_iter()
.map(|r| r.message_id)
.collect();
collected.sort();
assert_eq!(collected, ["m-each-1", "m-each-2"]);
}
#[test]
fn ttl_dedup_accepts_first_and_rejects_duplicate() {
let storage = TtlDedupStorage::new(Duration::from_secs(60));
assert!(drive_dedup_future(storage.first_seen("msg-1")).unwrap());
assert!(!drive_dedup_future(storage.first_seen("msg-1")).unwrap());
}
#[test]
fn ttl_dedup_accepts_after_expiry() {
let storage = TtlDedupStorage::new(Duration::from_nanos(1));
assert!(drive_dedup_future(storage.first_seen("msg-2")).unwrap());
std::thread::sleep(Duration::from_millis(10));
assert!(
drive_dedup_future(storage.first_seen("msg-2")).unwrap(),
"key should be accepted after TTL expiry"
);
}
#[test]
fn ttl_dedup_distinct_keys_are_independent() {
let storage = TtlDedupStorage::new(Duration::from_secs(60));
assert!(drive_dedup_future(storage.first_seen("key-a")).unwrap());
assert!(drive_dedup_future(storage.first_seen("key-b")).unwrap());
assert!(!drive_dedup_future(storage.first_seen("key-a")).unwrap());
assert!(!drive_dedup_future(storage.first_seen("key-b")).unwrap());
}
#[test]
fn ttl_dedup_expiry_check_is_independent_of_sweeping() {
let storage = TtlDedupStorage::with_sweep_interval(Duration::from_nanos(1), Duration::MAX);
assert!(drive_dedup_future(storage.first_seen("msg-nosweep")).unwrap());
std::thread::sleep(Duration::from_millis(10));
assert!(
drive_dedup_future(storage.first_seen("msg-nosweep")).unwrap(),
"an expired key must be re-accepted even when the sweep never fires"
);
}
#[test]
fn ttl_dedup_without_sweeping_never_releases_distinct_keys() {
let storage = TtlDedupStorage::with_sweep_interval(Duration::from_nanos(1), Duration::MAX);
for i in 0..64 {
assert!(drive_dedup_future(storage.first_seen(&format!("k{i}"))).unwrap());
}
assert_eq!(
storage.inner.lock().entries.len(),
64,
"every distinct key is retained when sweeping is disabled"
);
let sweeping =
TtlDedupStorage::with_sweep_interval(Duration::from_nanos(1), Duration::ZERO);
for i in 0..64 {
assert!(drive_dedup_future(sweeping.first_seen(&format!("k{i}"))).unwrap());
}
assert!(
sweeping.inner.lock().entries.len() < 64,
"a sweeping store reclaims expired entries"
);
}
#[test]
fn ttl_dedup_sweep_interval_constructor_respects_max_60s() {
let storage = TtlDedupStorage::new(Duration::from_secs(700));
assert!(drive_dedup_future(storage.first_seen("probe")).unwrap());
}
}