use super::{Contiguous, Many, fixed, variable};
use crate::journal::{Error, contiguous::Mutable};
use commonware_macros::boxed;
use commonware_runtime::{
ReadOptions, Runner as _, Spawner as _, Supervisor as _,
buffer::paged::CacheRef,
deterministic,
mocks::{DelayedSyncContext, PendingSyncs},
reschedule,
};
use commonware_utils::{NZU16, NZU64, NZUsize};
use futures::{FutureExt as _, StreamExt, future::BoxFuture};
use std::{
future::Future,
sync::atomic::{AtomicUsize, Ordering},
};
#[boxed]
pub(super) async fn run_contiguous_tests<F, J>(factory: F)
where
F: Fn(String, usize) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let counter = AtomicUsize::new(0);
let indexed_factory = |name: String| {
let idx = counter.fetch_add(1, Ordering::SeqCst);
factory(name, idx)
};
test_empty_journal_bounds(&indexed_factory).await;
test_bounds_with_items(&indexed_factory).await;
test_bounds_after_prune(&indexed_factory).await;
test_append_and_size(&indexed_factory).await;
test_sequential_appends(&indexed_factory).await;
test_replay_from_start(&indexed_factory).await;
test_replay_from_middle(&indexed_factory).await;
test_replay_from_unsealed_tail(&indexed_factory).await;
test_replay_with_small_buffer(&indexed_factory).await;
test_prune_retains_size(&indexed_factory).await;
test_through_trait(&indexed_factory).await;
test_replay_after_prune(&indexed_factory).await;
test_prune_then_append(&indexed_factory).await;
test_position_stability(&indexed_factory).await;
test_sync_behavior(&indexed_factory).await;
test_replay_on_empty(&indexed_factory).await;
test_replay_at_exact_size(&indexed_factory).await;
test_multiple_prunes(&indexed_factory).await;
test_prune_beyond_size(&indexed_factory).await;
test_persistence_basic(&indexed_factory).await;
test_persistence_after_prune(&indexed_factory).await;
test_read_by_position(&indexed_factory).await;
test_read_many(&indexed_factory).await;
test_read_out_of_range(&indexed_factory).await;
test_read_after_prune(&indexed_factory).await;
test_rewind_to_middle(&indexed_factory).await;
test_rewind_to_zero(&indexed_factory).await;
test_rewind_current_size(&indexed_factory).await;
test_rewind_invalid_forward(&indexed_factory).await;
test_rewind_invalid_pruned(&indexed_factory).await;
test_rewind_then_append(&indexed_factory).await;
test_rewind_zero_then_append(&indexed_factory).await;
test_rewind_after_prune(&indexed_factory).await;
test_section_boundary_behavior(&indexed_factory).await;
test_destroy_and_reinit(&indexed_factory).await;
test_append_many_empty(&indexed_factory).await;
test_append_many_basic(&indexed_factory).await;
test_append_many_across_sections(&indexed_factory).await;
test_append_many_then_append(&indexed_factory).await;
test_append_many_single_item(&indexed_factory).await;
}
async fn test_empty_journal_bounds<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let journal = factory("empty".into()).await.unwrap();
let bounds = journal.bounds();
assert_eq!(bounds.start, 0);
assert_eq!(bounds.end, 0);
assert!(bounds.is_empty());
journal.destroy().await.unwrap();
}
async fn test_bounds_with_items<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("bounds-with-items".into()).await.unwrap();
for i in 0..10 {
(journal, _) = journal.append(&(i * 100)).await.unwrap();
}
let bounds = journal.bounds();
assert_eq!(bounds.start, 0);
assert_eq!(bounds.end, 10);
assert!(!bounds.is_empty());
journal.destroy().await.unwrap();
}
async fn test_bounds_after_prune<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("bounds-after-prune".into()).await.unwrap();
for i in 0..30 {
(journal, _) = journal.append(&(i * 100)).await.unwrap();
}
let bounds = journal.bounds();
assert_eq!(bounds.start, 0);
assert_eq!(bounds.end, 30);
(journal, _) = journal.prune(10).await.unwrap();
let bounds = journal.bounds();
assert_eq!(bounds.start, 10);
assert_eq!(bounds.end, 30);
(journal, _) = journal.prune(25).await.unwrap();
let bounds = journal.bounds();
assert_eq!(bounds.start, 20);
assert_eq!(bounds.end, 30);
(journal, _) = journal.prune(30).await.unwrap();
let bounds = journal.bounds();
assert_eq!(bounds.start, 30);
assert_eq!(bounds.end, 30);
assert!(bounds.is_empty());
journal.sync().await.unwrap();
let journal = factory("bounds-after-prune".into()).await.unwrap();
let bounds = journal.bounds();
assert!(bounds.is_empty());
journal.destroy().await.unwrap();
}
async fn test_append_and_size<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("append-and-size".into()).await.unwrap();
let pos1;
(journal, pos1) = journal.append(&100).await.unwrap();
let pos2;
(journal, pos2) = journal.append(&200).await.unwrap();
let pos3;
(journal, pos3) = journal.append(&300).await.unwrap();
assert_eq!(pos1, 0);
assert_eq!(pos2, 1);
assert_eq!(pos3, 2);
assert_eq!(journal.bounds().end, 3);
assert_eq!(journal.read(0).await.unwrap(), 100);
assert_eq!(journal.read(1).await.unwrap(), 200);
assert_eq!(journal.read(2).await.unwrap(), 300);
journal.destroy().await.unwrap();
}
async fn test_sequential_appends<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("sequential-appends".into()).await.unwrap();
for i in 0..25u64 {
let pos;
(journal, pos) = journal.append(&(i * 10)).await.unwrap();
assert_eq!(pos, i);
}
assert_eq!(journal.bounds().end, 25);
for i in 0..25u64 {
assert_eq!(journal.read(i).await.unwrap(), i * 10);
}
journal.destroy().await.unwrap();
}
async fn test_replay_from_start<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("replay-from-start".into()).await.unwrap();
for i in 0..10u64 {
(journal, _) = journal.append(&(i * 10)).await.unwrap();
}
{
let stream = journal
.replay(0, NZUsize!(1024), ReadOptions::default())
.await
.unwrap();
futures::pin_mut!(stream);
let mut items = Vec::new();
while let Some(result) = stream.next().await {
items.push(result.unwrap());
}
assert_eq!(items.len(), 10);
for (i, (pos, value)) in items.iter().enumerate() {
assert_eq!(*pos, i as u64);
assert_eq!(*value, (i as u64) * 10);
}
}
journal.destroy().await.unwrap();
}
async fn test_replay_from_middle<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("replay-from-middle".into()).await.unwrap();
for i in 0..15u64 {
(journal, _) = journal.append(&(i * 10)).await.unwrap();
}
{
let stream = journal
.replay(7, NZUsize!(1024), ReadOptions::default())
.await
.unwrap();
futures::pin_mut!(stream);
let mut items = Vec::new();
while let Some(result) = stream.next().await {
items.push(result.unwrap());
}
assert_eq!(items.len(), 8);
for (i, (pos, value)) in items.iter().enumerate() {
assert_eq!(*pos, (i + 7) as u64);
assert_eq!(*value, ((i + 7) as u64) * 10);
}
}
journal.destroy().await.unwrap();
}
async fn test_replay_from_unsealed_tail<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("replay-from-unsealed-tail".into()).await.unwrap();
for i in 0..17u64 {
(journal, _) = journal.append(&(i * 10)).await.unwrap();
}
{
let stream = journal
.replay(13, NZUsize!(1024), ReadOptions::default())
.await
.unwrap();
futures::pin_mut!(stream);
let mut items = Vec::new();
while let Some(result) = stream.next().await {
items.push(result.unwrap());
}
assert_eq!(items.len(), 4);
for (i, (pos, value)) in items.iter().enumerate() {
let expected_pos = (i + 13) as u64;
assert_eq!(*pos, expected_pos);
assert_eq!(*value, expected_pos * 10);
}
}
journal.destroy().await.unwrap();
}
async fn test_replay_with_small_buffer<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("replay-with-small-buffer".into()).await.unwrap();
for i in 0..25u64 {
(journal, _) = journal.append(&(i * 10)).await.unwrap();
}
{
let stream = journal
.replay(0, NZUsize!(9), ReadOptions::default())
.await
.unwrap();
futures::pin_mut!(stream);
let mut items = Vec::new();
while let Some(result) = stream.next().await {
items.push(result.unwrap());
}
assert_eq!(items.len(), 25);
for (i, (pos, value)) in items.iter().enumerate() {
assert_eq!(*pos, i as u64);
assert_eq!(*value, (i as u64) * 10);
}
}
journal.destroy().await.unwrap();
}
async fn test_prune_retains_size<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("prune-retains-size".into()).await.unwrap();
for i in 0..20u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
let size_before = journal.bounds().end;
(journal, _) = journal.prune(10).await.unwrap();
let size_after = journal.bounds().end;
assert_eq!(size_before, size_after);
assert_eq!(size_after, 20);
(journal, _) = journal.prune(20).await.unwrap();
let size_after_all = journal.bounds().end;
assert_eq!(size_after, size_after_all);
journal.sync().await.unwrap();
let journal = factory("prune-retains-size".into()).await.unwrap();
let size_after_close = journal.bounds().end;
assert_eq!(size_after_close, size_after_all);
journal.destroy().await.unwrap();
}
async fn test_through_trait<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let journal = factory("through-trait".into()).await.unwrap();
let (journal, pos1) = Mutable::append(journal, &42).await.unwrap();
let (journal, pos2) = Mutable::append(journal, &100).await.unwrap();
assert_eq!(pos1, 0);
assert_eq!(pos2, 1);
let size = Contiguous::bounds(&journal).end;
assert_eq!(size, 2);
journal.destroy().await.unwrap();
}
async fn test_replay_after_prune<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("replay-after-prune".into()).await.unwrap();
for i in 0..20u64 {
(journal, _) = journal.append(&(i * 10)).await.unwrap();
}
let (journal, _) = journal.prune(10).await.unwrap();
{
let stream = journal
.replay(10, NZUsize!(1024), ReadOptions::default())
.await
.unwrap();
futures::pin_mut!(stream);
let mut items = Vec::new();
while let Some(result) = stream.next().await {
items.push(result.unwrap());
}
assert_eq!(items.len(), 10);
for (i, (pos, value)) in items.iter().enumerate() {
assert_eq!(*pos, (i + 10) as u64);
assert_eq!(*value, ((i + 10) as u64) * 10);
}
}
journal.destroy().await.unwrap();
}
async fn test_prune_then_append<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("prune-then-append".into()).await.unwrap();
for i in 0..10u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
(journal, _) = journal.prune(10).await.unwrap();
assert!(journal.bounds().is_empty());
let pos;
(journal, pos) = journal.append(&999).await.unwrap();
assert_eq!(pos, 10);
assert_eq!(journal.bounds().end, 11);
journal.destroy().await.unwrap();
}
async fn test_position_stability<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("position-stability".into()).await.unwrap();
for i in 0..20u64 {
(journal, _) = journal.append(&(i * 100)).await.unwrap();
}
(journal, _) = journal.prune(10).await.unwrap();
for i in 20..25u64 {
let pos;
(journal, pos) = journal.append(&(i * 100)).await.unwrap();
assert_eq!(pos, i);
}
assert_eq!(journal.read(10).await.unwrap(), 1000);
assert_eq!(journal.read(15).await.unwrap(), 1500);
assert_eq!(journal.read(20).await.unwrap(), 2000);
assert_eq!(journal.read(24).await.unwrap(), 2400);
{
let stream = journal
.replay(10, NZUsize!(1024), ReadOptions::default())
.await
.unwrap();
futures::pin_mut!(stream);
let mut items = Vec::new();
while let Some(result) = stream.next().await {
items.push(result.unwrap());
}
assert_eq!(items.len(), 15);
for (i, (pos, value)) in items.iter().enumerate() {
let expected_pos = (i + 10) as u64;
assert_eq!(*pos, expected_pos);
assert_eq!(*value, expected_pos * 100);
}
}
journal.destroy().await.unwrap();
}
async fn test_sync_behavior<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("sync-behavior".into()).await.unwrap();
for i in 0..5u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
journal = journal.sync().await.unwrap();
assert_eq!(journal.read(0).await.unwrap(), 0);
let pos;
(journal, pos) = journal.append(&100).await.unwrap();
assert_eq!(pos, 5);
assert_eq!(journal.read(5).await.unwrap(), 100);
assert_eq!(journal.bounds().end, 6);
journal.destroy().await.unwrap();
}
async fn test_replay_on_empty<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let journal = factory("replay-on-empty".into()).await.unwrap();
{
let stream = journal
.replay(0, NZUsize!(1024), ReadOptions::default())
.await
.unwrap();
futures::pin_mut!(stream);
let mut items = Vec::new();
while let Some(result) = stream.next().await {
items.push(result.unwrap());
}
assert_eq!(items.len(), 0);
}
journal.destroy().await.unwrap();
}
async fn test_replay_at_exact_size<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("replay-at-exact-size".into()).await.unwrap();
for i in 0..10u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
let bounds = journal.bounds();
{
let stream = journal
.replay(bounds.end, NZUsize!(1024), ReadOptions::default())
.await
.unwrap();
futures::pin_mut!(stream);
let mut items = Vec::new();
while let Some(result) = stream.next().await {
items.push(result.unwrap());
}
assert_eq!(items.len(), 0);
}
journal.destroy().await.unwrap();
}
async fn test_multiple_prunes<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("multiple-prunes".into()).await.unwrap();
for i in 0..20u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
let pruned1;
(journal, pruned1) = journal.prune(10).await.unwrap();
let pruned2;
(journal, pruned2) = journal.prune(10).await.unwrap();
assert!(pruned1);
assert!(!pruned2);
assert_eq!(journal.bounds().end, 20);
assert_eq!(journal.read(10).await.unwrap(), 10);
assert_eq!(journal.read(19).await.unwrap(), 19);
journal.destroy().await.unwrap();
}
async fn test_prune_beyond_size<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("prune-beyond-size".into()).await.unwrap();
for i in 0..10u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
(journal, _) = journal.prune(100).await.unwrap();
assert_eq!(journal.bounds().end, 10);
let pos;
(journal, pos) = journal.append(&999).await.unwrap();
assert_eq!(pos, 10);
assert_eq!(journal.read(10).await.unwrap(), 999);
journal.destroy().await.unwrap();
}
async fn test_persistence_basic<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let test_name = "persistence-basic".to_string();
{
let mut journal = factory(test_name.clone()).await.unwrap();
for i in 0..15u64 {
let pos;
(journal, pos) = journal.append(&(i * 10)).await.unwrap();
assert_eq!(pos, i);
}
assert_eq!(journal.bounds().end, 15);
journal.sync().await.unwrap();
}
{
let journal = factory(test_name.clone()).await.unwrap();
assert_eq!(journal.bounds().end, 15);
for i in 0..15u64 {
assert_eq!(journal.read(i).await.unwrap(), i * 10);
}
{
let stream = journal
.replay(0, NZUsize!(1024), ReadOptions::default())
.await
.unwrap();
futures::pin_mut!(stream);
let mut items = Vec::new();
while let Some(result) = stream.next().await {
items.push(result.unwrap());
}
assert_eq!(items.len(), 15);
for (i, (pos, value)) in items.iter().enumerate() {
assert_eq!(*pos, i as u64);
assert_eq!(*value, (i as u64) * 10);
}
}
journal.destroy().await.unwrap();
}
}
async fn test_persistence_after_prune<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let test_name = "persistence-after-prune".to_string();
{
let mut journal = factory(test_name.clone()).await.unwrap();
for i in 0..25u64 {
(journal, _) = journal.append(&(i * 100)).await.unwrap();
}
let pruned;
(journal, pruned) = journal.prune(10).await.unwrap();
assert!(pruned);
assert_eq!(journal.bounds().end, 25);
journal.sync().await.unwrap();
}
{
let mut journal = factory(test_name.clone()).await.unwrap();
assert_eq!(journal.bounds().end, 25);
for i in 0..10u64 {
assert!(matches!(journal.read(i).await, Err(Error::ItemPruned(_))));
}
for i in 10..25u64 {
assert_eq!(journal.read(i).await.unwrap(), i * 100);
}
{
let stream = journal
.replay(10, NZUsize!(1024), ReadOptions::default())
.await
.unwrap();
futures::pin_mut!(stream);
let mut items = Vec::new();
while let Some(result) = stream.next().await {
items.push(result.unwrap());
}
assert_eq!(items.len(), 15);
for (i, (pos, value)) in items.iter().enumerate() {
let expected_pos = (i + 10) as u64;
assert_eq!(*pos, expected_pos);
assert_eq!(*value, expected_pos * 100);
}
}
let pos;
(journal, pos) = journal.append(&999).await.unwrap();
assert_eq!(pos, 25);
assert_eq!(journal.read(25).await.unwrap(), 999);
journal.destroy().await.unwrap();
}
}
pub(super) async fn test_read_by_position<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("read-by-position".into()).await.unwrap();
for i in 0..1000u64 {
(journal, _) = journal.append(&(i * 100)).await.unwrap();
assert_eq!(journal.read(i).await.unwrap(), i * 100);
}
for i in 0..1000u64 {
assert_eq!(journal.read(i).await.unwrap(), i * 100);
}
journal.destroy().await.unwrap();
}
pub(super) async fn test_read_many<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("read-many".into()).await.unwrap();
for i in 0..15u64 {
(journal, _) = journal.append(&(i * 100)).await.unwrap();
}
let items = journal.read_many(&[1, 4, 12]).await.unwrap();
assert_eq!(items, vec![100, 400, 1200]);
journal.destroy().await.unwrap();
}
pub(super) async fn test_read_out_of_range<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("read-out-of-range".into()).await.unwrap();
(journal, _) = journal.append(&42).await.unwrap();
let result = journal.read(10).await;
assert!(matches!(result, Err(Error::ItemOutOfRange(_))));
journal.destroy().await.unwrap();
}
pub(super) async fn test_read_after_prune<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("read-after-prune".into()).await.unwrap();
for i in 0..20u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
let (journal, _) = journal.prune(10).await.unwrap();
let bounds = journal.bounds();
let result = journal.read(bounds.start - 1).await;
assert!(matches!(result, Err(Error::ItemPruned(_))));
journal.destroy().await.unwrap();
}
async fn test_rewind_to_middle<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("rewind-to-middle".into()).await.unwrap();
for i in 0..20u64 {
(journal, _) = journal.append(&(i * 100)).await.unwrap();
}
journal = journal.rewind(12).await.unwrap();
assert_eq!(journal.bounds().end, 12);
for i in 0..12u64 {
assert_eq!(journal.read(i).await.unwrap(), i * 100);
}
for i in 12..20u64 {
assert!(matches!(
journal.read(i).await,
Err(Error::ItemOutOfRange(_))
));
}
let pos;
(journal, pos) = journal.append(&999).await.unwrap();
assert_eq!(pos, 12);
assert_eq!(journal.read(12).await.unwrap(), 999);
journal.destroy().await.unwrap();
}
async fn test_rewind_to_zero<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("rewind-to-zero".into()).await.unwrap();
for i in 0..10u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
journal = journal.rewind(0).await.unwrap();
let bounds = journal.bounds();
assert_eq!(bounds.end, 0);
assert!(bounds.is_empty());
let pos;
(journal, pos) = journal.append(&42).await.unwrap();
assert_eq!(pos, 0);
journal.destroy().await.unwrap();
}
async fn test_rewind_current_size<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("rewind-current-size".into()).await.unwrap();
for i in 0..10u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
let journal = journal.rewind(10).await.unwrap();
assert_eq!(journal.bounds().end, 10);
journal.destroy().await.unwrap();
}
async fn test_rewind_invalid_forward<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("rewind-invalid-forward".into()).await.unwrap();
for i in 0..10u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
let result = journal.rewind(20).await;
assert!(matches!(result, Err(Error::InvalidRewind(20))));
let journal = factory("rewind-invalid-forward".into()).await.unwrap();
journal.destroy().await.unwrap();
}
async fn test_rewind_invalid_pruned<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("rewind-invalid-pruned".into()).await.unwrap();
for i in 0..20u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
let (journal, _) = journal.prune(10).await.unwrap();
let result = journal.rewind(5).await;
assert!(matches!(result, Err(Error::ItemPruned(5))));
let journal = factory("rewind-invalid-pruned".into()).await.unwrap();
journal.destroy().await.unwrap();
}
async fn test_rewind_then_append<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("rewind-then-append".into()).await.unwrap();
for i in 0..15u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
journal = journal.rewind(8).await.unwrap();
let pos1;
(journal, pos1) = journal.append(&888).await.unwrap();
let pos2;
(journal, pos2) = journal.append(&999).await.unwrap();
assert_eq!(pos1, 8);
assert_eq!(pos2, 9);
assert_eq!(journal.read(8).await.unwrap(), 888);
assert_eq!(journal.read(9).await.unwrap(), 999);
journal.destroy().await.unwrap();
}
async fn test_rewind_zero_then_append<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("rewind-zero-then-append".into()).await.unwrap();
for i in 0..10u64 {
(journal, _) = journal.append(&(i * 100)).await.unwrap();
}
journal = journal.rewind(0).await.unwrap();
let bounds = journal.bounds();
assert_eq!(bounds.end, 0);
assert!(bounds.is_empty());
let pos;
(journal, pos) = journal.append(&42).await.unwrap();
assert_eq!(pos, 0);
assert_eq!(journal.bounds().end, 1);
assert_eq!(journal.read(0).await.unwrap(), 42);
journal.destroy().await.unwrap();
}
async fn test_rewind_after_prune<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("rewind-after-prune".into()).await.unwrap();
for i in 0..30u64 {
(journal, _) = journal.append(&(i * 100)).await.unwrap();
}
(journal, _) = journal.prune(10).await.unwrap();
let bounds = journal.bounds();
assert_eq!(bounds.start, 10);
journal = journal.rewind(20).await.unwrap();
let bounds = journal.bounds();
assert_eq!(bounds.end, 20);
assert_eq!(bounds.start, 10);
for i in bounds.start..20 {
assert_eq!(journal.read(i).await.unwrap(), i * 100);
}
let pos;
(journal, pos) = journal.append(&999).await.unwrap();
assert_eq!(pos, 20);
assert_eq!(journal.read(20).await.unwrap(), 999);
assert_eq!(journal.bounds().start, 10);
let result = journal.rewind(5).await;
assert!(matches!(result, Err(Error::ItemPruned(5))));
let journal = factory("rewind-after-prune".into()).await.unwrap();
journal.destroy().await.unwrap();
}
async fn test_section_boundary_behavior<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("section-boundary".into()).await.unwrap();
for i in 0..10u64 {
let pos;
(journal, pos) = journal.append(&(i * 100)).await.unwrap();
assert_eq!(pos, i);
}
assert_eq!(journal.bounds().end, 10);
let pos;
(journal, pos) = journal.append(&999).await.unwrap();
assert_eq!(pos, 10);
assert_eq!(journal.bounds().end, 11);
(journal, _) = journal.prune(10).await.unwrap();
assert_eq!(journal.bounds().start, 10);
assert!(matches!(journal.read(9).await, Err(Error::ItemPruned(_))));
assert_eq!(journal.read(10).await.unwrap(), 999);
let pos;
(journal, pos) = journal.append(&888).await.unwrap();
assert_eq!(pos, 11);
assert_eq!(journal.bounds().end, 12);
journal = journal.rewind(10).await.unwrap();
let bounds = journal.bounds();
assert_eq!(bounds.end, 10);
assert!(bounds.is_empty());
let pos;
(journal, pos) = journal.append(&777).await.unwrap();
assert_eq!(pos, 10);
assert_eq!(journal.bounds().end, 11);
assert_eq!(journal.read(10).await.unwrap(), 777);
assert_eq!(journal.bounds().start, 10);
journal.destroy().await.unwrap();
}
async fn test_destroy_and_reinit<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let test_name = "destroy-and-reinit".to_string();
{
let mut journal = factory(test_name.clone()).await.unwrap();
for i in 0..20u64 {
(journal, _) = journal.append(&(i * 100)).await.unwrap();
}
let (journal, _) = journal.prune(10).await.unwrap();
assert_eq!(journal.bounds().end, 20);
assert!(!journal.bounds().is_empty());
journal.destroy().await.unwrap();
}
{
let journal = factory(test_name.clone()).await.unwrap();
let bounds = journal.bounds();
assert_eq!(bounds.end, 0);
assert!(bounds.is_empty());
{
let stream = journal
.replay(0, NZUsize!(1024), ReadOptions::default())
.await
.unwrap();
futures::pin_mut!(stream);
let mut items = Vec::new();
while let Some(result) = stream.next().await {
items.push(result.unwrap());
}
assert!(items.is_empty());
}
journal.destroy().await.unwrap();
}
}
async fn test_append_many_empty<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("append-many-empty".into()).await.unwrap();
(journal, _) = journal.append(&10).await.unwrap();
(journal, _) = journal.append(&20).await.unwrap();
assert_eq!(journal.bounds().end, 2);
assert!(matches!(
journal.append_many(Many::Flat(&[])).await,
Err(Error::EmptyAppend)
));
let journal = factory("append-many-empty".into()).await.unwrap();
journal.destroy().await.unwrap();
}
async fn test_append_many_basic<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("append-many-basic".into()).await.unwrap();
let pos;
(journal, pos) = journal
.append_many(Many::Flat(&[100, 200, 300]))
.await
.unwrap();
assert_eq!(pos, 2);
assert_eq!(journal.bounds().end, 3);
assert_eq!(journal.read(0).await.unwrap(), 100);
assert_eq!(journal.read(1).await.unwrap(), 200);
assert_eq!(journal.read(2).await.unwrap(), 300);
journal.destroy().await.unwrap();
}
async fn test_append_many_across_sections<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("append-many-sections".into()).await.unwrap();
let items: Vec<u64> = (0..25).map(|i| i * 10).collect();
let pos;
(journal, pos) = journal.append_many(Many::Flat(&items)).await.unwrap();
assert_eq!(pos, 24);
assert_eq!(journal.bounds().end, 25);
for i in 0..25u64 {
assert_eq!(journal.read(i).await.unwrap(), i * 10);
}
journal.destroy().await.unwrap();
}
async fn test_append_many_then_append<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("append-many-then-single".into()).await.unwrap();
(journal, _) = journal
.append_many(Many::Flat(&[10, 20, 30]))
.await
.unwrap();
let pos;
(journal, pos) = journal.append(&40).await.unwrap();
assert_eq!(pos, 3);
assert_eq!(journal.read(0).await.unwrap(), 10);
assert_eq!(journal.read(1).await.unwrap(), 20);
assert_eq!(journal.read(2).await.unwrap(), 30);
assert_eq!(journal.read(3).await.unwrap(), 40);
journal.destroy().await.unwrap();
}
async fn test_append_many_single_item<F, J>(factory: &F)
where
F: Fn(String) -> BoxFuture<'static, Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("append-many-single".into()).await.unwrap();
let pos;
(journal, pos) = journal.append_many(Many::Flat(&[42])).await.unwrap();
assert_eq!(pos, 0);
assert_eq!(journal.read(0).await.unwrap(), 42);
journal.destroy().await.unwrap();
}
#[boxed]
async fn test_start_sync_durability<F, Fut, J>(factory: F)
where
F: Fn(&'static str) -> Fut,
Fut: Future<Output = Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = factory("a").await.unwrap();
for i in 0..7u64 {
(journal, _) = journal.append(&(i * 10)).await.unwrap();
}
let handle;
(journal, handle) = journal.start_sync().await.unwrap();
handle.await.unwrap();
let size = journal.bounds().end;
drop(journal);
let journal = factory("b").await.unwrap();
assert_eq!(journal.bounds().end, size);
for i in 0..7u64 {
assert_eq!(journal.read(i).await.unwrap(), i * 10);
}
journal.destroy().await.unwrap();
}
#[test]
fn test_fixed_start_sync_durability() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = fixed::Config {
partition: "fixed-start-sync".into(),
items_per_blob: NZU64!(3),
page_cache: CacheRef::from_pooler(&context, NZU16!(44), NZUsize!(8)),
write_buffer: NZUsize!(2048),
replay_buffer: NZUsize!(2048),
};
test_start_sync_durability(|label| {
let cfg = cfg.clone();
fixed::Journal::<_, u64>::init(context.child(label), cfg)
})
.await;
});
}
#[test]
fn test_variable_start_sync_durability() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let cfg = variable::Config {
partition: "variable-start-sync".into(),
items_per_section: NZU64!(3),
compression: None,
codec_config: (),
page_cache: CacheRef::from_pooler(&context, NZU16!(44), NZUsize!(8)),
write_buffer: NZUsize!(2048),
replay_buffer: NZUsize!(2048),
};
test_start_sync_durability(|label| {
let cfg = cfg.clone();
variable::Journal::<_, u64>::init(context.child(label), cfg)
})
.await;
});
}
#[boxed]
async fn test_start_sync_overlaps_work<F, Fut, J>(
context: deterministic::Context,
pending: PendingSyncs,
make: F,
) where
F: Fn(DelayedSyncContext<deterministic::Context>) -> Fut,
Fut: Future<Output = Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = make(DelayedSyncContext {
inner: context.child("a"),
pending: pending.clone(),
})
.await
.unwrap();
for i in 0..4u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
let handle;
(journal, handle) = journal.start_sync().await.unwrap();
assert!(pending.starts() >= 1);
assert_eq!(pending.completions(), 0);
let waiter = context
.child("await_sync")
.spawn(|_| async move { handle.await.unwrap() });
while pending.entered() == 0 {
reschedule().await;
}
(journal, _) = journal.append(&999).await.unwrap();
assert_eq!(journal.read(0).await.unwrap(), 0);
assert_eq!(
pending.completions(),
0,
"the journal made progress while the sync was still in flight"
);
pending.unblock();
waiter.await.unwrap();
assert!(pending.completions() >= 1);
let handle;
(journal, handle) = journal.start_sync().await.unwrap();
handle.await.unwrap();
drop(journal);
let journal = make(DelayedSyncContext {
inner: context.child("b"),
pending: pending.clone(),
})
.await
.unwrap();
assert_eq!(journal.bounds().end, 5);
for i in 0..4u64 {
assert_eq!(journal.read(i).await.unwrap(), i);
}
assert_eq!(journal.read(4).await.unwrap(), 999);
journal.destroy().await.unwrap();
}
#[boxed]
async fn test_start_sync_waits_for_prior<F, Fut, J>(
context: deterministic::Context,
pending: PendingSyncs,
make: F,
) where
F: Fn(DelayedSyncContext<deterministic::Context>) -> Fut,
Fut: Future<Output = Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = make(DelayedSyncContext {
inner: context.child("a"),
pending: pending.clone(),
})
.await
.unwrap();
for i in 0..4u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
let first;
(journal, first) = journal.start_sync().await.unwrap();
let starts_before = pending.starts();
(journal, _) = journal.append(&999).await.unwrap();
let (journal, second) = {
let mut call = std::pin::pin!(journal.start_sync());
assert!(
call.as_mut().now_or_never().is_none(),
"start_sync proceeded while the prior sync was pending"
);
assert_eq!(
pending.starts(),
starts_before,
"a second backend sync started while the prior sync was pending"
);
pending.unblock();
call.await.unwrap()
};
first.await.unwrap();
second.await.unwrap();
journal.destroy().await.unwrap();
}
#[boxed]
async fn test_start_sync_overlaps_predecessor_and_tail<F, Fut, J>(
context: deterministic::Context,
pending: PendingSyncs,
make: F,
) where
F: FnOnce(DelayedSyncContext<deterministic::Context>) -> Fut,
Fut: Future<Output = Result<J, Error>>,
J: Mutable<Item = u64> + 'static,
{
let mut journal = make(DelayedSyncContext {
inner: context.child("a"),
pending: pending.clone(),
})
.await
.unwrap();
for i in 0..4u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
let starts_before = pending.starts();
assert!(starts_before > 0);
let handle;
(journal, handle) = journal.start_sync().await.unwrap();
assert!(
pending.starts() > starts_before,
"tail sync was not started while predecessor was in flight"
);
let tail = {
let mut parked = pending.lock();
assert!(
parked.len() >= 2,
"expected predecessor and tail syncs parked"
);
parked.pop().unwrap()
};
tail.release.send(Ok(())).unwrap();
futures::pin_mut!(handle);
assert!(
handle.as_mut().now_or_never().is_none(),
"sync handle completed without the predecessor sync"
);
pending.unblock();
handle.await.unwrap();
journal.destroy().await.unwrap();
}
#[boxed]
async fn test_start_sync_failure_propagates<F, Fut, J>(
context: deterministic::Context,
pending: PendingSyncs,
make: F,
) where
F: FnOnce(DelayedSyncContext<deterministic::Context>) -> Fut,
Fut: Future<Output = Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = make(DelayedSyncContext {
inner: context.child("a"),
pending: pending.clone(),
})
.await
.unwrap();
for i in 0..4u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
pending.arm_fail();
pending.unblock();
let handle;
(journal, handle) = journal.start_sync().await.unwrap();
assert!(
handle.await.is_err(),
"the sync handle surfaces the failure"
);
let starts_before = pending.starts();
assert!(
matches!(journal.commit().await, Err(Error::Runtime(_))),
"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"
);
}
fn fixed_overlap_cfg(context: &deterministic::Context, partition: &str) -> fixed::Config {
fixed::Config {
partition: partition.into(),
items_per_blob: NZU64!(10),
page_cache: CacheRef::from_pooler(context, NZU16!(44), NZUsize!(8)),
write_buffer: NZUsize!(2048),
replay_buffer: NZUsize!(2048),
}
}
fn variable_overlap_cfg(context: &deterministic::Context, partition: &str) -> variable::Config<()> {
variable::Config {
partition: partition.into(),
items_per_section: NZU64!(10),
compression: None,
codec_config: (),
page_cache: CacheRef::from_pooler(context, NZU16!(44), NZUsize!(8)),
write_buffer: NZUsize!(2048),
replay_buffer: NZUsize!(2048),
}
}
#[test]
fn test_fixed_start_sync_overlaps_predecessor_and_tail() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let cfg = fixed::Config {
partition: "fixed-start-sync-predecessor".into(),
items_per_blob: NZU64!(3),
page_cache: CacheRef::from_pooler(&context, NZU16!(44), NZUsize!(8)),
write_buffer: NZUsize!(2048),
replay_buffer: NZUsize!(2048),
};
test_start_sync_overlaps_predecessor_and_tail(context, pending, move |ctx| {
fixed::Journal::<_, u64>::init(ctx, cfg)
})
.await;
});
}
#[test]
fn test_variable_start_sync_overlaps_predecessor_and_tail() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let cfg = variable::Config {
partition: "variable-start-sync-predecessor".into(),
items_per_section: NZU64!(3),
compression: None,
codec_config: (),
page_cache: CacheRef::from_pooler(&context, NZU16!(44), NZUsize!(8)),
write_buffer: NZUsize!(2048),
replay_buffer: NZUsize!(2048),
};
test_start_sync_overlaps_predecessor_and_tail(context, pending, move |ctx| {
variable::Journal::<_, u64>::init(ctx, cfg)
})
.await;
});
}
#[test]
fn test_fixed_start_sync_overlaps_work() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let cfg = fixed_overlap_cfg(&context, "fixed-start-sync-overlap");
test_start_sync_overlaps_work(context, pending, move |ctx| {
let cfg = cfg.clone();
fixed::Journal::<_, u64>::init(ctx, cfg)
})
.await;
});
}
#[test]
fn test_variable_start_sync_overlaps_work() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let cfg = variable_overlap_cfg(&context, "variable-start-sync-overlap");
test_start_sync_overlaps_work(context, pending, move |ctx| {
let cfg = cfg.clone();
variable::Journal::<_, u64>::init(ctx, cfg)
})
.await;
});
}
#[test]
fn test_fixed_start_sync_waits_for_prior() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let cfg = fixed_overlap_cfg(&context, "fixed-start-sync-waits");
test_start_sync_waits_for_prior(context, pending, move |ctx| {
let cfg = cfg.clone();
fixed::Journal::<_, u64>::init(ctx, cfg)
})
.await;
});
}
#[test]
fn test_variable_start_sync_waits_for_prior() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let cfg = variable_overlap_cfg(&context, "variable-start-sync-waits");
test_start_sync_waits_for_prior(context, pending, move |ctx| {
let cfg = cfg.clone();
variable::Journal::<_, u64>::init(ctx, cfg)
})
.await;
});
}
#[test]
fn test_fixed_start_sync_failure_propagates() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let cfg = fixed_overlap_cfg(&context, "fixed-start-sync-fail");
test_start_sync_failure_propagates(context, pending, move |ctx| {
fixed::Journal::<_, u64>::init(ctx, cfg)
})
.await;
});
}
#[test]
fn test_variable_start_sync_failure_propagates() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let cfg = variable_overlap_cfg(&context, "variable-start-sync-fail");
test_start_sync_failure_propagates(context, pending, move |ctx| {
variable::Journal::<_, u64>::init(ctx, cfg)
})
.await;
});
}
#[boxed]
async fn test_prune_waits_for_pending_sync<F, Fut, J>(
context: deterministic::Context,
pending: PendingSyncs,
make: F,
) where
F: FnOnce(DelayedSyncContext<deterministic::Context>) -> Fut,
Fut: Future<Output = Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = make(DelayedSyncContext {
inner: context.child("a"),
pending: pending.clone(),
})
.await
.unwrap();
for i in 0..4u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
assert!(pending.starts() > 0);
let prune = journal.prune(3);
futures::pin_mut!(prune);
assert!(
prune.as_mut().now_or_never().is_none(),
"prune proceeded while the rollover sync was pending"
);
pending.unblock();
let (journal, _) = prune.await.unwrap();
assert_eq!(journal.bounds(), 3..4);
journal.destroy().await.unwrap();
}
#[boxed]
async fn test_rewind_surfaces_failed_sync<F, Fut, J>(
context: deterministic::Context,
pending: PendingSyncs,
make: F,
) where
F: FnOnce(DelayedSyncContext<deterministic::Context>) -> Fut,
Fut: Future<Output = Result<J, Error>>,
J: Mutable<Item = u64>,
{
let mut journal = make(DelayedSyncContext {
inner: context.child("a"),
pending: pending.clone(),
})
.await
.unwrap();
for i in 0..4u64 {
(journal, _) = journal.append(&i).await.unwrap();
}
assert!(pending.starts() > 0);
pending.arm_fail();
pending.unblock();
assert!(
matches!(journal.rewind(3).await, Err(Error::Runtime(_))),
"rewind must surface the failed rollover sync"
);
}
#[test]
fn test_fixed_prune_waits_for_pending_sync() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let cfg = fixed::Config {
partition: "fixed-prune-pending".into(),
items_per_blob: NZU64!(3),
page_cache: CacheRef::from_pooler(&context, NZU16!(44), NZUsize!(8)),
write_buffer: NZUsize!(2048),
replay_buffer: NZUsize!(2048),
};
test_prune_waits_for_pending_sync(context, pending, move |ctx| {
fixed::Journal::<_, u64>::init(ctx, cfg)
})
.await;
});
}
#[test]
fn test_variable_prune_waits_for_pending_sync() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let cfg = variable::Config {
partition: "variable-prune-pending".into(),
items_per_section: NZU64!(3),
compression: None,
codec_config: (),
page_cache: CacheRef::from_pooler(&context, NZU16!(44), NZUsize!(8)),
write_buffer: NZUsize!(2048),
replay_buffer: NZUsize!(2048),
};
test_prune_waits_for_pending_sync(context, pending, move |ctx| {
variable::Journal::<_, u64>::init(ctx, cfg)
})
.await;
});
}
#[test]
fn test_fixed_rewind_surfaces_failed_sync() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let cfg = fixed::Config {
partition: "fixed-rewind-fail".into(),
items_per_blob: NZU64!(3),
page_cache: CacheRef::from_pooler(&context, NZU16!(44), NZUsize!(8)),
write_buffer: NZUsize!(2048),
replay_buffer: NZUsize!(2048),
};
test_rewind_surfaces_failed_sync(context, pending, move |ctx| {
fixed::Journal::<_, u64>::init(ctx, cfg)
})
.await;
});
}
#[test]
fn test_variable_rewind_surfaces_failed_sync() {
let executor = deterministic::Runner::default();
executor.start(|context| async move {
let pending = PendingSyncs::default();
let cfg = variable::Config {
partition: "variable-rewind-fail".into(),
items_per_section: NZU64!(3),
compression: None,
codec_config: (),
page_cache: CacheRef::from_pooler(&context, NZU16!(44), NZUsize!(8)),
write_buffer: NZUsize!(2048),
replay_buffer: NZUsize!(2048),
};
test_rewind_surfaces_failed_sync(context, pending, move |ctx| {
variable::Journal::<_, u64>::init(ctx, cfg)
})
.await;
});
}