haematite 0.6.2

Content-addressed, branchable, actor-native storage engine
Documentation
use std::cell::RefCell;
use std::convert::Infallible;
use std::error::Error;
use std::sync::{Arc, Mutex};

use super::commit::{BranchCommitError, CommitDurability, CommitRequest, commit_branch};
use super::durable_record::{CreateExclusive, DurableRecordStore, EntryFence};
use super::handle::DEFAULT_SHARD_ID;
use super::lifecycle::create_branch;
use super::native_record_store::NativeDurableRecordStore;
use super::{
    BranchKind, BranchRefError, BranchRefRecord, BranchRefStore, BranchRegistry, BranchShardRef,
};
use crate::store::{MemoryStore, NodeStore};
use crate::tree::{Hash, LeafNode, Node};

fn record(name: &str, created: u64, seq: u64, head: u8) -> BranchRefRecord {
    BranchRefRecord {
        name: name.to_owned(),
        created,
        kind: BranchKind::Work,
        namespace_lineage: None,
        seq,
        timestamp: created + seq,
        shards: vec![BranchShardRef {
            shard_id: 0,
            fork_anchor: Hash::from_bytes([1; 32]),
            head: Hash::from_bytes([head; 32]),
        }],
        parents: Vec::new(),
    }
}

#[test]
fn native_record_store_contract_matches_branch_ref_store() -> Result<(), BranchRefError> {
    let directory = tempfile::tempdir().map_err(BranchRefError::Io)?;
    let mut backend = NativeDurableRecordStore::open(directory.path())?;
    assert!(backend.list_read_at_open()?.is_empty());

    let first = record("work", 7, 0, 2);
    assert_eq!(
        backend.create_exclusive(&first)?,
        CreateExclusive::Installed
    );
    backend.entry_fence(EntryFence::Present(&first))?;
    let loaded = backend.list_read_at_open()?;
    assert_eq!(loaded.len(), 1);
    assert_eq!(loaded.first(), Some(&first));

    let replacement = record("work", 7, 1, 3);
    assert_eq!(
        backend.cas_replace_install("work", 7, 0, &replacement)?,
        replacement
    );
    backend.entry_fence(EntryFence::Present(&replacement))?;
    assert!(backend.unlink("work")?);
    backend.entry_fence(EntryFence::Absent("work"))?;
    assert!(!backend.unlink("work")?);
    Ok(())
}

/// Lock a test mutex, adopting the inner state on poison (house idiom;
/// `unwrap_used`/`expect_used` are denied crate-wide).
fn locked<T>(mutex: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
    match mutex.lock() {
        Ok(guard) => guard,
        Err(poisoned) => poisoned.into_inner(),
    }
}

#[derive(Debug)]
struct TraceStore<B> {
    inner: B,
    trace: Arc<Mutex<Vec<&'static str>>>,
}

impl<B: DurableRecordStore<Error = BranchRefError>> DurableRecordStore for TraceStore<B> {
    type Error = BranchRefError;
    fn list_read_at_open(&mut self) -> Result<Vec<BranchRefRecord>, BranchRefError> {
        locked(&self.trace).push("list/read-at-open");
        self.inner.list_read_at_open()
    }

    fn create_exclusive(
        &mut self,
        record: &BranchRefRecord,
    ) -> Result<CreateExclusive, BranchRefError> {
        locked(&self.trace).push("create-exclusive");
        self.inner.create_exclusive(record)
    }

    fn cas_replace_install(
        &mut self,
        name: &str,
        expected_created: u64,
        expected_seq: u64,
        replacement: &BranchRefRecord,
    ) -> Result<BranchRefRecord, BranchRefError> {
        locked(&self.trace).push("cas-replace-install");
        self.inner
            .cas_replace_install(name, expected_created, expected_seq, replacement)
    }

    fn unlink(&mut self, name: &str) -> Result<bool, BranchRefError> {
        locked(&self.trace).push("unlink");
        self.inner.unlink(name)
    }

    fn entry_fence(&mut self, expected: EntryFence<'_>) -> Result<(), BranchRefError> {
        locked(&self.trace).push(match expected {
            EntryFence::Present(_) => "entry-fence-present",
            EntryFence::Absent(_) => "entry-fence-absent",
        });
        self.inner.entry_fence(expected)
    }
}

#[test]
fn native_record_store_trace_create_advance_failed_advance_remove() -> Result<(), BranchRefError> {
    let directory = tempfile::tempdir().map_err(BranchRefError::Io)?;
    let trace = Arc::new(Mutex::new(Vec::new()));
    let inner = NativeDurableRecordStore::open(directory.path())?;
    let backend = TraceStore {
        inner,
        trace: Arc::clone(&trace),
    };
    let mut refs = BranchRefStore::from_backend(backend)?;
    let first = record("trace", 11, 0, 4);
    refs.create(first)?;
    assert_eq!(
        refs.advance(
            "trace",
            11,
            0,
            &[(0, Hash::from_bytes([5; 32]))],
            Vec::new(),
            12,
        )?,
        1
    );
    assert!(matches!(
        refs.advance("trace", 11, 0, &[], Vec::new(), 13),
        Err(BranchRefError::StaleSeq { .. })
    ));
    assert!(refs.remove("trace")?.is_some());
    assert_eq!(
        locked(&trace).as_slice(),
        [
            "list/read-at-open",
            "create-exclusive",
            "entry-fence-present",
            "cas-replace-install",
            "entry-fence-present",
            "cas-replace-install",
            "unlink",
            "entry-fence-absent",
        ]
    );
    Ok(())
}

#[derive(Debug)]
struct FenceCountingStore {
    inner: MemoryStore,
    fences: RefCell<usize>,
}

impl NodeStore for FenceCountingStore {
    type Error = Infallible;

    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> {
        *self.fences.borrow_mut() += 1;
        Ok(())
    }
}

#[test]
fn native_record_store_trace_volatile_and_durable_noop_never_fence() -> Result<(), Box<dyn Error>> {
    let directory =
        tempfile::tempdir().map_err(|error| BranchCommitError::Ref(BranchRefError::Io(error)))?;
    let mut memory = MemoryStore::new();
    let root = memory.put(&Node::Leaf(LeafNode::new(Vec::new())?));
    let mut store = FenceCountingStore {
        inner: memory,
        fences: RefCell::new(0),
    };
    let registry = BranchRegistry::new();
    let mut refs = BranchRefStore::open(directory.path())?;
    let branch = create_branch(
        "no-fence",
        [(DEFAULT_SHARD_ID, root)],
        &mut refs,
        &registry,
        23,
    )?;
    branch.put(DEFAULT_SHARD_ID, b"key", b"value")?;
    commit_branch(
        &branch,
        &mut store,
        &registry,
        CommitRequest {
            durability: CommitDurability::Volatile,
            extra_parents: &[],
            timestamp: 24,
        },
    )?;
    commit_branch(
        &branch,
        &mut store,
        &registry,
        CommitRequest {
            durability: CommitDurability::Durable { refs: &mut refs },
            extra_parents: &[],
            timestamp: 25,
        },
    )?;
    assert_eq!(*store.fences.borrow(), 0);
    Ok(())
}

#[test]
fn native_installed_unfenced_trace_adopts_map_and_returns_error() -> Result<(), BranchRefError> {
    let directory = tempfile::tempdir().map_err(BranchRefError::Io)?;
    let mut refs = BranchRefStore::open(directory.path())?;
    let durable = record("unfenced", 19, 0, 6);
    super::persist::fail_next_parent_dir_sync();
    assert!(matches!(
        refs.create(durable.clone()),
        Err(BranchRefError::Io(_))
    ));
    assert_eq!(refs.get("unfenced"), Some(&durable));
    Ok(())
}