use std::error::Error;
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use beamr::module::ModuleRegistry;
use beamr::process::ExitReason;
use beamr::scheduler::{Scheduler, SchedulerConfig};
use super::RecordPromiseOutcome;
use super::handle::{RangeItem, ShardError, ShardHandle};
use crate::store::DiskStore;
use crate::sync::SyncNodeId;
use crate::sync::ballot::{Ballot, Stamp};
use crate::tree::{Hash, LeafNode, Node};
use crate::wal::DurableWal;
const TIMEOUT: Duration = Duration::from_secs(5);
type RangeEntries = Vec<(Vec<u8>, Vec<u8>)>;
fn test_scheduler() -> Result<Arc<Scheduler>, Box<dyn Error>> {
let scheduler = Scheduler::new(
SchedulerConfig {
thread_count: Some(1),
..SchedulerConfig::default()
},
Arc::new(ModuleRegistry::new()),
)
.map_err(|message| -> Box<dyn Error> { message.into() })?;
Ok(Arc::new(scheduler))
}
struct TestShard {
_dir: tempfile::TempDir,
store_dir: PathBuf,
wal_path: PathBuf,
handle: ShardHandle,
}
impl TestShard {
fn spawn(scheduler: &Arc<Scheduler>, name: &str) -> Result<Self, Box<dyn Error>> {
let dir = tempfile::tempdir()?;
let store_dir = dir.path().join(format!("{name}.store"));
let wal_path = dir.path().join(format!("{name}.wal"));
let mut store = DiskStore::new(&store_dir)?;
let _root = empty_root(&mut store)?;
drop(store);
let handle =
ShardHandle::spawn(Arc::clone(scheduler), store_dir.clone(), wal_path.clone())?;
Ok(Self {
_dir: dir,
store_dir,
wal_path,
handle,
})
}
fn respawn(&self, scheduler: &Arc<Scheduler>) -> Result<ShardHandle, Box<dyn Error>> {
Ok(ShardHandle::spawn(
Arc::clone(scheduler),
self.store_dir.clone(),
self.wal_path.clone(),
)?)
}
}
fn empty_root(store: &mut DiskStore) -> Result<Hash, Box<dyn Error>> {
let leaf = LeafNode::new(Vec::new())?;
Ok(store.put(&Node::Leaf(leaf))?)
}
fn put(handle: &ShardHandle, key: &[u8], value: &[u8]) -> Result<(), Box<dyn Error>> {
handle.put_with_ttl(key.to_vec(), value.to_vec(), None, TIMEOUT)?;
Ok(())
}
fn delete(handle: &ShardHandle, key: &[u8]) -> Result<(), Box<dyn Error>> {
handle.delete(key.to_vec(), crate::sync::ballot::Stamp::bottom(), TIMEOUT)?;
Ok(())
}
fn get(handle: &ShardHandle, key: &[u8]) -> Result<Option<Vec<u8>>, Box<dyn Error>> {
Ok(handle.get(key.to_vec(), TIMEOUT)?)
}
fn commit(handle: &ShardHandle) -> Result<Hash, Box<dyn Error>> {
Ok(handle.commit(TIMEOUT)?)
}
fn range(handle: &ShardHandle, from: &[u8], to: &[u8]) -> Result<RangeEntries, Box<dyn Error>> {
let items = handle.range(from.to_vec(), to.to_vec(), TIMEOUT)?;
let mut entries = Vec::new();
let mut saw_done = false;
for item in items {
match item {
RangeItem::Entry { key, value } => entries.push((key, value)),
RangeItem::Done => saw_done = true,
}
}
assert!(saw_done, "range result must terminate with Done");
Ok(entries)
}
#[test]
fn get_and_range_merge_tree_with_buffer_shadowing() -> Result<(), Box<dyn Error>> {
let scheduler = test_scheduler()?;
let shard = TestShard::spawn(&scheduler, "merge")?;
let handle = &shard.handle;
put(handle, b"a", b"tree-a")?;
put(handle, b"b", b"tree-b")?;
put(handle, b"d", b"tree-d")?;
let committed_root = commit(handle)?;
put(handle, b"b", b"buffer-b")?;
put(handle, b"c", b"buffer-c")?;
delete(handle, b"d")?;
assert_eq!(get(handle, b"b")?, Some(b"buffer-b".to_vec()));
assert_eq!(get(handle, b"a")?, Some(b"tree-a".to_vec()));
assert_eq!(get(handle, b"unknown")?, None);
assert_eq!(get(handle, b"d")?, None);
assert_eq!(
range(handle, b"a", b"e")?,
vec![
(b"a".to_vec(), b"tree-a".to_vec()),
(b"b".to_vec(), b"buffer-b".to_vec()),
(b"c".to_vec(), b"buffer-c".to_vec()),
]
);
let _ = committed_root;
scheduler.shutdown();
Ok(())
}
#[test]
fn put_and_delete_ack_after_wal_append_without_tree_mutation() -> Result<(), Box<dyn Error>> {
let scheduler = test_scheduler()?;
let shard = TestShard::spawn(&scheduler, "durable-first")?;
let handle = &shard.handle;
put(handle, b"event", b"payload")?;
assert_eq!(
DurableWal::read_file(&shard.wal_path)?.entries(),
&[crate::wal::WalEntry::put(
b"event".to_vec(),
b"payload".to_vec()
)]
);
assert_eq!(get(handle, b"event")?, Some(b"payload".to_vec()));
assert_eq!(
DurableWal::read_file(&shard.wal_path)?.committed_root(),
None
);
delete(handle, b"event")?;
assert_eq!(get(handle, b"event")?, None);
let tombstone =
crate::ttl::entry::encode_stamped_tombstone(crate::sync::ballot::Stamp::bottom());
assert_eq!(
DurableWal::read_file(&shard.wal_path)?.entries(),
&[
crate::wal::WalEntry::put(b"event".to_vec(), b"payload".to_vec()),
crate::wal::WalEntry::put(b"event".to_vec(), tombstone),
]
);
scheduler.shutdown();
Ok(())
}
#[test]
fn commit_persists_marker_and_is_idempotent() -> Result<(), Box<dyn Error>> {
let scheduler = test_scheduler()?;
let shard = TestShard::spawn(&scheduler, "commit")?;
let handle = &shard.handle;
put(handle, b"b", b"two")?;
put(handle, b"a", b"one")?;
put(handle, b"c", b"three")?;
let committed = commit(handle)?;
assert_eq!(
DurableWal::read_file(&shard.wal_path)?.committed_root(),
Some(committed)
);
assert_eq!(get(handle, b"a")?, Some(b"one".to_vec()));
assert_eq!(commit(handle)?, committed);
scheduler.shutdown();
Ok(())
}
#[test]
fn commit_root_is_history_independent_for_same_key_value_set() -> Result<(), Box<dyn Error>> {
let scheduler = test_scheduler()?;
let first = TestShard::spawn(&scheduler, "deterministic-a")?;
let second = TestShard::spawn(&scheduler, "deterministic-b")?;
put(&first.handle, b"alpha", b"1")?;
put(&first.handle, b"beta", b"2")?;
let first_root = commit(&first.handle)?;
put(&second.handle, b"beta", b"2")?;
put(&second.handle, b"alpha", b"1")?;
let second_root = commit(&second.handle)?;
assert_eq!(first_root, second_root);
scheduler.shutdown();
Ok(())
}
#[test]
fn respawn_replays_wal_and_leaves_sibling_running() -> Result<(), Box<dyn Error>> {
let scheduler = test_scheduler()?;
let failed = TestShard::spawn(&scheduler, "failed")?;
let sibling = TestShard::spawn(&scheduler, "sibling")?;
put(&sibling.handle, b"sibling-key", b"sibling-value")?;
put(&failed.handle, b"committed", b"tree-value")?;
let committed_root = commit(&failed.handle)?;
put(&failed.handle, b"buffered", b"wal-value")?;
assert_eq!(
DurableWal::read_file(&failed.wal_path)?.committed_root(),
Some(committed_root)
);
scheduler.exit_signal(0, failed.handle.pid(), ExitReason::Kill)?;
let recovered = failed.respawn(&scheduler)?;
assert_ne!(recovered.pid(), failed.handle.pid());
assert_eq!(
get(&sibling.handle, b"sibling-key")?,
Some(b"sibling-value".to_vec())
);
assert_eq!(get(&recovered, b"buffered")?, Some(b"wal-value".to_vec()));
assert_eq!(get(&recovered, b"committed")?, Some(b"tree-value".to_vec()));
scheduler.shutdown();
Ok(())
}
#[test]
fn boot_failure_keeps_scheduler_usable_and_fails_the_command() -> Result<(), Box<dyn Error>> {
let scheduler = test_scheduler()?;
let dir = tempfile::tempdir()?;
let store_dir = dir.path().join("not-a-dir.store");
std::fs::write(&store_dir, b"i am a file, not a directory")?;
let wal_path = dir.path().join("boot-fail.wal");
let handle = ShardHandle::spawn(Arc::clone(&scheduler), store_dir, wal_path)?;
let result = handle.get(b"any-key".to_vec(), TIMEOUT);
assert!(
result.is_err(),
"a command against a boot-failed shard must error, got Ok: {result:?}"
);
assert!(
matches!(
&result,
Err(ShardError::ReplyTimeout { .. }
| ShardError::ActorUnavailable { .. }
| ShardError::Spawn(_))
),
"boot-failure command should fail as ReplyTimeout / ActorUnavailable / \
Spawn (never ReplyDisconnected or a storage error), got {result:?}"
);
let healthy = TestShard::spawn(&scheduler, "healthy-after-boot-fail")?;
put(&healthy.handle, b"k", b"v")?;
assert_eq!(get(&healthy.handle, b"k")?, Some(b"v".to_vec()));
scheduler.shutdown();
Ok(())
}
fn ballot(counter: u64, node: &str) -> Ballot {
Ballot::new(counter, SyncNodeId::new(node))
}
fn promise(handle: &ShardHandle, ballot: Ballot) -> Result<(), Box<dyn Error>> {
match handle.record_promise(ballot, TIMEOUT)? {
RecordPromiseOutcome::Promised => Ok(()),
RecordPromiseOutcome::Rejected { promised } => {
Err(format!("expected Promised, got Rejected({promised:?})").into())
}
}
}
#[test]
fn fence_rejects_stale_epoch_write_applying_nothing() -> Result<(), Box<dyn Error>> {
let scheduler = test_scheduler()?;
let shard = TestShard::spawn(&scheduler, "fence-stale")?;
let handle = &shard.handle;
promise(handle, ballot(5, "X"))?;
let result = handle.apply_durable(
b"k".to_vec(),
None,
b"stale".to_vec(),
None,
Stamp::new(ballot(3, "Y"), 0),
TIMEOUT,
);
assert!(
matches!(
result,
Err(ShardError::Fenced {
ref promised,
ref attempted,
}) if *promised == ballot(5, "X") && *attempted == ballot(3, "Y")
),
"stale-epoch write must be Fenced{{promised:(5,X), attempted:(3,Y)}}, got {result:?}"
);
assert_eq!(get(handle, b"k")?, None, "fenced write must apply nothing");
assert_eq!(
DurableWal::read_file(&shard.wal_path)?.committed_root(),
None,
"a fenced write must not have committed anything"
);
scheduler.shutdown();
Ok(())
}
#[test]
fn fence_admits_ge_without_ever_raising_promised() -> Result<(), Box<dyn Error>> {
let scheduler = test_scheduler()?;
let shard = TestShard::spawn(&scheduler, "fence-r2")?;
let handle = &shard.handle;
promise(handle, ballot(5, "X"))?;
handle.apply_durable(
b"k".to_vec(),
None,
b"at-five".to_vec(),
None,
Stamp::new(ballot(5, "X"), 0),
TIMEOUT,
)?;
assert_eq!(get(handle, b"k")?, Some(b"at-five".to_vec()));
assert_eq!(
handle.read_promise_state(TIMEOUT)?.promised,
ballot(5, "X"),
"an admitted data write must NOT raise promised (R2)"
);
let current = Hash::of(b"at-five");
handle.apply_durable(
b"k".to_vec(),
Some(current),
b"at-seven".to_vec(),
None,
Stamp::new(ballot(7, "Z"), 0),
TIMEOUT,
)?;
assert_eq!(get(handle, b"k")?, Some(b"at-seven".to_vec()));
assert_eq!(
handle.read_promise_state(TIMEOUT)?.promised,
ballot(5, "X"),
"a higher-epoch data write STILL must not raise promised (R2)"
);
match handle.record_promise(ballot(6, "W"), TIMEOUT)? {
RecordPromiseOutcome::Promised => {}
RecordPromiseOutcome::Rejected { promised } => {
return Err(format!(
"record_promise((6,W)) was Rejected({promised:?}) — promised was silently \
raised by a data write, violating R2"
)
.into());
}
}
scheduler.shutdown();
Ok(())
}
fn export_committed(
scheduler: &Arc<Scheduler>,
name: &str,
key: &[u8],
value: &[u8],
stamp: Stamp,
) -> Result<(Option<Hash>, Vec<crate::sync::NodeTransfer>), Box<dyn Error>> {
let shard = TestShard::spawn(scheduler, name)?;
shard
.handle
.apply_durable(key.to_vec(), None, value.to_vec(), None, stamp, TIMEOUT)?;
let export = shard.handle.export_reachable(0, TIMEOUT)?;
Ok(export)
}
#[test]
fn merge_adopt_unions_forked_promiser_state() -> Result<(), Box<dyn Error>> {
let scheduler = test_scheduler()?;
let b = export_committed(
&scheduler,
"fork-b",
b"k3",
b"v3",
Stamp::new(ballot(2, "A"), 1),
)?;
let c = export_committed(
&scheduler,
"fork-c",
b"k2",
b"v2",
Stamp::new(ballot(2, "A"), 0),
)?;
let target = TestShard::spawn(&scheduler, "fork-target")?;
target.handle.merge_adopt(vec![b, c.clone()], TIMEOUT)?;
assert_eq!(
get(&target.handle, b"k2")?,
Some(b"v2".to_vec()),
"merge over ALL promisers must serve k2 (from C)"
);
assert_eq!(
get(&target.handle, b"k3")?,
Some(b"v3".to_vec()),
"merge over ALL promisers must serve k3 (from B)"
);
let single = TestShard::spawn(&scheduler, "fork-single")?;
single.handle.merge_adopt(vec![c], TIMEOUT)?;
assert_eq!(
get(&single.handle, b"k2")?,
Some(b"v2".to_vec()),
"single-root adoption serves the one promiser's key"
);
assert_eq!(
get(&single.handle, b"k3")?,
None,
"single-root adoption DROPS the forked write — proves the union is load-bearing"
);
scheduler.shutdown();
Ok(())
}
#[test]
fn merge_adopt_root_survives_crash_recovery() -> Result<(), Box<dyn Error>> {
let scheduler = test_scheduler()?;
let b = export_committed(
&scheduler,
"durable-b",
b"k3",
b"v3",
Stamp::new(ballot(2, "A"), 1),
)?;
let c = export_committed(
&scheduler,
"durable-c",
b"k2",
b"v2",
Stamp::new(ballot(2, "A"), 0),
)?;
let target = TestShard::spawn(&scheduler, "durable-target")?;
target.handle.merge_adopt(vec![b, c], TIMEOUT)?;
let recovered = target.respawn(&scheduler)?;
assert_eq!(
get(&recovered, b"k2")?,
Some(b"v2".to_vec()),
"merged k2 must survive crash recovery (fsync'd marker)"
);
assert_eq!(
get(&recovered, b"k3")?,
Some(b"v3".to_vec()),
"merged k3 must survive crash recovery"
);
scheduler.shutdown();
Ok(())
}
fn stamp_of(handle: &ShardHandle, key: &[u8]) -> Result<Option<Stamp>, Box<dyn Error>> {
let Some(raw) = handle.get_raw(key.to_vec(), TIMEOUT)? else {
return Ok(None);
};
let entry = crate::ttl::entry::StampedEntry::decode(&raw)?
.ok_or("committed value is not a stamped envelope")?;
Ok(Some(entry.stamp().clone()))
}
#[test]
fn batch_applies_all_keys_at_one_stamp_in_one_commit() -> Result<(), Box<dyn Error>> {
let scheduler = test_scheduler()?;
let shard = TestShard::spawn(&scheduler, "batch-atomic")?;
let handle = &shard.handle;
let batch_stamp = Stamp::new(ballot(4, "owner"), 7);
handle.apply_durable_batch(
vec![
(b"k1".to_vec(), None, b"v1".to_vec(), None),
(b"k2".to_vec(), None, b"v2".to_vec(), None),
(b"k3".to_vec(), None, b"v3".to_vec(), None),
],
batch_stamp.clone(),
TIMEOUT,
)?;
assert_eq!(get(handle, b"k1")?, Some(b"v1".to_vec()));
assert_eq!(get(handle, b"k2")?, Some(b"v2".to_vec()));
assert_eq!(get(handle, b"k3")?, Some(b"v3".to_vec()));
assert_eq!(stamp_of(handle, b"k1")?, Some(batch_stamp.clone()));
assert_eq!(stamp_of(handle, b"k2")?, Some(batch_stamp.clone()));
assert_eq!(stamp_of(handle, b"k3")?, Some(batch_stamp));
assert!(
DurableWal::read_file(&shard.wal_path)?
.committed_root()
.is_some(),
"a successful batch must have committed a root marker"
);
scheduler.shutdown();
Ok(())
}
#[test]
fn batch_fence_rejects_whole_batch_writing_nothing() -> Result<(), Box<dyn Error>> {
let scheduler = test_scheduler()?;
let shard = TestShard::spawn(&scheduler, "batch-fence")?;
let handle = &shard.handle;
promise(handle, ballot(5, "X"))?;
let result = handle.apply_durable_batch(
vec![
(b"k1".to_vec(), None, b"v1".to_vec(), None),
(b"k2".to_vec(), None, b"v2".to_vec(), None),
(b"k3".to_vec(), None, b"v3".to_vec(), None),
],
Stamp::new(ballot(3, "Y"), 0),
TIMEOUT,
);
assert!(
matches!(
result,
Err(ShardError::Fenced { ref promised, ref attempted })
if *promised == ballot(5, "X") && *attempted == ballot(3, "Y")
),
"a below-promised batch must be Fenced, got {result:?}"
);
assert_eq!(get(handle, b"k1")?, None, "fenced batch must write nothing");
assert_eq!(get(handle, b"k2")?, None, "fenced batch must write nothing");
assert_eq!(get(handle, b"k3")?, None, "fenced batch must write nothing");
assert_eq!(
DurableWal::read_file(&shard.wal_path)?.committed_root(),
None,
"a fenced batch must not have committed anything"
);
scheduler.shutdown();
Ok(())
}
#[test]
fn batch_cas_mismatch_rejects_whole_batch_no_partial_apply() -> Result<(), Box<dyn Error>> {
let scheduler = test_scheduler()?;
let shard = TestShard::spawn(&scheduler, "batch-cas")?;
let handle = &shard.handle;
handle.apply_durable(
b"k2".to_vec(),
None,
b"already-here".to_vec(),
None,
Stamp::new(ballot(1, "owner"), 0),
TIMEOUT,
)?;
let seeded = stamp_of(handle, b"k2")?;
let result = handle.apply_durable_batch(
vec![
(b"k1".to_vec(), None, b"v1".to_vec(), None),
(b"k2".to_vec(), None, b"v2".to_vec(), None),
(b"k3".to_vec(), None, b"v3".to_vec(), None),
],
Stamp::new(ballot(2, "owner"), 0),
TIMEOUT,
);
assert!(
matches!(
result,
Err(ShardError::CasHashMismatch { expected, actual })
if expected.is_none() && actual == Some(Hash::of(b"already-here"))
),
"a per-key CAS mismatch must reject the whole batch, got {result:?}"
);
assert_eq!(
get(handle, b"k1")?,
None,
"k1 must NOT be applied (no partial write)"
);
assert_eq!(
get(handle, b"k3")?,
None,
"k3 must NOT be applied (no partial write)"
);
assert_eq!(get(handle, b"k2")?, Some(b"already-here".to_vec()));
assert_eq!(
stamp_of(handle, b"k2")?,
seeded,
"k2 must be untouched by the rejected batch"
);
scheduler.shutdown();
Ok(())
}
#[test]
fn batch_survives_crash_recovery_with_shared_stamp() -> Result<(), Box<dyn Error>> {
let scheduler = test_scheduler()?;
let shard = TestShard::spawn(&scheduler, "batch-crash")?;
let batch_stamp = Stamp::new(ballot(6, "owner"), 2);
shard.handle.apply_durable_batch(
vec![
(b"k1".to_vec(), None, b"v1".to_vec(), None),
(b"k2".to_vec(), None, b"v2".to_vec(), None),
(b"k3".to_vec(), None, b"v3".to_vec(), None),
],
batch_stamp.clone(),
TIMEOUT,
)?;
scheduler.exit_signal(0, shard.handle.pid(), ExitReason::Kill)?;
let recovered = shard.respawn(&scheduler)?;
assert_ne!(recovered.pid(), shard.handle.pid());
assert_eq!(get(&recovered, b"k1")?, Some(b"v1".to_vec()));
assert_eq!(get(&recovered, b"k2")?, Some(b"v2".to_vec()));
assert_eq!(get(&recovered, b"k3")?, Some(b"v3".to_vec()));
assert_eq!(stamp_of(&recovered, b"k1")?, Some(batch_stamp.clone()));
assert_eq!(stamp_of(&recovered, b"k2")?, Some(batch_stamp.clone()));
assert_eq!(stamp_of(&recovered, b"k3")?, Some(batch_stamp));
scheduler.shutdown();
Ok(())
}