use std::time::Duration;
mod errors;
pub mod handle;
mod liveness;
pub mod native;
mod scan;
mod startup;
mod stream_index;
#[cfg(test)]
mod coop_smoke;
use errors::{AppendError, CasError, HashCasError};
pub use handle::{RangeItem, ShardError, ShardHandle};
use crate::branch::current_timestamp;
use crate::store::NodeStore;
use crate::sync::ballot::{Ballot, Stamp};
use crate::tree::{Cursor, Hash, LeafNode, Node, batch_mutate_owned};
use crate::ttl::entry::{
StampedEntry, encode_optional_ttl, encode_stamped_optional_ttl, encode_stamped_tombstone,
};
use crate::ttl::filter::{Visibility, is_expired_at, visible_value};
use crate::wal::{
DurableWal, LookupResult, Mutation, PromiseRecord, RecoveredWal, WalBuffer, WalError,
};
#[derive(Debug)]
pub struct ShardActor {
wal: DurableWal,
buffer: WalBuffer,
committed_root: Option<Hash>,
live_streams: stream_index::LiveStreamIndex,
stream_index_errors: stream_index::SequenceIndexErrors,
promised: Ballot,
owner_epoch: Option<Ballot>,
persisted_max_minted: u64,
}
enum ApplyKind {
Value {
value: Vec<u8>,
ttl: Option<Duration>,
},
Tombstone,
}
pub(super) enum GroupWrite {
Cas {
key: Vec<u8>,
expected: Option<u64>,
new: u64,
},
ApplyValue {
key: Vec<u8>,
expected: Option<Hash>,
value: Vec<u8>,
ttl: Option<Duration>,
stamp: Stamp,
},
ApplyTombstone {
key: Vec<u8>,
expected: Option<Hash>,
stamp: Stamp,
},
}
#[derive(Debug)]
pub(super) enum GroupOutcome {
Committed,
Rejected(ShardError),
CommitFailed(ShardError),
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum RecordPromiseOutcome {
Promised,
Rejected { promised: Ballot },
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct PromiseState {
pub promised: Ballot,
pub owner_epoch: Option<Ballot>,
pub persisted_max_minted: u64,
pub committed_root: Option<Hash>,
}
impl ShardActor {
#[cfg(test)]
#[must_use]
pub fn new(wal: DurableWal) -> Self {
Self {
wal,
buffer: WalBuffer::new(),
committed_root: None,
live_streams: stream_index::LiveStreamIndex::new(),
stream_index_errors: stream_index::SequenceIndexErrors::new(),
promised: Ballot::bottom(),
owner_epoch: None,
persisted_max_minted: 0,
}
}
pub fn from_recovered<S>(
mut wal: DurableWal,
recovered: RecoveredWal,
store: &S,
) -> Result<Self, WalError>
where
S: NodeStore + ?Sized,
{
let committed_root = recovered.committed_root();
let stream_index = stream_index::rebuild(store, committed_root).map_err(tree_error)?;
let promise = recovered
.promise()
.cloned()
.unwrap_or_else(PromiseRecord::initial);
wal.seed_promise(promise.clone());
let buffer = recovered.into_buffer();
let PromiseRecord {
promised,
owner_epoch,
persisted_max_minted,
} = promise;
Ok(Self {
wal,
buffer,
committed_root,
live_streams: stream_index.live,
stream_index_errors: stream_index.errors,
promised,
owner_epoch,
persisted_max_minted,
})
}
#[must_use]
pub const fn committed_root(&self) -> Option<Hash> {
self.committed_root
}
#[cfg(test)]
pub fn put<K, V>(&mut self, key: K, value: V) -> Result<(), WalError>
where
K: Into<Vec<u8>>,
V: Into<Vec<u8>>,
{
let key = key.into();
let value = value.into();
self.put_encoded(key, value)
}
pub fn put_with_ttl<K, V>(
&mut self,
key: K,
value: V,
ttl: Option<Duration>,
) -> Result<(), WalError>
where
K: Into<Vec<u8>>,
V: Into<Vec<u8>>,
{
let key = key.into();
let value = encode_ttl_value(value.into(), ttl)?;
self.put_encoded(key, value)
}
fn put_encoded(&mut self, key: Vec<u8>, value: Vec<u8>) -> Result<(), WalError> {
let mutation = Mutation::Put {
key: key.clone(),
value: value.clone(),
};
self.wal.append_mutation(&mutation)?;
self.buffer.put(key, value);
Ok(())
}
pub fn delete<K>(&mut self, key: K, stamp: Stamp) -> Result<(), WalError>
where
K: Into<Vec<u8>>,
{
let key = key.into();
let encoded = encode_stamped_tombstone(stamp);
self.put_encoded(key, encoded)
}
pub fn delete_if_expired<S>(&mut self, key: &[u8], store: &S) -> Result<bool, WalError>
where
S: NodeStore + ?Sized,
{
let current = match self.buffer.get(key) {
LookupResult::BufferedValue(value) => Some(value),
LookupResult::BufferedDelete => None,
LookupResult::NotBuffered => match self.committed_root {
Some(root) => Cursor::new(store, root).get(key).map_err(tree_error)?,
None => None,
},
};
let Some(value) = current else {
return Ok(false);
};
if StampedEntry::decode(&value)
.map_err(tree_error)?
.is_some_and(|entry| entry.is_tombstone())
{
return Ok(false);
}
if is_expired_at(&value, current_timestamp()).map_err(tree_error)? {
self.buffer_remove(key.to_vec())?;
Ok(true)
} else {
Ok(false)
}
}
fn buffer_remove(&mut self, key: Vec<u8>) -> Result<(), WalError> {
let mutation = Mutation::Delete { key: key.clone() };
self.wal.append_mutation(&mutation)?;
self.buffer.delete(key);
Ok(())
}
pub fn get<K, S>(&self, key: K, store: &S) -> Result<Option<Vec<u8>>, WalError>
where
K: AsRef<[u8]>,
S: NodeStore + ?Sized,
{
let key = key.as_ref();
match self.buffer.get(key) {
LookupResult::BufferedValue(value) => visible_ttl_value(&value),
LookupResult::BufferedDelete => Ok(None),
LookupResult::NotBuffered => self.committed_root.map_or_else(
|| Ok(None),
|root| {
visible_optional_ttl_value(
Cursor::new(store, root).get(key).map_err(tree_error)?,
)
},
),
}
}
#[doc(hidden)]
pub fn get_raw<K, S>(&self, key: K, store: &S) -> Result<Option<Vec<u8>>, WalError>
where
K: AsRef<[u8]>,
S: NodeStore + ?Sized,
{
let key = key.as_ref();
match self.buffer.get(key) {
LookupResult::BufferedValue(value) => Ok(Some(value)),
LookupResult::BufferedDelete => Ok(None),
LookupResult::NotBuffered => self.committed_root.map_or_else(
|| Ok(None),
|root| Cursor::new(store, root).get(key).map_err(tree_error),
),
}
}
pub fn commit<S>(&mut self, store: &mut S) -> Result<Hash, WalError>
where
S: NodeStore + ?Sized,
{
let baseline_root = match self.committed_root {
Some(root) => root,
None => store_empty_root(store)?,
};
let batch = buffered_batch(&self.buffer);
let new_root = batch_mutate_owned(store, baseline_root, batch).map_err(tree_error)?;
store.sync_dirty_dirs().map_err(tree_error)?;
self.wal.commit(new_root)?;
stream_index::apply_committed_buffer(
&mut self.live_streams,
&mut self.stream_index_errors,
&self.buffer,
);
self.buffer = WalBuffer::new();
self.committed_root = Some(new_root);
Ok(new_root)
}
fn append<S>(
&mut self,
key: &[u8],
entries: Vec<Vec<u8>>,
expected_seq: u64,
ttl: Option<Duration>,
store: &mut S,
) -> Result<u64, AppendError>
where
S: NodeStore + ?Sized,
{
if entries.is_empty() {
return Ok(expected_seq);
}
let seq_key = sequence_key(key);
let actual = self.read_sequence(&seq_key, store)?;
if actual != expected_seq {
return Err(AppendError::SequenceConflict {
expected: expected_seq,
actual,
});
}
let entry_count = u64::try_from(entries.len())
.map_err(|_| WalError::TreeError("too many append entries".to_owned()))?;
let new_seq = actual
.checked_add(entry_count)
.ok_or_else(|| WalError::TreeError("append sequence overflow".to_owned()))?;
let mut mutations = Vec::with_capacity(entries.len().saturating_add(1));
for (offset, entry) in entries.into_iter().enumerate() {
let offset = u64::try_from(offset)
.map_err(|_| WalError::TreeError("too many append entries".to_owned()))?;
let seq = actual
.checked_add(offset.saturating_add(1))
.ok_or_else(|| WalError::TreeError("append sequence overflow".to_owned()))?;
let value = encode_ttl_value(entry, ttl)?;
mutations.push(Mutation::Put {
key: event_key(key, seq),
value,
});
}
mutations.push(Mutation::Put {
key: seq_key,
value: new_seq.to_be_bytes().to_vec(),
});
let previous_buffer = self.buffer.clone();
for mutation in mutations {
buffer_mutation(&mut self.buffer, mutation);
}
match self.commit(store) {
Ok(_root) => Ok(new_seq),
Err(error) => {
self.buffer = previous_buffer;
Err(AppendError::from(error))
}
}
}
fn read_value<S>(&self, key: &[u8], store: &S) -> Result<Option<u64>, WalError>
where
S: NodeStore + ?Sized,
{
self.get(key, store)?.map_or(Ok(None), |bytes| {
bytes
.as_slice()
.try_into()
.map(|raw| Some(u64::from_be_bytes(raw)))
.map_err(|_| WalError::TreeError("invalid scalar value".to_owned()))
})
}
fn cas<S>(
&mut self,
key: &[u8],
expected: Option<u64>,
new: u64,
store: &mut S,
) -> Result<(), CasError>
where
S: NodeStore + ?Sized,
{
let prior = self.stage_cas(key, expected, new, store)?;
match self.commit(store) {
Ok(_root) => Ok(()),
Err(error) => {
self.buffer.restore_entry(key, prior);
Err(CasError::from(error))
}
}
}
fn stage_cas<S>(
&mut self,
key: &[u8],
expected: Option<u64>,
new: u64,
store: &S,
) -> Result<Option<Mutation>, CasError>
where
S: NodeStore + ?Sized,
{
let actual = self.read_value(key, store)?;
if actual != expected {
return Err(CasError::Mismatch { expected, actual });
}
let prior = self.buffer.snapshot_entry(key);
self.buffer.put(key, new.to_be_bytes());
Ok(prior)
}
fn apply_durable<S>(
&mut self,
key: &[u8],
expected: Option<Hash>,
value: Vec<u8>,
ttl: Option<Duration>,
stamp: Stamp,
store: &mut S,
) -> Result<(), HashCasError>
where
S: NodeStore + ?Sized,
{
self.apply_durable_kind(key, expected, ApplyKind::Value { value, ttl }, stamp, store)
}
fn apply_durable_tombstone<S>(
&mut self,
key: &[u8],
expected: Option<Hash>,
stamp: Stamp,
store: &mut S,
) -> Result<(), HashCasError>
where
S: NodeStore + ?Sized,
{
self.apply_durable_kind(key, expected, ApplyKind::Tombstone, stamp, store)
}
fn apply_durable_kind<S>(
&mut self,
key: &[u8],
expected: Option<Hash>,
kind: ApplyKind,
stamp: Stamp,
store: &mut S,
) -> Result<(), HashCasError>
where
S: NodeStore + ?Sized,
{
let prior = self.stage_apply_kind(key, expected, kind, stamp, store)?;
match self.commit(store) {
Ok(_root) => Ok(()),
Err(error) => {
self.buffer.restore_entry(key, prior);
Err(HashCasError::from(error))
}
}
}
fn stage_apply_kind<S>(
&mut self,
key: &[u8],
expected: Option<Hash>,
kind: ApplyKind,
stamp: Stamp,
store: &S,
) -> Result<Option<Mutation>, HashCasError>
where
S: NodeStore + ?Sized,
{
if stamp.epoch < self.promised {
return Err(HashCasError::Fenced {
promised: self.promised.clone(),
attempted: stamp.epoch,
});
}
let actual = self.current_value_hash(key, store)?;
if actual != expected {
return Err(HashCasError::HashMismatch { expected, actual });
}
let prior = self.buffer.snapshot_entry(key);
let encoded = match kind {
ApplyKind::Value { value, ttl } => {
encode_stamped_optional_ttl(value, stamp, ttl).map_err(tree_error)?
}
ApplyKind::Tombstone => encode_stamped_tombstone(stamp),
};
self.buffer.put(key, encoded);
Ok(prior)
}
fn apply_durable_batch<S>(
&mut self,
items: Vec<handle::BatchItem>,
stamp: Stamp,
store: &mut S,
) -> Result<(), HashCasError>
where
S: NodeStore + ?Sized,
{
if stamp.epoch < self.promised {
return Err(HashCasError::Fenced {
promised: self.promised.clone(),
attempted: stamp.epoch,
});
}
for (key, expected, _value, _ttl) in &items {
let actual = self.current_value_hash(key, store)?;
if actual != *expected {
return Err(HashCasError::HashMismatch {
expected: *expected,
actual,
});
}
}
if items.is_empty() {
return Ok(());
}
let previous_buffer = self.buffer.clone();
for (key, _expected, value, ttl) in items {
let encoded = match encode_stamped_optional_ttl(value, stamp.clone(), ttl) {
Ok(encoded) => encoded,
Err(error) => {
self.buffer = previous_buffer;
return Err(tree_error(error).into());
}
};
self.buffer.put(key, encoded);
}
match self.commit(store) {
Ok(_root) => Ok(()),
Err(error) => {
self.buffer = previous_buffer;
Err(HashCasError::from(error))
}
}
}
pub(super) fn apply_group<S>(
&mut self,
writes: Vec<GroupWrite>,
store: &mut S,
) -> Vec<GroupOutcome>
where
S: NodeStore + ?Sized,
{
let mut outcomes: Vec<Option<GroupOutcome>> = Vec::with_capacity(writes.len());
let mut staged: Vec<(usize, Vec<u8>, Option<Mutation>)> = Vec::with_capacity(writes.len());
for write in writes {
let index = outcomes.len();
match self.stage_group_write(write, store) {
Ok((key, prior)) => {
staged.push((index, key, prior));
outcomes.push(None); }
Err(error) => outcomes.push(Some(GroupOutcome::Rejected(error))),
}
}
if staged.is_empty() {
return finalise_group_outcomes(outcomes);
}
match self.commit(store) {
Ok(_root) => {
for (index, _key, _prior) in staged {
outcomes[index] = Some(GroupOutcome::Committed);
}
}
Err(error) => {
let message = error.to_string();
for (index, key, prior) in staged.into_iter().rev() {
self.buffer.restore_entry(&key, prior);
outcomes[index] = Some(GroupOutcome::CommitFailed(ShardError::Wal(
WalError::TreeError(message.clone()),
)));
}
}
}
finalise_group_outcomes(outcomes)
}
fn stage_group_write<S>(
&mut self,
write: GroupWrite,
store: &S,
) -> Result<(Vec<u8>, Option<Mutation>), ShardError>
where
S: NodeStore + ?Sized,
{
match write {
GroupWrite::Cas { key, expected, new } => {
let prior = self.stage_cas(&key, expected, new, store)?;
Ok((key, prior))
}
GroupWrite::ApplyValue {
key,
expected,
value,
ttl,
stamp,
} => {
let prior = self.stage_apply_kind(
&key,
expected,
ApplyKind::Value { value, ttl },
stamp,
store,
)?;
Ok((key, prior))
}
GroupWrite::ApplyTombstone {
key,
expected,
stamp,
} => {
let prior =
self.stage_apply_kind(&key, expected, ApplyKind::Tombstone, stamp, store)?;
Ok((key, prior))
}
}
}
#[must_use]
pub fn promise_state(&self) -> PromiseState {
PromiseState {
promised: self.promised().clone(),
owner_epoch: self.owner_epoch().cloned(),
persisted_max_minted: self.persisted_max_minted(),
committed_root: self.committed_root(),
}
}
#[must_use]
pub const fn promised(&self) -> &Ballot {
&self.promised
}
#[must_use]
pub const fn owner_epoch(&self) -> Option<&Ballot> {
self.owner_epoch.as_ref()
}
#[must_use]
pub const fn persisted_max_minted(&self) -> u64 {
self.persisted_max_minted
}
fn promise_snapshot(&self) -> PromiseRecord {
PromiseRecord {
promised: self.promised.clone(),
owner_epoch: self.owner_epoch.clone(),
persisted_max_minted: self.persisted_max_minted,
}
}
pub fn record_promise(&mut self, ballot: Ballot) -> Result<RecordPromiseOutcome, WalError> {
if ballot <= self.promised {
return Ok(RecordPromiseOutcome::Rejected {
promised: self.promised.clone(),
});
}
let snapshot = PromiseRecord {
promised: ballot.clone(),
owner_epoch: self.owner_epoch.clone(),
persisted_max_minted: self.persisted_max_minted,
};
self.wal.append_promise(&snapshot)?;
self.promised = ballot;
Ok(RecordPromiseOutcome::Promised)
}
pub fn record_owner_epoch(&mut self, ballot: Ballot) -> Result<(), WalError> {
let mut snapshot = self.promise_snapshot();
snapshot.owner_epoch = Some(ballot.clone());
self.wal.append_promise(&snapshot)?;
self.owner_epoch = Some(ballot);
Ok(())
}
pub fn reserve_minted(&mut self, counter: u64) -> Result<u64, WalError> {
let reserved = self.persisted_max_minted.max(counter);
if reserved == self.persisted_max_minted {
return Ok(reserved);
}
let mut snapshot = self.promise_snapshot();
snapshot.persisted_max_minted = reserved;
self.wal.append_promise(&snapshot)?;
self.persisted_max_minted = reserved;
Ok(reserved)
}
fn current_value_hash<S>(&self, key: &[u8], store: &S) -> Result<Option<Hash>, WalError>
where
S: NodeStore + ?Sized,
{
Ok(self.get(key, store)?.map(|value| Hash::of(&value)))
}
#[must_use]
pub const fn buffer(&self) -> &WalBuffer {
&self.buffer
}
pub(super) fn scan_sequences(&self) -> Result<Vec<handle::StreamSeq>, ShardError> {
scan::scan_sequences(&self.live_streams, &self.stream_index_errors)
}
#[cfg(test)]
fn scan_sequences_full_walk<S>(&self, store: &S) -> Result<Vec<handle::StreamSeq>, ShardError>
where
S: NodeStore + ?Sized,
{
stream_index::full_walk_with_buffer(store, self.committed_root, &self.buffer)
}
fn read_sequence<S>(&self, seq_key: &[u8], store: &S) -> Result<u64, WalError>
where
S: NodeStore + ?Sized,
{
self.get(seq_key, store)?.map_or(Ok(0), |bytes| {
bytes
.as_slice()
.try_into()
.map(u64::from_be_bytes)
.map_err(|_| WalError::TreeError("invalid sequence metadata".to_owned()))
})
}
pub fn export_reachable<S>(
&self,
shard_id: crate::branch::ShardId,
store: &S,
) -> Result<(Option<Hash>, Vec<crate::sync::NodeTransfer>), WalError>
where
S: NodeStore + ?Sized,
{
let source_root = self.committed_root;
let missing =
crate::sync::find_missing_nodes(store, &EmptyTarget, shard_id, source_root, None)
.map_err(|error| WalError::TreeError(error.to_string()))?;
Ok((source_root, missing.transfers))
}
pub fn merge_adopt<S>(
&mut self,
promisers: &[(Option<Hash>, Vec<crate::sync::NodeTransfer>)],
store: &mut S,
) -> Result<Option<Hash>, WalError>
where
S: NodeStore + ?Sized,
{
let mut acc = self.committed_root;
for (promiser_root, transfers) in promisers {
for transfer in transfers {
let actual = transfer.node.hash();
if actual != transfer.hash {
return Err(WalError::TreeError(format!(
"handoff node hash mismatch: expected {:?}, actual {actual:?}",
transfer.hash
)));
}
let stored = store.put(&transfer.node).map_err(tree_error)?;
if stored != transfer.hash {
return Err(WalError::TreeError(format!(
"handoff store wrote {stored:?}, expected {:?}",
transfer.hash
)));
}
}
acc = crate::sync::merge_committed_union(acc, *promiser_root, store)
.map_err(|error| WalError::TreeError(error.to_string()))?;
}
let Some(root) = acc else {
return Ok(self.committed_root);
};
self.buffer = WalBuffer::new();
let stream_index = stream_index::rebuild(store, Some(root)).map_err(tree_error)?;
store.sync_dirty_dirs().map_err(tree_error)?;
self.wal.commit(root)?;
self.committed_root = Some(root);
self.live_streams = stream_index.live;
self.stream_index_errors = stream_index.errors;
Ok(Some(root))
}
}
fn finalise_group_outcomes(outcomes: Vec<Option<GroupOutcome>>) -> Vec<GroupOutcome> {
outcomes
.into_iter()
.map(|slot| {
slot.unwrap_or_else(|| {
GroupOutcome::CommitFailed(ShardError::Wal(WalError::TreeError(
"group-commit outcome was not resolved".to_owned(),
)))
})
})
.collect()
}
fn buffer_mutation(buffer: &mut WalBuffer, mutation: Mutation) {
match mutation {
Mutation::Put { key, value } => buffer.put(key, value),
Mutation::Delete { key } => buffer.delete(key),
}
}
fn event_key(key: &[u8], seq: u64) -> Vec<u8> {
let mut encoded = Vec::with_capacity(key.len().saturating_add(9));
encoded.extend_from_slice(key);
encoded.push(0);
encoded.extend_from_slice(&seq.to_be_bytes());
encoded
}
pub(super) const SEQ_SUFFIX: &[u8] = &[0xff, b's', b'e', b'q'];
fn sequence_key(key: &[u8]) -> Vec<u8> {
let mut encoded = Vec::with_capacity(key.len().saturating_add(SEQ_SUFFIX.len()));
encoded.extend_from_slice(key);
encoded.extend_from_slice(SEQ_SUFFIX);
encoded
}
pub(super) fn decode_sequence_key(encoded: &[u8]) -> Option<&[u8]> {
encoded
.len()
.checked_sub(SEQ_SUFFIX.len())
.and_then(|split| {
let (stream_key, suffix) = encoded.split_at(split);
(suffix == SEQ_SUFFIX).then_some(stream_key)
})
}
fn buffered_batch(buffer: &WalBuffer) -> Vec<(Vec<u8>, Option<Vec<u8>>)> {
buffer
.iter()
.map(|mutation| match mutation {
Mutation::Put { key, value } => (key.clone(), Some(value.clone())),
Mutation::Delete { key } => (key.clone(), None),
})
.collect()
}
fn store_empty_root<S>(store: &mut S) -> Result<Hash, WalError>
where
S: NodeStore + ?Sized,
{
let node = Node::Leaf(LeafNode::new(Vec::new()).map_err(tree_error)?);
store.put(&node).map_err(tree_error)
}
fn encode_ttl_value(value: Vec<u8>, ttl: Option<Duration>) -> Result<Vec<u8>, WalError> {
encode_optional_ttl(value, ttl).map_err(tree_error)
}
fn visible_optional_ttl_value(value: Option<Vec<u8>>) -> Result<Option<Vec<u8>>, WalError> {
value.map_or(Ok(None), |value| visible_ttl_value(&value))
}
fn visible_ttl_value(value: &[u8]) -> Result<Option<Vec<u8>>, WalError> {
match visible_value(value).map_err(tree_error)? {
Visibility::Live(value) => Ok(Some(value)),
Visibility::Expired => Ok(None),
}
}
fn tree_error(error: impl std::fmt::Display) -> WalError {
WalError::TreeError(error.to_string())
}
struct EmptyTarget;
impl crate::sync::TargetNodeReader for EmptyTarget {
fn read_target_node(
&self,
_hash: Hash,
) -> Result<Option<crate::sync::TargetNodeSummary>, crate::sync::SyncError> {
Ok(None)
}
}
#[cfg(test)]
#[path = "actor/tests.rs"]
mod tests;
#[cfg(test)]
#[path = "actor/stream_index_tests.rs"]
mod stream_index_tests;
#[cfg(test)]
mod storage_tests {
use super::ShardActor;
use crate::store::MemoryStore;
use crate::tree::{Hash, LeafNode, Node, batch_mutate};
use crate::wal::{DurableWal, FsyncPolicy, LookupResult, WalEntry, WalError, WalRecovery};
use std::path::{Path, PathBuf};
#[derive(Debug)]
struct TempWal {
dir: tempfile::TempDir,
path: PathBuf,
}
impl TempWal {
fn path(&self) -> &Path {
debug_assert!(self.path.starts_with(self.dir.path()));
&self.path
}
}
fn temp_path(name: &str) -> Result<TempWal, WalError> {
let dir = tempfile::tempdir()?;
let path = dir.path().join(name);
Ok(TempWal { dir, path })
}
fn empty_root(store: &mut MemoryStore) -> Result<Hash, WalError> {
let leaf =
LeafNode::new(Vec::new()).map_err(|error| WalError::TreeError(error.to_string()))?;
Ok(store.put(&Node::Leaf(leaf)))
}
fn test_stamp(counter: u64, node: &str, seq: u64) -> crate::sync::ballot::Stamp {
use crate::sync::ballot::{Ballot, Stamp};
use crate::sync::topology::SyncNodeId;
Stamp::new(Ballot::new(counter, SyncNodeId::new(node)), seq)
}
#[test]
fn put_returns_ok_only_after_entry_is_written_to_wal() -> Result<(), WalError> {
let temp = temp_path("actor-put.wal")?;
let path = temp.path();
let wal = DurableWal::new(path, FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
actor.put(b"event".to_vec(), b"payload".to_vec())?;
assert_eq!(
actor.buffer().get(b"event"),
LookupResult::BufferedValue(b"payload".to_vec())
);
assert_eq!(
DurableWal::read_file(path)?.entries(),
&[WalEntry::put(b"event".to_vec(), b"payload".to_vec())]
);
Ok(())
}
#[test]
fn delete_writes_a_stamped_tombstone_to_wal_and_reads_as_absent() -> Result<(), WalError> {
let temp = temp_path("actor-delete.wal")?;
let path = temp.path();
let wal = DurableWal::new(path, FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
let store = MemoryStore::new();
let stamp = test_stamp(2, "owner", 0);
actor.delete(b"event".to_vec(), stamp.clone())?;
let tombstone = crate::ttl::entry::encode_stamped_tombstone(stamp);
assert_eq!(
actor.buffer().get(b"event"),
LookupResult::BufferedValue(tombstone.clone())
);
assert_eq!(
DurableWal::read_file(path)?.entries(),
&[WalEntry::put(b"event".to_vec(), tombstone)]
);
assert_eq!(actor.get(b"event", &store)?, None);
Ok(())
}
#[test]
fn delete_if_expired_removes_only_expired_values() -> Result<(), WalError> {
let temp = temp_path("actor-delete-if-expired.wal")?;
let wal = DurableWal::new(temp.path(), FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
let store = MemoryStore::new();
actor.put(b"live".to_vec(), b"keep".to_vec())?;
assert!(!actor.delete_if_expired(b"live", &store)?);
assert_eq!(actor.get(b"live", &store)?, Some(b"keep".to_vec()));
actor.put_with_ttl(
b"gone".to_vec(),
b"stale".to_vec(),
Some(std::time::Duration::ZERO),
)?;
assert!(actor.delete_if_expired(b"gone", &store)?);
assert_eq!(actor.get(b"gone", &store)?, None);
assert!(!actor.delete_if_expired(b"missing", &store)?);
Ok(())
}
#[test]
fn r_tomb_sweep_never_removes_a_tombstone() -> Result<(), WalError> {
let temp = temp_path("actor-rtomb.wal")?;
let wal = DurableWal::new(temp.path(), FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
let store = MemoryStore::new();
let stamp = test_stamp(5, "owner", 9);
actor.delete(b"tomb".to_vec(), stamp.clone())?;
assert_eq!(actor.get(b"tomb", &store)?, None, "tombstone reads as None");
assert!(
!actor.delete_if_expired(b"tomb", &store)?,
"R-TOMB: the sweep must NEVER remove a tombstone"
);
let raw = actor
.get_raw(b"tomb", &store)?
.ok_or_else(|| WalError::TreeError("tombstone vanished from storage".to_owned()))?;
let decoded = crate::ttl::entry::StampedEntry::decode(&raw)
.map_err(|error| WalError::TreeError(error.to_string()))?
.ok_or_else(|| WalError::TreeError("tombstone is not a stamped entry".to_owned()))?;
assert!(
decoded.is_tombstone(),
"the swept-over entry is still a tombstone"
);
assert_eq!(
decoded.stamp(),
&stamp,
"the tombstone's stamp is intact after the sweep"
);
assert_eq!(actor.get(b"tomb", &store)?, None, "still reads as None");
actor.put_with_ttl(
b"expired".to_vec(),
b"stale".to_vec(),
Some(std::time::Duration::ZERO),
)?;
assert!(
actor.delete_if_expired(b"expired", &store)?,
"an actually-expired value is still swept"
);
assert_eq!(actor.get(b"expired", &store)?, None);
Ok(())
}
#[test]
fn from_recovered_accepts_put_get_delete_and_appends_after_replayed_entries()
-> Result<(), WalError> {
let temp = temp_path("actor-resume.wal")?;
let mut store = MemoryStore::new();
let committed_root = empty_root(&mut store)?;
let mut wal = DurableWal::new(temp.path(), FsyncPolicy::CommitOnly)?;
wal.commit(committed_root)?;
wal.append(&WalEntry::put(b"replayed".to_vec(), b"before".to_vec()))?;
drop(wal);
let recovered = WalRecovery::recover_path(temp.path(), &store)?;
let wal = DurableWal::new(temp.path(), FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::from_recovered(wal, recovered, &store)?;
assert_eq!(actor.committed_root(), Some(committed_root));
assert_eq!(actor.get(b"replayed", &store)?, Some(b"before".to_vec()));
actor.put(b"new".to_vec(), b"after".to_vec())?;
let stamp = test_stamp(1, "owner", 0);
actor.delete(b"replayed".to_vec(), stamp.clone())?;
assert_eq!(actor.get(b"new", &store)?, Some(b"after".to_vec()));
assert_eq!(actor.get(b"replayed", &store)?, None);
let tombstone = crate::ttl::entry::encode_stamped_tombstone(stamp);
assert_eq!(
DurableWal::read_file(temp.path())?.entries(),
&[
WalEntry::put(b"replayed".to_vec(), b"before".to_vec()),
WalEntry::put(b"new".to_vec(), b"after".to_vec()),
WalEntry::put(b"replayed".to_vec(), tombstone),
]
);
Ok(())
}
#[test]
fn commit_after_recovery_truncates_wal_updates_root_and_tree_reads() -> Result<(), WalError> {
let temp = temp_path("actor-commit-after-recovery.wal")?;
let mut store = MemoryStore::new();
let committed_root = empty_root(&mut store)?;
let mut wal = DurableWal::new(temp.path(), FsyncPolicy::CommitOnly)?;
wal.commit(committed_root)?;
wal.append(&WalEntry::put(b"event".to_vec(), b"payload".to_vec()))?;
drop(wal);
let recovered = WalRecovery::recover_path(temp.path(), &store)?;
let wal = DurableWal::new(temp.path(), FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::from_recovered(wal, recovered, &store)?;
let new_root = actor.commit(&mut store)?;
let contents = DurableWal::read_file(temp.path())?;
assert_eq!(contents.committed_root(), Some(new_root));
assert_eq!(contents.entries(), &[]);
assert_eq!(actor.committed_root(), Some(new_root));
assert!(actor.buffer().is_empty());
assert_eq!(actor.get(b"event", &store)?, Some(b"payload".to_vec()));
assert_ne!(new_root, committed_root);
Ok(())
}
#[test]
fn recovered_actor_matches_uncrashed_actor_after_same_commit() -> Result<(), WalError> {
let crashed = temp_path("actor-crashed.wal")?;
let uncrashed = temp_path("actor-uncrashed.wal")?;
let mut crashed_store = MemoryStore::new();
let mut uncrashed_store = MemoryStore::new();
let crashed_root = empty_root(&mut crashed_store)?;
let uncrashed_root = empty_root(&mut uncrashed_store)?;
let mut crashed_wal = DurableWal::new(crashed.path(), FsyncPolicy::CommitOnly)?;
crashed_wal.commit(crashed_root)?;
crashed_wal.append(&WalEntry::put(b"k".to_vec(), b"v1".to_vec()))?;
drop(crashed_wal);
let recovered = WalRecovery::recover_path(crashed.path(), &crashed_store)?;
let crashed_wal = DurableWal::new(crashed.path(), FsyncPolicy::CommitOnly)?;
let mut recovered_actor =
ShardActor::from_recovered(crashed_wal, recovered, &crashed_store)?;
recovered_actor.put(b"k".to_vec(), b"v2".to_vec())?;
let recovered_root = recovered_actor.commit(&mut crashed_store)?;
let uncrashed_wal = DurableWal::new(uncrashed.path(), FsyncPolicy::CommitOnly)?;
let mut uncrashed_actor = ShardActor::new(uncrashed_wal);
let uncrashed_root = batch_mutate(
&mut uncrashed_store,
uncrashed_root,
&[(b"k".to_vec(), Some(b"v2".to_vec()))],
)
.map_err(|error| WalError::TreeError(error.to_string()))?;
uncrashed_actor.put(b"k".to_vec(), b"v2".to_vec())?;
let committed_uncrashed_root = uncrashed_actor.commit(&mut uncrashed_store)?;
assert_eq!(
recovered_actor.get(b"k", &crashed_store)?,
Some(b"v2".to_vec())
);
assert_eq!(committed_uncrashed_root, uncrashed_root);
assert_eq!(recovered_root, committed_uncrashed_root);
Ok(())
}
}
#[cfg(test)]
mod node_dir_fsync_tests {
use super::ShardActor;
use crate::store::{DiskStore, MemoryStore, NodeStore};
use crate::tree::{Hash, Node};
use crate::wal::{DurableWal, FsyncPolicy, WalError, WalRecovery};
use std::cell::RefCell;
use std::path::PathBuf;
use std::sync::Arc;
#[derive(Debug)]
struct CrashWindowStore {
inner: DiskStore,
lossy: bool,
pending: RefCell<Vec<PathBuf>>,
dir: PathBuf,
}
impl CrashWindowStore {
fn new(dir: PathBuf, lossy: bool) -> Result<Self, WalError> {
let inner =
DiskStore::new(&dir).map_err(|error| WalError::TreeError(error.to_string()))?;
Ok(Self {
inner,
lossy,
pending: RefCell::new(Vec::new()),
dir,
})
}
fn node_path(&self, hash: &Hash) -> PathBuf {
let hex = hash.to_string();
let (prefix, file_name) = hex.split_at(2);
self.dir.join(prefix).join(file_name)
}
}
impl NodeStore for CrashWindowStore {
type Error = crate::store::StoreError;
fn get(&self, hash: &Hash) -> Result<Option<Arc<Node>>, Self::Error> {
self.inner.get(hash)
}
fn put(&mut self, node: &Node) -> Result<Hash, Self::Error> {
let hash = self.inner.put(node)?;
self.pending.borrow_mut().push(self.node_path(&hash));
Ok(hash)
}
fn sync_dirty_dirs(&self) -> Result<(), Self::Error> {
let pending = std::mem::take(&mut *self.pending.borrow_mut());
if self.lossy {
for path in pending {
match std::fs::remove_file(&path) {
Ok(()) => {}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(error) => return Err(crate::store::StoreError::Io(error)),
}
}
Ok(())
} else {
self.inner.sync_dirty_dirs()
}
}
}
fn commit_then_recover(
store: &mut CrashWindowStore,
wal_path: &std::path::Path,
nodes_dir: &std::path::Path,
) -> Result<(Hash, DiskStore, Result<ShardActor, WalError>), WalError> {
let wal = DurableWal::new(wal_path, FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
actor.put(b"durable-key".to_vec(), b"durable-value".to_vec())?;
let committed_root = actor.commit(store)?;
drop(actor);
let cold = DiskStore::new(nodes_dir).map_err(|e| WalError::TreeError(e.to_string()))?;
let actor = match WalRecovery::recover_path(wal_path, &cold) {
Ok(recovered) => {
let wal = DurableWal::new(wal_path, FsyncPolicy::CommitOnly)?;
ShardActor::from_recovered(wal, recovered, &cold)
}
Err(error) => Err(error),
};
Ok((committed_root, cold, actor))
}
#[test]
fn lossy_dir_barrier_makes_recovery_reject_missing_committed_root() -> Result<(), WalError> {
let dir = tempfile::tempdir()?;
let wal = tempfile::tempdir()?;
let wal_path = wal.path().join("shard.wal");
let nodes_dir = dir.path().join("nodes");
let mut store = CrashWindowStore::new(nodes_dir.clone(), true)?;
let (committed_root, _cold, recovered) =
commit_then_recover(&mut store, &wal_path, &nodes_dir)?;
match recovered {
Err(WalError::MissingCommittedRoot { root }) => {
assert_eq!(
root, committed_root,
"recovery must name the marker's now-unreachable root"
);
Ok(())
}
Err(other) => Err(other),
Ok(_actor) => Err(WalError::TreeError(
"expected MissingCommittedRoot when the dir barrier loses node files, \
but recovery succeeded"
.to_owned(),
)),
}
}
#[test]
fn durable_dir_barrier_lets_recovery_read_committed_value() -> Result<(), WalError> {
let dir = tempfile::tempdir()?;
let wal = tempfile::tempdir()?;
let wal_path = wal.path().join("shard.wal");
let nodes_dir = dir.path().join("nodes");
let mut store = CrashWindowStore::new(nodes_dir.clone(), false)?;
let (committed_root, cold, recovered) =
commit_then_recover(&mut store, &wal_path, &nodes_dir)?;
let actor = recovered?;
assert_eq!(actor.committed_root(), Some(committed_root));
assert_eq!(
actor.get(b"durable-key", &cold)?,
Some(b"durable-value".to_vec()),
"the committed value must be readable from disk after recovery"
);
Ok(())
}
#[test]
fn barrier_runs_strictly_before_the_wal_marker() -> Result<(), WalError> {
#[derive(Debug)]
struct OrderingStore {
inner: MemoryStore,
wal_path: PathBuf,
marker_present_at_barrier: RefCell<Option<bool>>,
}
impl NodeStore for OrderingStore {
type Error = std::convert::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> {
let present = DurableWal::read_file(&self.wal_path)
.ok()
.and_then(|contents| contents.committed_root())
.is_some();
*self.marker_present_at_barrier.borrow_mut() = Some(present);
Ok(())
}
}
let wal_dir = tempfile::tempdir()?;
let wal_path = wal_dir.path().join("shard.wal");
let mut store = OrderingStore {
inner: MemoryStore::new(),
wal_path: wal_path.clone(),
marker_present_at_barrier: RefCell::new(None),
};
let wal = DurableWal::new(&wal_path, FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
actor.put(b"k".to_vec(), b"v".to_vec())?;
actor.commit(&mut store)?;
let observed = *store.marker_present_at_barrier.borrow();
assert_eq!(
observed,
Some(false),
"the dir-sync barrier must run while the WAL marker is still absent \
(strictly before wal.commit)"
);
assert!(
DurableWal::read_file(&wal_path)?.committed_root().is_some(),
"the marker must be written by commit (after the barrier)"
);
Ok(())
}
}
#[cfg(test)]
mod promise_recovery_tests {
use super::{Ballot, RecordPromiseOutcome, ShardActor};
use crate::store::MemoryStore;
use crate::sync::topology::SyncNodeId;
use crate::wal::{DurableWal, FsyncPolicy, WalError, WalRecovery};
use std::path::{Path, PathBuf};
struct TempWal {
_dir: tempfile::TempDir,
path: PathBuf,
}
fn temp_wal() -> Result<TempWal, WalError> {
let dir = tempfile::tempdir()?;
let path = dir.path().join("shard.wal");
Ok(TempWal { _dir: dir, path })
}
fn ballot(counter: u64, node: &str) -> Ballot {
Ballot::new(counter, SyncNodeId::from(node))
}
fn reopen(path: &Path) -> Result<ShardActor, WalError> {
let store = MemoryStore::new();
let recovered = WalRecovery::recover_path(path, &store)?;
let wal = DurableWal::new(path, FsyncPolicy::CommitOnly)?;
ShardActor::from_recovered(wal, recovered, &store)
}
#[test]
fn promise_is_durable_and_monotonic_across_crash() -> Result<(), WalError> {
let temp = temp_wal()?;
{
let wal = DurableWal::new(&temp.path, FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
assert_eq!(
actor.record_promise(ballot(5, "X"))?,
RecordPromiseOutcome::Promised
);
drop(actor);
}
let mut recovered = reopen(&temp.path)?;
assert_eq!(
recovered.promised(),
&ballot(5, "X"),
"a returned record_promise must survive a crash"
);
let outcome = recovered.record_promise(ballot(3, "Y"))?;
assert_eq!(
outcome,
RecordPromiseOutcome::Rejected {
promised: ballot(5, "X")
},
"promised must never regress below a persisted ballot after restart"
);
assert_eq!(recovered.promised(), &ballot(5, "X"), "promised unchanged");
assert_eq!(
recovered.record_promise(ballot(6, "A"))?,
RecordPromiseOutcome::Promised
);
drop(recovered);
let again = reopen(&temp.path)?;
assert_eq!(again.promised(), &ballot(6, "A"), "higher ballot persisted");
Ok(())
}
#[test]
fn reserved_minted_counter_never_regresses_across_crash() -> Result<(), WalError> {
let temp = temp_wal()?;
{
let wal = DurableWal::new(&temp.path, FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
assert_eq!(actor.reserve_minted(7)?, 7);
drop(actor);
}
let mut recovered = reopen(&temp.path)?;
assert!(
recovered.persisted_max_minted() >= 7,
"reserved minted counter must survive a crash"
);
assert_eq!(
recovered.reserve_minted(4)?,
7,
"lower request keeps the floor"
);
assert_eq!(recovered.persisted_max_minted(), 7);
let next = recovered.persisted_max_minted() + 1;
assert_eq!(recovered.reserve_minted(next)?, 8);
assert!(
recovered.persisted_max_minted() >= next,
"next reserved counter must strictly exceed the prior persisted floor"
);
drop(recovered);
let again = reopen(&temp.path)?;
assert_eq!(
again.persisted_max_minted(),
8,
"advance persisted across crash"
);
Ok(())
}
#[test]
fn owner_epoch_is_durable_across_crash() -> Result<(), WalError> {
let temp = temp_wal()?;
{
let wal = DurableWal::new(&temp.path, FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
actor.record_owner_epoch(ballot(4, "owner"))?;
drop(actor);
}
let recovered = reopen(&temp.path)?;
assert_eq!(recovered.owner_epoch(), Some(&ballot(4, "owner")));
Ok(())
}
#[test]
fn all_three_values_co_persist_across_crash() -> Result<(), WalError> {
let temp = temp_wal()?;
{
let wal = DurableWal::new(&temp.path, FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
assert_eq!(
actor.record_promise(ballot(2, "P"))?,
RecordPromiseOutcome::Promised
);
actor.record_owner_epoch(ballot(2, "P"))?;
assert_eq!(actor.reserve_minted(5)?, 5);
drop(actor);
}
let recovered = reopen(&temp.path)?;
assert_eq!(recovered.promised(), &ballot(2, "P"));
assert_eq!(recovered.owner_epoch(), Some(&ballot(2, "P")));
assert_eq!(recovered.persisted_max_minted(), 5);
Ok(())
}
#[test]
fn promise_survives_commit_truncation_and_crash() -> Result<(), WalError> {
let temp = temp_wal()?;
let mut store = MemoryStore::new();
{
let wal = DurableWal::new(&temp.path, FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
assert_eq!(
actor.record_promise(ballot(9, "Z"))?,
RecordPromiseOutcome::Promised
);
actor.put(b"k".to_vec(), b"v".to_vec())?;
let _root = actor.commit(&mut store)?;
drop(actor);
}
let recovered = WalRecovery::recover_path(&temp.path, &store)?;
let wal = DurableWal::new(&temp.path, FsyncPolicy::CommitOnly)?;
let actor = ShardActor::from_recovered(wal, recovered, &store)?;
assert_eq!(
actor.promised(),
&ballot(9, "Z"),
"promise must survive a commit truncation + crash (re-emit after marker)"
);
Ok(())
}
}
#[cfg(test)]
mod group_commit_tests {
use super::{GroupOutcome, GroupWrite, ShardActor};
use crate::shard::actor::handle::ShardError;
use crate::store::{DiskStore, MemoryStore, NodeStore};
use crate::sync::ballot::{Ballot, Stamp};
use crate::sync::topology::SyncNodeId;
use crate::tree::{Hash, LeafNode, Node};
use crate::wal::{DurableWal, FsyncPolicy, WalError, WalRecovery};
use std::cell::Cell;
use std::path::{Path, PathBuf};
use std::sync::Arc;
#[derive(Debug)]
struct CommitCountingStore {
inner: MemoryStore,
commits: Cell<usize>,
}
impl CommitCountingStore {
fn new() -> Self {
Self {
inner: MemoryStore::new(),
commits: Cell::new(0),
}
}
fn commit_count(&self) -> usize {
self.commits.get()
}
}
impl NodeStore for CommitCountingStore {
type Error = std::convert::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.commits.set(self.commits.get().saturating_add(1));
Ok(())
}
}
fn stamp(counter: u64, node: &str, seq: u64) -> Stamp {
Stamp::new(Ballot::new(counter, SyncNodeId::new(node)), seq)
}
fn apply_value(key: &[u8], expected: Option<Hash>, value: &[u8], stamp: Stamp) -> GroupWrite {
GroupWrite::ApplyValue {
key: key.to_vec(),
expected,
value: value.to_vec(),
ttl: None,
stamp,
}
}
#[test]
fn group_of_writes_coalesces_into_one_commit_all_readable() -> Result<(), WalError> {
let dir = tempfile::tempdir()?;
let wal = DurableWal::new(dir.path().join("group.wal"), FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
let mut store = CommitCountingStore::new();
let writes = vec![
apply_value(b"k1", None, b"first", stamp(1, "owner", 0)),
apply_value(b"k2", None, b"v2", stamp(1, "owner", 1)),
GroupWrite::Cas {
key: b"counter".to_vec(),
expected: None,
new: 7,
},
apply_value(
b"k1",
Some(Hash::of(b"first")),
b"second",
stamp(1, "owner", 2),
),
];
let outcomes = actor.apply_group(writes, &mut store);
assert_eq!(outcomes.len(), 4);
for outcome in &outcomes {
assert!(
matches!(outcome, GroupOutcome::Committed),
"every write in a clean group must commit, got {outcome:?}"
);
}
assert_eq!(
store.commit_count(),
1,
"N grouped writes must produce exactly ONE commit/fsync"
);
assert_eq!(actor.get(b"k1", &store)?, Some(b"second".to_vec()));
assert_eq!(actor.get(b"k2", &store)?, Some(b"v2".to_vec()));
assert_eq!(actor.read_value(b"counter", &store)?, Some(7));
assert!(
actor.committed_root().is_some(),
"the group committed a root"
);
assert!(actor.buffer().is_empty(), "commit cleared the buffer");
Ok(())
}
#[test]
fn partial_failure_commits_survivors_and_leaves_failed_key_unchanged() -> Result<(), WalError> {
let dir = tempfile::tempdir()?;
let wal = DurableWal::new(dir.path().join("partial.wal"), FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
let mut store = CommitCountingStore::new();
let seeded = actor.apply_group(
vec![apply_value(b"mid", None, b"seeded", stamp(1, "owner", 0))],
&mut store,
);
assert!(matches!(seeded.as_slice(), [GroupOutcome::Committed]));
let commits_after_seed = store.commit_count();
let writes = vec![
apply_value(b"ok1", None, b"a", stamp(2, "owner", 0)),
apply_value(b"mid", None, b"should-not-apply", stamp(2, "owner", 1)),
apply_value(b"ok2", None, b"b", stamp(2, "owner", 2)),
];
let outcomes = actor.apply_group(writes, &mut store);
assert!(
matches!(outcomes[0], GroupOutcome::Committed),
"survivor ok1 must commit, got {:?}",
outcomes[0]
);
assert!(
matches!(
&outcomes[1],
GroupOutcome::Rejected(ShardError::CasHashMismatch { expected, actual })
if expected.is_none() && *actual == Some(Hash::of(b"seeded"))
),
"only the failed CAS write is Rejected, got {:?}",
outcomes[1]
);
assert!(
matches!(outcomes[2], GroupOutcome::Committed),
"survivor ok2 must commit, got {:?}",
outcomes[2]
);
assert_eq!(
store.commit_count(),
commits_after_seed.saturating_add(1),
"the survivors share exactly one group commit"
);
assert_eq!(actor.get(b"ok1", &store)?, Some(b"a".to_vec()));
assert_eq!(actor.get(b"ok2", &store)?, Some(b"b".to_vec()));
assert_eq!(
actor.get(b"mid", &store)?,
Some(b"seeded".to_vec()),
"a failed CAS must leave the buffer as if that write never happened"
);
Ok(())
}
#[test]
fn all_writes_fail_means_no_commit_at_all() -> Result<(), WalError> {
let dir = tempfile::tempdir()?;
let wal = DurableWal::new(dir.path().join("allfail.wal"), FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
let mut store = CommitCountingStore::new();
actor.apply_group(
vec![
apply_value(b"a", None, b"x", stamp(1, "owner", 0)),
apply_value(b"b", None, b"y", stamp(1, "owner", 1)),
],
&mut store,
);
let commits_after_seed = store.commit_count();
let outcomes = actor.apply_group(
vec![
apply_value(b"a", None, b"nope", stamp(2, "owner", 0)),
apply_value(b"b", None, b"nope", stamp(2, "owner", 1)),
],
&mut store,
);
assert!(outcomes.iter().all(|o| matches!(
o,
GroupOutcome::Rejected(ShardError::CasHashMismatch { .. })
)));
assert_eq!(
store.commit_count(),
commits_after_seed,
"an all-rejected group must not commit (no fsync)"
);
Ok(())
}
#[derive(Debug)]
struct CrashWindowStore {
inner: DiskStore,
lossy: bool,
pending: std::cell::RefCell<Vec<PathBuf>>,
dir: PathBuf,
}
impl CrashWindowStore {
fn new(dir: PathBuf, lossy: bool) -> Result<Self, WalError> {
let inner =
DiskStore::new(&dir).map_err(|error| WalError::TreeError(error.to_string()))?;
Ok(Self {
inner,
lossy,
pending: std::cell::RefCell::new(Vec::new()),
dir,
})
}
fn node_path(&self, hash: &Hash) -> PathBuf {
let hex = hash.to_string();
let (prefix, file_name) = hex.split_at(2);
self.dir.join(prefix).join(file_name)
}
}
impl NodeStore for CrashWindowStore {
type Error = crate::store::StoreError;
fn get(&self, hash: &Hash) -> Result<Option<Arc<Node>>, Self::Error> {
self.inner.get(hash)
}
fn put(&mut self, node: &Node) -> Result<Hash, Self::Error> {
let hash = self.inner.put(node)?;
self.pending.borrow_mut().push(self.node_path(&hash));
Ok(hash)
}
fn sync_dirty_dirs(&self) -> Result<(), Self::Error> {
let pending = std::mem::take(&mut *self.pending.borrow_mut());
if self.lossy {
for path in pending {
match std::fs::remove_file(&path) {
Ok(()) => {}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(error) => return Err(crate::store::StoreError::Io(error)),
}
}
Ok(())
} else {
self.inner.sync_dirty_dirs()
}
}
}
fn empty_root(store: &mut CrashWindowStore) -> Result<Hash, WalError> {
let leaf =
LeafNode::new(Vec::new()).map_err(|error| WalError::TreeError(error.to_string()))?;
store
.put(&Node::Leaf(leaf))
.map_err(|error| WalError::TreeError(error.to_string()))
}
fn group_then_recover(
store: &mut CrashWindowStore,
wal_path: &Path,
nodes_dir: &Path,
) -> Result<(DiskStore, Result<ShardActor, WalError>), WalError> {
let wal = DurableWal::new(wal_path, FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
let outcomes = actor.apply_group(
vec![
apply_value(b"g1", None, b"v1", stamp(1, "owner", 0)),
apply_value(b"g2", None, b"v2", stamp(1, "owner", 1)),
apply_value(b"g3", None, b"v3", stamp(1, "owner", 2)),
],
store,
);
assert!(
outcomes
.iter()
.all(|o| matches!(o, GroupOutcome::Committed)),
"the group must have committed before the crash"
);
drop(actor);
let cold = DiskStore::new(nodes_dir).map_err(|e| WalError::TreeError(e.to_string()))?;
let actor = match WalRecovery::recover_path(wal_path, &cold) {
Ok(recovered) => {
let wal = DurableWal::new(wal_path, FsyncPolicy::CommitOnly)?;
ShardActor::from_recovered(wal, recovered, &cold)
}
Err(error) => Err(error),
};
Ok((cold, actor))
}
#[test]
fn group_commit_with_lossy_barrier_is_rejected_never_partial() -> Result<(), WalError> {
let dir = tempfile::tempdir()?;
let wal = tempfile::tempdir()?;
let wal_path = wal.path().join("group.wal");
let nodes_dir = dir.path().join("nodes");
let mut store = CrashWindowStore::new(nodes_dir.clone(), true)?;
empty_root(&mut store)?;
let (_cold, recovered) = group_then_recover(&mut store, &wal_path, &nodes_dir)?;
match recovered {
Err(WalError::MissingCommittedRoot { .. }) => Ok(()),
Err(other) => Err(other),
Ok(_actor) => Err(WalError::TreeError(
"expected MissingCommittedRoot when the group's node dir entries are lost, \
but recovery succeeded"
.to_owned(),
)),
}
}
#[test]
fn group_commit_with_durable_barrier_recovers_all_survivors() -> Result<(), WalError> {
let dir = tempfile::tempdir()?;
let wal = tempfile::tempdir()?;
let wal_path = wal.path().join("group.wal");
let nodes_dir = dir.path().join("nodes");
let mut store = CrashWindowStore::new(nodes_dir.clone(), false)?;
empty_root(&mut store)?;
let (cold, recovered) = group_then_recover(&mut store, &wal_path, &nodes_dir)?;
let actor = recovered?;
assert_eq!(actor.get(b"g1", &cold)?, Some(b"v1".to_vec()));
assert_eq!(actor.get(b"g2", &cold)?, Some(b"v2".to_vec()));
assert_eq!(actor.get(b"g3", &cold)?, Some(b"v3".to_vec()));
Ok(())
}
#[test]
fn group_commit_failure_rolls_back_all_survivors() -> Result<(), WalError> {
#[derive(Debug)]
struct FailingStore {
inner: MemoryStore,
puts: Cell<usize>,
fail_after: usize,
}
impl NodeStore for FailingStore {
type Error = crate::store::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> {
let count = self.puts.get().saturating_add(1);
self.puts.set(count);
if count > self.fail_after {
return Err(crate::store::StoreError::Io(std::io::Error::other(
"injected commit failure",
)));
}
Ok(self.inner.put(node))
}
}
let dir = tempfile::tempdir()?;
let wal = DurableWal::new(dir.path().join("fail.wal"), FsyncPolicy::CommitOnly)?;
let mut actor = ShardActor::new(wal);
let mut store = FailingStore {
inner: MemoryStore::new(),
puts: Cell::new(0),
fail_after: 0,
};
let outcomes = actor.apply_group(
vec![
apply_value(b"x", None, b"vx", stamp(1, "owner", 0)),
apply_value(b"y", None, b"vy", stamp(1, "owner", 1)),
],
&mut store,
);
assert_eq!(outcomes.len(), 2);
for outcome in &outcomes {
assert!(
matches!(outcome, GroupOutcome::CommitFailed(_)),
"a failed group commit must tell every survivor CommitFailed, got {outcome:?}"
);
}
assert!(
actor.buffer().is_empty(),
"a failed group commit rolls back every staged survivor's key"
);
assert_eq!(actor.committed_root(), None, "nothing was committed");
Ok(())
}
}