use std::error::Error;
use std::fmt;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
use super::ShardActor;
use super::expiry_index::{ArmError, DeadlineScheduler, Generation, WallClock};
use crate::store::MemoryStore;
use crate::sync::ballot::{Ballot, Stamp};
use crate::sync::topology::SyncNodeId;
use crate::ttl::filter::{Visibility, visible_value_at};
use crate::wal::{DurableWal, FsyncPolicy, WalRecovery};
#[derive(Default)]
pub(super) struct RaceClock(AtomicU64);
impl RaceClock {
pub(super) fn at(now: u64) -> Self {
Self(AtomicU64::new(now))
}
pub(super) fn set(&self, now: u64) {
self.0.store(now, Ordering::SeqCst);
}
}
impl fmt::Debug for RaceClock {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter
.debug_tuple("RaceClock")
.field(&self.0.load(Ordering::SeqCst))
.finish()
}
}
impl WallClock for RaceClock {
fn now(&self) -> u64 {
self.0.load(Ordering::SeqCst)
}
}
#[derive(Default)]
pub(super) struct Arms(pub(super) Vec<(Duration, Generation)>);
impl DeadlineScheduler for Arms {
fn schedule(&mut self, delay: Duration, generation: Generation) -> Result<(), ArmError> {
self.0.push((delay, generation));
Ok(())
}
}
type ActorFixture = (tempfile::TempDir, ShardActor, MemoryStore);
pub(super) fn actor(clock: Arc<RaceClock>) -> Result<ActorFixture, Box<dyn Error>> {
let dir = tempfile::tempdir()?;
let wal = DurableWal::new(dir.path().join("race.wal"), FsyncPolicy::CommitOnly)?;
Ok((
dir,
ShardActor::new_with_clock(wal, clock),
MemoryStore::new(),
))
}
pub(super) fn arm(actor: &mut ShardActor, arms: &mut Arms) -> Result<Generation, Box<dyn Error>> {
actor.arm_expiry(arms)?;
actor
.current_expiry_generation()
.ok_or_else(|| "missing current generation".into())
}
pub(super) fn fire(
actor: &mut ShardActor,
generation: Generation,
store: &MemoryStore,
) -> Result<bool, Box<dyn Error>> {
let Some(due) = actor.begin_expiry_deadline(generation) else {
return Ok(false);
};
for (_deadline, key) in due {
actor.inspect_expiry_key();
actor.delete_if_expired(&key, store)?;
}
actor.finish_expiry_deadline();
Ok(true)
}
pub(super) fn stamp(epoch: u64, sequence: u64) -> Stamp {
Stamp::new(Ballot::new(epoch, SyncNodeId::new("expiry-race")), sequence)
}
#[test]
fn refresh_minimum_to_later_makes_old_delivery_stale_and_value_survives()
-> Result<(), Box<dyn Error>> {
let clock = Arc::new(RaceClock::at(1_000));
let (_dir, mut actor, store) = actor(clock.clone())?;
let mut arms = Arms::default();
actor.put_with_ttl(b"k", b"old", Some(Duration::from_nanos(10)), &store)?;
let stale = arm(&mut actor, &mut arms)?;
actor.put_with_ttl(b"k", b"fresh", Some(Duration::from_nanos(100)), &store)?;
let current = arm(&mut actor, &mut arms)?;
clock.set(1_010);
assert!(!fire(&mut actor, stale, &store)?);
let refreshed = actor
.get_raw(b"k", &store)?
.ok_or("refreshed value missing")?;
assert_eq!(
visible_value_at(&refreshed, 1_010)?,
Visibility::Live(b"fresh".to_vec())
);
clock.set(1_100);
assert!(fire(&mut actor, current, &store)?);
assert!(actor.get_raw(b"k", &store)?.is_none());
assert_eq!(arms.0.len(), 2);
Ok(())
}
#[test]
fn refresh_non_minimum_to_earlier_advances_generation_once() -> Result<(), Box<dyn Error>> {
let clock = Arc::new(RaceClock::at(100));
let (_dir, mut actor, store) = actor(clock)?;
let mut arms = Arms::default();
actor.put_with_ttl(b"minimum", b"v", Some(Duration::from_nanos(50)), &store)?;
let old = arm(&mut actor, &mut arms)?;
actor.put_with_ttl(b"other", b"v", Some(Duration::from_nanos(100)), &store)?;
actor.arm_expiry(&mut arms)?;
actor.put_with_ttl(b"other", b"earlier", Some(Duration::from_nanos(25)), &store)?;
let new = arm(&mut actor, &mut arms)?;
assert_ne!(old, new);
assert_eq!(arms.0.len(), 2);
assert_eq!(arms.0[1].0, Duration::from_nanos(25));
assert_eq!(actor.expiry_metrics().minimum_changes, 2);
Ok(())
}
#[test]
fn delete_current_and_one_equal_minimum_obey_rearm_rules() -> Result<(), Box<dyn Error>> {
let clock = Arc::new(RaceClock::at(500));
let (_dir, mut actor, store) = actor(clock.clone())?;
let mut arms = Arms::default();
for key in [b"a".as_slice(), b"b".as_slice()] {
actor.put_with_ttl(key, key, Some(Duration::from_nanos(10)), &store)?;
}
actor.put_with_ttl(b"next", b"v", Some(Duration::from_nanos(20)), &store)?;
let equal_generation = arm(&mut actor, &mut arms)?;
actor.delete(b"a", Stamp::bottom(), &store)?;
actor.arm_expiry(&mut arms)?;
assert_eq!(
arms.0.len(),
1,
"deleting one equal-minimum peer must not re-arm"
);
clock.set(510);
assert!(fire(&mut actor, equal_generation, &store)?);
assert!(actor.get_raw(b"b", &store)?.is_none());
let next_generation = arm(&mut actor, &mut arms)?;
assert_eq!(arms.0.len(), 2);
actor.delete(b"next", Stamp::bottom(), &store)?;
actor.arm_expiry(&mut arms)?;
assert_eq!(actor.current_expiry_generation(), None);
assert_eq!(arms.0.len(), 2, "empty index must not arm");
assert!(!fire(&mut actor, next_generation, &store)?);
Ok(())
}
#[test]
fn restart_rebuilds_before_arm_and_stale_delivery_moves_zero_work_counters()
-> Result<(), Box<dyn Error>> {
let clock = Arc::new(RaceClock::at(1));
let dir = tempfile::tempdir()?;
let path = dir.path().join("restart.wal");
let mut store = MemoryStore::new();
let wal = DurableWal::new(&path, FsyncPolicy::CommitOnly)?;
let mut first = ShardActor::new_with_global_clock(wal, clock);
let mut first_arms = Arms::default();
first.put_with_ttl(b"k", b"v", Some(Duration::from_secs(3_600)), &store)?;
first.commit(&mut store)?;
let old = arm(&mut first, &mut first_arms)?;
drop(first);
let recovered = WalRecovery::recover_path(&path, &store)?;
let wal = DurableWal::new(&path, FsyncPolicy::CommitOnly)?;
let mut restarted =
ShardActor::from_recovered(wal, recovered, &store, crate::tree::TreePolicy::V1_DEFAULT)?;
assert_eq!(restarted.expiry_metrics().rebuild_entries, 1);
let mut restarted_arms = Arms::default();
let current = arm(&mut restarted, &mut restarted_arms)?;
assert_ne!(old, current);
let before = restarted.expiry_metrics();
assert!(!fire(&mut restarted, old, &store)?);
let after = restarted.expiry_metrics();
assert_eq!(after.stale_drops, before.stale_drops + 1);
assert_eq!(after.deadline_deliveries, before.deadline_deliveries + 1);
assert_eq!(after.actor_wakes, before.actor_wakes + 1);
assert_eq!(after.index_mutations, before.index_mutations);
assert_eq!(after.inspected_keys, before.inspected_keys);
assert_eq!(after.delete_attempts, before.delete_attempts);
assert_eq!(after.deletes, before.deletes);
assert_eq!(after.physical_arms, before.physical_arms);
Ok(())
}
#[test]
fn merge_adopt_invalidates_old_delivery_and_rebuilt_minimum_wins() -> Result<(), Box<dyn Error>> {
let clock = Arc::new(RaceClock::at(1_000));
let source_dir = tempfile::tempdir()?;
let source_wal = DurableWal::new(
source_dir.path().join("source.wal"),
FsyncPolicy::CommitOnly,
)?;
let mut source = ShardActor::new_with_clock(source_wal, clock.clone());
let mut source_store = MemoryStore::new();
source.apply_durable(
b"adopted",
None,
b"v".to_vec(),
Some(Duration::from_nanos(10)),
stamp(1, 0),
&mut source_store,
)?;
let export = source.export_reachable(0, &source_store)?;
let target_dir = tempfile::tempdir()?;
let target_wal = DurableWal::new(
target_dir.path().join("target.wal"),
FsyncPolicy::CommitOnly,
)?;
let mut target = ShardActor::new_with_clock(target_wal, clock);
let mut target_store = MemoryStore::new();
let mut arms = Arms::default();
target.apply_durable(
b"local",
None,
b"v".to_vec(),
Some(Duration::from_nanos(50)),
stamp(1, 1),
&mut target_store,
)?;
let stale = arm(&mut target, &mut arms)?;
target.merge_adopt(&[export], &mut target_store)?;
let before = target.expiry_metrics();
assert!(!fire(&mut target, stale, &target_store)?);
let after = target.expiry_metrics();
assert_eq!(after.inspected_keys, before.inspected_keys);
assert_eq!(after.deletes, before.deletes);
let current = arm(&mut target, &mut arms)?;
assert_ne!(stale, current);
assert_eq!(arms.0.len(), 2);
assert_eq!(arms.0[1].0, Duration::from_nanos(10));
Ok(())
}