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::ShardHandle;
use crate::store::DiskStore;
use crate::tree::{LeafNode, Node};
use crate::wal::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(())
}