use std::collections::BTreeMap;
use std::path::Path;
use std::sync::Arc;
use std::time::{Duration, Instant};
use beamr::module::ModuleRegistry;
use beamr::scheduler::{Scheduler, SchedulerConfig};
use super::recover::{collect_tree, recover_view};
use super::{SweepError, SweepHandle, SweepStats};
use crate::shard::actor::{ShardActor, ShardHandle};
use crate::store::DiskStore;
use crate::tree::{LeafNode, Node};
use crate::wal::{DurableWal, FsyncPolicy, Mutation};
const TIMEOUT: Duration = Duration::from_secs(5);
#[test]
fn sweep_stats_default_is_zero() {
assert_eq!(SweepStats::default().scanned, 0);
assert_eq!(SweepStats::default().expired, 0);
assert_eq!(SweepStats::default().deleted, 0);
}
fn physically_present(store_dir: &Path, wal_path: &Path, key: &[u8]) -> Result<bool, SweepError> {
let (store, root, buffer) = recover_view(store_dir, wal_path)?;
let mut merged = BTreeMap::new();
if let Some(root) = root {
collect_tree(&store, root, &mut merged)?;
}
for mutation in &buffer {
match mutation {
Mutation::Put { key, value } => {
merged.insert(key.clone(), value.clone());
}
Mutation::Delete { key } => {
merged.remove(key);
}
}
}
Ok(merged.contains_key(key))
}
#[test]
fn periodic_tick_physically_deletes_expired_entry() -> Result<(), Box<dyn std::error::Error>> {
let scheduler = Arc::new(Scheduler::new(
SchedulerConfig::default(),
Arc::new(ModuleRegistry::new()),
)?);
let dir = tempfile::tempdir()?;
let store_dir = dir.path().join("sweep.store");
let wal_path = dir.path().join("sweep.wal");
let mut store = DiskStore::new(&store_dir)?;
let _root = store.put(&Node::Leaf(LeafNode::new(Vec::new())?))?;
drop(store);
let shard = ShardHandle::spawn(Arc::clone(&scheduler), store_dir.clone(), wal_path.clone())?;
let key = b"ttl-victim".to_vec();
shard.put_with_ttl(
key.clone(),
b"doomed".to_vec(),
Some(Duration::from_millis(100)),
TIMEOUT,
)?;
assert!(
physically_present(&store_dir, &wal_path, &key)?,
"entry must be physically present before any sweep runs"
);
let interval = Duration::from_millis(20);
let sweep = SweepHandle::spawn(
Arc::clone(&scheduler),
store_dir.clone(),
wal_path.clone(),
shard,
interval,
TIMEOUT,
)?;
let deadline = Instant::now() + Duration::from_secs(10);
while physically_present(&store_dir, &wal_path, &key)? && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(20));
}
assert!(
!physically_present(&store_dir, &wal_path, &key)?,
"periodic sweep must physically remove the expired entry from the store/WAL"
);
assert!(sweep.shutdown(TIMEOUT).is_ok());
scheduler.shutdown();
Ok(())
}
fn merged_view(
store_dir: &Path,
wal_path: &Path,
) -> Result<BTreeMap<Vec<u8>, Vec<u8>>, Box<dyn std::error::Error>> {
let (store, root, buffer) = recover_view(store_dir, wal_path)?;
let mut merged = BTreeMap::new();
if let Some(root) = root {
collect_tree(&store, root, &mut merged)?;
}
for mutation in &buffer {
match mutation {
Mutation::Put { key, value } => {
merged.insert(key.clone(), value.clone());
}
Mutation::Delete { key } => {
merged.remove(key);
}
}
}
Ok(merged)
}
#[test]
fn recover_view_after_crash_before_sweep_is_intact() -> Result<(), Box<dyn std::error::Error>> {
let dir = tempfile::tempdir()?;
let store_dir = dir.path().join("sweep.store");
let wal_path = dir.path().join("sweep.wal");
let mut store = DiskStore::new(&store_dir)?;
let wal = DurableWal::new(&wal_path, FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
actor.put_with_ttl(b"live".to_vec(), b"keep".to_vec(), None)?;
actor.put_with_ttl(
b"doomed".to_vec(),
b"gone".to_vec(),
Some(Duration::from_nanos(1)),
)?;
let _root = actor.commit(&mut store)?;
drop(actor);
let merged = merged_view(&store_dir, &wal_path)?;
assert!(
merged.contains_key(b"live".as_slice()),
"the live entry must survive the crash"
);
assert!(
merged.contains_key(b"doomed".as_slice()),
"an expired entry not yet swept must be fully present after a crash, never torn"
);
Ok(())
}
#[test]
fn recover_view_after_staged_sweep_delete_is_consistent() -> Result<(), Box<dyn std::error::Error>>
{
let dir = tempfile::tempdir()?;
let store_dir = dir.path().join("sweep.store");
let wal_path = dir.path().join("sweep.wal");
let mut store = DiskStore::new(&store_dir)?;
let wal = DurableWal::new(&wal_path, FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
actor.put_with_ttl(b"live".to_vec(), b"keep".to_vec(), None)?;
actor.put_with_ttl(
b"doomed".to_vec(),
b"gone".to_vec(),
Some(Duration::from_nanos(1)),
)?;
actor.commit(&mut store)?;
let removed = actor.delete_if_expired(b"doomed", &store)?;
assert!(removed, "the expired key must be a sweep candidate");
drop(actor);
let merged = merged_view(&store_dir, &wal_path)?;
assert!(
merged.contains_key(b"live".as_slice()),
"the live entry must remain fully intact"
);
assert!(
!merged.contains_key(b"doomed".as_slice()),
"a staged sweep delete must recover as fully absent, never torn"
);
Ok(())
}