haematite 0.7.0

Content-addressed, branchable, actor-native storage engine
Documentation
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(())
}