use super::operation::Operation;
pub use crate::qmdb::compact::Config;
use crate::{
Context,
journal::contiguous::variable::{self, Config as JournalConfig},
merkle::{Family, Location, Proof, batch, compact as compact_merkle},
qmdb::{
self, Error,
any::value::ValueEncoding,
batch_chain::{self, Bounds, Commitment},
compact::{
batch as compact_batch,
witness::{self, VerifiedWitness},
},
sync::{CompactTarget, FeedbackTx, Request, Response, Source},
},
};
use commonware_codec::{Encode, EncodeShared, Read};
use commonware_cryptography::{Digest, Hasher};
use commonware_macros::boxed;
use commonware_parallel::Strategy;
use commonware_runtime::Handle;
use std::sync::{Arc, Weak};
pub struct Db<F, E, V, H, C, S: Strategy>
where
F: Family,
E: Context,
V: ValueEncoding,
H: Hasher,
Operation<F, V>: EncodeShared,
Operation<F, V>: Read<Cfg = C>,
C: Clone + Send + Sync + 'static,
{
merkle: compact_merkle::Merkle<F, H::Digest, S>,
root: H::Digest,
last_commit_loc: Location<F>,
last_commit_metadata: Option<V::Value>,
inactivity_floor_loc: Location<F>,
commit_codec_config: C,
witness: witness::Store<E, F, H::Digest>,
}
impl<F, E, V, H, C, S: Strategy> std::fmt::Debug for Db<F, E, V, H, C, S>
where
F: Family,
E: Context,
V: ValueEncoding,
H: Hasher,
Operation<F, V>: EncodeShared,
Operation<F, V>: Read<Cfg = C>,
C: Clone + Send + Sync + 'static,
{
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Db")
.field("size", &self.size())
.field("inactivity_floor_loc", &self.inactivity_floor_loc())
.finish_non_exhaustive()
}
}
#[allow(clippy::type_complexity)]
pub struct UnmerkleizedBatch<F, H, V, S: Strategy>
where
F: Family,
V: ValueEncoding,
H: Hasher,
Operation<F, V>: EncodeShared,
{
merkle_batch: compact_merkle::UnmerkleizedBatch<F, H::Digest, S>,
appends: Vec<V::Value>,
parent: Option<Arc<MerkleizedBatch<F, H::Digest, V, S>>>,
base: batch_chain::Commitment<F, H::Digest>,
}
#[derive(Clone)]
pub struct MerkleizedBatch<F: Family, D: Digest, V: ValueEncoding, S: Strategy>
where
Operation<F, V>: EncodeShared,
{
pub(super) merkle_batch: Arc<batch::MerkleizedBatch<F, D, S>>,
operations: Arc<Vec<Operation<F, V>>>,
pub(super) commit_metadata: Option<V::Value>,
pub(super) parent: Option<Weak<Self>>,
pub(super) bounds: batch_chain::Bounds<F, D>,
}
impl<F: Family, D: Digest, V: ValueEncoding, S: Strategy> MerkleizedBatch<F, D, V, S>
where
Operation<F, V>: EncodeShared,
{
pub(super) fn ancestors(&self) -> impl Iterator<Item = Arc<Self>> + use<F, D, V, S> {
batch_chain::ancestors(self.parent.clone(), |batch| batch.parent.as_ref())
}
pub(super) const fn commitment(&self) -> Commitment<F, D> {
self.bounds.tip
}
pub const fn root(&self) -> D {
self.bounds.tip.root
}
pub const fn bounds(&self) -> &Bounds<F, D> {
&self.bounds
}
pub fn operations(&self) -> (Location<F>, Arc<Vec<Operation<F, V>>>) {
(self.bounds.base.size, Arc::clone(&self.operations))
}
pub fn proof<E, C, H>(&self, db: &Db<F, E, V, H, C, S>) -> Result<Proof<F, D>, Error<F>>
where
E: Context,
H: Hasher<Digest = D>,
C: Clone + Send + Sync + 'static,
Operation<F, V>: Read<Cfg = C>,
{
let inactive_peaks = F::inactive_peaks(self.bounds.tip.size, self.bounds.inactivity_floor);
let hasher = qmdb::hasher::<H>();
db.merkle
.with_mem(|base| {
self.merkle_batch.range_proof(
base,
&hasher,
self.bounds.base.size..self.bounds.tip.size,
inactive_peaks,
)
})
.map_err(Into::into)
}
pub fn pinned_nodes<E, C, H>(&self, db: &Db<F, E, V, H, C, S>) -> Result<Vec<D>, Error<F>>
where
E: Context,
H: Hasher<Digest = D>,
C: Clone + Send + Sync + 'static,
Operation<F, V>: Read<Cfg = C>,
{
db.merkle
.with_mem(|base| {
F::nodes_to_pin(self.bounds.base.size)
.map(|pos| {
self.merkle_batch
.get_node(pos)
.or_else(|| base.get_node(pos))
.ok_or(crate::merkle::Error::ElementPruned(pos))
})
.collect::<Result<Vec<_>, _>>()
})
.map_err(Into::into)
}
pub fn new_batch<H>(self: &Arc<Self>) -> UnmerkleizedBatch<F, H, V, S>
where
H: Hasher<Digest = D>,
{
UnmerkleizedBatch {
merkle_batch: compact_merkle::UnmerkleizedBatch::wrap(self.merkle_batch.new_batch()),
appends: Vec::new(),
parent: Some(Arc::clone(self)),
base: self.commitment(),
}
}
}
impl<F, H, V, S> UnmerkleizedBatch<F, H, V, S>
where
F: Family,
V: ValueEncoding,
H: Hasher,
S: Strategy,
Operation<F, V>: EncodeShared,
{
pub(super) fn new<E, C>(
db: &Db<F, E, V, H, C, S>,
base: batch_chain::Commitment<F, H::Digest>,
) -> Self
where
E: Context,
C: Clone + Send + Sync + 'static,
Operation<F, V>: Read<Cfg = C>,
{
Self {
merkle_batch: db.merkle.new_batch(),
appends: Vec::new(),
parent: None,
base,
}
}
fn db(&self) -> Commitment<F, H::Digest> {
self.parent
.as_ref()
.map_or(self.base, |parent| parent.bounds.db)
}
pub fn append(mut self, value: V::Value) -> Self {
self.appends.push(value);
self
}
#[tracing::instrument(
name = "qmdb.keyless.compact.batch.merkleize",
level = "info",
skip_all
)]
pub async fn merkleize<E, C>(
self,
db: &Db<F, E, V, H, C, S>,
metadata: Option<V::Value>,
inactivity_floor: Location<F>,
) -> Arc<MerkleizedBatch<F, H::Digest, V, S>>
where
F: Family,
E: Context,
C: Clone + Send + Sync + 'static,
Operation<F, V>: Read<Cfg = C>,
{
let live_ancestors: Vec<_> =
batch_chain::parent_and_ancestors(self.parent.as_ref(), |parent| parent.ancestors())
.collect();
let boundary = batch_chain::effective_boundary(
self.db(),
live_ancestors.last().map(|oldest| oldest.bounds.base),
);
let mut ops: Vec<Operation<F, V>> = Vec::with_capacity(self.appends.len() + 1);
for value in self.appends {
ops.push(Operation::Append(value));
}
ops.push(Operation::Commit(metadata.clone(), inactivity_floor));
let operations = Arc::new(ops);
let total_size = self.base.size + operations.len() as u64;
let inactive_peaks = F::inactive_peaks(total_size, inactivity_floor);
let (merkle, root) = compact_batch::merkleize_ops::<F, H, S, _>(
&db.merkle,
self.merkle_batch,
Arc::clone(&operations),
inactive_peaks,
)
.await
.expect("inactive_peaks computed from batch size");
let ancestors = batch_chain::collect_ancestor_bounds(
live_ancestors,
|batch| batch.bounds.inactivity_floor,
|batch| batch.commitment(),
);
Arc::new(MerkleizedBatch {
merkle_batch: merkle,
operations,
commit_metadata: metadata,
parent: self.parent.as_ref().map(Arc::downgrade),
bounds: batch_chain::Bounds {
base: self.base,
db: boundary,
tip: Commitment::new(total_size, root),
ancestors,
inactivity_floor,
},
})
}
}
impl<F, E, V, H, C, S> Db<F, E, V, H, C, S>
where
F: Family,
E: Context,
V: ValueEncoding,
H: Hasher,
S: Strategy,
Operation<F, V>: EncodeShared,
Operation<F, V>: Read<Cfg = C>,
C: Clone + Send + Sync + 'static,
{
fn encode_commit_op(metadata: Option<V::Value>, inactivity_floor_loc: Location<F>) -> Vec<u8> {
Operation::<F, V>::Commit(metadata, inactivity_floor_loc)
.encode()
.to_vec()
}
pub(crate) fn init_from_sync(
strategy: S,
journal: witness::Journal<E, F, H::Digest>,
commit_codec_config: C,
last_commit_loc: Location<F>,
pinned_nodes: Vec<H::Digest>,
last_commit_op: Operation<F, V>,
) -> Result<Self, Error<F>> {
let Operation::Commit(last_commit_metadata, inactivity_floor_loc) = last_commit_op else {
return Err(Error::UnexpectedData(last_commit_loc));
};
witness::validate_inactivity_floor(inactivity_floor_loc, last_commit_loc)?;
let op_bytes = Self::encode_commit_op(last_commit_metadata.clone(), inactivity_floor_loc);
let merkle =
compact_merkle::Merkle::from_compact_state(strategy, last_commit_loc, pinned_nodes)?;
let hasher = qmdb::hasher::<H>();
merkle.append_leaf(&hasher, &op_bytes)?;
let imported = witness::build_witness::<F, H, S>(&merkle, inactivity_floor_loc, op_bytes)?;
merkle.prune_to_frontier();
let witness = witness::Store::from_import(journal, imported);
let root = witness.with(|w| w.root);
Ok(Self {
merkle,
root,
last_commit_loc,
last_commit_metadata,
inactivity_floor_loc,
commit_codec_config,
witness,
})
}
#[boxed]
pub(crate) async fn init_from_merkle(
mut merkle: compact_merkle::Merkle<F, H::Digest, S>,
witness_context: E,
witness_config: JournalConfig<()>,
commit_codec_config: C,
) -> Result<Self, Error<F>>
where
F: Family,
Operation<F, V>: Read<Cfg = C>,
{
let journal: witness::Journal<E, F, H::Digest> =
variable::Journal::init(witness_context, witness_config).await?;
let (witness, last_commit_op) = witness::init::<E, F, H, S, Operation<F, V>>(
journal,
&mut merkle,
&commit_codec_config,
Operation::<F, V>::Commit(None, Location::new(0))
.encode()
.to_vec(),
)
.await?;
let Operation::Commit(last_commit_metadata, inactivity_floor_loc) = last_commit_op else {
return Err(Error::DataCorrupted("last operation was not a commit"));
};
let last_commit_loc = witness.with(|w| w.size()) - 1;
let root = witness.with(|w| w.root);
Ok(Self {
merkle,
root,
last_commit_loc,
last_commit_metadata,
inactivity_floor_loc,
commit_codec_config,
witness,
})
}
pub const fn root(&self) -> H::Digest {
self.root
}
pub const fn strategy(&self) -> &S {
self.merkle.strategy()
}
pub const fn last_commit_loc(&self) -> Location<F> {
self.last_commit_loc
}
pub const fn inactivity_floor_loc(&self) -> Location<F> {
self.inactivity_floor_loc
}
pub fn size(&self) -> Location<F> {
self.last_commit_loc + 1
}
pub fn get_metadata(&self) -> Option<V::Value> {
self.last_commit_metadata.clone()
}
pub fn target(&self) -> CompactTarget<F, H::Digest> {
self.witness.with(VerifiedWitness::target)
}
pub(crate) fn commitment(&self) -> batch_chain::Commitment<F, H::Digest> {
batch_chain::Commitment::new(self.last_commit_loc + 1, self.root())
}
pub fn new_batch(&self) -> UnmerkleizedBatch<F, H, V, S> {
UnmerkleizedBatch::new(self, self.commitment())
}
pub fn to_batch(&self) -> Arc<MerkleizedBatch<F, H::Digest, V, S>>
where
F: Family,
{
Arc::new(MerkleizedBatch {
merkle_batch: self.merkle.to_batch(),
operations: Arc::new(Vec::new()),
commit_metadata: self.last_commit_metadata.clone(),
parent: None,
bounds: batch_chain::Bounds::from_db(self.commitment(), self.inactivity_floor_loc),
})
}
pub fn validate_batch(
&self,
batch: &MerkleizedBatch<F, H::Digest, V, S>,
) -> Result<(), Error<F>> {
batch
.bounds
.validate_apply_to(self.commitment(), self.inactivity_floor_loc)
}
#[tracing::instrument(name = "qmdb.keyless.compact.db.apply_batch", level = "info", skip_all)]
pub async fn apply_batch(
mut self,
batch: Arc<MerkleizedBatch<F, H::Digest, V, S>>,
) -> Result<(Self, core::ops::Range<Location<F>>), Error<F>> {
self.validate_batch(&batch)?;
let start_loc = self.last_commit_loc + 1;
self.merkle.apply_batch(&batch.merkle_batch)?;
self.root = batch.root();
self.last_commit_loc = batch.bounds.tip.size - 1;
self.last_commit_metadata = batch.commit_metadata.clone();
self.inactivity_floor_loc = batch.bounds.inactivity_floor;
let last_commit_metadata = self.last_commit_metadata.clone();
let inactivity_floor_loc = self.inactivity_floor_loc;
self.witness = self
.witness
.apply::<H, S>(&self.merkle, inactivity_floor_loc, || {
Self::encode_commit_op(last_commit_metadata, inactivity_floor_loc)
})
.await?;
Ok((self, start_loc..batch.bounds.tip.size))
}
#[tracing::instrument(name = "qmdb.keyless.compact.db.start_sync", level = "info", skip_all)]
pub async fn start_sync(mut self) -> Result<(Self, Handle<()>), Error<F>> {
let last_commit_metadata = self.last_commit_metadata.clone();
let inactivity_floor_loc = self.inactivity_floor_loc;
let handle;
(self.witness, handle) = self
.witness
.start_sync::<H, S>(&self.merkle, inactivity_floor_loc, || {
Self::encode_commit_op(last_commit_metadata, inactivity_floor_loc)
})
.await?;
Ok((self, handle))
}
#[tracing::instrument(name = "qmdb.keyless.compact.db.commit", level = "info", skip_all)]
pub async fn commit(mut self) -> Result<Self, Error<F>> {
let last_commit_metadata = self.last_commit_metadata.clone();
let inactivity_floor_loc = self.inactivity_floor_loc;
self.witness = self
.witness
.commit::<H, S>(&self.merkle, inactivity_floor_loc, || {
Self::encode_commit_op(last_commit_metadata, inactivity_floor_loc)
})
.await?;
Ok(self)
}
#[tracing::instrument(name = "qmdb.keyless.compact.db.sync", level = "info", skip_all)]
pub async fn sync(mut self) -> Result<Self, Error<F>> {
let last_commit_metadata = self.last_commit_metadata.clone();
let inactivity_floor_loc = self.inactivity_floor_loc;
self.witness = self
.witness
.sync::<H, S>(&self.merkle, inactivity_floor_loc, || {
Self::encode_commit_op(last_commit_metadata, inactivity_floor_loc)
})
.await?;
Ok(self)
}
#[tracing::instrument(name = "qmdb.keyless.compact.db.rewind", level = "info", skip_all)]
pub async fn rewind(mut self, target: Location<F>) -> Result<Self, Error<F>>
where
F: Family,
{
if self.size() == target
&& self.witness.with(|w| w.size()) == target
&& !self.witness.has_uncommitted_state()
{
self.witness.wait_for_sync().await?;
return Ok(self);
}
let last_commit_op;
(self.witness, last_commit_op) = self
.witness
.rewind::<H, S, Operation<F, V>>(&self.merkle, target, &self.commit_codec_config)
.await?;
let Operation::Commit(last_commit_metadata, inactivity_floor_loc) = last_commit_op else {
return Err(Error::DataCorrupted("last operation was not a commit"));
};
self.last_commit_metadata = last_commit_metadata;
self.inactivity_floor_loc = inactivity_floor_loc;
self.last_commit_loc = target - 1;
self.root = self.witness.with(|w| w.root);
Ok(self)
}
#[tracing::instrument(name = "qmdb.keyless.compact.db.prune", level = "info", skip_all)]
pub async fn prune(mut self, pruning_boundary: Location<F>) -> Result<Self, Error<F>> {
self.witness = self.witness.prune(pruning_boundary).await?;
Ok(self)
}
#[boxed]
pub async fn destroy(self) -> Result<(), Error<F>> {
self.witness.destroy().await?;
Ok(())
}
}
impl<F, E, V, H, C, S> Source for Db<F, E, V, H, C, S>
where
F: Family,
E: Context,
V: ValueEncoding,
H: Hasher,
Operation<F, V>: EncodeShared + Read<Cfg = C>,
C: Clone + Send + Sync + 'static,
S: Strategy,
{
type Family = F;
type Digest = H::Digest;
type Op = Operation<F, V>;
type Error = qmdb::Error<F>;
async fn serve(
&self,
request: Request<F>,
) -> Result<(Response<F, Self::Op, H::Digest>, FeedbackTx), Self::Error> {
Ok((
self.witness
.compact_state(&self.commit_codec_config, request)?,
None,
))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{
merkle::{mmb, mmr},
qmdb::{
any::value::FixedEncoding, compact::witness, verify_proof,
verify_proof_and_pinned_nodes,
},
};
use commonware_cryptography::{Sha256, sha256::Digest};
use commonware_macros::test_traced;
use commonware_parallel::Sequential;
use commonware_runtime::{
BufferPooler, Runner as _, Spawner as _, Supervisor as _,
buffer::paged::CacheRef,
deterministic,
mocks::{DelayedSyncContext, PendingSyncs, drive_pending_syncs, fail_pending_syncs},
reschedule,
};
use commonware_utils::{NZU16, NZU64, NZUsize, sequence::U64};
use core::future::Future;
use futures::FutureExt as _;
use std::num::{NonZeroU16, NonZeroUsize};
type TestDb<F> = Db<F, deterministic::Context, FixedEncoding<U64>, Sha256, (), Sequential>;
const WITNESS_PAGE_SIZE: NonZeroU16 = NZU16!(77);
const WITNESS_PAGE_CACHE_SIZE: NonZeroUsize = NZUsize!(9);
fn witness_config(partition: &str, pooler: &impl BufferPooler) -> JournalConfig<()> {
JournalConfig {
partition: format!("{partition}-witness"),
items_per_section: NZU64!(64),
compression: None,
codec_config: (),
page_cache: CacheRef::from_pooler(pooler, WITNESS_PAGE_SIZE, WITNESS_PAGE_CACHE_SIZE),
write_buffer: NZUsize!(1024),
replay_buffer: NZUsize!(1024),
}
}
async fn open_db<F: Family>(context: deterministic::Context, partition: &str) -> TestDb<F> {
let witness_cfg = witness_config(partition, &context);
let merkle = crate::merkle::compact::Merkle::new(Sequential);
Db::init_from_merkle(merkle, context.child("witness"), witness_cfg, ())
.await
.unwrap()
}
async fn compact_operations_and_proof_inner<F: Family>(context: deterministic::Context) {
let db = open_db::<F>(context.child("db"), "keyless-operations-and-proof").await;
let mut seed = db.new_batch();
for value in 1u64..=6 {
seed = seed.append(U64::new(value));
}
let seed = seed
.merkleize(&db, Some(U64::new(7)), Location::new(0))
.await;
let (db, _) = db.apply_batch(seed).await.unwrap();
let db = db.sync().await.unwrap();
let floor = db.size();
assert_eq!(floor, Location::new(8));
assert!(matches!(
db.to_batch().proof(&db),
Err(Error::Merkle(crate::merkle::Error::Empty))
));
let mut parent = db.new_batch();
for value in 8u64..=12 {
parent = parent.append(U64::new(value));
}
let parent = parent.merkleize(&db, Some(U64::new(13)), floor).await;
let child = parent
.new_batch::<Sha256>()
.append(U64::new(14))
.merkleize(&db, Some(U64::new(15)), floor)
.await;
let (child_start, child_ops) = child.operations();
let (_, child_ops_again) = child.operations();
let child_end = child.bounds().tip.size;
assert!(Arc::ptr_eq(&child_ops, &child_ops_again));
assert_eq!(child_start, parent.bounds().tip.size);
assert_eq!(*child_start + child_ops.len() as u64, *child_end);
assert_eq!(child_end, Location::new(16));
assert!(matches!(
child_ops.as_slice(),
[Operation::Append(value), Operation::Commit(Some(metadata), operation_floor)]
if value == &U64::new(14)
&& metadata == &U64::new(15)
&& operation_floor == &floor
));
let child_root = child.root();
let child_proof = child.proof(&db).unwrap();
let child_pins = child.pinned_nodes(&db).unwrap();
assert_eq!(child_proof.leaves, child_end);
assert_eq!(
child_proof.inactive_peaks,
F::inactive_peaks(child_end, floor),
);
assert!(verify_proof::<Sha256, _, _>(
&child_proof,
child_start,
&child_ops,
&child_root
));
assert!(verify_proof_and_pinned_nodes::<Sha256, _, _>(
&child_proof,
child_start,
&child_ops,
&child_pins,
&child_root
));
assert!(child_pins.len() > 1);
let mut reordered_child_pins = child_pins.clone();
reordered_child_pins.swap(0, 1);
assert!(!verify_proof_and_pinned_nodes::<Sha256, _, _>(
&child_proof,
child_start,
&child_ops,
&reordered_child_pins,
&child_root
));
let (db, _) = db.apply_batch(parent).await.unwrap();
let child_proof_after = child.proof(&db).unwrap();
assert_eq!(child.pinned_nodes(&db).unwrap(), child_pins);
assert!(verify_proof_and_pinned_nodes::<Sha256, _, _>(
&child_proof_after,
child_start,
&child_ops,
&child_pins,
&child_root
));
let (db, child_range) = db.apply_batch(child).await.unwrap();
assert_eq!(child_range, child_start..child_end);
let commit_floor = db.size();
let commit_only = db
.new_batch()
.merkleize(&db, Some(U64::new(16)), commit_floor)
.await;
let (commit_start, commit_ops) = commit_only.operations();
let commit_end = commit_only.bounds().tip.size;
let commit_root = commit_only.root();
let commit_proof = commit_only.proof(&db).unwrap();
let commit_pins = commit_only.pinned_nodes(&db).unwrap();
assert_eq!(commit_start, commit_floor);
assert!(matches!(
commit_ops.as_slice(),
[Operation::Commit(Some(metadata), operation_floor)]
if metadata == &U64::new(16) && operation_floor == &commit_floor
));
assert_eq!(*commit_start + commit_ops.len() as u64, *commit_end);
assert_eq!(commit_proof.leaves, commit_end);
assert_eq!(
commit_proof.inactive_peaks,
F::inactive_peaks(commit_end, commit_floor)
);
assert!(verify_proof_and_pinned_nodes::<Sha256, _, _>(
&commit_proof,
commit_start,
&commit_ops,
&commit_pins,
&commit_root
));
let (db, commit_range) = db.apply_batch(commit_only).await.unwrap();
assert_eq!(commit_range, commit_start..commit_end);
let db = db.sync().await.unwrap();
let late_parent = db
.new_batch()
.append(U64::new(17))
.merkleize(&db, Some(U64::new(18)), db.size())
.await;
let late = late_parent
.new_batch::<Sha256>()
.append(U64::new(19))
.merkleize(&db, Some(U64::new(20)), db.size())
.await;
let (db, _) = db.apply_batch(Arc::clone(&late_parent)).await.unwrap();
assert!(matches!(
late_parent.proof(&db),
Err(Error::Merkle(crate::merkle::Error::ElementPruned(_)))
));
assert!(matches!(
late_parent.pinned_nodes(&db),
Err(Error::Merkle(crate::merkle::Error::ElementPruned(_)))
));
drop(late_parent);
let (late_start, late_ops) = late.operations();
let late_root = late.root();
let late_proof = late.proof(&db).unwrap();
let late_pins = late.pinned_nodes(&db).unwrap();
assert!(verify_proof_and_pinned_nodes::<Sha256, _, _>(
&late_proof,
late_start,
&late_ops,
&late_pins,
&late_root
));
let (db, _) = db.apply_batch(late).await.unwrap();
db.destroy().await.unwrap();
}
#[test_traced]
fn test_compact_operations_and_proof_mmr() {
deterministic::Runner::default().start(compact_operations_and_proof_inner::<mmr::Family>);
}
#[test_traced]
fn test_compact_operations_and_proof_mmb() {
deterministic::Runner::default().start(compact_operations_and_proof_inner::<mmb::Family>);
}
async fn open_witness_journal(
context: deterministic::Context,
partition: &str,
) -> witness::Journal<deterministic::Context, mmr::Family, Digest> {
let cfg = witness_config(partition, &context);
witness::Journal::init(context, cfg).await.unwrap()
}
#[test_traced("INFO")]
fn test_serve_refuses_requests_outside_witness() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "keyless-serve-refusal").await;
let floor = db.inactivity_floor_loc();
let batch = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, Some(U64::new(11)), floor)
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let n = db.target().size;
let boundary = |size: Location<mmr::Family>, start: Location<mmr::Family>| {
Request::Boundary { size, start }
};
let operations =
|size: Location<mmr::Family>, start: Location<mmr::Family>| Request::Operations {
size,
start,
max_ops: NZU64!(1),
};
let beyond = n + 1;
assert!(matches!(
db.serve(boundary(beyond, n)).await,
Err(Error::Merkle(crate::merkle::Error::RangeOutOfBounds(_)))
));
assert!(matches!(
db.serve(operations(Location::new(0), Location::new(0)))
.await,
Err(Error::Merkle(crate::merkle::Error::RangeOutOfBounds(_)))
));
assert!(matches!(
db.serve(operations(n - 1, n - 2)).await,
Err(Error::Journal(crate::journal::Error::ItemPruned(_)))
));
assert!(matches!(
db.serve(boundary(n, n)).await,
Err(Error::Merkle(crate::merkle::Error::RangeOutOfBounds(_)))
));
assert!(matches!(
db.serve(boundary(n, n - 2)).await,
Err(Error::Journal(crate::journal::Error::ItemPruned(_)))
));
let (response, feedback_tx) = db
.serve(Request::Operations {
size: n,
start: n - 1,
max_ops: NZU64!(5),
})
.await
.unwrap();
assert!(feedback_tx.is_none());
let Response::Operations { operations, .. } = response else {
panic!("operations request should get an operations response");
};
assert_eq!(operations.len(), 1);
let (response, _) = db.serve(boundary(n, n - 1)).await.unwrap();
assert!(matches!(response, Response::Boundary { .. }));
});
}
type DelayedDb = Db<
mmr::Family,
DelayedSyncContext<deterministic::Context>,
FixedEncoding<U64>,
Sha256,
(),
Sequential,
>;
fn open_delayed_db(
context: &deterministic::Context,
label: &'static str,
partition: &str,
pending: &PendingSyncs,
) -> impl Future<Output = Result<DelayedDb, Error<mmr::Family>>> {
let witness_cfg = witness_config(partition, context);
let merkle = crate::merkle::compact::Merkle::new(Sequential);
let context = DelayedSyncContext {
inner: context.child(label),
pending: pending.clone(),
};
DelayedDb::init_from_merkle(merkle, context.child("witness"), witness_cfg, ())
}
async fn apply_append(db: DelayedDb, seed: u64) -> DelayedDb {
let floor = db.inactivity_floor_loc();
let batch = db
.new_batch()
.append(U64::new(seed))
.merkleize(&db, Some(U64::new(seed)), floor)
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
db
}
async fn fail_dropped_watermark_sync(mut db: DelayedDb, pending: &PendingSyncs) -> DelayedDb {
let first;
(db, first) = db.start_sync().await.unwrap();
drive_pending_syncs(pending, first).await.unwrap();
let dropped;
(db, dropped) = db.start_sync().await.unwrap();
assert_eq!(
pending.lock().len(),
1,
"expected only the recovery-watermark sync"
);
fail_pending_syncs(pending);
drop(dropped);
db
}
#[test_traced]
fn test_compact_apply_overlaps_start_sync() {
deterministic::Runner::default().start(|ctx| async move {
let partition = "keyless-start-sync-overlap";
let pending = PendingSyncs::default();
let open = open_delayed_db(&ctx, "delayed", partition, &pending);
let mut db = drive_pending_syncs(&pending, open).await.unwrap();
db = apply_append(db, 1).await;
let starts_before = pending.starts();
let entered_before = pending.entered();
let completions_before = pending.completions();
let handle;
(db, handle) = db.start_sync().await.unwrap();
assert!(pending.starts() > starts_before);
assert_eq!(pending.completions(), completions_before);
let first_target = db.target();
let waiter = ctx
.child("await_sync")
.spawn(|_| async move { handle.await.unwrap() });
while pending.entered() == entered_before {
reschedule().await;
}
db = apply_append(db, 2).await;
assert_ne!(db.root(), first_target.root);
let second_target = db.target();
assert_ne!(second_target, first_target);
assert_eq!(
pending.completions(),
completions_before,
"the database made progress while the sync was still in flight"
);
pending.unblock();
waiter.await.unwrap();
let handle;
(db, handle) = db.start_sync().await.unwrap();
handle.await.unwrap();
drop(db);
let db = open_delayed_db(&ctx, "reopen", partition, &pending)
.await
.unwrap();
assert_eq!(db.target(), second_target);
assert_eq!(db.get_metadata(), Some(U64::new(2)));
let db = db.rewind(first_target.size).await.unwrap();
assert_eq!(db.target(), first_target);
db.destroy().await.unwrap();
});
}
#[test_traced]
fn test_compact_start_sync_recovery() {
deterministic::Runner::default().start(|ctx| async move {
let partition = "keyless-start-sync-recovery";
let pending = PendingSyncs::default();
pending.unblock();
let mut db = open_delayed_db(&ctx, "delayed", partition, &pending)
.await
.unwrap();
db = apply_append(db, 1).await;
let handle;
(db, handle) = db.start_sync().await.unwrap();
handle.await.unwrap();
let root = db.root();
drop(db);
let db = open_delayed_db(&ctx, "reopen", partition, &pending)
.await
.unwrap();
assert_eq!(db.root(), root);
assert_eq!(db.get_metadata(), Some(U64::new(1)));
db.destroy().await.unwrap();
});
}
#[test_traced]
fn test_compact_start_sync_failure_propagates() {
deterministic::Runner::default().start(|ctx| async move {
let pending = PendingSyncs::default();
pending.unblock();
let mut db = open_delayed_db(&ctx, "delayed", "keyless-start-sync-fail", &pending)
.await
.unwrap();
db = apply_append(db, 1).await;
pending.arm_fail();
let handle;
(db, handle) = db.start_sync().await.unwrap();
assert!(
handle.await.is_err(),
"the sync handle surfaces the failure"
);
let starts_before = pending.starts();
assert!(
db.commit().await.is_err(),
"the next durability op surfaces the failed in-flight sync"
);
assert_eq!(
pending.starts(),
starts_before,
"the surfaced error is the retained failure, not a fresh sync's"
);
});
}
#[test_traced]
fn test_compact_start_sync_then_noop_sync_drains() {
deterministic::Runner::default().start(|ctx| async move {
let partition = "keyless-start-sync-noop-drain";
let pending = PendingSyncs::default();
let open = open_delayed_db(&ctx, "delayed", partition, &pending);
let mut db = drive_pending_syncs(&pending, open).await.unwrap();
db = apply_append(db, 1).await;
let starts_before = pending.starts();
let completions_before = pending.completions();
let handle;
(db, handle) = db.start_sync().await.unwrap();
assert!(pending.starts() > starts_before);
assert_eq!(pending.completions(), completions_before);
let root = db.root();
let db = {
let mut sync = std::pin::pin!(db.sync());
assert!(
sync.as_mut().now_or_never().is_none(),
"sync proceeded while the started sync was pending"
);
pending.unblock();
sync.await.unwrap()
};
handle.await.unwrap();
assert!(pending.completions() > completions_before);
drop(db);
let db = open_delayed_db(&ctx, "reopen", partition, &pending)
.await
.unwrap();
assert_eq!(db.root(), root);
db.destroy().await.unwrap();
});
}
#[test_traced]
fn test_compact_start_sync_then_noop_sync_fails_without_new_work() {
deterministic::Runner::default().start(|ctx| async move {
let pending = PendingSyncs::default();
let open = open_delayed_db(
&ctx,
"delayed",
"keyless-start-sync-noop-sync-fail",
&pending,
);
let mut db = drive_pending_syncs(&pending, open).await.unwrap();
db = apply_append(db, 1).await;
db = fail_dropped_watermark_sync(db, &pending).await;
let starts_before = pending.starts();
assert!(
drive_pending_syncs(&pending, db.sync()).await.is_err(),
"sync absorbed the retained metadata failure"
);
assert_eq!(
pending.starts(),
starts_before,
"sync started new journal work before returning the retained failure"
);
});
}
#[test_traced]
fn test_compact_start_sync_then_noop_commit_waits() {
deterministic::Runner::default().start(|ctx| async move {
let pending = PendingSyncs::default();
let open = open_delayed_db(&ctx, "delayed", "keyless-start-sync-noop-commit", &pending);
let mut db = drive_pending_syncs(&pending, open).await.unwrap();
db = apply_append(db, 1).await;
let handle;
(db, handle) = db.start_sync().await.unwrap();
let starts_before = pending.starts();
let db = {
let mut commit = std::pin::pin!(db.commit());
assert!(
commit.as_mut().now_or_never().is_none(),
"commit proceeded while the started sync was pending"
);
pending.unblock();
commit.await.unwrap()
};
handle.await.unwrap();
assert_eq!(
pending.starts(),
starts_before,
"a successful pipelined sync still triggered journal work"
);
db.destroy().await.unwrap();
});
}
#[test_traced]
fn test_compact_start_sync_noop_second_call() {
deterministic::Runner::default().start(|ctx| async move {
let partition = "keyless-start-sync-noop-second";
let pending = PendingSyncs::default();
let open = open_delayed_db(&ctx, "delayed", partition, &pending);
let mut db = drive_pending_syncs(&pending, open).await.unwrap();
db = apply_append(db, 1).await;
let h1;
(db, h1) = db.start_sync().await.unwrap();
pending.unblock();
h1.await.unwrap();
let h2;
(db, h2) = db.start_sync().await.unwrap();
h2.await.unwrap();
let root = db.root();
drop(db);
let journal = open_witness_journal(ctx.child("probe"), partition).await;
assert_eq!(journal.size(), 2);
drop(journal);
let db = open_delayed_db(&ctx, "reopen", partition, &pending)
.await
.unwrap();
assert_eq!(db.root(), root);
db.destroy().await.unwrap();
});
}
#[test_traced]
fn test_compact_start_sync_rewind_fast_path_drains() {
deterministic::Runner::default().start(|ctx| async move {
let partition = "keyless-start-sync-rewind-drain";
let pending = PendingSyncs::default();
let open = open_delayed_db(&ctx, "delayed", partition, &pending);
let mut db = drive_pending_syncs(&pending, open).await.unwrap();
db = apply_append(db, 1).await;
let handle;
(db, handle) = db.start_sync().await.unwrap();
let root = db.root();
let size = db.size();
let starts_before = pending.starts();
let db = {
let mut rewind = std::pin::pin!(db.rewind(size));
assert!(
rewind.as_mut().now_or_never().is_none(),
"rewind proceeded while the started sync was pending"
);
pending.unblock();
rewind.await.unwrap()
};
handle.await.unwrap();
assert_eq!(
pending.starts(),
starts_before,
"the fast path started journal work instead of adopting the proven sync"
);
assert_eq!(db.root(), root);
drop(db);
let db = open_delayed_db(&ctx, "reopen", partition, &pending)
.await
.unwrap();
assert_eq!(db.root(), root);
db.destroy().await.unwrap();
});
}
#[test_traced]
fn test_compact_start_sync_rewind_fast_path_fails() {
deterministic::Runner::default().start(|ctx| async move {
let pending = PendingSyncs::default();
pending.unblock();
let mut db =
open_delayed_db(&ctx, "delayed", "keyless-start-sync-rewind-fail", &pending)
.await
.unwrap();
db = apply_append(db, 1).await;
pending.arm_fail();
let handle;
(db, handle) = db.start_sync().await.unwrap();
assert!(handle.await.is_err());
let size = db.size();
assert!(
db.rewind(size).await.is_err(),
"rewind reported an unproven tip as durable"
);
});
}
#[test_traced]
fn test_compact_start_sync_metadata_failure_resurfaces_on_commit() {
deterministic::Runner::default().start(|ctx| async move {
let pending = PendingSyncs::default();
let open = open_delayed_db(&ctx, "delayed", "keyless-start-sync-meta-fail", &pending);
let mut db = drive_pending_syncs(&pending, open).await.unwrap();
db = apply_append(db, 1).await;
let handle;
(db, handle) = db.start_sync().await.unwrap();
{
let mut parked = pending.lock();
assert_eq!(
parked.len(),
2,
"expected the data and offsets syncs parked"
);
let offsets = parked.pop().unwrap();
let data = parked.pop().unwrap();
data.release.send(Ok(())).unwrap();
offsets
.release
.send(Err(commonware_runtime::Error::Io(
std::io::Error::other("injected sync failure").into(),
)))
.unwrap();
}
assert!(
handle.await.is_err(),
"the sync handle surfaces the failure"
);
pending.unblock();
let db = apply_append(db, 2).await;
assert!(
db.commit().await.is_err(),
"commit absorbed the retained metadata failure"
);
});
}
#[test_traced]
fn test_compact_start_sync_retains_dropped_metadata_failure() {
deterministic::Runner::default().start(|ctx| async move {
let pending = PendingSyncs::default();
let open = open_delayed_db(
&ctx,
"delayed",
"keyless-start-sync-dropped-meta-fail",
&pending,
);
let mut db = drive_pending_syncs(&pending, open).await.unwrap();
db = apply_append(db, 1).await;
db = fail_dropped_watermark_sync(db, &pending).await;
let next;
(db, next) = db.start_sync().await.unwrap();
assert!(
drive_pending_syncs(&pending, next).await.is_err(),
"a later start_sync masked the retained metadata failure"
);
drop(db);
});
}
#[test_traced]
fn test_compact_start_sync_proven_skips_journal() {
deterministic::Runner::default().start(|ctx| async move {
let pending = PendingSyncs::default();
pending.unblock();
let mut db = open_delayed_db(&ctx, "delayed", "keyless-start-sync-proven", &pending)
.await
.unwrap();
db = apply_append(db, 1).await;
let handle;
(db, handle) = db.start_sync().await.unwrap();
handle.await.unwrap();
let starts_before = pending.starts();
let db = db.commit().await.unwrap();
let size = db.size();
let db = db.rewind(size).await.unwrap();
assert_eq!(
pending.starts(),
starts_before,
"a proven pipelined sync still triggered journal work"
);
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_start_sync_persists_import() {
deterministic::Runner::default().start(|context| async move {
let dst = "keyless-import-start-sync-dst";
let src = "keyless-import-start-sync-src";
let meta_a = U64::new(11);
let meta_b = U64::new(22);
let target_b = {
let source = open_db::<mmr::Family>(context.child("src"), src).await;
let batch = source
.new_batch()
.append(U64::new(2))
.merkleize(&source, Some(meta_b.clone()), Location::new(0))
.await;
let (source, _) = source.apply_batch(batch).await.unwrap();
let source = source.sync().await.unwrap();
source.target()
};
let (_, size_b, pinned_b) = {
let journal = open_witness_journal(context.child("src_tip"), src).await;
witness::tests::tip(&journal).await
};
assert_eq!(size_b, target_b.size);
{
let seeded = open_db::<mmr::Family>(context.child("seed"), dst).await;
let batch = seeded
.new_batch()
.append(U64::new(1))
.merkleize(&seeded, Some(meta_a), Location::new(0))
.await;
let (seeded, _) = seeded.apply_batch(batch).await.unwrap();
let seeded = seeded.sync().await.unwrap();
assert_ne!(seeded.target(), target_b);
}
{
let journal = open_witness_journal(context.child("import"), dst).await;
let imported = TestDb::<mmr::Family>::init_from_sync(
Sequential,
journal,
(),
size_b - 1,
pinned_b,
Operation::Commit(Some(meta_b.clone()), Location::new(0)),
)
.unwrap();
assert_eq!(imported.target(), target_b);
let (_imported, handle) = imported.start_sync().await.unwrap();
handle.await.unwrap();
}
let db = open_db::<mmr::Family>(context.child("reopen"), dst).await;
assert_eq!(db.target(), target_b);
assert_eq!(db.root(), target_b.root);
assert_eq!(db.get_metadata(), Some(meta_b));
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_stale_batch_rejected() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "keyless-stale").await;
let floor = db.inactivity_floor_loc();
let batch_a = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, Some(U64::new(11)), floor)
.await;
let batch_b = db
.new_batch()
.append(U64::new(2))
.merkleize(&db, Some(U64::new(22)), floor)
.await;
let expected_root = batch_a.root();
let (db, _) = db.apply_batch(batch_a).await.unwrap();
assert_eq!(db.root(), expected_root);
assert!(matches!(
db.apply_batch(batch_b).await,
Err(Error::StaleBatch)
));
});
}
#[test_traced("INFO")]
fn test_compact_delayed_merkleize_after_ancestor_apply() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "keyless-delayed-child").await;
let floor = db.inactivity_floor_loc();
let a = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, None, floor)
.await;
let b = a
.new_batch::<Sha256>()
.append(U64::new(2))
.merkleize(&db, None, floor)
.await;
let c = b.new_batch::<Sha256>().append(U64::new(3));
let (db, _) = db.apply_batch(a).await.unwrap();
let c = c.merkleize(&db, None, floor).await;
let expected_root = c.root();
let (db, _) = db.apply_batch(c).await.unwrap();
assert_eq!(db.root(), expected_root);
});
}
#[test_traced("INFO")]
fn test_compact_to_batch_reflects_live_state() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "keyless-to-batch-live").await;
let floor = db.inactivity_floor_loc();
let pre_apply_root = db.root();
let pre_snapshot = db.to_batch();
assert_eq!(
pre_snapshot.root(),
pre_apply_root,
"snapshot before any mutation should match the live root"
);
let batch = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, Some(U64::new(11)), floor)
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let live_root = db.root();
assert_ne!(
live_root, pre_apply_root,
"applying a non-empty batch must change the live root"
);
let snapshot = db.to_batch();
assert_eq!(
snapshot.root(),
live_root,
"to_batch().root() must match the live db.root() even before sync"
);
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_stale_batch_chained() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "keyless-chained-stale").await;
let floor = db.inactivity_floor_loc();
let common_parent = db
.new_batch()
.append(U64::new(10))
.merkleize(&db, Some(U64::new(110)), floor)
.await;
let sibling_a = common_parent
.new_batch::<Sha256>()
.append(U64::new(11))
.merkleize(&db, Some(U64::new(111)), floor)
.await;
let sibling_b = common_parent
.new_batch::<Sha256>()
.append(U64::new(12))
.merkleize(&db, Some(U64::new(112)), floor)
.await;
let (db, _) = db.apply_batch(sibling_a).await.unwrap();
assert!(matches!(
db.validate_batch(&sibling_b),
Err(Error::StaleBatch)
));
let parent_a = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, Some(U64::new(11)), floor)
.await;
let parent_b = db
.new_batch()
.append(U64::new(2))
.merkleize(&db, Some(U64::new(22)), floor)
.await;
let child_b = parent_b
.new_batch::<Sha256>()
.append(U64::new(3))
.merkleize(&db, Some(U64::new(33)), floor)
.await;
let (db, _) = db.apply_batch(parent_a).await.unwrap();
assert!(matches!(
db.validate_batch(&child_b),
Err(Error::StaleBatch)
));
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_stale_parent_after_child_applied() {
deterministic::Runner::default().start(|context| async move {
let db =
open_db::<mmr::Family>(context.child("db"), "keyless-child-before-parent").await;
let floor = db.inactivity_floor_loc();
let parent = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, Some(U64::new(11)), floor)
.await;
let child = parent
.new_batch::<Sha256>()
.append(U64::new(2))
.merkleize(&db, Some(U64::new(22)), floor)
.await;
let (db, _) = db.apply_batch(child).await.unwrap();
assert!(matches!(
db.apply_batch(parent).await,
Err(Error::StaleBatch)
));
});
}
#[test_traced("INFO")]
fn test_compact_sequential_commit_parent_then_child() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "keyless-parent-child").await;
let floor = db.inactivity_floor_loc();
let parent = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, Some(U64::new(11)), floor)
.await;
let child = parent
.new_batch::<Sha256>()
.append(U64::new(2))
.merkleize(&db, Some(U64::new(22)), floor)
.await;
let expected_root = child.root();
let (db, _) = db.apply_batch(parent).await.unwrap();
let (db, _) = db.apply_batch(child).await.unwrap();
let db = db.sync().await.unwrap();
assert_eq!(db.root(), expected_root);
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_floor_regressed() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "keyless-floor-regressed").await;
let advance_floor = db.new_batch().append(U64::new(1));
let advance_floor = advance_floor.merkleize(&db, None, Location::new(1)).await;
let (db, _) = db.apply_batch(advance_floor).await.unwrap();
let db = db.sync().await.unwrap();
let target = db.target();
let regressed = db
.new_batch()
.append(U64::new(2))
.merkleize(&db, None, Location::new(0))
.await;
assert!(matches!(
db.apply_batch(regressed).await,
Err(Error::FloorRegressed(new, current))
if new == Location::new(0) && current == Location::new(1)
));
let db =
open_db::<mmr::Family>(context.child("reopen"), "keyless-floor-regressed").await;
assert_eq!(db.target(), target);
});
}
#[test_traced("INFO")]
fn test_compact_ancestor_floor_regressed() {
deterministic::Runner::default().start(|context| async move {
let db =
open_db::<mmr::Family>(context.child("db"), "keyless-ancestor-floor-regressed")
.await;
let parent = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, None, Location::new(2))
.await;
let child = parent
.new_batch::<Sha256>()
.append(U64::new(2))
.merkleize(&db, None, Location::new(1))
.await;
let target = db.target();
assert!(matches!(
db.apply_batch(child).await,
Err(Error::FloorRegressed(new, prev))
if new == Location::new(1) && prev == Location::new(2)
));
let db =
open_db::<mmr::Family>(context.child("reopen"), "keyless-ancestor-floor-regressed")
.await;
assert_eq!(db.target(), target);
});
}
#[test_traced("INFO")]
fn test_compact_rewind_restores_commit_metadata_and_floor() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "keyless-rewind-meta").await;
let v1 = U64::new(1);
let meta1 = U64::new(11);
let floor1 = Location::new(0);
let batch = db
.new_batch()
.append(v1)
.merkleize(&db, Some(meta1.clone()), floor1)
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let root_after_first = db.root();
let size_after_first = db.size();
let v2 = U64::new(2);
let meta2 = U64::new(22);
let floor2 = Location::new(1);
let batch = db
.new_batch()
.append(v2)
.merkleize(&db, Some(meta2.clone()), floor2)
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
assert_eq!(db.get_metadata(), Some(meta2));
assert_eq!(db.inactivity_floor_loc(), floor2);
let db = db.rewind(size_after_first).await.unwrap();
assert_eq!(db.root(), root_after_first);
assert_eq!(db.get_metadata(), Some(meta1));
assert_eq!(db.inactivity_floor_loc(), floor1);
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_rewind_persists_across_reopen() {
deterministic::Runner::default().start(|context| async move {
let partition = "keyless-rewind-reopen";
let meta1 = U64::new(11);
let floor1 = Location::new(0);
let meta2 = U64::new(22);
let floor2 = Location::new(1);
let root_after_first = {
let db = open_db::<mmr::Family>(context.child("first"), partition).await;
let batch = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, Some(meta1.clone()), floor1)
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let root = db.root();
let size_after_first = db.size();
let batch = db
.new_batch()
.append(U64::new(2))
.merkleize(&db, Some(meta2), floor2)
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let _db = db.rewind(size_after_first).await.unwrap();
root
};
let db = open_db::<mmr::Family>(context.child("second"), partition).await;
assert_eq!(db.root(), root_after_first);
assert_eq!(db.get_metadata(), Some(meta1));
assert_eq!(db.inactivity_floor_loc(), floor1);
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_commit_persists_across_reopen() {
deterministic::Runner::default().start(|context| async move {
let partition = "keyless-commit-reopen";
let meta1 = U64::new(11);
let meta2 = U64::new(22);
let root_after_second = {
let db = open_db::<mmr::Family>(context.child("first"), partition).await;
let batch = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, Some(meta1), Location::new(0))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.commit().await.unwrap();
let batch = db
.new_batch()
.append(U64::new(2))
.merkleize(&db, Some(meta2.clone()), Location::new(1))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.commit().await.unwrap();
db.root()
};
let db = open_db::<mmr::Family>(context.child("second"), partition).await;
assert_eq!(db.root(), root_after_second);
assert_eq!(db.get_metadata(), Some(meta2));
assert_eq!(db.inactivity_floor_loc(), Location::new(1));
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_rewind_to_committed_entry_after_reopen() {
deterministic::Runner::default().start(|context| async move {
let partition = "keyless-commit-rewind-reopen";
let meta1 = U64::new(11);
let meta2 = U64::new(22);
let (root_a, size_a) = {
let db = open_db::<mmr::Family>(context.child("first"), partition).await;
let batch = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, Some(meta1.clone()), Location::new(0))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.commit().await.unwrap();
let root_a = db.root();
let size_a = db.size();
let batch = db
.new_batch()
.append(U64::new(2))
.merkleize(&db, Some(meta2), Location::new(1))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let _db = db.commit().await.unwrap();
(root_a, size_a)
};
let db = open_db::<mmr::Family>(context.child("second"), partition).await;
let db = db.rewind(size_a).await.unwrap();
assert_eq!(db.root(), root_a);
assert_eq!(db.get_metadata(), Some(meta1));
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_sync_after_commit() {
deterministic::Runner::default().start(|context| async move {
let partition = "keyless-sync-after-commit";
let meta = U64::new(11);
let root = {
let db = open_db::<mmr::Family>(context.child("first"), partition).await;
let batch = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, Some(meta.clone()), Location::new(0))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.commit().await.unwrap();
let db = db.sync().await.unwrap();
db.root()
};
let db = open_db::<mmr::Family>(context.child("second"), partition).await;
assert_eq!(db.root(), root);
assert_eq!(db.get_metadata(), Some(meta));
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_import_persists_with_commit() {
deterministic::Runner::default().start(|context| async move {
let dst = "keyless-import-commit-dst";
let src = "keyless-import-commit-src";
let meta_a = U64::new(11);
let meta_b = U64::new(22);
let (target_b, pinned_b) = {
let source = open_db::<mmr::Family>(context.child("src"), src).await;
let batch = source
.new_batch()
.append(U64::new(2))
.merkleize(&source, Some(meta_b.clone()), Location::new(0))
.await;
let (source, _) = source.apply_batch(batch).await.unwrap();
let source = source.sync().await.unwrap();
let target = source.target();
let (response, _) = source
.serve(Request::Boundary {
size: target.size,
start: target.size - 1,
})
.await
.unwrap();
let Response::Boundary { pinned_nodes, .. } = response else {
panic!("boundary request should get a boundary response");
};
(target, pinned_nodes)
};
{
let seeded = open_db::<mmr::Family>(context.child("seed"), dst).await;
let batch = seeded
.new_batch()
.append(U64::new(1))
.merkleize(&seeded, Some(meta_a), Location::new(0))
.await;
let (seeded, _) = seeded.apply_batch(batch).await.unwrap();
let seeded = seeded.sync().await.unwrap();
assert_ne!(seeded.target(), target_b);
}
{
let journal = open_witness_journal(context.child("import"), dst).await;
let imported = TestDb::<mmr::Family>::init_from_sync(
Sequential,
journal,
(),
target_b.size - 1,
pinned_b,
Operation::Commit(Some(meta_b.clone()), Location::new(0)),
)
.unwrap();
assert_eq!(imported.target(), target_b);
let _imported = imported.commit().await.unwrap();
}
let db = open_db::<mmr::Family>(context.child("reopen"), dst).await;
assert_eq!(db.target(), target_b);
assert_eq!(db.root(), target_b.root);
assert_eq!(db.get_metadata(), Some(meta_b));
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_reopen_rejects_tampered_witness() {
deterministic::Runner::default().start(|context| async move {
let partition = "keyless-witness-tamper";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.append(U64::new(7))
.merkleize(&db, Some(U64::new(11)), Location::new(1))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
drop(db);
let journal = open_witness_journal(context.child("tamper"), partition).await;
let (op_bytes, size, mut pinned_nodes) = witness::tests::tip(&journal).await;
pinned_nodes.push(Sha256::fill(0xff));
witness::tests::overwrite_tip(journal, op_bytes, size, pinned_nodes).await;
let merkle = crate::merkle::compact::Merkle::new(Sequential);
let reopened = TestDb::<mmr::Family>::init_from_merkle(
merkle,
context.child("reopen_witness"),
witness_config(partition, &context),
(),
)
.await;
assert!(matches!(reopened, Err(Error::DataCorrupted(_))));
});
}
#[test_traced("INFO")]
fn test_compact_rewind_rejects_corrupt_target_entry() {
deterministic::Runner::default().start(|context| async move {
let partition = "keyless-corrupt-rewind-target";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, None, Location::new(1))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let rewind_target = db.target().size;
let batch = db
.new_batch()
.append(U64::new(2))
.merkleize(&db, None, Location::new(1))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let tip_target = db.target();
drop(db);
let mut journal = open_witness_journal(context.child("corrupt"), partition).await;
journal = witness::tests::corrupt_entry(journal, 1, |entry| {
entry.pinned_nodes.push(Sha256::fill(0xff));
})
.await;
drop(journal);
let merkle = crate::merkle::compact::Merkle::new(Sequential);
let reopened = TestDb::<mmr::Family>::init_from_merkle(
merkle,
context.child("reopen"),
witness_config(partition, &context),
(),
)
.await
.unwrap();
assert_eq!(reopened.target(), tip_target);
assert!(matches!(
reopened.rewind(rewind_target).await,
Err(Error::DataCorrupted(_))
));
let merkle = crate::merkle::compact::Merkle::new(Sequential);
let reopened = TestDb::<mmr::Family>::init_from_merkle(
merkle,
context.child("reopen2"),
witness_config(partition, &context),
(),
)
.await
.unwrap();
assert_eq!(reopened.target(), tip_target);
reopened.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_reopen_rejects_interrupted_import() {
deterministic::Runner::default().start(|context| async move {
let partition = "keyless-interrupted-import";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.append(U64::new(7))
.merkleize(&db, Some(U64::new(11)), Location::new(1))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
drop(db);
let journal = open_witness_journal(context.child("clear"), partition).await;
let size = journal.size();
let journal = journal.clear_to_size(size.max(1)).await.unwrap();
drop(journal);
let merkle = crate::merkle::compact::Merkle::new(Sequential);
let reopened = TestDb::<mmr::Family>::init_from_merkle(
merkle,
context.child("reopen_witness"),
witness_config(partition, &context),
(),
)
.await;
assert!(matches!(reopened, Err(Error::Journal(_))));
});
}
#[test_traced("INFO")]
fn test_compact_reopen_rejects_commit_floor_beyond_tip() {
deterministic::Runner::default().start(|context| async move {
let partition = "keyless-invalid-persisted-floor";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.append(U64::new(7))
.merkleize(&db, Some(U64::new(11)), Location::new(1))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
drop(db);
let oversized_floor = Location::new(10);
let journal = open_witness_journal(context.child("tamper"), partition).await;
let (_, size, pinned_nodes) = witness::tests::tip(&journal).await;
let bad_op = Operation::<mmr::Family, FixedEncoding<U64>>::Commit(
Some(U64::new(11)),
oversized_floor,
)
.encode()
.to_vec();
witness::tests::overwrite_tip(journal, bad_op, size, pinned_nodes).await;
let merkle = crate::merkle::compact::Merkle::new(Sequential);
let reopened = TestDb::<mmr::Family>::init_from_merkle(
merkle,
context.child("reopen_witness"),
witness_config(partition, &context),
(),
)
.await;
assert!(matches!(
reopened,
Err(Error::DataCorrupted("invalid compact witness"))
));
});
}
#[test_traced("INFO")]
fn test_compact_reopen_rejects_tampered_pinned_nodes() {
deterministic::Runner::default().start(|context| async move {
let partition = "keyless-pinned-nodes-tamper";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.append(U64::new(7))
.merkleize(&db, Some(U64::new(11)), Location::new(1))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let tampered_target = db.target();
drop(db);
let journal = open_witness_journal(context.child("tamper"), partition).await;
let (op_bytes, size, mut pinned_nodes) = witness::tests::tip(&journal).await;
pinned_nodes[0] = Sha256::fill(0xff);
witness::tests::overwrite_tip(journal, op_bytes, size, pinned_nodes).await;
let merkle = crate::merkle::compact::Merkle::new(Sequential);
let reopened = TestDb::<mmr::Family>::init_from_merkle(
merkle,
context.child("reopen_witness"),
witness_config(partition, &context),
(),
)
.await
.unwrap();
assert_ne!(reopened.target(), tampered_target);
reopened.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_rewind_to_current_is_noop() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "keyless-rewind-noop").await;
let batch = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, Some(U64::new(11)), Location::new(0))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let root = db.root();
let size = db.size();
let db = db.rewind(size).await.unwrap();
assert_eq!(db.root(), root);
assert_eq!(db.size(), size);
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_prune_past_tip_keeps_tip() {
deterministic::Runner::default().start(|context| async move {
let partition = "keyless-prune-past-tip";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, Some(U64::new(11)), Location::new(0))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let target = db.target();
let boundary = db.size() + 100;
let db = db.prune(boundary).await.unwrap();
assert_eq!(db.target(), target);
drop(db);
let reopened = open_db::<mmr::Family>(context.child("reopen"), partition).await;
assert_eq!(reopened.target(), target);
reopened.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_rewind_beyond_history() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "keyless-rewind-beyond").await;
assert!(matches!(
db.rewind(Location::new(0)).await,
Err(Error::Merkle(crate::merkle::Error::RewindBeyondHistory))
));
let db = open_db::<mmr::Family>(context.child("reopen"), "keyless-rewind-beyond").await;
let beyond_tip = db.size() + 100;
assert!(matches!(
db.rewind(beyond_tip).await,
Err(Error::Merkle(crate::merkle::Error::RewindBeyondHistory))
));
});
}
#[test_traced("INFO")]
fn test_compact_rewind_between_commits() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "keyless-rewind-between").await;
let floor = db.inactivity_floor_loc();
let batch = db
.new_batch()
.append(U64::new(1))
.append(U64::new(2))
.merkleize(&db, Some(U64::new(11)), floor)
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let root_a = db.root();
let size_a = db.size();
assert_eq!(size_a, Location::new(4));
let batch = db
.new_batch()
.append(U64::new(3))
.merkleize(&db, Some(U64::new(22)), floor)
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let root_b = db.root();
let mut db = db;
for target in [2u64, 3, 5] {
assert!(matches!(
db.rewind(Location::new(target)).await,
Err(Error::Merkle(crate::merkle::Error::RewindBeyondHistory))
));
db = open_db::<mmr::Family>(
context.child("reopen").with_attribute("target", target),
"keyless-rewind-between",
)
.await;
}
assert_eq!(db.root(), root_b);
let db = db.rewind(size_a).await.unwrap();
assert_eq!(db.root(), root_a);
assert_eq!(db.get_metadata(), Some(U64::new(11)));
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_reopen_drops_unsynced_witness() {
deterministic::Runner::default().start(|context| async move {
let partition = "keyless-witness-unsynced";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, Some(U64::new(11)), Location::new(1))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let target_a = db.target();
drop(db);
let journal = open_witness_journal(context.child("crash"), partition).await;
let (op_bytes, mut size, pinned_nodes) = witness::tests::tip(&journal).await;
size += 2;
witness::tests::append_unsynced(journal, op_bytes, size, pinned_nodes).await;
let reopened = open_db::<mmr::Family>(context.child("reopen"), partition).await;
assert_eq!(reopened.target(), target_a);
reopened.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_rewind_multiple_commits() {
deterministic::Runner::default().start(|context| async move {
let partition = "keyless-rewind-multi";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, Some(U64::new(11)), Location::new(0))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let root_a = db.root();
let size_a = db.size();
let target_a = db.target();
let mut db = db;
for i in [2u64, 3] {
let batch = db
.new_batch()
.append(U64::new(i))
.merkleize(&db, Some(U64::new(i * 11)), Location::new(0))
.await;
(db, _) = db.apply_batch(batch).await.unwrap();
db = db.sync().await.unwrap();
}
assert_ne!(db.root(), root_a);
let db = db.rewind(size_a).await.unwrap();
assert_eq!(db.root(), root_a);
assert_eq!(db.size(), size_a);
assert_eq!(db.get_metadata(), Some(U64::new(11)));
assert_eq!(db.target(), target_a);
drop(db);
let db = open_db::<mmr::Family>(context.child("reopen"), partition).await;
assert_eq!(db.root(), root_a);
assert_eq!(db.target(), target_a);
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_prune_then_rewind() {
deterministic::Runner::default().start(|context| async move {
let mut witness_cfg = witness_config("keyless-prune-rewind", &context);
witness_cfg.items_per_section = NZU64!(1);
let merkle = crate::merkle::compact::Merkle::new(Sequential);
let mut db: TestDb<mmr::Family> =
Db::init_from_merkle(merkle, context.child("witness"), witness_cfg.clone(), ())
.await
.unwrap();
let mut sizes = Vec::new();
for i in [1u64, 2, 3] {
let batch = db
.new_batch()
.append(U64::new(i))
.merkleize(&db, Some(U64::new(i * 11)), Location::new(0))
.await;
(db, _) = db.apply_batch(batch).await.unwrap();
db = db.sync().await.unwrap();
sizes.push(db.size());
}
let db = db.prune(sizes[1]).await.unwrap();
assert!(matches!(
db.rewind(sizes[0]).await,
Err(Error::Merkle(crate::merkle::Error::RewindBeyondHistory))
));
let merkle = crate::merkle::compact::Merkle::new(Sequential);
let db: TestDb<mmr::Family> = Db::init_from_merkle(
merkle,
context.child("witness").with_attribute("index", 2),
witness_cfg,
(),
)
.await
.unwrap();
let db = db.rewind(sizes[1]).await.unwrap();
assert_eq!(db.size(), sizes[1]);
assert_eq!(db.get_metadata(), Some(U64::new(22)));
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_rewind_preserves_pre_advance_batch() {
deterministic::Runner::default().start(|context| async move {
let db =
open_db::<mmr::Family>(context.child("db"), "keyless-rewind-preserves-pre-advance")
.await;
let batch = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, None, Location::new(0))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let size_after_first = db.size();
let held = db
.new_batch()
.append(U64::new(2))
.merkleize(&db, None, Location::new(0))
.await;
let batch = db
.new_batch()
.append(U64::new(3))
.merkleize(&db, None, Location::new(0))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let db = db.rewind(size_after_first).await.unwrap();
let (db, _) = db.apply_batch(held).await.unwrap();
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_noop_commit_after_commit() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "keyless-noop-after-commit").await;
let batch = db
.new_batch()
.append(U64::new(1))
.append(U64::new(2))
.merkleize(&db, Some(U64::new(11)), Location::new(0))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let root_after_first = db.root();
assert_eq!(db.size(), Location::new(4));
let db = db.sync().await.unwrap();
assert_eq!(db.size(), Location::new(4));
assert_eq!(db.root(), root_after_first);
assert_eq!(db.target().root, db.root());
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_noop_commit_after_reopen() {
deterministic::Runner::default().start(|context| async move {
let partition = "keyless-noop-after-reopen";
let root_before_drop = {
let db = open_db::<mmr::Family>(context.child("first"), partition).await;
let batch = db
.new_batch()
.append(U64::new(1))
.append(U64::new(2))
.merkleize(&db, Some(U64::new(11)), Location::new(0))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let root = db.root();
assert_eq!(db.size(), Location::new(4));
root
};
let db = open_db::<mmr::Family>(context.child("second"), partition).await;
assert_eq!(db.root(), root_before_drop);
assert_eq!(db.size(), Location::new(4));
let db = db.sync().await.unwrap();
assert_eq!(db.size(), Location::new(4));
assert_eq!(db.root(), root_before_drop);
assert_eq!(db.target().root, db.root());
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_noop_commit_after_rewind() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "keyless-noop-after-rewind").await;
let batch = db
.new_batch()
.append(U64::new(1))
.append(U64::new(2))
.merkleize(&db, Some(U64::new(11)), Location::new(0))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let root_after_first = db.root();
let batch = db
.new_batch()
.append(U64::new(3))
.merkleize(&db, Some(U64::new(22)), Location::new(1))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let db = db.rewind(Location::new(4)).await.unwrap();
assert_eq!(db.size(), Location::new(4));
assert_eq!(db.root(), root_after_first);
let db = db.sync().await.unwrap();
assert_eq!(db.size(), Location::new(4));
assert_eq!(db.root(), root_after_first);
assert_eq!(db.target().root, db.root());
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_rewind_makes_post_advance_batch_stale() {
deterministic::Runner::default().start(|context| async move {
let db =
open_db::<mmr::Family>(context.child("db"), "keyless-rewind-makes-stale").await;
let batch = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, None, Location::new(0))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let size_after_first = db.size();
let batch = db
.new_batch()
.append(U64::new(2))
.merkleize(&db, None, Location::new(0))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let held = db
.new_batch()
.append(U64::new(3))
.merkleize(&db, None, Location::new(0))
.await;
let db = db.rewind(size_after_first).await.unwrap();
assert!(matches!(db.apply_batch(held).await, Err(Error::StaleBatch)));
});
}
#[test_traced("INFO")]
fn test_compact_floor_beyond_size() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "keyless-floor-beyond").await;
let batch = db.new_batch().merkleize(&db, None, Location::new(2)).await;
assert!(matches!(
db.apply_batch(batch).await,
Err(Error::FloorBeyondSize(floor, tip))
if floor == Location::new(2) && tip == Location::new(1)
));
});
}
#[test_traced("INFO")]
fn test_compact_ancestor_floor_beyond_size() {
deterministic::Runner::default().start(|context| async move {
let db =
open_db::<mmr::Family>(context.child("db"), "keyless-ancestor-floor-beyond").await;
let parent = db
.new_batch()
.append(U64::new(1))
.merkleize(&db, None, Location::new(3))
.await;
let child = parent
.new_batch::<Sha256>()
.append(U64::new(2))
.merkleize(&db, None, Location::new(0))
.await;
assert!(matches!(
db.apply_batch(child).await,
Err(Error::FloorBeyondSize(floor, commit))
if floor == Location::new(3) && commit == Location::new(2)
));
});
}
}