use std::error::Error;
use crate::db::{Database, DatabaseConfig};
use crate::shard::commit_state::CommitSnapshot;
type BoxError = Box<dyn Error>;
fn single_shard_db() -> Result<(tempfile::TempDir, Database), BoxError> {
let dir = tempfile::tempdir()?;
let db = Database::create(DatabaseConfig {
data_dir: dir.path().join("db"),
shard_count: 1,
distributed: None,
executor_threads: None,
})?;
Ok((dir, db))
}
fn is_dirty(snapshot: &CommitSnapshot) -> bool {
snapshot.recovering
|| snapshot.unreconciled
|| snapshot.committed_gen != snapshot.dirty_gen
|| snapshot.root.is_none()
}
#[test]
fn put_dirties_the_shard_and_commit_cleans_it() -> Result<(), BoxError> {
let (_dir, db) = single_shard_db()?;
db.put(b"k".to_vec(), b"v".to_vec())?;
let dirty = db.commit_state_snapshot(0);
assert!(is_dirty(&dirty), "a buffered put must classify dirty");
db.commit()?;
let clean = db.commit_state_snapshot(0);
assert!(!is_dirty(&clean), "commit must classify clean");
assert!(clean.root.is_some(), "commit publishes the root");
assert_eq!(clean.committed_gen, clean.dirty_gen);
Ok(())
}
#[test]
fn delete_dirties_the_shard() -> Result<(), BoxError> {
let (_dir, db) = single_shard_db()?;
db.put(b"k".to_vec(), b"v".to_vec())?;
db.commit()?;
assert!(!is_dirty(&db.commit_state_snapshot(0)));
db.delete(b"k".to_vec())?;
assert!(
is_dirty(&db.commit_state_snapshot(0)),
"a buffered tombstone must classify dirty"
);
Ok(())
}
#[test]
fn append_advances_and_leaves_clean() -> Result<(), BoxError> {
let (_dir, db) = single_shard_db()?;
db.append(b"stream".to_vec(), vec![b"e".to_vec()], 0)?;
let snapshot = db.commit_state_snapshot(0);
assert!(!is_dirty(&snapshot), "append commits => clean");
assert!(snapshot.root.is_some());
assert_eq!(snapshot.advance_gen, 1, "append advanced the root once");
Ok(())
}
#[test]
fn cas_advances_and_leaves_clean() -> Result<(), BoxError> {
let (_dir, db) = single_shard_db()?;
db.cas(b"c".to_vec(), None, 1)?;
let snapshot = db.commit_state_snapshot(0);
assert!(!is_dirty(&snapshot), "cas commits => clean");
assert_eq!(snapshot.advance_gen, 1);
Ok(())
}
#[test]
fn empty_commit_stays_clean_without_advancing() -> Result<(), BoxError> {
let (_dir, db) = single_shard_db()?;
db.put(b"k".to_vec(), b"v".to_vec())?;
db.commit()?;
let advance_after_first = db.commit_state_snapshot(0).advance_gen;
db.commit()?; let snapshot = db.commit_state_snapshot(0);
assert!(!is_dirty(&snapshot));
assert_eq!(
snapshot.advance_gen, advance_after_first,
"an empty commit does not bump advance_gen"
);
Ok(())
}
#[test]
fn materialise_runs_restart_protocol_once() -> Result<(), BoxError> {
let (_dir, db) = single_shard_db()?;
db.put(b"k".to_vec(), b"v".to_vec())?; assert_eq!(
db.commit_state_snapshot(0).incarnation,
1,
"first spawn bumps the incarnation to 1"
);
Ok(())
}
#[test]
fn recovery_replay_reopens_dirty_then_commits() -> Result<(), BoxError> {
let dir = tempfile::tempdir()?;
let data_dir = dir.path().join("db");
{
let db = Database::create(DatabaseConfig {
data_dir: data_dir.clone(),
shard_count: 1,
distributed: None,
executor_threads: None,
})?;
db.put(b"k".to_vec(), b"v".to_vec())?; }
let reopened = Database::open(&data_dir)?;
assert_eq!(reopened.get(b"k")?, Some(b"v".to_vec()));
let recovered = reopened.commit_state_snapshot(0);
assert!(
is_dirty(&recovered),
"a recovered non-empty buffer must classify dirty"
);
assert_eq!(recovered.incarnation, 1, "the reopened actor spawned once");
reopened.commit()?;
assert!(
!is_dirty(&reopened.commit_state_snapshot(0)),
"committing the recovered buffer cleans the shard"
);
Ok(())
}
#[test]
fn cas_rejection_from_clean_stays_clean() -> Result<(), BoxError> {
let (_dir, db) = single_shard_db()?;
db.cas(b"c".to_vec(), None, 1)?; assert!(!is_dirty(&db.commit_state_snapshot(0)), "clean baseline");
assert!(
db.cas(b"c".to_vec(), Some(99), 2).is_err(),
"CAS is rejected"
);
assert!(
!is_dirty(&db.commit_state_snapshot(0)),
"a rejected CAS restores an empty buffer => shard stays CLEAN"
);
db.commit()?;
assert!(
db.cas(b"c".to_vec(), Some(1), 5).is_ok(),
"counter still 1 — the rejected CAS left it untouched"
);
Ok(())
}
#[test]
fn cas_rejection_from_dirty_stays_dirty_and_preserves_prior() -> Result<(), BoxError> {
let (_dir, db) = single_shard_db()?;
db.put(b"k".to_vec(), b"v".to_vec())?; assert!(is_dirty(&db.commit_state_snapshot(0)), "dirty baseline");
assert!(
db.cas(b"c".to_vec(), Some(99), 2).is_err(),
"CAS on a missing counter is rejected"
);
assert!(
is_dirty(&db.commit_state_snapshot(0)),
"a rejected CAS leaves the pre-existing buffered put => shard stays DIRTY"
);
db.commit()?;
assert_eq!(
db.get(b"k")?,
Some(b"v".to_vec()),
"the prior buffered put is not lost by the rejected CAS"
);
Ok(())
}
#[test]
fn append_rejection_from_clean_stays_clean() -> Result<(), BoxError> {
let (_dir, db) = single_shard_db()?;
db.append(b"s".to_vec(), vec![b"e0".to_vec()], 0)?; assert!(!is_dirty(&db.commit_state_snapshot(0)), "clean baseline");
assert!(
db.append(b"s".to_vec(), vec![b"e1".to_vec(), b"e2".to_vec()], 0)
.is_err(),
"a stale-seq append is rejected"
);
assert!(
!is_dirty(&db.commit_state_snapshot(0)),
"a rejected append restores the prior buffer => shard stays CLEAN"
);
assert_eq!(
db.read_events(b"s")?,
vec![b"e0".to_vec()],
"no partial entries leaked from the rejected multi-entry append"
);
Ok(())
}
#[test]
fn append_rejection_from_dirty_stays_dirty_and_preserves_prior() -> Result<(), BoxError> {
let (_dir, db) = single_shard_db()?;
db.append(b"s".to_vec(), vec![b"e0".to_vec()], 0)?; db.put(b"k".to_vec(), b"v".to_vec())?; assert!(is_dirty(&db.commit_state_snapshot(0)), "dirty baseline");
assert!(
db.append(b"s".to_vec(), vec![b"e1".to_vec(), b"e2".to_vec()], 0)
.is_err(),
"stale-seq append is rejected"
);
assert!(
is_dirty(&db.commit_state_snapshot(0)),
"a rejected append leaves the buffered put => shard stays DIRTY"
);
db.commit()?;
assert_eq!(
db.get(b"k")?,
Some(b"v".to_vec()),
"the prior buffered put is not lost by the rejected append"
);
assert_eq!(
db.read_events(b"s")?,
vec![b"e0".to_vec()],
"the stream is unchanged by the rejected append"
);
Ok(())
}
#[test]
fn restart_marker_durable_reopens_clean_and_converged() -> Result<(), BoxError> {
let dir = tempfile::tempdir()?;
let data_dir = dir.path().join("db");
let committed_root;
{
let db = Database::create(DatabaseConfig {
data_dir: data_dir.clone(),
shard_count: 1,
distributed: None,
executor_threads: None,
})?;
db.put(b"k".to_vec(), b"v".to_vec())?;
committed_root = *db.commit()?.get(&0).ok_or("shard 0 committed root")?; }
let reopened = Database::open(&data_dir)?;
assert_eq!(
reopened.get(b"k")?,
Some(b"v".to_vec()),
"the committed data is intact after reopen (marker recovered)"
);
let recovered = reopened.commit_state_snapshot(0);
assert!(
!is_dirty(&recovered),
"a shard with a durable marker recovers CLEAN (no stale-root window)"
);
assert_eq!(
recovered.root,
Some(committed_root),
"the recovered cell carries the durable marker root"
);
assert_eq!(recovered.incarnation, 1, "the reopened actor spawned once");
let jobs_before = reopened.executor().dispatched_jobs();
let roots = reopened.commit()?;
assert_eq!(
*roots.get(&0).ok_or("shard 0 root")?,
committed_root,
"the next commit returns the recovered marker root"
);
assert_eq!(
reopened.executor().dispatched_jobs(),
jobs_before,
"convergence is a clean skip — no job dispatched"
);
Ok(())
}
#[test]
fn advance_gen_is_monotonic_across_commits() -> Result<(), BoxError> {
let (_dir, db) = single_shard_db()?;
for round in 0..5_u32 {
db.put(
format!("k{round}").into_bytes(),
round.to_be_bytes().to_vec(),
)?;
db.commit()?;
db.commit()?; assert_eq!(
db.commit_state_snapshot(0).advance_gen,
u64::from(round) + 1,
"advance_gen bumps exactly once per advancing commit"
);
}
Ok(())
}