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},
},
operation::Key,
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 core::marker::PhantomData;
use std::{
collections::BTreeMap,
sync::{Arc, Weak},
};
pub struct Db<F, E, K, V, H, C, S: Strategy>
where
F: Family,
E: Context,
K: Key,
V: ValueEncoding,
H: Hasher,
Operation<F, K, V>: EncodeShared,
Operation<F, K, 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>,
_key: PhantomData<K>,
}
impl<F, E, K, V, H, C, S: Strategy> std::fmt::Debug for Db<F, E, K, V, H, C, S>
where
F: Family,
E: Context,
K: Key,
V: ValueEncoding,
H: Hasher,
Operation<F, K, V>: EncodeShared,
Operation<F, K, 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, K, V, S: Strategy>
where
F: Family,
K: Key,
V: ValueEncoding,
H: Hasher,
Operation<F, K, V>: EncodeShared,
{
merkle_batch: compact_merkle::UnmerkleizedBatch<F, H::Digest, S>,
mutations: BTreeMap<K, V::Value>,
parent: Option<Arc<MerkleizedBatch<F, H::Digest, K, V, S>>>,
base: batch_chain::Commitment<F, H::Digest>,
}
#[derive(Clone)]
pub struct MerkleizedBatch<F: Family, D: Digest, K: Key, V: ValueEncoding, S: Strategy>
where
Operation<F, K, V>: EncodeShared,
{
pub(super) merkle_batch: Arc<batch::MerkleizedBatch<F, D, S>>,
operations: Arc<Vec<Operation<F, K, V>>>,
pub(super) commit_metadata: Option<V::Value>,
pub(super) parent: Option<Weak<Self>>,
pub(super) bounds: batch_chain::Bounds<F, D>,
pub(super) _key: PhantomData<K>,
}
impl<F: Family, D: Digest, K: Key, V: ValueEncoding, S: Strategy> MerkleizedBatch<F, D, K, V, S>
where
Operation<F, K, V>: EncodeShared,
{
pub(super) fn ancestors(&self) -> impl Iterator<Item = Arc<Self>> + use<F, D, K, 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
}
#[allow(clippy::type_complexity)]
pub fn operations(&self) -> (Location<F>, Arc<Vec<Operation<F, K, V>>>) {
(self.bounds.base.size, Arc::clone(&self.operations))
}
pub fn proof<E, C, H>(&self, db: &Db<F, E, K, V, H, C, S>) -> Result<Proof<F, D>, Error<F>>
where
E: Context,
H: Hasher<Digest = D>,
C: Clone + Send + Sync + 'static,
Operation<F, K, 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, K, V, H, C, S>) -> Result<Vec<D>, Error<F>>
where
E: Context,
H: Hasher<Digest = D>,
C: Clone + Send + Sync + 'static,
Operation<F, K, 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, K, V, S>
where
H: Hasher<Digest = D>,
{
UnmerkleizedBatch {
merkle_batch: compact_merkle::UnmerkleizedBatch::wrap(self.merkle_batch.new_batch()),
mutations: BTreeMap::new(),
parent: Some(Arc::clone(self)),
base: self.commitment(),
}
}
}
impl<F, H, K, V, S> UnmerkleizedBatch<F, H, K, V, S>
where
F: Family,
K: Key,
V: ValueEncoding,
H: Hasher,
S: Strategy,
Operation<F, K, V>: EncodeShared,
{
pub(super) fn new<E, C>(
db: &Db<F, E, K, V, H, C, S>,
base: batch_chain::Commitment<F, H::Digest>,
) -> Self
where
E: Context,
C: Clone + Send + Sync + 'static,
Operation<F, K, V>: Read<Cfg = C>,
{
Self {
merkle_batch: db.merkle.new_batch(),
mutations: BTreeMap::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 set(mut self, key: K, value: V::Value) -> Self {
self.mutations.insert(key, value);
self
}
#[tracing::instrument(
name = "qmdb.immutable.compact.batch.merkleize",
level = "info",
skip_all
)]
pub async fn merkleize<E, C>(
self,
db: &Db<F, E, K, V, H, C, S>,
metadata: Option<V::Value>,
inactivity_floor: Location<F>,
) -> Arc<MerkleizedBatch<F, H::Digest, K, V, S>>
where
F: Family,
E: Context,
C: Clone + Send + Sync + 'static,
Operation<F, K, 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, K, V>> = Vec::with_capacity(self.mutations.len() + 1);
for (key, value) in self.mutations {
ops.push(Operation::Set(key, 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,
},
_key: PhantomData,
})
}
}
impl<F, E, K, V, H, C, S> Db<F, E, K, V, H, C, S>
where
F: Family,
E: Context,
K: Key,
V: ValueEncoding,
H: Hasher,
S: Strategy,
Operation<F, K, V>: EncodeShared,
Operation<F, K, 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, K, 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, K, 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,
_key: PhantomData,
})
}
#[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, K, 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, K, V>>(
journal,
&mut merkle,
&commit_codec_config,
Operation::<F, K, 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,
_key: PhantomData,
})
}
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, K, V, S> {
UnmerkleizedBatch::new(self, self.commitment())
}
pub fn to_batch(&self) -> Arc<MerkleizedBatch<F, H::Digest, K, 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),
_key: PhantomData,
})
}
pub fn validate_batch(
&self,
batch: &MerkleizedBatch<F, H::Digest, K, V, S>,
) -> Result<(), Error<F>> {
batch
.bounds
.validate_apply_to(self.commitment(), self.inactivity_floor_loc)
}
#[tracing::instrument(
name = "qmdb.immutable.compact.db.apply_batch",
level = "info",
skip_all
)]
pub async fn apply_batch(
mut self,
batch: Arc<MerkleizedBatch<F, H::Digest, K, 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.immutable.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.immutable.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.immutable.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.immutable.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, K, 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.immutable.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, K, V, H, C, S> Source for Db<F, E, K, V, H, C, S>
where
F: Family,
E: Context,
K: Key,
V: ValueEncoding,
H: Hasher,
Operation<F, K, V>: EncodeShared + Read<Cfg = C>,
C: Clone + Send + Sync + 'static,
S: Strategy,
{
type Family = F;
type Digest = H::Digest;
type Op = Operation<F, K, 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 _, Supervisor as _,
buffer::paged::CacheRef,
deterministic,
mocks::{DelayedSyncContext, PendingSyncs, drive_pending_syncs},
};
use commonware_utils::{NZU16, NZU64, NZUsize};
use core::future::Future;
use futures::FutureExt as _;
use std::num::{NonZeroU16, NonZeroUsize};
type TestDb<F> =
Db<F, deterministic::Context, Digest, FixedEncoding<Digest>, 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"), "immutable-operations-and-proof").await;
let value = |key: u8| Sha256::fill(key.wrapping_add(100));
let mut seed = db.new_batch();
for key in 1u8..=6 {
seed = seed.set(Sha256::fill(key), value(key));
}
let seed = seed
.merkleize(&db, Some(Sha256::fill(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 key in 8u8..=12 {
parent = parent.set(Sha256::fill(key), value(key));
}
let parent = parent.merkleize(&db, Some(Sha256::fill(13)), floor).await;
let child = parent
.new_batch::<Sha256>()
.set(Sha256::fill(14), value(14))
.merkleize(&db, Some(Sha256::fill(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::Set(key, set_value), Operation::Commit(Some(metadata), operation_floor)]
if key == &Sha256::fill(14)
&& set_value == &value(14)
&& metadata == &Sha256::fill(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(Sha256::fill(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 == &Sha256::fill(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()
.set(Sha256::fill(17), value(17))
.merkleize(&db, Some(Sha256::fill(18)), db.size())
.await;
let late = late_parent
.new_batch::<Sha256>()
.set(Sha256::fill(19), value(19))
.merkleize(&db, Some(Sha256::fill(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()
}
type DelayedDb = Db<
mmr::Family,
DelayedSyncContext<deterministic::Context>,
Digest,
FixedEncoding<Digest>,
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_set(db: DelayedDb, key: Digest, value: Digest, metadata: Digest) -> DelayedDb {
let floor = db.inactivity_floor_loc();
let batch = db
.new_batch()
.set(key, value)
.merkleize(&db, Some(metadata), floor)
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
db
}
#[test_traced]
fn test_compact_start_sync_recovery() {
deterministic::Runner::default().start(|ctx| async move {
let partition = "immutable-start-sync-recovery";
let pending = PendingSyncs::default();
pending.unblock();
let mut db = open_delayed_db(&ctx, "delayed", partition, &pending)
.await
.unwrap();
let metadata = Sha256::fill(9u8);
db = apply_set(db, Sha256::fill(1u8), Sha256::fill(2u8), metadata).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(metadata));
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", "immutable-start-sync-fail", &pending)
.await
.unwrap();
db = apply_set(db, Sha256::fill(1u8), Sha256::fill(2u8), Sha256::fill(9u8)).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_rewind_fast_path_drains() {
deterministic::Runner::default().start(|ctx| async move {
let partition = "immutable-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_set(db, Sha256::fill(1u8), Sha256::fill(2u8), Sha256::fill(9u8)).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",
"immutable-start-sync-rewind-fail",
&pending,
)
.await
.unwrap();
db = apply_set(db, Sha256::fill(1u8), Sha256::fill(2u8), Sha256::fill(9u8)).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("INFO")]
fn test_compact_stale_batch_rejected() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "immutable-stale").await;
let key1 = Sha256::hash(&[&[1]]);
let key2 = Sha256::hash(&[&[2]]);
let value1 = Sha256::fill(10u8);
let value2 = Sha256::fill(20u8);
let batch_a = db
.new_batch()
.set(key1, value1)
.merkleize(&db, None, Location::new(0))
.await;
let batch_b = db
.new_batch()
.set(key2, value2)
.merkleize(&db, None, Location::new(0))
.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"), "immutable-delayed-child").await;
let key1 = Sha256::hash(&[&[1]]);
let key2 = Sha256::hash(&[&[2]]);
let key3 = Sha256::hash(&[&[3]]);
let value1 = Sha256::fill(10u8);
let value2 = Sha256::fill(20u8);
let value3 = Sha256::fill(30u8);
let a = db
.new_batch()
.set(key1, value1)
.merkleize(&db, None, Location::new(0))
.await;
let b = a
.new_batch::<Sha256>()
.set(key2, value2)
.merkleize(&db, None, Location::new(0))
.await;
let c = b.new_batch::<Sha256>().set(key3, value3);
let (db, _) = db.apply_batch(a).await.unwrap();
let c = c.merkleize(&db, None, Location::new(0)).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"), "immutable-to-batch-live").await;
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 key = Sha256::hash(&[&[1]]);
let value = Sha256::fill(10u8);
let batch = db
.new_batch()
.set(key, value)
.merkleize(&db, None, Location::new(0))
.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"), "immutable-chained-stale").await;
let common_parent = db
.new_batch()
.set(Sha256::hash(&[&[10]]), Sha256::fill(10u8))
.merkleize(&db, None, Location::new(0))
.await;
let sibling_a = common_parent
.new_batch::<Sha256>()
.set(Sha256::hash(&[&[11]]), Sha256::fill(11u8))
.merkleize(&db, None, Location::new(0))
.await;
let sibling_b = common_parent
.new_batch::<Sha256>()
.set(Sha256::hash(&[&[12]]), Sha256::fill(12u8))
.merkleize(&db, None, Location::new(0))
.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()
.set(Sha256::hash(&[&[1]]), Sha256::fill(1u8))
.merkleize(&db, None, Location::new(0))
.await;
let parent_b = db
.new_batch()
.set(Sha256::hash(&[&[2]]), Sha256::fill(2u8))
.merkleize(&db, None, Location::new(0))
.await;
let child_b = parent_b
.new_batch::<Sha256>()
.set(Sha256::hash(&[&[3]]), Sha256::fill(3u8))
.merkleize(&db, None, Location::new(0))
.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"), "immutable-child-before-parent").await;
let parent = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(1u8))
.merkleize(&db, None, Location::new(0))
.await;
let child = parent
.new_batch::<Sha256>()
.set(Sha256::hash(&[&[2]]), Sha256::fill(2u8))
.merkleize(&db, None, Location::new(0))
.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"), "immutable-parent-child").await;
let parent = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(1u8))
.merkleize(&db, None, Location::new(0))
.await;
let child = parent
.new_batch::<Sha256>()
.set(Sha256::hash(&[&[2]]), Sha256::fill(2u8))
.merkleize(&db, None, Location::new(0))
.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"), "immutable-floor-regressed").await;
let advance_floor = db.new_batch().set(Sha256::hash(&[&[1]]), Sha256::fill(1u8));
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()
.set(Sha256::hash(&[&[2]]), Sha256::fill(2u8))
.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"), "immutable-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"), "immutable-regressed-ancestor-floor")
.await;
let parent = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(1u8))
.merkleize(&db, None, Location::new(1))
.await;
let child = parent
.new_batch::<Sha256>()
.set(Sha256::hash(&[&[2]]), Sha256::fill(2u8))
.merkleize(&db, None, Location::new(0))
.await;
let target = db.target();
assert!(matches!(
db.apply_batch(child).await,
Err(Error::FloorRegressed(new, prev))
if new == Location::new(0) && prev == Location::new(1)
));
let db = open_db::<mmr::Family>(
context.child("reopen"),
"immutable-regressed-ancestor-floor",
)
.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"), "immutable-rewind-meta").await;
let k1 = Sha256::hash(&[&[1]]);
let v1 = Sha256::fill(11u8);
let meta1 = Sha256::fill(0xaa);
let floor1 = Location::new(0);
let batch = db
.new_batch()
.set(k1, v1)
.merkleize(&db, Some(meta1), 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 k2 = Sha256::hash(&[&[2]]);
let v2 = Sha256::fill(22u8);
let meta2 = Sha256::fill(0xbb);
let floor2 = Location::new(1);
let batch = db
.new_batch()
.set(k2, v2)
.merkleize(&db, Some(meta2), 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 = "immutable-rewind-reopen";
let meta1 = Sha256::fill(0xaa);
let floor1 = Location::new(0);
let meta2 = Sha256::fill(0xbb);
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()
.set(Sha256::hash(&[&[1]]), Sha256::fill(11u8))
.merkleize(&db, Some(meta1), 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()
.set(Sha256::hash(&[&[2]]), Sha256::fill(22u8))
.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 = "immutable-commit-reopen";
let meta1 = Sha256::fill(0xaa);
let meta2 = Sha256::fill(0xbb);
let root_after_second = {
let db = open_db::<mmr::Family>(context.child("first"), partition).await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(11u8))
.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()
.set(Sha256::hash(&[&[2]]), Sha256::fill(22u8))
.merkleize(&db, Some(meta2), 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 = "immutable-commit-rewind-reopen";
let meta1 = Sha256::fill(0xaa);
let meta2 = Sha256::fill(0xbb);
let (root_a, size_a) = {
let db = open_db::<mmr::Family>(context.child("first"), partition).await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(11u8))
.merkleize(&db, Some(meta1), 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()
.set(Sha256::hash(&[&[2]]), Sha256::fill(22u8))
.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 = "immutable-sync-after-commit";
let meta = Sha256::fill(0xaa);
let root = {
let db = open_db::<mmr::Family>(context.child("first"), partition).await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(11u8))
.merkleize(&db, Some(meta), 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_reopen_rejects_tampered_witness() {
deterministic::Runner::default().start(|context| async move {
let partition = "immutable-witness-tamper";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[7]]), Sha256::fill(7u8))
.merkleize(&db, Some(Sha256::fill(0xaa)), 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 = "immutable-corrupt-rewind-target";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(1u8))
.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()
.set(Sha256::hash(&[&[2]]), Sha256::fill(2u8))
.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 = "immutable-interrupted-import";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[7]]), Sha256::fill(7u8))
.merkleize(&db, Some(Sha256::fill(0xaa)), 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 = "immutable-invalid-persisted-floor";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[7]]), Sha256::fill(7u8))
.merkleize(&db, Some(Sha256::fill(0xaa)), 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, Digest, FixedEncoding<Digest>>::Commit(
Some(Sha256::fill(0xaa)),
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 = "immutable-pinned-nodes-tamper";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[7]]), Sha256::fill(7u8))
.merkleize(&db, Some(Sha256::fill(0xaa)), 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_reopen_drops_unsynced_witness() {
deterministic::Runner::default().start(|context| async move {
let partition = "immutable-witness-unsynced";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(1u8))
.merkleize(&db, Some(Sha256::fill(0xa1)), 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_beyond_history() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "immutable-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"), "immutable-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"), "immutable-rewind-between").await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(1u8))
.set(Sha256::hash(&[&[2]]), Sha256::fill(2u8))
.merkleize(&db, Some(Sha256::fill(0xa1)), 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();
assert_eq!(size_a, Location::new(4));
let batch = db
.new_batch()
.set(Sha256::hash(&[&[3]]), Sha256::fill(3u8))
.merkleize(&db, Some(Sha256::fill(0xb1)), Location::new(0))
.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),
"immutable-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(Sha256::fill(0xa1)));
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_rewind_multiple_commits() {
deterministic::Runner::default().start(|context| async move {
let partition = "immutable-rewind-multi";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(1u8))
.merkleize(&db, Some(Sha256::fill(0xa1)), 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 [2u8, 3] {
let batch = db
.new_batch()
.set(Sha256::hash(&[&[i]]), Sha256::fill(i))
.merkleize(&db, Some(Sha256::fill(i)), 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(Sha256::fill(0xa1)));
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_rewind_to_current_is_noop() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(context.child("db"), "immutable-rewind-noop").await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(1u8))
.merkleize(&db, Some(Sha256::fill(0xa1)), 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_then_rewind() {
deterministic::Runner::default().start(|context| async move {
let mut witness_cfg = witness_config("immutable-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 [1u8, 2, 3] {
let batch = db
.new_batch()
.set(Sha256::hash(&[&[i]]), Sha256::fill(i))
.merkleize(&db, Some(Sha256::fill(i)), 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(Sha256::fill(2)));
db.destroy().await.unwrap();
});
}
#[test_traced("INFO")]
fn test_compact_prune_past_tip_keeps_tip() {
deterministic::Runner::default().start(|context| async move {
let partition = "immutable-prune-past-tip";
let db = open_db::<mmr::Family>(context.child("db"), partition).await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(1u8))
.merkleize(&db, Some(Sha256::fill(0xa1)), 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_preserves_pre_advance_batch() {
deterministic::Runner::default().start(|context| async move {
let db = open_db::<mmr::Family>(
context.child("db"),
"immutable-rewind-preserves-pre-advance",
)
.await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(1u8))
.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()
.set(Sha256::hash(&[&[2]]), Sha256::fill(2u8))
.merkleize(&db, None, Location::new(0))
.await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[3]]), Sha256::fill(3u8))
.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"), "immutable-noop-after-commit").await;
let k1 = Sha256::hash(&[&[1]]);
let v1 = Sha256::fill(11u8);
let k2 = Sha256::hash(&[&[2]]);
let v2 = Sha256::fill(22u8);
let batch = db
.new_batch()
.set(k1, v1)
.set(k2, v2)
.merkleize(&db, Some(Sha256::fill(0xaa)), 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 size_after_first = db.size();
let db = db.sync().await.unwrap();
assert_eq!(db.size(), size_after_first);
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 = "immutable-noop-after-reopen";
let (root_before_drop, size_before_drop) = {
let db = open_db::<mmr::Family>(context.child("first"), partition).await;
let k1 = Sha256::hash(&[&[1]]);
let v1 = Sha256::fill(11u8);
let k2 = Sha256::hash(&[&[2]]);
let v2 = Sha256::fill(22u8);
let batch = db
.new_batch()
.set(k1, v1)
.set(k2, v2)
.merkleize(&db, Some(Sha256::fill(0xaa)), Location::new(0))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
(db.root(), db.size())
};
let db = open_db::<mmr::Family>(context.child("second"), partition).await;
assert_eq!(db.root(), root_before_drop);
assert_eq!(db.size(), size_before_drop);
let db = db.sync().await.unwrap();
assert_eq!(db.size(), size_before_drop);
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"), "immutable-noop-after-rewind").await;
let k1 = Sha256::hash(&[&[1]]);
let v1 = Sha256::fill(11u8);
let k2 = Sha256::hash(&[&[2]]);
let v2 = Sha256::fill(22u8);
let batch = db
.new_batch()
.set(k1, v1)
.set(k2, v2)
.merkleize(&db, Some(Sha256::fill(0xaa)), 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 size_after_first = db.size();
let k3 = Sha256::hash(&[&[3]]);
let v3 = Sha256::fill(33u8);
let batch = db
.new_batch()
.set(k3, v3)
.merkleize(&db, Some(Sha256::fill(0xbb)), Location::new(1))
.await;
let (db, _) = db.apply_batch(batch).await.unwrap();
let db = db.sync().await.unwrap();
let db = db.rewind(size_after_first).await.unwrap();
assert_eq!(db.size(), size_after_first);
assert_eq!(db.root(), root_after_first);
let db = db.sync().await.unwrap();
assert_eq!(db.size(), size_after_first);
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"), "immutable-rewind-makes-stale").await;
let batch = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(1u8))
.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()
.set(Sha256::hash(&[&[2]]), Sha256::fill(2u8))
.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()
.set(Sha256::hash(&[&[3]]), Sha256::fill(3u8))
.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"), "immutable-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"), "immutable-ancestor-floor-beyond")
.await;
let parent = db
.new_batch()
.set(Sha256::hash(&[&[1]]), Sha256::fill(1u8))
.merkleize(&db, None, Location::new(3))
.await;
let child = parent
.new_batch::<Sha256>()
.set(Sha256::hash(&[&[2]]), Sha256::fill(2u8))
.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)
));
});
}
}