use crate::translator::Translator;
use commonware_runtime::buffer::paged::CacheRef;
use std::num::{NonZeroU64, NonZeroUsize};
mod storage;
pub use storage::Archive;
#[derive(Clone)]
pub struct Config<T: Translator, C> {
pub translator: T,
pub metadata_partition: String,
pub key_partition: String,
pub key_page_cache: CacheRef,
pub value_partition: String,
pub compression: Option<u8>,
pub codec_config: C,
pub items_per_section: NonZeroU64,
pub key_write_buffer: NonZeroUsize,
pub value_write_buffer: NonZeroUsize,
pub replay_buffer: NonZeroUsize,
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{
archive::{Archive as _, Error, Identifier, MultiArchive as _},
journal::{Error as JournalError, segmented::glob::corrupt_frame},
translator::{FourCap, TwoCap},
};
use commonware_codec::{DecodeExt, Error as CodecError, FixedSize};
use commonware_cryptography::Crc32;
use commonware_macros::{test_group, test_traced};
use commonware_runtime::{
Blob as _, BufferPooler, Error as RError, Metrics as _, ReadOptions, Runner, Spawner as _,
Storage as _, Supervisor as _, WriteOptions, deterministic,
mocks::{
DelayedSyncContext, PendingSyncs, drive_pending_syncs, fail_pending_syncs,
release_next_pending_syncs, release_pending_syncs,
},
telemetry::metrics::has_metric_value,
};
use commonware_utils::{NZU16, NZU64, NZUsize, sequence::FixedBytes};
use rand::RngExt as _;
use std::{
collections::BTreeMap,
num::{NonZeroU16, NonZeroU64},
sync::{
Arc,
atomic::{AtomicUsize, Ordering},
},
};
fn test_key(key: &str) -> FixedBytes<64> {
let mut buf = [0u8; 64];
let key = key.as_bytes();
assert!(key.len() <= buf.len());
buf[..key.len()].copy_from_slice(key);
FixedBytes::decode(buf.as_ref()).unwrap()
}
const DEFAULT_ITEMS_PER_SECTION: u64 = 65536;
const DEFAULT_WRITE_BUFFER: usize = 1024;
const DEFAULT_REPLAY_BUFFER: usize = 4096;
const PAGE_SIZE: NonZeroU16 = NZU16!(1024);
const PAGE_CACHE_SIZE: NonZeroUsize = NZUsize!(10);
fn test_config<E: BufferPooler>(
context: &E,
partition_prefix: &str,
items_per_section: NonZeroU64,
) -> Config<FourCap, ()> {
Config {
translator: FourCap,
metadata_partition: format!("{partition_prefix}-metadata"),
key_partition: format!("{partition_prefix}-index"),
key_page_cache: CacheRef::from_pooler(context, PAGE_SIZE, PAGE_CACHE_SIZE),
value_partition: format!("{partition_prefix}-value"),
codec_config: (),
compression: None,
key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
replay_buffer: NZUsize!(DEFAULT_REPLAY_BUFFER),
items_per_section,
}
}
const I32_VALUE_FRAME_SIZE: u64 =
(i32::SIZE + crate::journal::segmented::glob::CHECKSUM_SIZE) as u64;
#[test_traced]
fn test_put_after_start_sync_is_accepted_before_handle_completes() {
let executor = deterministic::Runner::default();
let (_, checkpoint) = executor.start_and_recover(|context| async move {
let pending = PendingSyncs::default();
let context = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&context, "test", NZU64!(DEFAULT_ITEMS_PER_SECTION));
let archive = Archive::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let (mut archive, handle) = archive
.put_start_sync(1, test_key("aaa"), 10)
.await
.expect("Failed to start sync");
let pending_after_start = pending.lock().len();
assert!(
pending_after_start > 0,
"put_start_sync should return while the sync handle is still pending"
);
archive = archive
.put(2, test_key("bbb"), 20)
.await
.expect("archive should remain usable before sync completion");
assert_eq!(
pending.lock().len(),
pending_after_start,
"put should not issue a new storage sync while accepting later data"
);
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), Some(10));
release_pending_syncs(&pending);
handle.await.expect("sync handle should complete");
let (_archive, follow_up) = archive
.start_sync()
.await
.expect("Failed to start next sync");
assert!(
!pending.lock().is_empty(),
"the later put must remain pending for a future sync"
);
release_pending_syncs(&pending);
follow_up.await.expect("follow-up sync should complete");
});
deterministic::Runner::from(checkpoint).start(|context| async move {
let cfg = test_config(&context, "test", NZU64!(DEFAULT_ITEMS_PER_SECTION));
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("reopen"), cfg)
.await
.expect("Failed to reopen archive");
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), Some(10));
assert_eq!(archive.get(Identifier::Index(2)).await.unwrap(), Some(20));
});
}
#[test_traced]
fn test_duplicate_put_start_sync_observes_in_flight_sync() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let context = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&context, "test", NZU64!(DEFAULT_ITEMS_PER_SECTION));
let archive = Archive::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let (archive, first) = archive
.put_start_sync(1, test_key("aaa"), 10)
.await
.expect("Failed to start sync");
assert_eq!(pending.lock().len(), 2);
let (archive, second) = archive
.put_start_sync(1, test_key("duplicate"), 99)
.await
.expect("Failed to start duplicate sync");
assert_eq!(
pending.lock().len(),
2,
"duplicate put_start_sync must not issue a new storage sync"
);
let started = Arc::new(AtomicUsize::new(0));
let completed = Arc::new(AtomicUsize::new(0));
let started_clone = started.clone();
let completed_clone = completed.clone();
let waiter = context.inner.child("duplicate").spawn(|_| async move {
started_clone.fetch_add(1, Ordering::Relaxed);
second.await.expect("duplicate sync handle should complete");
completed_clone.fetch_add(1, Ordering::Relaxed);
});
while started.load(Ordering::Relaxed) == 0 {
commonware_runtime::reschedule().await;
}
commonware_runtime::reschedule().await;
assert_eq!(
completed.load(Ordering::Relaxed),
0,
"duplicate put_start_sync must observe the original in-flight sync"
);
release_pending_syncs(&pending);
first.await.expect("first sync handle should complete");
while completed.load(Ordering::Relaxed) == 0 {
commonware_runtime::reschedule().await;
}
waiter.await.expect("duplicate waiter failed");
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), Some(10));
});
}
#[test_traced]
fn test_below_floor_put_start_sync_covers_prior_pending_write() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let context = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&context, "test", NZU64!(1));
let archive = Archive::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let archive = archive.prune(1).await.expect("Failed to set prune floor");
let archive = archive
.put(2, test_key("pending"), 20)
.await
.expect("Failed to buffer retained write");
assert!(pending.lock().is_empty());
let (archive, handle) = archive
.put_start_sync(0, test_key("pruned"), 0)
.await
.expect("Failed to request sync through below-floor put");
assert_eq!(
pending.lock().len(),
2,
"the sync combinator must cover writes accepted before its below-floor put"
);
release_pending_syncs(&pending);
handle.await.expect("covering sync should complete");
assert_eq!(archive.get(Identifier::Index(2)).await.unwrap(), Some(20));
assert_eq!(archive.get(Identifier::Index(0)).await.unwrap(), None);
});
}
#[test_traced]
fn test_below_floor_put_multi_sync_covers_prior_pending_write() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let context = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&context, "test", NZU64!(1));
let archive = Archive::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let archive = archive.prune(1).await.expect("Failed to set prune floor");
let archive = archive
.put_multi(2, test_key("pending"), 20)
.await
.expect("Failed to buffer retained write");
pending.arm();
let completed = Arc::new(AtomicUsize::new(0));
let completed_clone = completed.clone();
let task = context.inner.child("put_multi_sync").spawn(|_| async move {
let result = archive.put_multi_sync(0, test_key("pruned"), 0).await;
completed_clone.store(1, Ordering::Relaxed);
result
});
while pending.calls() == 0 && completed.load(Ordering::Relaxed) == 0 {
commonware_runtime::reschedule().await;
}
assert_eq!(
completed.load(Ordering::Relaxed),
0,
"put_multi_sync must wait for writes accepted before its below-floor put"
);
assert!(pending.calls() > 0);
release_pending_syncs(&pending);
let archive = task
.await
.expect("put_multi_sync task failed")
.expect("put_multi_sync failed");
assert_eq!(archive.get_all(2).await.unwrap(), Some(vec![20]));
assert_eq!(archive.get_all(0).await.unwrap(), None);
});
}
#[test_traced]
fn test_overlapping_put_start_sync_waits_for_in_flight_sync() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let context = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&context, "test", NZU64!(DEFAULT_ITEMS_PER_SECTION));
let archive = Archive::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let (archive, first) = archive
.put_start_sync(1, test_key("aaa"), 10)
.await
.expect("Failed to start sync");
let pending_after_first = pending.lock().len();
assert!(pending_after_first > 0);
let started = Arc::new(AtomicUsize::new(0));
let completed = Arc::new(AtomicUsize::new(0));
let started_clone = started.clone();
let completed_clone = completed.clone();
let waiter = context.inner.child("second").spawn(|_| async move {
started_clone.fetch_add(1, Ordering::Relaxed);
let (archive, second) = archive
.put_start_sync(2, test_key("bbb"), 20)
.await
.expect("Failed to start second sync");
completed_clone.fetch_add(1, Ordering::Relaxed);
(archive, second)
});
while started.load(Ordering::Relaxed) == 0 {
commonware_runtime::reschedule().await;
}
commonware_runtime::reschedule().await;
assert_eq!(completed.load(Ordering::Relaxed), 0);
assert_eq!(
pending.lock().len(),
pending_after_first,
"second put_start_sync must not start new syncs before the first completes"
);
release_pending_syncs(&pending);
first.await.expect("first sync handle should complete");
while completed.load(Ordering::Relaxed) == 0 {
commonware_runtime::reschedule().await;
}
let (archive, second) = waiter.await.expect("second put task failed");
assert!(!pending.lock().is_empty());
release_pending_syncs(&pending);
second.await.expect("second sync handle should complete");
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), Some(10));
assert_eq!(archive.get(Identifier::Index(2)).await.unwrap(), Some(20));
});
}
#[test_traced]
fn test_sync_after_put_start_sync_waits_for_in_flight_sync() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let context = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&context, "test", NZU64!(DEFAULT_ITEMS_PER_SECTION));
let archive = Archive::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let (archive, first) = archive
.put_start_sync(1, test_key("aaa"), 10)
.await
.expect("Failed to start sync");
assert!(!pending.lock().is_empty());
let started = Arc::new(AtomicUsize::new(0));
let completed = Arc::new(AtomicUsize::new(0));
let started_clone = started.clone();
let completed_clone = completed.clone();
let waiter = context.inner.child("sync").spawn(|_| async move {
started_clone.fetch_add(1, Ordering::Relaxed);
let archive = archive.sync().await.expect("sync should complete");
completed_clone.fetch_add(1, Ordering::Relaxed);
archive
});
while started.load(Ordering::Relaxed) == 0 {
commonware_runtime::reschedule().await;
}
commonware_runtime::reschedule().await;
assert_eq!(
completed.load(Ordering::Relaxed),
0,
"shutdown sync must wait for the in-flight put_start_sync handle"
);
release_pending_syncs(&pending);
first.await.expect("first sync handle should complete");
while completed.load(Ordering::Relaxed) == 0 {
commonware_runtime::reschedule().await;
}
let archive = waiter.await.expect("sync task failed");
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), Some(10));
});
}
#[test_traced]
fn test_destroy_after_put_start_sync_waits_for_in_flight_sync() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let context = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&context, "test", NZU64!(DEFAULT_ITEMS_PER_SECTION));
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let (archive, first) = archive
.put_start_sync(1, test_key("aaa"), 10)
.await
.expect("Failed to start sync");
assert!(!pending.lock().is_empty());
let started = Arc::new(AtomicUsize::new(0));
let completed = Arc::new(AtomicUsize::new(0));
let started_clone = started.clone();
let completed_clone = completed.clone();
let waiter = context.inner.child("destroy").spawn(|_| async move {
started_clone.fetch_add(1, Ordering::Relaxed);
archive.destroy().await.expect("destroy should complete");
completed_clone.fetch_add(1, Ordering::Relaxed);
});
while started.load(Ordering::Relaxed) == 0 {
commonware_runtime::reschedule().await;
}
commonware_runtime::reschedule().await;
assert_eq!(
completed.load(Ordering::Relaxed),
0,
"destroy must wait for the in-flight put_start_sync handle"
);
release_pending_syncs(&pending);
first.await.expect("first sync handle should complete");
while completed.load(Ordering::Relaxed) == 0 {
commonware_runtime::reschedule().await;
}
waiter.await.expect("destroy task failed");
});
}
#[test_traced]
fn test_prune_after_put_start_sync_waits_for_in_flight_sync() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let context = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&context, "test", NZU64!(1));
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let (archive, first) = archive
.put_start_sync(1, test_key("aaa"), 10)
.await
.expect("Failed to start sync");
assert!(!pending.lock().is_empty());
let started = Arc::new(AtomicUsize::new(0));
let completed = Arc::new(AtomicUsize::new(0));
let started_clone = started.clone();
let completed_clone = completed.clone();
let waiter = context.inner.child("prune").spawn(|_| async move {
started_clone.fetch_add(1, Ordering::Relaxed);
let archive = archive.prune(2).await.expect("prune should complete");
completed_clone.fetch_add(1, Ordering::Relaxed);
archive
});
while started.load(Ordering::Relaxed) == 0 {
commonware_runtime::reschedule().await;
}
commonware_runtime::reschedule().await;
assert_eq!(
completed.load(Ordering::Relaxed),
0,
"prune must wait for in-flight syncs on pruned sections"
);
release_pending_syncs(&pending);
first
.await
.expect("sync handle should complete despite pruning");
while completed.load(Ordering::Relaxed) == 0 {
commonware_runtime::reschedule().await;
}
let archive = waiter.await.expect("prune task failed");
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), None);
});
}
#[test_traced]
fn test_prune_surfaces_failed_in_flight_sync() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let context = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&context, "test", NZU64!(1));
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let (archive, first) = archive
.put_start_sync(1, test_key("aaa"), 10)
.await
.expect("Failed to start sync");
fail_pending_syncs(&pending);
let err = archive
.prune(2)
.await
.expect_err("prune must surface a failed in-flight sync");
assert!(matches!(
err,
Error::Journal(JournalError::Runtime(RError::Io(_)))
));
let err = first.await.expect_err("first sync handle should fail");
assert!(matches!(err, RError::Io(_)));
});
}
#[test_traced]
fn test_put_start_sync_after_prune_drops_pruned_sync_requests() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let context = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&context, "test", NZU64!(1));
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let (archive, first) = archive
.put_start_sync(1, test_key("aaa"), 10)
.await
.expect("Failed to start sync");
release_pending_syncs(&pending);
first.await.expect("first sync handle should complete");
let archive = archive.prune(2).await.expect("Failed to prune");
let (archive, second) = archive
.put_start_sync(2, test_key("bbb"), 20)
.await
.expect("put_start_sync after prune should succeed");
release_pending_syncs(&pending);
second.await.expect("second sync handle should complete");
let archive = drive_pending_syncs(&pending, archive.sync())
.await
.expect("sync after prune should succeed");
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), None);
assert_eq!(archive.get(Identifier::Index(2)).await.unwrap(), Some(20));
});
}
#[test_traced]
fn test_overlapping_put_start_sync_restarts_after_all_handles_complete() {
let executor = deterministic::Runner::default();
let (_, checkpoint) = executor.start_and_recover(|context| async move {
let pending = PendingSyncs::default();
let context = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&context, "test", NZU64!(1));
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let (archive, first) = archive
.put_start_sync(1, test_key("aaa"), 10)
.await
.expect("Failed to start first sync");
assert_eq!(pending.lock().len(), 2);
let (_archive, second) = archive
.put_start_sync(2, test_key("bbb"), 20)
.await
.expect("Failed to start second sync");
assert_eq!(
pending.lock().len(),
4,
"different sections should be able to have independent in-flight syncs"
);
release_pending_syncs(&pending);
first.await.expect("first sync handle should complete");
second.await.expect("second sync handle should complete");
});
deterministic::Runner::from(checkpoint).start(|context| async move {
let cfg = test_config(&context, "test", NZU64!(1));
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("reopen"), cfg)
.await
.expect("Failed to reopen archive");
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), Some(10));
assert_eq!(archive.get(Identifier::Index(2)).await.unwrap(), Some(20));
});
}
#[test_traced]
fn test_overlapping_put_start_sync_restarts_only_completed_handles() {
let executor = deterministic::Runner::default();
let (_, checkpoint) = executor.start_and_recover(|context| async move {
let pending = PendingSyncs::default();
let context = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&context, "test", NZU64!(1));
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let (archive, first) = archive
.put_start_sync(1, test_key("aaa"), 10)
.await
.expect("Failed to start first sync");
let (archive, second) = archive
.put_start_sync(2, test_key("bbb"), 20)
.await
.expect("Failed to start second sync");
assert_eq!(pending.lock().len(), 4);
release_next_pending_syncs(&pending, 2);
first.await.expect("first sync handle should complete");
drop(second);
drop(archive);
});
deterministic::Runner::from(checkpoint).start(|context| async move {
let cfg = test_config(&context, "test", NZU64!(1));
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("reopen"), cfg)
.await
.expect("Failed to reopen archive");
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), Some(10));
assert_eq!(archive.get(Identifier::Index(2)).await.unwrap(), None);
});
}
#[test_traced]
fn test_failed_start_sync_is_returned_by_next_start_sync_handle() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let context = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&context, "test", NZU64!(DEFAULT_ITEMS_PER_SECTION));
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let (archive, first) = archive
.put_start_sync(1, test_key("aaa"), 10)
.await
.expect("Failed to start sync");
assert_eq!(pending.lock().len(), 2);
fail_pending_syncs(&pending);
let archive = archive
.put(2, test_key("bbb"), 20)
.await
.expect("write should be accepted before observing the failed sync");
let (_archive, second) = archive
.start_sync()
.await
.expect("start_sync should return a handle for the failed sync");
let err = second
.await
.expect_err("next start_sync handle should observe failed in-flight sync");
assert!(matches!(err, RError::Io(_)));
let err = first.await.expect_err("first sync handle should fail");
assert!(matches!(err, RError::Io(_)));
});
}
#[test_traced]
fn test_archive_truncates_at_first_invalid_value() {
deterministic::Runner::default().start(|context| async move {
for (name, bad_position, retained) in [("first", 0, 0), ("middle", 1, 1)] {
let cfg = test_config(&context, &format!("invalid-{name}"), NZU64!(4));
let mut archive = Archive::init(context.child(name), cfg.clone())
.await
.unwrap();
for (index, value) in [10, 20, 30].into_iter().enumerate() {
archive = archive
.put(index as u64, test_key(&format!("key-{index}")), value)
.await
.unwrap();
}
archive = archive.sync().await.unwrap();
drop(archive);
corrupt_frame(
&context,
&cfg.value_partition,
&0u64.to_be_bytes(),
bad_position,
I32_VALUE_FRAME_SIZE,
)
.await;
let archive =
Archive::<_, _, FixedBytes<64>, i32>::init(context.child(name), cfg.clone())
.await
.unwrap();
assert_eq!(
archive.ranges().collect::<Vec<_>>(),
if retained == 0 {
Vec::new()
} else {
vec![(0, retained - 1)]
}
);
for (index, value) in [10, 20, 30].into_iter().enumerate() {
let expected = (index < retained as usize).then_some(value);
assert_eq!(
archive.get(Identifier::Index(index as u64)).await.unwrap(),
expected
);
}
drop(archive);
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child(name), cfg)
.await
.unwrap();
assert_eq!(archive.last_index(), retained.checked_sub(1));
archive.destroy().await.unwrap();
}
});
}
#[test_traced]
fn test_archive_completes_interrupted_rewind_to_empty_section() {
deterministic::Runner::default().start(|context| async move {
let cfg = test_config(&context, "empty-rewind", NZU64!(4));
let archive = Archive::init(context.child("seed"), cfg.clone())
.await
.unwrap();
let archive = archive.put(0, test_key("zero"), 10).await.unwrap();
let archive = archive.sync().await.unwrap();
drop(archive);
context.remove(&cfg.metadata_partition, None).await.unwrap();
let (index, _) = context
.open(&cfg.key_partition, &0u64.to_be_bytes())
.await
.unwrap();
index.resize(0).await.unwrap();
index.sync().await.unwrap();
drop(index);
let (_, value_size) = context
.open(&cfg.value_partition, &0u64.to_be_bytes())
.await
.unwrap();
assert_eq!(value_size, I32_VALUE_FRAME_SIZE);
let archive =
Archive::<_, _, FixedBytes<64>, i32>::init(context.child("repair"), cfg.clone())
.await
.unwrap();
assert_eq!(archive.last_index(), None);
drop(archive);
let (_, value_size) = context
.open(&cfg.value_partition, &0u64.to_be_bytes())
.await
.unwrap();
assert_eq!(value_size, 0, "startup must finish the value truncation");
context.remove(&cfg.metadata_partition, None).await.unwrap();
let pending = PendingSyncs::default();
let delayed = DelayedSyncContext {
inner: context.child("clean_restart"),
pending: pending.clone(),
};
pending.arm();
let completed = Arc::new(AtomicUsize::new(0));
let completed_clone = completed.clone();
let cfg_clone = cfg.clone();
let task = context.child("clean_restart_task").spawn(|_| async move {
let result = Archive::init(delayed.child("archive"), cfg_clone).await;
completed_clone.store(1, Ordering::Relaxed);
result
});
while pending.calls() == 0 && completed.load(Ordering::Relaxed) == 0 {
commonware_runtime::reschedule().await;
}
if pending.calls() != 0 {
pending.unblock();
let _ = task.await;
panic!("clean empty-section restart must not issue durability operations");
}
pending.unblock();
let archive = task.await.unwrap().unwrap();
let archive = archive.put(0, test_key("new"), 20).await.unwrap();
let archive = archive.sync().await.unwrap();
drop(archive);
let (_, value_size) = context
.open(&cfg.value_partition, &0u64.to_be_bytes())
.await
.unwrap();
assert_eq!(value_size, I32_VALUE_FRAME_SIZE);
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("reopen"), cfg)
.await
.unwrap();
assert_eq!(archive.get(Identifier::Index(0)).await.unwrap(), Some(20));
archive.destroy().await.unwrap();
});
}
#[test_traced]
fn test_validation_marker_skips_previously_validated_values() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = test_config(&context, "marker-skip", NZU64!(4));
let mut archive = Archive::init(context.child("seed"), cfg.clone())
.await
.unwrap();
archive = archive.put(0, test_key("zero"), 10).await.unwrap();
archive = archive.put(1, test_key("one"), 20).await.unwrap();
archive = archive.sync().await.unwrap();
drop(archive);
context.remove(&cfg.metadata_partition, None).await.unwrap();
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(
context.child("first_open"),
cfg.clone(),
)
.await
.unwrap();
assert_eq!(archive.ranges().collect::<Vec<_>>(), vec![(0, 1)]);
drop(archive);
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(
context.child("second_open"),
cfg.clone(),
)
.await
.unwrap();
assert_eq!(archive.get(Identifier::Index(0)).await.unwrap(), Some(10));
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), Some(20));
let archive = archive.put(2, test_key("two"), 30).await.unwrap();
let archive = archive.sync().await.unwrap();
drop(archive);
let archive =
Archive::<_, _, FixedBytes<64>, i32>::init(context.child("third_open"), cfg)
.await
.unwrap();
assert_eq!(archive.ranges().collect::<Vec<_>>(), vec![(0, 2)]);
assert_eq!(archive.get(Identifier::Index(0)).await.unwrap(), Some(10));
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), Some(20));
assert_eq!(archive.get(Identifier::Index(2)).await.unwrap(), Some(30));
});
}
#[test_traced]
fn test_validation_marker_skips_covered_interior_values() {
deterministic::Runner::default().start(|context| async move {
let cfg = test_config(&context, "marker-covered-interior", NZU64!(4));
let archive = Archive::init(context.child("seed"), cfg.clone())
.await
.unwrap();
let archive = archive.put(0, test_key("zero"), 10).await.unwrap();
let archive = archive.put(1, test_key("one"), 20).await.unwrap();
let archive = archive.sync().await.unwrap();
let archive = archive.sync().await.unwrap();
drop(archive);
corrupt_frame(
&context,
&cfg.value_partition,
&0u64.to_be_bytes(),
0,
I32_VALUE_FRAME_SIZE,
)
.await;
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("reopen"), cfg)
.await
.unwrap();
assert!(archive.get(Identifier::Index(0)).await.is_err());
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), Some(20));
});
}
#[test_traced]
fn test_validation_marker_damage_never_mutates() {
#[derive(Clone, Copy)]
enum Damage {
MissingIndex,
MissingValues,
TruncatedIndex,
TruncatedValues,
CorruptIndex,
CorruptValues,
}
deterministic::Runner::default().start(|context| async move {
for (name, damage) in [
("missing_index", Damage::MissingIndex),
("missing_values", Damage::MissingValues),
("truncated_index", Damage::TruncatedIndex),
("truncated_values", Damage::TruncatedValues),
("corrupt_index", Damage::CorruptIndex),
("corrupt_values", Damage::CorruptValues),
] {
let case = context.child(name);
let cfg = test_config(&case, name, NZU64!(4));
let archive = Archive::init(case.child("seed"), cfg.clone())
.await
.unwrap();
let archive = archive.put_sync(0, test_key("zero"), 10).await.unwrap();
let archive = archive.sync().await.unwrap();
drop(archive);
let (_, index_size) = context
.open(&cfg.key_partition, &0u64.to_be_bytes())
.await
.unwrap();
let (_, value_size) = context
.open(&cfg.value_partition, &0u64.to_be_bytes())
.await
.unwrap();
let damage_index = matches!(
damage,
Damage::MissingIndex | Damage::TruncatedIndex | Damage::CorruptIndex
);
let damaged_partition = if damage_index {
&cfg.key_partition
} else {
&cfg.value_partition
};
let damaged_size = match damage {
Damage::MissingIndex | Damage::MissingValues => {
context
.remove(damaged_partition, Some(&0u64.to_be_bytes()))
.await
.unwrap();
None
}
Damage::TruncatedIndex | Damage::TruncatedValues => {
let (blob, size) = context
.open(damaged_partition, &0u64.to_be_bytes())
.await
.unwrap();
let size = if damage_index { size - 1 } else { 0 };
blob.resize(size).await.unwrap();
blob.sync().await.unwrap();
Some(size)
}
Damage::CorruptIndex | Damage::CorruptValues => {
let (blob, size) = context
.open(damaged_partition, &0u64.to_be_bytes())
.await
.unwrap();
let byte = blob
.read_at(0, 1, ReadOptions::default())
.await
.unwrap()
.coalesce();
let byte = byte.as_ref()[0];
blob.write_at(0, vec![byte ^ 0xFF], WriteOptions::SYNC)
.await
.unwrap();
Some(size)
}
};
if matches!(damage, Damage::CorruptValues) {
for child in ["first", "second"] {
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(
case.child(child),
cfg.clone(),
)
.await
.expect("marked value damage must not fail startup");
assert!(archive.get(Identifier::Index(0)).await.is_err());
drop(archive);
let (_, size) = context
.open(&cfg.value_partition, &0u64.to_be_bytes())
.await
.unwrap();
assert_eq!(size, value_size, "adoption must preserve the damaged frame");
}
continue;
}
for child in ["first", "second"] {
let result =
Archive::<_, _, FixedBytes<64>, i32>::init(case.child(child), cfg.clone())
.await;
assert!(
matches!(result, Err(Error::Journal(JournalError::Corruption(_)))),
"damaged marked section must remain visible as corruption"
);
let (surviving_partition, surviving_size) = if damage_index {
(&cfg.value_partition, value_size)
} else {
(&cfg.key_partition, index_size)
};
let (_, size) = context
.open(surviving_partition, &0u64.to_be_bytes())
.await
.unwrap();
assert_eq!(
size, surviving_size,
"failed startup must preserve the surviving journal section"
);
if let Some(damaged_size) = damaged_size {
let (_, size) = context
.open(damaged_partition, &0u64.to_be_bytes())
.await
.unwrap();
assert_eq!(
size, damaged_size,
"failed startup must not normalize the damaged journal section"
);
}
}
}
});
}
#[test_traced]
fn test_validation_floor_rejection_precedes_index_suffix_repair() {
deterministic::Runner::default().start(|context| async move {
let cfg = test_config(&context, "floor-order", NZU64!(4));
let archive =
Archive::<_, _, FixedBytes<64>, i32>::init(context.child("seed"), cfg.clone())
.await
.unwrap();
let archive = archive.put_sync(0, test_key("zero"), 10).await.unwrap();
let archive = archive.sync().await.unwrap();
drop(archive);
let (index, index_size) = context
.open(&cfg.key_partition, &0u64.to_be_bytes())
.await
.unwrap();
index
.write_at(index_size, vec![0xA5; 7], WriteOptions::SYNC)
.await
.unwrap();
let expected_size = index_size + 7;
let byte = index
.read_at(0, 1, ReadOptions::default())
.await
.unwrap()
.coalesce();
index
.write_at(0, vec![byte.as_ref()[0] ^ 0xFF], WriteOptions::SYNC)
.await
.unwrap();
let expected = index
.read_at(0, expected_size as usize, ReadOptions::default())
.await
.unwrap()
.coalesce();
drop(index);
for child in ["first", "second"] {
let result =
Archive::<_, _, FixedBytes<64>, i32>::init(context.child(child), cfg.clone())
.await;
assert!(matches!(
result,
Err(Error::Journal(JournalError::Corruption(_)))
));
let (index, actual_size) = context
.open(&cfg.key_partition, &0u64.to_be_bytes())
.await
.unwrap();
assert_eq!(actual_size, expected_size);
let actual = index
.read_at(0, actual_size as usize, ReadOptions::default())
.await
.unwrap()
.coalesce();
assert_eq!(actual.as_ref(), expected.as_ref());
}
});
}
#[test_traced]
fn test_validation_marker_survives_torn_index_tail_rewrite() {
let executor = deterministic::Runner::default();
let (_, checkpoint) = executor.start_and_recover(|context| async move {
let cfg = test_config(&context, "marker-torn-tail", NZU64!(4));
let archive = Archive::init(context.child("seed"), cfg.clone())
.await
.unwrap();
let archive = archive.put_sync(0, test_key("zero"), 10).await.unwrap();
let archive = archive.sync().await.unwrap();
let page_size = usize::from(PAGE_SIZE.get());
let physical_page_size = page_size + 12;
let record_size = u64::SIZE + FixedBytes::<64>::SIZE + u64::SIZE + u32::SIZE;
assert!(record_size < page_size);
let (index, size) = context
.open(&cfg.key_partition, &0u64.to_be_bytes())
.await
.unwrap();
assert_eq!(size, physical_page_size as u64);
let old_page = index
.read_at(0, physical_page_size, ReadOptions::default())
.await
.unwrap()
.coalesce();
let old_page = old_page.as_ref().to_vec();
drop(index);
let old_len =
u16::from_be_bytes(old_page[page_size..page_size + 2].try_into().unwrap()) as usize;
let old_crc =
u32::from_be_bytes(old_page[page_size + 2..page_size + 6].try_into().unwrap());
assert_eq!(old_len, record_size);
assert_eq!(old_crc, Crc32::checksum(&old_page[..old_len]));
let archive = archive.put_sync(1, test_key("one"), 20).await.unwrap();
drop(archive);
let (index, size) = context
.open(&cfg.key_partition, &0u64.to_be_bytes())
.await
.unwrap();
assert_eq!(size, physical_page_size as u64);
let new_page = index
.read_at(0, physical_page_size, ReadOptions::default())
.await
.unwrap()
.coalesce();
let new_page = new_page.as_ref().to_vec();
assert_eq!(
&new_page[page_size..page_size + 6],
&old_page[page_size..page_size + 6],
);
let new_len =
u16::from_be_bytes(new_page[page_size + 6..page_size + 8].try_into().unwrap())
as usize;
let new_crc =
u32::from_be_bytes(new_page[page_size + 8..page_size + 12].try_into().unwrap());
assert_eq!(new_len, 2 * record_size);
assert_eq!(new_crc, Crc32::checksum(&new_page[..new_len]));
index
.write_at(0, old_page.clone(), WriteOptions::SYNC)
.await
.unwrap();
let torn_prefix = page_size + 6 + 2;
index
.write_at(0, new_page[..torn_prefix].to_vec(), WriteOptions::SYNC)
.await
.unwrap();
let torn_page = index
.read_at(0, physical_page_size, ReadOptions::default())
.await
.unwrap()
.coalesce();
let torn_page = torn_page.as_ref();
assert_eq!(
&torn_page[page_size..page_size + 6],
&old_page[page_size..page_size + 6],
);
assert_eq!(
u16::from_be_bytes(torn_page[page_size + 6..page_size + 8].try_into().unwrap(),)
as usize,
new_len,
);
let torn_crc =
u32::from_be_bytes(torn_page[page_size + 8..page_size + 12].try_into().unwrap());
assert_ne!(torn_crc, Crc32::checksum(&torn_page[..new_len]));
drop(index);
let (_, value_size) = context
.open(&cfg.value_partition, &0u64.to_be_bytes())
.await
.unwrap();
assert_eq!(value_size, 2 * I32_VALUE_FRAME_SIZE);
});
deterministic::Runner::from(checkpoint).start(|context| async move {
let pending = PendingSyncs::default();
let delayed = DelayedSyncContext {
inner: context.child("delayed"),
pending: pending.clone(),
};
pending.arm();
let cfg = test_config(&delayed, "marker-torn-tail", NZU64!(4));
let archive = drive_pending_syncs(
&pending,
Archive::<_, _, FixedBytes<64>, i32>::init(delayed.child("reopen"), cfg.clone()),
)
.await
.unwrap();
assert_eq!(pending.calls(), 1);
assert_eq!(archive.last_index(), Some(0));
assert_eq!(archive.ranges().collect::<Vec<_>>(), vec![(0, 0)]);
assert_eq!(archive.get(Identifier::Index(0)).await.unwrap(), Some(10));
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), None);
drop(archive);
let (_, value_size) = context
.open(&cfg.value_partition, &0u64.to_be_bytes())
.await
.unwrap();
assert_eq!(value_size, I32_VALUE_FRAME_SIZE);
});
}
#[test_traced]
fn test_startup_publishes_validated_marker_without_data_resync() {
deterministic::Runner::default().start(|context| async move {
let cfg = test_config(&context, "startup-order", NZU64!(4));
let archive = Archive::init(context.child("seed"), cfg.clone())
.await
.unwrap();
let archive = archive.put(0, test_key("zero"), 10).await.unwrap();
let archive = archive.sync().await.unwrap();
drop(archive);
let pending = PendingSyncs::default();
let delayed = DelayedSyncContext {
inner: context.child("delayed"),
pending: pending.clone(),
};
pending.arm();
let completed = Arc::new(AtomicUsize::new(0));
let completed_clone = completed.clone();
let reopen_cfg = cfg.clone();
let task = context.child("startup").spawn(|_| async move {
let result =
Archive::<_, _, FixedBytes<64>, i32>::init(delayed.child("reopen"), reopen_cfg)
.await;
completed_clone.store(1, Ordering::Relaxed);
result
});
while pending.calls() == 0 && completed.load(Ordering::Relaxed) == 0 {
commonware_runtime::reschedule().await;
}
commonware_runtime::reschedule().await;
let calls = pending.calls();
let finished = completed.load(Ordering::Relaxed);
if calls != 1 || finished != 1 {
pending.unblock();
let _ = task.await;
panic!(
"clean startup must return after starting one marker sync, calls={calls}, \
finished={finished}"
);
}
pending.unblock();
let archive = task.await.unwrap().unwrap().sync().await.unwrap();
assert_eq!(archive.get(Identifier::Index(0)).await.unwrap(), Some(10));
drop(archive);
Archive::<_, _, FixedBytes<64>, i32>::init(context.child("marker_reopen"), cfg)
.await
.unwrap()
.destroy()
.await
.unwrap();
});
}
#[test_traced]
fn test_startup_marker_failure_fails_next_sync() {
deterministic::Runner::default().start(|context| async move {
let cfg = test_config(&context, "startup-marker-failure", NZU64!(4));
let archive = Archive::init(context.child("seed"), cfg.clone())
.await
.unwrap();
let archive = archive.put(0, test_key("zero"), 10).await.unwrap();
let archive = archive.sync().await.unwrap();
drop(archive);
let pending = PendingSyncs::default();
let delayed = DelayedSyncContext {
inner: context.child("delayed"),
pending: pending.clone(),
};
pending.arm();
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(delayed.child("reopen"), cfg)
.await
.unwrap();
assert_eq!(pending.lock().len(), 1);
fail_pending_syncs(&pending);
assert!(matches!(archive.sync().await, Err(Error::Metadata(_))));
});
}
#[test_traced]
fn test_start_sync_publishes_closed_section_boundary() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = test_config(&context, "marker-lag", NZU64!(4));
let pending = PendingSyncs::default();
let delayed = DelayedSyncContext {
inner: context.child("delayed"),
pending: pending.clone(),
};
let archive = Archive::init(delayed.child("archive"), cfg.clone())
.await
.unwrap();
let (archive, first) = archive
.put_start_sync(0, test_key("zero"), 10)
.await
.unwrap();
assert_eq!(pending.lock().len(), 2);
release_pending_syncs(&pending);
first.await.unwrap();
let archive = archive.put(4, test_key("four"), 40).await.unwrap();
let (archive, second) = archive.start_sync().await.unwrap();
assert_eq!(pending.lock().len(), 3);
release_pending_syncs(&pending);
second.await.unwrap();
drop(archive);
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(
context.child("first_reopen"),
cfg.clone(),
)
.await
.unwrap();
assert_eq!(archive.ranges().collect::<Vec<_>>(), vec![(0, 0), (4, 4)]);
drop(archive);
Archive::<_, _, FixedBytes<64>, i32>::init(context.child("second_reopen"), cfg)
.await
.unwrap()
.destroy()
.await
.unwrap();
});
}
#[test_traced]
fn test_sync_publishes_closed_section_boundary() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = test_config(&context, "blocking-marker-lag", NZU64!(4));
let pending = PendingSyncs::default();
let delayed = DelayedSyncContext {
inner: context.child("delayed"),
pending: pending.clone(),
};
let archive = Archive::init(delayed.child("archive"), cfg.clone())
.await
.unwrap();
let (archive, first) = archive
.put_start_sync(0, test_key("zero"), 10)
.await
.unwrap();
assert_eq!(pending.lock().len(), 2);
release_pending_syncs(&pending);
first.await.unwrap();
let archive = archive.put(4, test_key("four"), 40).await.unwrap();
pending.arm();
let completed = Arc::new(AtomicUsize::new(0));
let completed_clone = completed.clone();
let task = delayed.inner.child("sync").spawn(|_| async move {
let archive = archive.sync().await.unwrap();
completed_clone.store(1, Ordering::Relaxed);
archive
});
while pending.calls() < 2 {
commonware_runtime::reschedule().await;
}
commonware_runtime::reschedule().await;
let parked_syncs = pending.lock().len();
if parked_syncs != 3 {
pending.unblock();
let _ = task.await;
panic!(
"blocking sync must not serialize the derived marker behind data: \
parked {parked_syncs} durability operations"
);
}
let metadata = pending.lock().remove(2);
metadata.release.send(Ok(())).unwrap();
commonware_runtime::reschedule().await;
assert_eq!(
completed.load(Ordering::Relaxed),
0,
"publishing the previous marker must not complete the current data sync"
);
release_pending_syncs(&pending);
let archive = task.await.unwrap();
drop(archive);
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(
delayed.inner.child("first_reopen"),
cfg.clone(),
)
.await
.unwrap();
assert_eq!(archive.get(Identifier::Index(0)).await.unwrap(), Some(10));
assert_eq!(archive.get(Identifier::Index(4)).await.unwrap(), Some(40));
drop(archive);
Archive::<_, _, FixedBytes<64>, i32>::init(delayed.inner.child("second_reopen"), cfg)
.await
.unwrap()
.destroy()
.await
.unwrap();
});
}
#[test_traced]
fn test_sync_delays_immediately_ready_durable_boundary() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
pending.unblock();
let immediate = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&immediate, "ready-marker-lag", NZU64!(4));
let archive = Archive::init(immediate.child("archive"), cfg)
.await
.unwrap();
let initial_starts = pending.starts();
let archive = archive.put_sync(0, test_key("zero"), 10).await.unwrap();
assert_eq!(pending.starts() - initial_starts, 2);
let archive = archive.sync().await.unwrap();
assert_eq!(pending.starts() - initial_starts, 3);
archive.destroy().await.unwrap();
});
}
#[test_traced]
fn test_sync_batches_markers_by_active_section() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
pending.unblock();
let immediate = DelayedSyncContext {
inner: context,
pending: pending.clone(),
};
let cfg = test_config(&immediate, "section-marker-batch", NZU64!(4));
let archive = Archive::init(immediate.child("archive"), cfg)
.await
.unwrap();
let initial_starts = pending.starts();
let archive = archive.put_sync(0, test_key("zero"), 10).await.unwrap();
let archive = archive.put_sync(1, test_key("one"), 20).await.unwrap();
assert_eq!(pending.starts() - initial_starts, 4);
let archive = archive.put_sync(4, test_key("four"), 40).await.unwrap();
assert_eq!(pending.starts() - initial_starts, 7);
let archive = archive.sync().await.unwrap();
assert_eq!(pending.starts() - initial_starts, 8);
let archive = archive.sync().await.unwrap();
assert_eq!(pending.starts() - initial_starts, 8);
archive.destroy().await.unwrap();
});
}
#[test_traced]
fn test_start_sync_withholds_marker_for_unproven_boundary() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = test_config(&context, "marker-unproven-boundary", NZU64!(4));
let pending = PendingSyncs::default();
let delayed = DelayedSyncContext {
inner: context.child("delayed"),
pending: pending.clone(),
};
let archive = Archive::init(delayed.child("archive"), cfg.clone())
.await
.unwrap();
let (archive, first) = archive
.put_start_sync(0, test_key("zero"), 10)
.await
.unwrap();
release_pending_syncs(&pending);
first.await.unwrap();
let archive = archive.put(4, test_key("four"), 40).await.unwrap();
let (archive, second) = archive.start_sync().await.unwrap();
assert_eq!(pending.lock().len(), 3);
let marker = pending.lock().pop().expect("marker sync parked");
marker
.release
.send(Ok(()))
.expect("marker sync receiver dropped");
let archive = archive.put(8, test_key("eight"), 80).await.unwrap();
let before = pending.starts();
let (archive, third) = archive.start_sync().await.unwrap();
assert_eq!(pending.starts() - before, 2);
release_pending_syncs(&pending);
second.await.unwrap();
third.await.unwrap();
let archive = drive_pending_syncs(&pending, archive.sync()).await.unwrap();
drop(archive);
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("reopen"), cfg)
.await
.unwrap();
assert_eq!(
archive.ranges().collect::<Vec<_>>(),
vec![(0, 0), (4, 4), (8, 8)]
);
archive.destroy().await.unwrap();
});
}
#[test_traced]
fn test_sync_publishes_previous_durable_sections_across_section_changes() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = test_config(&context, "cross-section-marker-lag", NZU64!(1));
let pending = PendingSyncs::default();
let delayed = DelayedSyncContext {
inner: context.child("delayed"),
pending: pending.clone(),
};
let archive = Archive::init(delayed.child("archive"), cfg.clone())
.await
.unwrap();
let archive = drive_pending_syncs(&pending, archive.put_sync(0, test_key("zero"), 10))
.await
.unwrap();
let archive = drive_pending_syncs(&pending, archive.put_sync(1, test_key("one"), 20))
.await
.unwrap();
let archive = drive_pending_syncs(&pending, archive.put_sync(2, test_key("two"), 30))
.await
.unwrap();
drop(archive);
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(
context.child("first_reopen"),
cfg.clone(),
)
.await
.unwrap();
assert_eq!(archive.get(Identifier::Index(0)).await.unwrap(), Some(10));
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), Some(20));
assert_eq!(archive.get(Identifier::Index(2)).await.unwrap(), Some(30));
drop(archive);
Archive::<_, _, FixedBytes<64>, i32>::init(context.child("second_reopen"), cfg)
.await
.unwrap()
.destroy()
.await
.unwrap();
});
}
#[test_traced]
fn test_empty_sync_publishes_final_durable_boundary() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = test_config(&context, "empty-sync-marker", NZU64!(1));
let pending = PendingSyncs::default();
let delayed = DelayedSyncContext {
inner: context.child("delayed"),
pending: pending.clone(),
};
let archive = Archive::init(delayed.child("archive"), cfg.clone())
.await
.unwrap();
let archive = drive_pending_syncs(&pending, archive.put_sync(0, test_key("zero"), 10))
.await
.unwrap();
assert_eq!(pending.starts(), 2);
let archive = drive_pending_syncs(&pending, archive.sync()).await.unwrap();
assert_eq!(pending.starts(), 3);
let archive = drive_pending_syncs(&pending, archive.sync()).await.unwrap();
assert_eq!(pending.starts(), 3);
drop(archive);
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("reopen"), cfg)
.await
.unwrap();
assert_eq!(archive.get(Identifier::Index(0)).await.unwrap(), Some(10));
});
}
#[test_traced]
fn test_sync_recreates_settled_section_barrier() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = test_config(&context, "marker-recreate", NZU64!(2));
let pending = PendingSyncs::default();
let delayed = DelayedSyncContext {
inner: context.child("delayed"),
pending: pending.clone(),
};
let archive = Archive::init(delayed.child("archive"), cfg.clone())
.await
.unwrap();
let archive = drive_pending_syncs(&pending, archive.put_sync(0, test_key("zero"), 10))
.await
.unwrap();
assert_eq!(pending.starts(), 2);
let archive = drive_pending_syncs(&pending, archive.put_sync(2, test_key("two"), 30))
.await
.unwrap();
assert_eq!(pending.starts(), 5);
let archive = drive_pending_syncs(&pending, archive.put_sync(1, test_key("one"), 20))
.await
.unwrap();
assert_eq!(pending.starts(), 8);
let archive = drive_pending_syncs(&pending, archive.sync()).await.unwrap();
assert_eq!(pending.starts(), 9);
let archive = drive_pending_syncs(&pending, archive.sync()).await.unwrap();
assert_eq!(pending.starts(), 9);
drop(archive);
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("reopen"), cfg)
.await
.unwrap();
assert_eq!(archive.get(Identifier::Index(0)).await.unwrap(), Some(10));
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), Some(20));
assert_eq!(archive.get(Identifier::Index(2)).await.unwrap(), Some(30));
});
}
#[test_traced]
fn test_prune_clears_validation_marker_before_section_reuse() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = test_config(&context, "marker-reuse", NZU64!(2));
let archive = Archive::init(context.child("seed"), cfg.clone())
.await
.unwrap();
let archive = archive.put_sync(0, test_key("old"), 10).await.unwrap();
let archive = archive.sync().await.unwrap();
let archive = archive.prune(2).await.unwrap();
drop(archive);
let pending = PendingSyncs::default();
let delayed = DelayedSyncContext {
inner: context.child("delayed"),
pending: pending.clone(),
};
let archive = Archive::init(delayed.child("reuse"), cfg.clone())
.await
.unwrap();
let archive = archive.put(0, test_key("new"), 20).await.unwrap();
let (archive, handle) = archive.start_sync().await.unwrap();
assert_eq!(pending.lock().len(), 2);
release_pending_syncs(&pending);
handle.await.unwrap();
drop(archive);
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("reopen"), cfg)
.await
.unwrap();
assert_eq!(archive.get(Identifier::Index(0)).await.unwrap(), Some(20));
});
}
#[test_traced]
fn test_archive_compression_then_none() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = Config {
translator: FourCap,
metadata_partition: "test-metadata".into(),
key_partition: "test-index".into(),
key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
value_partition: "test-value".into(),
codec_config: (),
compression: Some(3),
key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
replay_buffer: NZUsize!(DEFAULT_REPLAY_BUFFER),
items_per_section: NZU64!(DEFAULT_ITEMS_PER_SECTION),
};
let mut archive = Archive::init(context.child("first"), cfg.clone())
.await
.expect("Failed to initialize archive");
let index = 1u64;
let key = test_key("testkey");
let data = 1;
archive = archive
.put(index, key.clone(), data)
.await
.expect("Failed to put data");
let archive = archive.sync().await.expect("Failed to sync archive");
drop(archive);
let cfg = Config {
translator: FourCap,
metadata_partition: "test-metadata".into(),
key_partition: "test-index".into(),
key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
value_partition: "test-value".into(),
codec_config: (),
compression: None,
key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
replay_buffer: NZUsize!(DEFAULT_REPLAY_BUFFER),
items_per_section: NZU64!(DEFAULT_ITEMS_PER_SECTION),
};
let archive =
Archive::<_, _, FixedBytes<64>, i32>::init(context.child("second"), cfg.clone())
.await
.unwrap();
let result: Result<Option<i32>, _> = archive.get(Identifier::Index(index)).await;
assert!(matches!(
result,
Err(Error::Journal(JournalError::Codec(CodecError::ExtraData(
_
))))
));
});
}
#[test_traced]
fn test_archive_overlapping_key_basic() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = Config {
translator: FourCap,
metadata_partition: "test-metadata".into(),
key_partition: "test-index".into(),
key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
value_partition: "test-value".into(),
codec_config: (),
compression: None,
key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
replay_buffer: NZUsize!(DEFAULT_REPLAY_BUFFER),
items_per_section: NZU64!(DEFAULT_ITEMS_PER_SECTION),
};
let mut archive = Archive::init(context.child("storage"), cfg.clone())
.await
.expect("Failed to initialize archive");
let index1 = 1u64;
let key1 = test_key("keys1");
let data1 = 1;
let index2 = 2u64;
let key2 = test_key("keys2");
let data2 = 2;
archive = archive
.put(index1, key1.clone(), data1)
.await
.expect("Failed to put data");
archive = archive
.put(index2, key2.clone(), data2)
.await
.expect("Failed to put data");
let retrieved = archive
.get(Identifier::Key(&key1))
.await
.expect("Failed to get data")
.expect("Data not found");
assert_eq!(retrieved, data1);
let retrieved = archive
.get(Identifier::Key(&key2))
.await
.expect("Failed to get data")
.expect("Data not found");
assert_eq!(retrieved, data2);
let buffer = context.encode();
assert!(has_metric_value(&buffer, "items_tracked", 2));
assert!(buffer.contains("unnecessary_reads_total 1"));
assert!(buffer.contains("gets_total 2"));
});
}
#[test_traced]
fn test_archive_overlapping_key_multiple_sections() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = Config {
translator: FourCap,
metadata_partition: "test-metadata".into(),
key_partition: "test-index".into(),
key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
value_partition: "test-value".into(),
codec_config: (),
compression: None,
key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
replay_buffer: NZUsize!(DEFAULT_REPLAY_BUFFER),
items_per_section: NZU64!(DEFAULT_ITEMS_PER_SECTION),
};
let mut archive = Archive::init(context.child("storage"), cfg.clone())
.await
.expect("Failed to initialize archive");
let index1 = 1u64;
let key1 = test_key("keys1");
let data1 = 1;
let index2 = 2_000_000u64;
let key2 = test_key("keys2");
let data2 = 2;
archive = archive
.put(index1, key1.clone(), data1)
.await
.expect("Failed to put data");
archive = archive
.put(index2, key2.clone(), data2)
.await
.expect("Failed to put data");
let retrieved = archive
.get(Identifier::Key(&key1))
.await
.expect("Failed to get data")
.expect("Data not found");
assert_eq!(retrieved, data1);
let retrieved = archive
.get(Identifier::Key(&key2))
.await
.expect("Failed to get data")
.expect("Data not found");
assert_eq!(retrieved, data2);
});
}
#[test_traced]
fn test_archive_prune_keys() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = Config {
translator: FourCap,
metadata_partition: "test-metadata".into(),
key_partition: "test-index".into(),
key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
value_partition: "test-value".into(),
codec_config: (),
compression: None,
key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
replay_buffer: NZUsize!(DEFAULT_REPLAY_BUFFER),
items_per_section: NZU64!(1), };
let mut archive = Archive::init(context.child("storage"), cfg.clone())
.await
.expect("Failed to initialize archive");
let keys = vec![
(1u64, test_key("key1-blah"), 1),
(2u64, test_key("key2-blah"), 2),
(3u64, test_key("key3-blah"), 3),
(4u64, test_key("key3-bleh"), 3),
(5u64, test_key("key4-blah"), 4),
];
for (index, key, data) in &keys {
archive = archive
.put(*index, key.clone(), *data)
.await
.expect("Failed to put data");
}
let buffer = context.encode();
assert!(has_metric_value(&buffer, "items_tracked", 5));
archive = archive.prune(3).await.expect("Failed to prune");
for (index, key, data) in keys {
let retrieved = archive
.get(Identifier::Key(&key))
.await
.expect("Failed to get data");
if index < 3 {
assert!(retrieved.is_none());
} else {
assert_eq!(retrieved.expect("Data not found"), data);
}
}
let buffer = context.encode();
assert!(has_metric_value(&buffer, "items_tracked", 3));
assert!(has_metric_value(&buffer, "indices_pruned_total", 2));
assert!(has_metric_value(&buffer, "pruned_total", 0));
archive = archive.prune(2).await.expect("Failed to prune");
archive = archive.prune(3).await.expect("Failed to prune");
archive = archive
.put(6, test_key("key2-blfh"), 5)
.await
.expect("Failed to put data");
let buffer = context.encode();
assert!(has_metric_value(&buffer, "items_tracked", 4)); assert!(has_metric_value(&buffer, "indices_pruned_total", 2));
assert!(has_metric_value(&buffer, "pruned_total", 1));
let archive = archive
.put(1, test_key("key1-blah"), 1)
.await
.expect("Failed to put below floor");
assert_eq!(
archive
.get(Identifier::Key(&test_key("key1-blah")))
.await
.expect("Failed to get data"),
None
);
let (archive, handle) = archive
.put_start_sync(1, test_key("key1-blah"), 1)
.await
.expect("Failed to put_start_sync below floor");
handle.await.expect("handle must resolve");
let archive = archive
.put_sync(2, test_key("key2-blfh"), 2)
.await
.expect("Failed to put_sync below floor");
assert_eq!(archive.get(Identifier::Index(1)).await.unwrap(), None);
assert_eq!(archive.get(Identifier::Index(2)).await.unwrap(), None);
});
}
fn test_archive_keys_and_restart(num_keys: usize) -> String {
let executor = deterministic::Runner::default();
executor.start(|mut context| async move {
let items_per_section = 256u64;
let cfg = Config {
translator: TwoCap,
metadata_partition: "test-metadata".into(),
key_partition: "test-index".into(),
key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
value_partition: "test-value".into(),
codec_config: (),
compression: None,
key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
replay_buffer: NZUsize!(DEFAULT_REPLAY_BUFFER),
items_per_section: NZU64!(items_per_section),
};
let mut archive = Archive::init(
context.child("init").with_attribute("index", 1),
cfg.clone(),
)
.await
.expect("Failed to initialize archive");
let mut keys = BTreeMap::new();
while keys.len() < num_keys {
let index = keys.len() as u64;
let mut key = [0u8; 64];
context.fill(&mut key);
let key = FixedBytes::<64>::decode(key.as_ref()).unwrap();
let mut data = [0u8; 1024];
context.fill(&mut data);
let data = FixedBytes::<1024>::decode(data.as_ref()).unwrap();
archive = archive
.put(index, key.clone(), data.clone())
.await
.expect("Failed to put data");
keys.insert(key, (index, data));
}
for (key, (index, data)) in &keys {
let retrieved = archive
.get(Identifier::Index(*index))
.await
.expect("Failed to get data")
.expect("Data not found");
assert_eq!(&retrieved, data);
let retrieved = archive
.get(Identifier::Key(key))
.await
.expect("Failed to get data")
.expect("Data not found");
assert_eq!(&retrieved, data);
}
let buffer = context.encode();
assert!(has_metric_value(&buffer, "items_tracked", num_keys));
assert!(has_metric_value(&buffer, "pruned_total", 0));
let archive = archive.sync().await.expect("Failed to sync archive");
drop(archive);
let cfg = Config {
translator: TwoCap,
metadata_partition: "test-metadata".into(),
key_partition: "test-index".into(),
key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
value_partition: "test-value".into(),
codec_config: (),
compression: None,
key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
replay_buffer: NZUsize!(DEFAULT_REPLAY_BUFFER),
items_per_section: NZU64!(items_per_section),
};
let mut archive = Archive::<_, _, _, FixedBytes<1024>>::init(
context.child("init").with_attribute("index", 2),
cfg.clone(),
)
.await
.expect("Failed to initialize archive");
for (key, (index, data)) in &keys {
let retrieved = archive
.get(Identifier::Index(*index))
.await
.expect("Failed to get data")
.expect("Data not found");
assert_eq!(&retrieved, data);
let retrieved = archive
.get(Identifier::Key(key))
.await
.expect("Failed to get data")
.expect("Data not found");
assert_eq!(&retrieved, data);
}
let min = (keys.len() / 2) as u64;
archive = archive.prune(min).await.expect("Failed to prune");
let min = (min / items_per_section) * items_per_section;
let mut removed = 0;
for (key, (index, data)) in keys {
if index >= min {
let retrieved = archive
.get(Identifier::Key(&key))
.await
.expect("Failed to get data")
.expect("Data not found");
assert_eq!(retrieved, data);
let (current_end, start_next) = archive.next_gap(index);
assert_eq!(current_end.unwrap(), num_keys as u64 - 1);
assert!(start_next.is_none());
} else {
let retrieved = archive
.get(Identifier::Key(&key))
.await
.expect("Failed to get data");
assert!(retrieved.is_none());
removed += 1;
let (current_end, start_next) = archive.next_gap(index);
assert!(current_end.is_none());
assert_eq!(start_next.unwrap(), min);
}
}
let buffer = context.encode();
assert!(has_metric_value(
&buffer,
"items_tracked",
num_keys - removed
));
assert!(has_metric_value(&buffer, "indices_pruned_total", removed));
assert!(has_metric_value(&buffer, "pruned_total", 0));
context.auditor().state()
})
}
#[test_group("slow")]
#[test_traced]
fn test_archive_many_keys_and_restart() {
test_archive_keys_and_restart(100_000);
}
#[test_group("slow")]
#[test_traced]
fn test_determinism() {
let state1 = test_archive_keys_and_restart(5_000);
let state2 = test_archive_keys_and_restart(5_000);
assert_eq!(state1, state2);
}
#[test_traced]
fn test_archive_key_lookup_skips_pruned_duplicates() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = Config {
translator: FourCap,
metadata_partition: "test-metadata".into(),
key_partition: "test-index".into(),
key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
value_partition: "test-value".into(),
codec_config: (),
compression: None,
key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
replay_buffer: NZUsize!(DEFAULT_REPLAY_BUFFER),
items_per_section: NZU64!(1),
};
let mut archive = Archive::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let key = test_key("dupe-key");
archive = archive.put(2, key.clone(), 20).await.unwrap();
archive = archive.put(5, key.clone(), 50).await.unwrap();
assert!(archive.get(Identifier::Key(&key)).await.unwrap().is_some());
assert!(archive.has(Identifier::Key(&key)).await.unwrap());
archive = archive.prune(3).await.unwrap();
let got = archive.get(Identifier::Key(&key)).await.unwrap();
assert_eq!(
got,
Some(50),
"key lookup must skip the pruned entry and return the surviving one"
);
assert!(archive.has(Identifier::Key(&key)).await.unwrap());
let archive = archive.prune(6).await.unwrap();
assert_eq!(archive.get(Identifier::Key(&key)).await.unwrap(), None);
assert!(!archive.has(Identifier::Key(&key)).await.unwrap());
});
}
#[test_traced]
fn test_get_all_after_prune() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = Config {
translator: FourCap,
metadata_partition: "test-metadata".into(),
key_partition: "test-index".into(),
key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
value_partition: "test-value".into(),
codec_config: (),
compression: None,
key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
replay_buffer: NZUsize!(DEFAULT_REPLAY_BUFFER),
items_per_section: NZU64!(1),
};
let mut archive = Archive::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
archive = archive.put_multi(1, test_key("aaa"), 10).await.unwrap();
archive = archive.put_multi(1, test_key("bbb"), 20).await.unwrap();
archive = archive.put_multi(3, test_key("ccc"), 30).await.unwrap();
let archive = archive.prune(3).await.unwrap();
let all = archive.get_all(1).await.unwrap();
assert_eq!(all, None);
let all = archive.get_all(3).await.unwrap();
assert_eq!(all, Some(vec![30]));
});
}
#[test_traced]
fn test_has_at() {
let executor = deterministic::Runner::default();
let (_, checkpoint) = executor.start_and_recover(|context| async move {
let cfg = test_config(&context, "test", NZU64!(2));
let mut archive = Archive::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
assert!(!archive.has_at(1, &test_key("aaaa1")).await.unwrap());
archive = archive.put_multi(1, test_key("aaaa1"), 10).await.unwrap();
assert!(archive.has_at(1, &test_key("aaaa1")).await.unwrap());
assert!(!archive.has_at(2, &test_key("aaaa1")).await.unwrap());
assert!(!archive.has_at(1, &test_key("aaaa2")).await.unwrap());
archive = archive.put_multi(1, test_key("aaaa2"), 20).await.unwrap();
assert!(archive.has_at(1, &test_key("aaaa1")).await.unwrap());
assert!(archive.has_at(1, &test_key("aaaa2")).await.unwrap());
assert!(!archive.has_at(1, &test_key("bbbb")).await.unwrap());
archive = archive.put_multi(3, test_key("cccc"), 30).await.unwrap();
archive.sync().await.unwrap();
});
deterministic::Runner::from(checkpoint).start(|context| async move {
let cfg = test_config(&context, "test", NZU64!(2));
let archive = Archive::<_, _, FixedBytes<64>, i32>::init(context.child("reopen"), cfg)
.await
.expect("Failed to reopen archive");
assert!(archive.has_at(1, &test_key("aaaa1")).await.unwrap());
assert!(archive.has_at(1, &test_key("aaaa2")).await.unwrap());
assert!(!archive.has_at(1, &test_key("bbbb")).await.unwrap());
let archive = archive.prune(2).await.unwrap();
assert!(!archive.has_at(1, &test_key("aaaa1")).await.unwrap());
assert!(!archive.has_at(1, &test_key("aaaa2")).await.unwrap());
assert!(archive.has_at(3, &test_key("cccc")).await.unwrap());
archive.destroy().await.unwrap();
});
}
#[test_traced]
fn test_has_key() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = test_config(&context, "test", NZU64!(2));
let mut archive = Archive::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
let key = test_key("aaaa1");
assert!(!archive.has(Identifier::Key(&key)).await.unwrap());
archive = archive.put(1, key.clone(), 10).await.unwrap();
assert!(archive.has(Identifier::Key(&key)).await.unwrap());
let collision = test_key("aaaa2");
assert!(!archive.has(Identifier::Key(&collision)).await.unwrap());
archive = archive.put(2, collision.clone(), 20).await.unwrap();
assert!(archive.has(Identifier::Key(&collision)).await.unwrap());
archive = archive.put(4, test_key("cccc"), 30).await.unwrap();
let archive = archive.prune(4).await.unwrap();
assert!(!archive.has(Identifier::Key(&key)).await.unwrap());
assert!(!archive.has(Identifier::Key(&collision)).await.unwrap());
assert!(
archive
.has(Identifier::Key(&test_key("cccc")))
.await
.unwrap()
);
archive.destroy().await.unwrap();
});
}
#[test_traced]
fn test_put_multi_prune() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = Config {
translator: FourCap,
metadata_partition: "test-metadata".into(),
key_partition: "test-index".into(),
key_page_cache: CacheRef::from_pooler(&context, PAGE_SIZE, PAGE_CACHE_SIZE),
value_partition: "test-value".into(),
codec_config: (),
compression: None,
key_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
value_write_buffer: NZUsize!(DEFAULT_WRITE_BUFFER),
replay_buffer: NZUsize!(DEFAULT_REPLAY_BUFFER),
items_per_section: NZU64!(1),
};
let mut archive = Archive::init(context.child("storage"), cfg)
.await
.expect("Failed to initialize archive");
archive = archive.put_multi(1, test_key("aaa"), 10).await.unwrap();
archive = archive.put_multi(1, test_key("bbb"), 20).await.unwrap();
archive = archive.put_multi(3, test_key("ccc"), 30).await.unwrap();
let buffer = context.encode();
assert!(has_metric_value(&buffer, "items_tracked", 2));
let archive = archive.prune(3).await.unwrap();
assert_eq!(
archive
.get(Identifier::Key(&test_key("aaa")))
.await
.unwrap(),
None
);
assert_eq!(
archive
.get(Identifier::Key(&test_key("bbb")))
.await
.unwrap(),
None
);
assert_eq!(
archive
.get(Identifier::Key(&test_key("ccc")))
.await
.unwrap(),
Some(30)
);
let buffer = context.encode();
assert!(has_metric_value(&buffer, "items_tracked", 1));
assert!(has_metric_value(&buffer, "indices_pruned_total", 1));
let archive = archive
.put_multi(2, test_key("ddd"), 40)
.await
.expect("Failed to put below floor");
assert_eq!(
archive
.get(Identifier::Key(&test_key("ddd")))
.await
.expect("Failed to get data"),
None
);
let (archive, handle) = archive
.put_multi_start_sync(2, test_key("ddd"), 41)
.await
.expect("Failed to put_multi_start_sync below floor");
handle.await.expect("handle must resolve");
assert_eq!(archive.get_all(2).await.expect("Failed to get data"), None);
let archive = archive
.put_multi_sync(2, test_key("ddd"), 42)
.await
.expect("Failed to put_multi_sync below floor");
assert_eq!(archive.get_all(2).await.expect("Failed to get data"), None);
});
}
}