use std::path::{Path, PathBuf};
use std::sync::Arc;
use commonware_cryptography::{Blake3, Hasher as CHasher};
use commonware_parallel::Sequential;
use commonware_runtime::{Runner as _, Supervisor as _, buffer::paged::CacheRef};
use commonware_storage::merkle::Bagging;
use commonware_storage::merkle::mmr::{
Location as MmrLocation, Proof as MmrProof, StandardHasher,
full::{Config as JConfig, Mmr as JournaledMmr},
mem::Mmr as MemMmr,
};
use commonware_utils::{NZU16, NZU64, NZUsize};
use crate::hash::{HASH_LEN, Hash};
use crate::protocol::async_shim::Executor;
use crate::refs::validate_ref_name;
pub mod tokio_executor;
pub use tokio_executor::TokioExecutor;
const HISTORY_BAGGING: Bagging = Bagging::ForwardFold;
#[cfg(test)]
thread_local! {
static SYNC_CALL_COUNT: std::cell::Cell<u64> = const { std::cell::Cell::new(0) };
}
#[cfg(test)]
fn record_sync_call() {
SYNC_CALL_COUNT.with(|c| c.set(c.get() + 1));
}
#[cfg(test)]
pub(crate) fn reset_sync_call_count() {
SYNC_CALL_COUNT.with(|c| c.set(0));
}
#[cfg(test)]
pub(crate) fn sync_call_count() -> u64 {
SYNC_CALL_COUNT.with(std::cell::Cell::get)
}
fn history_hasher() -> StandardHasher<Blake3> {
StandardHasher::new(HISTORY_BAGGING)
}
#[derive(Copy, Clone, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct Position(pub u64);
impl Position {
#[must_use]
pub const fn as_u64(self) -> u64 {
self.0
}
}
pub type InclusionProof = MmrProof<<Blake3 as CHasher>::Digest>;
#[derive(Debug, thiserror::Error)]
pub enum HistoryError {
#[error("mmr error: {0}")]
Mmr(String),
#[error("invalid branch name for history journal: {0:?}")]
InvalidBranch(String),
#[error("history journal is corrupt: {0}")]
Corrupted(String),
#[error("failed to bootstrap commonware runtime: {0}")]
RuntimeBootstrap(String),
#[error("history directory I/O: {0}")]
Io(#[from] std::io::Error),
}
pub const HISTORY_DIR: &str = crate::layout::HISTORY_DIR_NAME;
pub const JOURNAL_PARTITION_SUFFIX: &str = "__journal";
pub const METADATA_PARTITION_SUFFIX: &str = "__metadata";
pub struct CommitHistory<X: Executor = TokioExecutor> {
backend: Backend<X>,
hasher: StandardHasher<Blake3>,
}
enum Backend<X: Executor> {
Mem {
mmr: MemMmr<<Blake3 as CHasher>::Digest>,
},
Journaled(Box<JournaledBackend<X>>),
}
struct JournaledBackend<X: Executor> {
mmr: JournaledMmr<commonware_runtime::tokio::Context, <Blake3 as CHasher>::Digest, Sequential>,
executor: Arc<X>,
ctx: commonware_runtime::tokio::Context,
common_dir: PathBuf,
branch: String,
}
impl<X: Executor> core::fmt::Debug for CommitHistory<X> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
match &self.backend {
Backend::Mem { mmr } => f
.debug_struct("CommitHistory::Mem")
.field("leaves", &u64::from(mmr.leaves()))
.field("size", &u64::from(mmr.size()))
.finish_non_exhaustive(),
Backend::Journaled(b) => f
.debug_struct("CommitHistory::Journaled")
.field("branch", &b.branch)
.field("leaves", &u64::from(b.mmr.leaves()))
.field("size", &u64::from(b.mmr.size()))
.finish_non_exhaustive(),
}
}
}
impl CommitHistory<TokioExecutor> {
#[must_use]
pub fn open() -> Self {
let hasher = history_hasher();
let mmr = MemMmr::new();
Self {
backend: Backend::Mem { mmr },
hasher,
}
}
}
impl<X: Executor + 'static> CommitHistory<X> {
pub fn open_at(
executor: Arc<X>,
layout: &crate::layout::RepoLayout,
branch: &str,
) -> Result<Self, HistoryError> {
Self::open_at_common_dir(executor, layout.common_dir(), branch)
}
fn open_at_common_dir(
executor: Arc<X>,
common_dir: &Path,
branch: &str,
) -> Result<Self, HistoryError> {
if !validate_ref_name(branch) {
return Err(HistoryError::InvalidBranch(branch.to_string()));
}
let history_dir = common_dir.join(HISTORY_DIR);
std::fs::create_dir_all(&history_dir)?;
let ctx = bootstrap_commonware_context(&history_dir)?;
Self::init_journaled(executor, ctx, common_dir, branch)
}
fn init_journaled(
executor: Arc<X>,
ctx: commonware_runtime::tokio::Context,
common_dir: &Path,
branch: &str,
) -> Result<Self, HistoryError> {
let sanitized = sanitize_branch(branch);
let journal_partition = format!("{sanitized}{JOURNAL_PARTITION_SUFFIX}");
let metadata_partition = format!("{sanitized}{METADATA_PARTITION_SUFFIX}");
let cfg = JConfig {
journal_partition,
metadata_partition,
items_per_blob: NZU64!(4096),
write_buffer: NZUsize!(4096),
strategy: Sequential,
page_cache: CacheRef::from_pooler(&ctx, NZU16!(4096), NZUsize!(8)),
};
let hasher = history_hasher();
let mmr = {
let hasher_inner = history_hasher();
let ctx_for_init = ctx.child("mmr_init");
executor
.block_on(async move {
JournaledMmr::<_, <Blake3 as CHasher>::Digest, Sequential>::init(
ctx_for_init,
&hasher_inner,
cfg,
)
.await
})
.map_err(|e| HistoryError::Corrupted(e.to_string()))?
};
Ok(Self {
backend: Backend::Journaled(Box::new(JournaledBackend {
mmr,
executor,
ctx,
common_dir: common_dir.to_path_buf(),
branch: branch.to_string(),
})),
hasher,
})
}
#[must_use]
pub fn common_dir(&self) -> Option<&Path> {
match &self.backend {
Backend::Mem { .. } => None,
Backend::Journaled(b) => Some(b.common_dir.as_path()),
}
}
#[must_use]
pub fn branch(&self) -> Option<&str> {
match &self.backend {
Backend::Mem { .. } => None,
Backend::Journaled(b) => Some(b.branch.as_str()),
}
}
pub fn reopen(&mut self) -> Result<(), HistoryError> {
let Backend::Journaled(b) = &self.backend else {
return Ok(());
};
let ctx = b.ctx.child("mmr_reopen");
let fresh = Self::init_journaled(b.executor.clone(), ctx, &b.common_dir, &b.branch)?;
*self = fresh;
Ok(())
}
pub fn append(&mut self, commit_hash: &Hash) -> Result<Position, HistoryError> {
let pos = self.append_no_sync(commit_hash)?;
self.sync()?;
Ok(pos)
}
fn append_no_sync(&mut self, commit_hash: &Hash) -> Result<Position, HistoryError> {
let leaf = digest_from_hash(commit_hash);
match &mut self.backend {
Backend::Mem { mmr } => {
let leaf_loc = mmr.leaves();
let batch = mmr
.new_batch()
.add(&self.hasher, &leaf)
.merkleize(mmr, &self.hasher);
mmr.apply_batch(&batch)
.map_err(|e| HistoryError::Mmr(e.to_string()))?;
Ok(Position(u64::from(leaf_loc)))
}
Backend::Journaled(b) => {
let leaf_loc = b.mmr.leaves();
let batch = b.mmr.new_batch().add(&self.hasher, &leaf);
let batch = b.mmr.with_mem(|mem| batch.merkleize(mem, &self.hasher));
b.mmr
.apply_batch(&batch)
.map_err(|e| HistoryError::Mmr(e.to_string()))?;
Ok(Position(u64::from(leaf_loc)))
}
}
}
fn sync(&mut self) -> Result<(), HistoryError> {
if let Backend::Journaled(b) = &mut self.backend {
#[cfg(test)]
record_sync_call();
let sync_fut = b.mmr.sync();
b.executor
.block_on(sync_fut)
.map_err(|e| HistoryError::Mmr(e.to_string()))?;
}
Ok(())
}
#[must_use]
pub fn root(&self) -> Hash {
let digest = match &self.backend {
Backend::Mem { mmr } => mmr.root(&self.hasher, 0),
Backend::Journaled(b) => b.mmr.root(&self.hasher, 0),
}
.expect("0 inactive peaks is always a valid root request");
let mut out = [0u8; HASH_LEN];
out.copy_from_slice(digest.as_ref());
out
}
#[must_use]
pub fn len(&self) -> u64 {
match &self.backend {
Backend::Mem { mmr } => u64::from(mmr.leaves()),
Backend::Journaled(b) => u64::from(b.mmr.leaves()),
}
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn prove(&self, position: Position) -> Result<InclusionProof, HistoryError> {
let loc = MmrLocation::new(position.0);
match &self.backend {
Backend::Mem { mmr } => mmr
.proof(&self.hasher, loc, 0)
.map_err(|e| HistoryError::Mmr(e.to_string())),
Backend::Journaled(b) => {
let hasher = self.hasher.clone();
let proof_fut = b.mmr.proof(&hasher, loc, 0);
b.executor
.block_on(proof_fut)
.map_err(|e| HistoryError::Mmr(e.to_string()))
}
}
}
pub fn destroy(self) -> Result<(), HistoryError> {
match self.backend {
Backend::Mem { .. } => Ok(()),
Backend::Journaled(b) => {
let JournaledBackend {
mmr,
executor,
ctx: _ctx,
..
} = *b;
executor
.block_on(mmr.destroy())
.map_err(|e| HistoryError::Mmr(e.to_string()))
}
}
}
}
impl Default for CommitHistory<TokioExecutor> {
fn default() -> Self {
Self::open()
}
}
#[must_use]
pub fn verify_inclusion(
commit_hash: &Hash,
position: Position,
proof: &InclusionProof,
root: &Hash,
) -> bool {
let leaf = digest_from_hash(commit_hash);
let root_digest = digest_from_hash(root);
let loc = MmrLocation::new(position.0);
let hasher = history_hasher();
proof.verify_element_inclusion(&hasher, leaf.as_ref(), loc, &root_digest)
}
pub fn rebuild_from_chain<X, F, E>(
history: &mut CommitHistory<X>,
tip: Hash,
mut parent_of: F,
) -> Result<u64, RebuildError<E>>
where
X: Executor + 'static,
F: FnMut(&Hash) -> Result<Option<Hash>, E>,
E: core::fmt::Display,
{
let mut chain = Vec::new();
let mut cursor = Some(tip);
while let Some(h) = cursor {
chain.push(h);
cursor = parent_of(&h).map_err(RebuildError::Walker)?;
}
chain.reverse();
let count = chain.len() as u64;
for h in &chain {
history.append_no_sync(h).map_err(RebuildError::History)?;
}
if count > 0 {
history.sync().map_err(RebuildError::History)?;
}
Ok(count)
}
#[derive(Debug, thiserror::Error)]
pub enum RebuildError<E: core::fmt::Display> {
#[error("parent-chain walker failed: {0}")]
Walker(E),
#[error(transparent)]
History(HistoryError),
}
fn digest_from_hash(h: &Hash) -> <Blake3 as CHasher>::Digest {
<<Blake3 as CHasher>::Digest as From<[u8; HASH_LEN]>>::from(*h)
}
use crate::refs::sanitize_ref_name as sanitize_branch;
#[cfg(test)]
mod bootstrap_probe {
use std::cell::RefCell;
use std::sync::Arc;
use std::sync::atomic::AtomicUsize;
thread_local! {
static COUNTER: RefCell<Option<Arc<AtomicUsize>>> = const { RefCell::new(None) };
}
struct Clear;
impl Drop for Clear {
fn drop(&mut self) {
COUNTER.with(|c| *c.borrow_mut() = None);
}
}
#[must_use]
pub(super) fn track(counter: Arc<AtomicUsize>) -> impl Drop {
COUNTER.with(|c| *c.borrow_mut() = Some(counter));
Clear
}
pub(super) fn current() -> Option<Arc<AtomicUsize>> {
COUNTER.with(|c| c.borrow().clone())
}
}
fn bootstrap_commonware_context(
storage_directory: &Path,
) -> Result<commonware_runtime::tokio::Context, HistoryError> {
let dir = storage_directory.to_path_buf();
#[cfg(test)]
let probe = bootstrap_probe::current();
std::thread::spawn(move || {
#[cfg(test)]
if let Some(counter) = &probe {
counter.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
}
let cfg = commonware_runtime::tokio::Config::new().with_storage_directory(dir);
let runner = commonware_runtime::tokio::Runner::new(cfg);
runner
.start(|ctx| async move { commonware_runtime::Supervisor::child(&ctx, "mkit_history") })
})
.join()
.map_err(|_| HistoryError::RuntimeBootstrap("bootstrap thread panicked".to_string()))
}
#[cfg(test)]
mod tests {
use super::*;
use tempfile::TempDir;
fn synth(i: u64) -> Hash {
crate::hash::hash(&i.to_be_bytes())
}
fn fresh_mkit_dir() -> (TempDir, crate::layout::RepoLayout) {
let tmp = TempDir::new().unwrap();
let mkit_dir = crate::layout::RepoLayout::single(tmp.path());
std::fs::create_dir_all(mkit_dir.common_dir()).unwrap();
(tmp, mkit_dir)
}
fn fresh_executor() -> Arc<TokioExecutor> {
Arc::new(TokioExecutor::new().expect("tokio runtime"))
}
#[test]
fn mem_empty_history_root_is_well_defined() {
let h1 = CommitHistory::open();
let h2 = CommitHistory::open();
assert_eq!(h1.root(), h2.root(), "empty root must be deterministic");
assert!(h1.is_empty());
assert_eq!(h1.len(), 0);
}
#[test]
fn mem_append_returns_dense_positions() {
let mut h = CommitHistory::open();
for i in 0..16u64 {
let pos = h.append(&synth(i)).unwrap();
assert_eq!(pos, Position(i), "positions must be dense and 0-based");
}
assert_eq!(h.len(), 16);
}
#[test]
fn mem_prove_and_verify_position_712_of_1000() {
let mut h = CommitHistory::open();
let commits: Vec<Hash> = (0..1000u64).map(synth).collect();
for c in &commits {
h.append(c).unwrap();
}
assert_eq!(h.len(), 1000);
let target = Position(712);
let proof = h.prove(target).unwrap();
let root = h.root();
assert!(
verify_inclusion(&commits[712], target, &proof, &root),
"honest proof must verify"
);
}
#[test]
fn mem_tampered_proof_fails_verification() {
let mut h = CommitHistory::open();
for i in 0..256u64 {
h.append(&synth(i)).unwrap();
}
let target = Position(42);
let mut proof = h.prove(target).unwrap();
let root = h.root();
let commit = synth(42);
assert!(verify_inclusion(&commit, target, &proof, &root));
assert!(
!proof.digests.is_empty(),
"non-trivial proof must carry at least one sibling"
);
let mut bytes: [u8; HASH_LEN] = [0u8; HASH_LEN];
bytes.copy_from_slice(proof.digests[0].as_ref());
bytes[0] ^= 0x01;
proof.digests[0] = <<Blake3 as CHasher>::Digest as From<[u8; HASH_LEN]>>::from(bytes);
assert!(
!verify_inclusion(&commit, target, &proof, &root),
"tampered proof must fail"
);
}
#[test]
fn verify_inclusion_rejects_wrong_position() {
let mut h = CommitHistory::open();
let commits: Vec<Hash> = (0..64u64).map(synth).collect();
for c in &commits {
h.append(c).unwrap();
}
let target = Position(42);
let proof = h.prove(target).unwrap();
let root = h.root();
assert!(verify_inclusion(&commits[42], target, &proof, &root));
assert!(!verify_inclusion(&commits[42], Position(41), &proof, &root));
assert!(!verify_inclusion(&commits[42], Position(0), &proof, &root));
}
#[test]
fn verify_inclusion_rejects_mismatched_leaf_count() {
let mut h = CommitHistory::open();
let commits: Vec<Hash> = (0..64u64).map(synth).collect();
for c in &commits {
h.append(c).unwrap();
}
let target = Position(42);
let mut proof = h.prove(target).unwrap();
let root = h.root();
assert!(verify_inclusion(&commits[42], target, &proof, &root));
proof.leaves = MmrLocation::new(63);
assert!(!verify_inclusion(&commits[42], target, &proof, &root));
}
#[test]
fn verify_inclusion_rejects_truncated_or_over_long_digests() {
let mut h = CommitHistory::open();
let commits: Vec<Hash> = (0..64u64).map(synth).collect();
for c in &commits {
h.append(c).unwrap();
}
let target = Position(42);
let proof = h.prove(target).unwrap();
let root = h.root();
assert!(verify_inclusion(&commits[42], target, &proof, &root));
assert!(
!proof.digests.is_empty(),
"non-trivial proof must carry at least one digest"
);
let mut truncated = proof.clone();
truncated.digests.pop();
assert!(!verify_inclusion(&commits[42], target, &truncated, &root));
let mut over_long = proof;
over_long.digests.push(over_long.digests[0]);
assert!(!verify_inclusion(&commits[42], target, &over_long, &root));
}
#[test]
fn open_at_rejects_invalid_branch() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let err = CommitHistory::open_at(exec, &mkit_dir, "../escape").expect_err("invalid branch");
assert!(matches!(err, HistoryError::InvalidBranch(_)));
}
#[test]
fn open_at_empty_root_matches_mem() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let h_disk = CommitHistory::open_at(exec, &mkit_dir, "main").unwrap();
let h_mem = CommitHistory::open();
assert_eq!(
h_disk.root(),
h_mem.root(),
"empty journaled root must match empty mem root"
);
assert!(h_disk.is_empty());
}
#[test]
fn open_at_round_trip_100_commits() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let commits: Vec<Hash> = (0..100u64).map(synth).collect();
let (root_before, len_before) = {
let mut h = CommitHistory::open_at(exec.clone(), &mkit_dir, "main").unwrap();
for c in &commits {
h.append(c).unwrap();
}
(h.root(), h.len())
};
let h = CommitHistory::open_at(exec.clone(), &mkit_dir, "main").unwrap();
assert_eq!(h.len(), len_before, "leaf count must survive reopen");
assert_eq!(h.root(), root_before, "root must survive reopen");
}
#[test]
fn open_at_prove_after_reopen() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let commits: Vec<Hash> = (0..64u64).map(synth).collect();
let root = {
let mut h = CommitHistory::open_at(exec.clone(), &mkit_dir, "main").unwrap();
for c in &commits {
h.append(c).unwrap();
}
h.root()
};
let h = CommitHistory::open_at(exec.clone(), &mkit_dir, "main").unwrap();
let target = Position(17);
let proof = h.prove(target).unwrap();
assert!(verify_inclusion(&commits[17], target, &proof, &root));
}
#[test]
fn reopen_reuses_the_bootstrapped_commonware_context_instead_of_rebootstrapping() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let counter = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let _guard = bootstrap_probe::track(counter.clone());
let mut h = CommitHistory::open_at(exec, &mkit_dir, "main").unwrap();
assert_eq!(
counter.load(std::sync::atomic::Ordering::SeqCst),
1,
"open_at must bootstrap the commonware Context exactly once"
);
h.reopen().unwrap();
assert_eq!(
counter.load(std::sync::atomic::Ordering::SeqCst),
1,
"reopen() must reuse the already-bootstrapped Context rather than \
spawning a second bootstrap thread"
);
}
#[test]
fn open_at_distinct_branches_have_distinct_roots() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let mut main = CommitHistory::open_at(exec.clone(), &mkit_dir, "main").unwrap();
let mut dev = CommitHistory::open_at(exec.clone(), &mkit_dir, "dev").unwrap();
main.append(&synth(0)).unwrap();
dev.append(&synth(1)).unwrap();
assert_ne!(main.root(), dev.root());
}
#[test]
fn open_at_branch_with_slash_is_isolated() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let mut a = CommitHistory::open_at(exec.clone(), &mkit_dir, "feat/v1").unwrap();
let mut b = CommitHistory::open_at(exec.clone(), &mkit_dir, "feat/v2").unwrap();
a.append(&synth(0)).unwrap();
assert_eq!(a.len(), 1);
assert_eq!(b.len(), 0, "sibling branch must not see appended leaf");
b.append(&synth(1)).unwrap();
assert_ne!(a.root(), b.root());
}
#[test]
fn destroy_removes_on_disk_partition() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let commits: Vec<Hash> = (0..5u64).map(synth).collect();
let mut h = CommitHistory::open_at(exec, &mkit_dir, "feature").unwrap();
for c in &commits {
h.append(c).unwrap();
}
assert_eq!(h.len(), 5);
let journal_blobs = mkit_dir.history_dir().join("feature__journal-blobs");
let journal_metadata = mkit_dir.history_dir().join("feature__journal-metadata");
assert!(
journal_blobs.exists(),
"journal blob dir must exist after appends"
);
h.destroy().unwrap();
assert!(
!journal_blobs.exists(),
"destroy must remove the on-disk journal blob partition"
);
assert!(
!journal_metadata.exists(),
"destroy must remove the on-disk journal metadata partition"
);
}
#[test]
fn destroyed_journal_reopens_empty_not_resuming_dead_incarnation() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let commits: Vec<Hash> = (0..3u64).map(synth).collect();
let mut h = CommitHistory::open_at(exec.clone(), &mkit_dir, "feature").unwrap();
for c in &commits {
h.append(c).unwrap();
}
assert_eq!(h.len(), 3);
h.destroy().unwrap();
let reopened = CommitHistory::open_at(exec, &mkit_dir, "feature").unwrap();
assert_eq!(
reopened.len(),
0,
"a destroyed journal must reopen with zero leaves, not resume the dead incarnation"
);
assert_eq!(
reopened.root(),
CommitHistory::open().root(),
"a fresh reopen after destroy must match a genuinely empty MMR's root"
);
}
#[test]
fn destroy_of_one_branch_does_not_touch_a_sibling_branch() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let mut a = CommitHistory::open_at(exec.clone(), &mkit_dir, "a").unwrap();
let mut b = CommitHistory::open_at(exec.clone(), &mkit_dir, "b").unwrap();
a.append(&synth(0)).unwrap();
b.append(&synth(1)).unwrap();
let b_root = b.root();
drop(b);
a.destroy().unwrap();
let b_reopened = CommitHistory::open_at(exec, &mkit_dir, "b").unwrap();
assert_eq!(
b_reopened.len(),
1,
"destroying branch 'a' must not affect sibling branch 'b'"
);
assert_eq!(b_reopened.root(), b_root);
}
#[test]
fn open_at_root_matches_live_mem_root() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let commits: Vec<Hash> = (0..50u64).map(synth).collect();
let mut journaled = CommitHistory::open_at(exec, &mkit_dir, "main").unwrap();
let mut mem = CommitHistory::open();
for c in &commits {
journaled.append(c).unwrap();
mem.append(c).unwrap();
}
assert_eq!(
journaled.root(),
mem.root(),
"journaled and mem MMRs must produce the same root for the same leaf sequence"
);
}
#[test]
fn journaled_root_changes_on_append() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let mut h = CommitHistory::open_at(exec, &mkit_dir, "main").unwrap();
let root_before = h.root();
h.append(&synth(0)).unwrap();
let root_after = h.root();
assert_ne!(
root_before, root_after,
"appending a leaf must change the journaled MMR root"
);
}
#[test]
fn journaled_wrong_commit_fails_verification() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let mut h = CommitHistory::open_at(exec, &mkit_dir, "main").unwrap();
let h_a = synth(0);
let h_other = synth(99);
h.append(&h_a).unwrap();
let proof = h.prove(Position(0)).unwrap();
let root = h.root();
assert!(
verify_inclusion(&h_a, Position(0), &proof, &root),
"honest proof must verify with the appended commit"
);
assert!(
!verify_inclusion(&h_other, Position(0), &proof, &root),
"swapping in a different commit hash must fail verification"
);
}
#[test]
fn journaled_wrong_root_fails_verification() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let h_a = synth(0);
let mut a = CommitHistory::open_at(exec.clone(), &mkit_dir, "main").unwrap();
a.append(&h_a).unwrap();
let proof = a.prove(Position(0)).unwrap();
let root_a = a.root();
assert!(
verify_inclusion(&h_a, Position(0), &proof, &root_a),
"sanity: honest proof verifies against its own root"
);
let mut b = CommitHistory::open_at(exec, &mkit_dir, "dev").unwrap();
b.append(&synth(1)).unwrap();
b.append(&synth(2)).unwrap();
let root_b = b.root();
assert_ne!(root_a, root_b, "distinct branches must have distinct roots");
assert!(
!verify_inclusion(&h_a, Position(0), &proof, &root_b),
"proof from one branch must not verify against another branch's root"
);
}
#[test]
fn truncated_journal_rolls_forward_or_surfaces_corruption() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let commits: Vec<Hash> = (0..32u64).map(synth).collect();
let len_before = {
let mut h = CommitHistory::open_at(exec.clone(), &mkit_dir, "main").unwrap();
for c in &commits {
h.append(c).unwrap();
}
h.len()
};
assert_eq!(len_before, 32);
let journal_root = mkit_dir.history_dir().join("main__journal-blobs");
assert!(
journal_root.exists(),
"journal blob dir must exist after appends"
);
let mut largest: Option<(std::path::PathBuf, u64)> = None;
for entry in std::fs::read_dir(&journal_root).unwrap() {
let entry = entry.unwrap();
let meta = entry.metadata().unwrap();
if meta.is_file() {
let path = entry.path();
let len = meta.len();
if largest.as_ref().is_none_or(|(_, l)| len > *l) {
largest = Some((path, len));
}
}
}
let (blob_path, blob_len) = largest.expect("at least one blob present after 32 appends");
assert!(
blob_len > 5,
"blob must be large enough to truncate (got {blob_len} bytes)"
);
let truncated_len = blob_len - 5;
let f = std::fs::OpenOptions::new()
.write(true)
.open(&blob_path)
.unwrap();
f.set_len(truncated_len).unwrap();
drop(f);
let reopened = CommitHistory::open_at(exec, &mkit_dir, "main");
match reopened {
Ok(h) => {
assert!(
h.len() <= len_before,
"rolled-forward leaf count must not exceed the pre-truncation count"
);
if !h.is_empty() {
let n = usize::try_from(h.len()).expect("leaf count fits in usize");
let mut reference = CommitHistory::open();
for c in &commits[..n] {
reference.append(c).unwrap();
}
assert_eq!(
h.root(),
reference.root(),
"rolled-forward root must equal a clean replay of the surviving leaves"
);
}
}
Err(HistoryError::Corrupted(_)) => {
}
Err(other) => panic!("unexpected error on truncated journal: {other:?}"),
}
}
#[test]
fn rebuild_from_chain_matches_live_appends() {
let (_tmp, mkit_dir_a) = fresh_mkit_dir();
let (_tmp_b, mkit_dir_b) = fresh_mkit_dir();
let exec = fresh_executor();
let commits: Vec<Hash> = (0..10u64).map(synth).collect();
let reference_root = {
let mut h = CommitHistory::open_at(exec.clone(), &mkit_dir_a, "main").unwrap();
for c in &commits {
h.append(c).unwrap();
}
h.root()
};
let mut h = CommitHistory::open_at(exec.clone(), &mkit_dir_b, "main").unwrap();
let tip = commits[commits.len() - 1];
let count = rebuild_from_chain::<TokioExecutor, _, std::io::Error>(&mut h, tip, |hash| {
let idx = commits
.iter()
.position(|c| c == hash)
.expect("walker called with unknown hash");
if idx == 0 {
Ok(None)
} else {
Ok(Some(commits[idx - 1]))
}
})
.unwrap();
assert_eq!(count, commits.len() as u64);
assert_eq!(
h.root(),
reference_root,
"rebuilt root must match live-append root"
);
}
#[test]
fn rebuild_from_chain_batches_backfill_into_a_single_sync() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let commits: Vec<Hash> = (0..500u64).map(synth).collect();
let mut h = CommitHistory::open_at(exec, &mkit_dir, "main").unwrap();
reset_sync_call_count();
let tip = commits[commits.len() - 1];
let count = rebuild_from_chain::<TokioExecutor, _, std::io::Error>(&mut h, tip, |hash| {
let idx = commits
.iter()
.position(|c| c == hash)
.expect("walker called with unknown hash");
if idx == 0 {
Ok(None)
} else {
Ok(Some(commits[idx - 1]))
}
})
.unwrap();
assert_eq!(count, 500, "sanity: all 500 commits were backfilled");
assert_eq!(h.len(), 500, "sanity: all 500 leaves landed in the MMR");
assert_eq!(
sync_call_count(),
1,
"backfilling 500 commits must fsync exactly once for the \
whole batch, not once per commit"
);
}
#[test]
fn append_still_syncs_once_per_call() {
let (_tmp, mkit_dir) = fresh_mkit_dir();
let exec = fresh_executor();
let mut h = CommitHistory::open_at(exec, &mkit_dir, "main").unwrap();
reset_sync_call_count();
h.append(&synth(0)).unwrap();
h.append(&synth(1)).unwrap();
h.append(&synth(2)).unwrap();
assert_eq!(
sync_call_count(),
3,
"each live `append` must still fsync on its own — only the \
backfill path batches"
);
}
#[test]
fn sanitize_branch_round_trip_invariants() {
assert_ne!(sanitize_branch("feat/v1.0"), sanitize_branch("feat_v1_0"));
assert_eq!(sanitize_branch("main"), "main");
assert_eq!(sanitize_branch("release-2026"), "release-2026");
assert_eq!(sanitize_branch("a/b"), "a_2fb");
assert_eq!(sanitize_branch("v1.0"), "v1_2e0");
assert_eq!(sanitize_branch("a_b"), "a_5fb");
}
}