haematite 0.7.0

Content-addressed, branchable, actor-native storage engine
Documentation
use std::cell::Cell;
use std::error::Error;
use std::sync::Arc;
use std::time::Duration;

use super::expiry_race_tests::{Arms, RaceClock, actor, arm, fire, stamp};
use super::handle::BatchItem;
use super::{GroupOutcome, GroupWrite, ShardActor};
use crate::store::{MemoryStore, NodeStore, StoreError};
use crate::sync::ballot::Ballot;
use crate::sync::topology::SyncNodeId;
use crate::tree::{Hash, Node};
use crate::ttl::filter::{Visibility, visible_value_at};
use crate::wal::{DurableWal, FsyncPolicy};

#[derive(Debug)]
struct SyncFailStore {
    inner: MemoryStore,
    fail_sync: Cell<bool>,
}

impl NodeStore for SyncFailStore {
    type Error = StoreError;

    fn get(&self, hash: &Hash) -> Result<Option<Arc<Node>>, Self::Error> {
        Ok(self.inner.get(hash))
    }

    fn put(&mut self, node: &Node) -> Result<Hash, Self::Error> {
        Ok(self.inner.put(node))
    }

    fn sync_dirty_dirs(&self) -> Result<(), Self::Error> {
        if self.fail_sync.get() {
            Err(StoreError::Io(std::io::Error::other(
                "injected merge barrier failure",
            )))
        } else {
            Ok(())
        }
    }
}

#[test]
fn merge_adopt_failure_preserves_buffer_index_and_oracle() -> 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"source",
        None,
        b"v".to_vec(),
        Some(Duration::from_nanos(5)),
        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 store = SyncFailStore {
        inner: MemoryStore::new(),
        fail_sync: Cell::new(false),
    };
    target.apply_durable(
        b"old",
        None,
        b"v".to_vec(),
        Some(Duration::from_nanos(20)),
        stamp(1, 1),
        &mut store,
    )?;
    target.put_with_ttl(b"buffered", b"v", Some(Duration::from_nanos(30)), &store)?;
    let mut arms = Arms::default();
    let generation = arm(&mut target, &mut arms)?;
    let before = target.expiry_metrics();
    store.fail_sync.set(true);
    assert!(target.merge_adopt(&[export], &mut store).is_err());

    assert!(matches!(
        target.buffer().get(b"buffered"),
        crate::wal::LookupResult::BufferedValue(_)
    ));
    assert!(target.get_raw(b"buffered", &store)?.is_some());
    assert_eq!(target.current_expiry_generation(), Some(generation));
    assert_eq!(target.expiry_metrics(), before);
    target.arm_expiry(&mut arms)?;
    assert_eq!(
        arms.0.len(),
        1,
        "failed adoption must not speculatively re-arm"
    );
    Ok(())
}

#[test]
fn durable_and_batch_rejections_move_zero_index_or_arm_counters() -> Result<(), Box<dyn Error>> {
    let clock = Arc::new(RaceClock::at(100));
    let (_dir, mut actor, mut store) = actor(clock)?;
    let mut arms = Arms::default();
    actor.put_with_ttl(b"existing", b"old", Some(Duration::from_nanos(50)), &store)?;
    arm(&mut actor, &mut arms)?;
    actor.record_promise(Ballot::new(9, SyncNodeId::new("new-owner")))?;
    let before_fence = actor.expiry_metrics();
    assert!(
        actor
            .apply_durable(
                b"fenced",
                None,
                b"v".to_vec(),
                Some(Duration::from_nanos(1)),
                stamp(1, 0),
                &mut store
            )
            .is_err()
    );
    assert_eq!(actor.expiry_metrics(), before_fence);

    let before_cas = actor.expiry_metrics();
    assert!(
        actor
            .apply_durable(
                b"existing",
                Some(Hash::of(b"wrong")),
                b"new".to_vec(),
                Some(Duration::from_nanos(1)),
                stamp(9, 1),
                &mut store
            )
            .is_err()
    );
    assert_eq!(actor.expiry_metrics(), before_cas);

    let items: Vec<BatchItem> = vec![
        (
            b"would-pass".to_vec(),
            None,
            b"v".to_vec(),
            Some(Duration::from_nanos(1)),
        ),
        (
            b"existing".to_vec(),
            None,
            b"bad".to_vec(),
            Some(Duration::from_nanos(2)),
        ),
    ];
    let before_batch = actor.expiry_metrics();
    assert!(
        actor
            .apply_durable_batch(items, stamp(9, 2), &mut store)
            .is_err()
    );
    assert_eq!(actor.expiry_metrics(), before_batch);
    assert!(actor.get_raw(b"would-pass", &store)?.is_none());
    Ok(())
}

#[test]
fn grouped_partial_rejection_indexes_survivors_in_order() -> Result<(), Box<dyn Error>> {
    let base = crate::branch::current_timestamp();
    let clock = Arc::new(RaceClock::at(base));
    let (_dir, mut actor, mut store) = actor(clock.clone())?;
    actor.apply_durable(
        b"rejected",
        None,
        b"seed".to_vec(),
        Some(Duration::from_secs(300)),
        stamp(1, 0),
        &mut store,
    )?;
    let before = actor.expiry_metrics();
    let outcomes = actor.apply_group(
        vec![
            GroupWrite::ApplyValue {
                key: b"first".to_vec(),
                expected: None,
                value: b"v".to_vec(),
                ttl: Some(Duration::from_secs(100)),
                stamp: stamp(2, 0),
            },
            GroupWrite::ApplyValue {
                key: b"rejected".to_vec(),
                expected: None,
                value: b"bad".to_vec(),
                ttl: Some(Duration::from_secs(1)),
                stamp: stamp(2, 1),
            },
            GroupWrite::ApplyValue {
                key: b"second".to_vec(),
                expected: None,
                value: b"v".to_vec(),
                ttl: Some(Duration::from_secs(200)),
                stamp: stamp(2, 2),
            },
        ],
        &mut store,
    );
    assert!(matches!(
        outcomes.as_slice(),
        [
            GroupOutcome::Committed,
            GroupOutcome::Rejected(_),
            GroupOutcome::Committed
        ]
    ));
    assert_eq!(
        actor.expiry_metrics().index_mutations - before.index_mutations,
        2
    );
    let mut arms = Arms::default();
    let first = arm(&mut actor, &mut arms)?;
    clock.set(base.saturating_add(100_000_000_000));
    assert!(fire(&mut actor, first, &store)?);
    let second = arm(&mut actor, &mut arms)?;
    assert_eq!(arms.0[1].0, Duration::from_secs(100));
    clock.set(base.saturating_add(200_000_000_000));
    assert!(fire(&mut actor, second, &store)?);
    assert_eq!(actor.get(b"rejected", &store)?, Some(b"seed".to_vec()));
    Ok(())
}

#[test]
fn overdue_delivery_is_immediate_once_and_backward_jump_allows_resurrection()
-> Result<(), Box<dyn Error>> {
    let clock = Arc::new(RaceClock::at(100));
    let (_dir, mut actor, store) = actor(clock.clone())?;
    let mut arms = Arms::default();
    actor.put_with_ttl(b"k", b"v", Some(Duration::ZERO), &store)?;
    let generation = arm(&mut actor, &mut arms)?;
    assert_eq!(arms.0, vec![(Duration::ZERO, generation)]);
    assert!(fire(&mut actor, generation, &store)?);
    actor.arm_expiry(&mut arms)?;
    assert_eq!(arms.0.len(), 1);

    let encoded = crate::ttl::entry::encode_optional_ttl_at(
        b"v".to_vec(),
        Some(Duration::from_nanos(10)),
        100,
    )?;
    assert_eq!(visible_value_at(&encoded, 110)?, Visibility::Expired);
    clock.set(90);
    assert_eq!(
        visible_value_at(&encoded, 90)?,
        Visibility::Live(b"v".to_vec())
    );
    Ok(())
}

#[derive(Debug)]
struct AlwaysFailStore;

impl NodeStore for AlwaysFailStore {
    type Error = StoreError;

    fn get(&self, _hash: &Hash) -> Result<Option<Arc<Node>>, Self::Error> {
        Ok(None)
    }

    fn put(&mut self, _node: &Node) -> Result<Hash, Self::Error> {
        Err(StoreError::Io(std::io::Error::other(
            "injected group commit failure",
        )))
    }
}

#[test]
fn grouped_commit_failure_rolls_back_every_expiry_effect() -> Result<(), Box<dyn Error>> {
    let clock = Arc::new(RaceClock::at(1_000));
    let dir = tempfile::tempdir()?;
    let wal = DurableWal::new(
        dir.path().join("group-failure.wal"),
        FsyncPolicy::CommitOnly,
    )?;
    let mut actor = ShardActor::new_with_clock(wal, clock);
    let before = actor.expiry_metrics();
    let outcomes = actor.apply_group(
        vec![
            GroupWrite::ApplyValue {
                key: b"a".to_vec(),
                expected: None,
                value: b"v".to_vec(),
                ttl: Some(Duration::from_nanos(1)),
                stamp: stamp(1, 0),
            },
            GroupWrite::ApplyValue {
                key: b"b".to_vec(),
                expected: None,
                value: b"v".to_vec(),
                ttl: Some(Duration::from_nanos(2)),
                stamp: stamp(1, 1),
            },
        ],
        &mut AlwaysFailStore,
    );
    assert!(
        outcomes
            .iter()
            .all(|outcome| matches!(outcome, GroupOutcome::CommitFailed(_)))
    );
    assert_eq!(actor.expiry_metrics(), before);
    assert!(actor.buffer().is_empty());
    let mut arms = Arms::default();
    actor.arm_expiry(&mut arms)?;
    assert!(arms.0.is_empty());
    Ok(())
}