#![cfg(feature = "import")]
use mnesis::Version;
use thiserror::Error;
use mnesis_store::PendingBatch;
use mnesis_store::envelope::PersistedEnvelope;
use mnesis_store::error::AppendError;
use mnesis_store::store::{RawEventStore, Store};
use mnesis_store::stream_id::StreamKey;
#[cfg(test)]
#[allow(
clippy::unwrap_used,
clippy::expect_used,
clippy::panic,
reason = "test code asserts exact values"
)]
mod tests {
use super::*;
use bytes::Bytes;
use futures::StreamExt;
use mnesis_store::envelope::{PendingEnvelope, pending_envelope};
use mnesis_store::export::EventExporter;
use mnesis_store::import::*;
use mnesis_store::value::SchemaVersion;
use proptest::prelude::*;
use static_assertions::assert_impl_all;
use std::error::Error as _;
use std::sync::Arc;
use tokio::sync::Barrier;
fn v(n: u64) -> Version {
Version::new(n).expect("test version must be nonzero")
}
fn persisted(version: u64, payload: &[u8]) -> PersistedEnvelope {
let mut buf = Vec::new();
buf.extend_from_slice(b"E");
buf.extend_from_slice(payload);
let et_end = 1u32;
let pl_end = et_end + u32::try_from(payload.len()).expect("payload fits u32");
PersistedEnvelope::try_new(
v(version),
Bytes::from(buf),
SchemaVersion::INITIAL,
0..et_end,
et_end..pl_end,
None,
)
.expect("valid persisted envelope")
}
fn evt(version: u64) -> ImportBlock {
ImportBlock::Event(persisted(version, b"p"))
}
fn section(origin: &str, blocks: Vec<ImportBlock>) -> StreamSection {
StreamSection {
origin: Bytes::copy_from_slice(origin.as_bytes()),
blocks,
}
}
fn report(stream: &str, outcome: StreamOutcome) -> StreamReport {
StreamReport {
stream: StreamKey::from_slice(stream.as_bytes()),
outcome,
}
}
assert_impl_all!(StreamOutcome: Send, Sync, Copy, Clone);
assert_impl_all!(Atomicity: Send, Sync, Copy, Clone);
assert_impl_all!(AbortReason: Send, Sync, Clone);
assert_impl_all!(StreamReport: Send, Sync, Clone);
assert_impl_all!(ImportReport: Send, Sync, Clone);
assert_impl_all!(ImportError<std::io::Error>: Send, Sync);
#[test]
fn atomicity_equality_and_copy_semantics() {
let policy = Atomicity::WholeChunk;
let copied = policy; assert_eq!(policy, copied);
assert_eq!(Atomicity::WholeChunk, Atomicity::WholeChunk);
assert_eq!(Atomicity::PerStream, Atomicity::PerStream);
assert_ne!(Atomicity::WholeChunk, Atomicity::PerStream);
}
#[test]
fn complete_is_complete_and_reaches_its_version() {
let outcome = StreamOutcome::Complete { version: v(5) };
assert!(outcome.is_complete());
assert_eq!(outcome.reached(), Some(v(5)));
}
#[test]
fn corrupt_untouched_is_not_complete_and_reaches_none() {
let outcome = StreamOutcome::Corrupt { reached: None };
assert!(!outcome.is_complete());
assert_eq!(outcome.reached(), None);
}
#[test]
fn corrupt_with_prefix_reaches_that_prefix() {
let outcome = StreamOutcome::Corrupt {
reached: Some(v(3)),
};
assert!(!outcome.is_complete());
assert_eq!(outcome.reached(), Some(v(3)));
}
#[test]
fn mismatch_untouched_reaches_none_not_got() {
let outcome = StreamOutcome::Mismatch {
reached: None,
got: v(8),
};
assert!(!outcome.is_complete());
assert_eq!(outcome.reached(), None);
}
#[test]
fn mismatch_with_prefix_reaches_prefix_not_got() {
let outcome = StreamOutcome::Mismatch {
reached: Some(v(3)),
got: v(8),
};
assert!(!outcome.is_complete());
assert_eq!(outcome.reached(), Some(v(3)));
}
#[test]
fn empty_report_is_vacuously_all_complete_with_no_unfinished() {
let report: ImportReport = ImportReport::new(Vec::new());
assert!(report.streams().is_empty());
assert!(report.all_complete());
assert_eq!(report.unfinished().count(), 0);
}
#[test]
fn all_complete_report_reports_no_unfinished() {
let streams = vec![
report("a", StreamOutcome::Complete { version: v(1) }),
report("b", StreamOutcome::Complete { version: v(2) }),
];
let report = ImportReport::new(streams.clone());
assert_eq!(report.streams(), streams.as_slice());
assert!(report.all_complete());
assert_eq!(report.unfinished().count(), 0);
}
#[test]
fn mixed_report_filters_exactly_the_non_complete_streams() {
let streams = vec![
report("done-1", StreamOutcome::Complete { version: v(1) }),
report(
"corrupt",
StreamOutcome::Corrupt {
reached: Some(v(2)),
},
),
report("done-2", StreamOutcome::Complete { version: v(3) }),
report(
"mismatch",
StreamOutcome::Mismatch {
reached: None,
got: v(9),
},
),
];
let report = ImportReport::new(streams);
let all_ids: Vec<&[u8]> = report
.streams()
.iter()
.map(|s| s.stream.as_bytes())
.collect();
assert_eq!(
all_ids,
[
b"done-1".as_ref(),
b"corrupt".as_ref(),
b"done-2".as_ref(),
b"mismatch".as_ref()
]
);
let unfinished_ids: Vec<&[u8]> = report.unfinished().map(|s| s.stream.as_bytes()).collect();
assert_eq!(unfinished_ids, [b"corrupt".as_ref(), b"mismatch".as_ref()]);
assert!(!report.all_complete());
}
#[test]
fn single_incomplete_stream_blocks_all_complete() {
let report = ImportReport::new(vec![report(
"only",
StreamOutcome::Corrupt { reached: None },
)]);
assert!(!report.all_complete());
assert_eq!(report.unfinished().count(), 1);
}
#[test]
fn abort_reason_display_strings_are_exact() {
assert_eq!(AbortReason::Corrupt.to_string(), "block failed checksum");
assert_eq!(
AbortReason::Mismatch {
expected: v(4),
got: v(7),
}
.to_string(),
"version mismatch (expected 4, got 7)",
);
}
#[test]
fn import_error_aborted_display_formats_stream_and_reason() {
let err: ImportError<std::io::Error> = ImportError::Aborted {
stream: StreamKey::from_slice(b"phone:task-123"),
reason: AbortReason::Mismatch {
expected: v(2),
got: v(5),
},
};
assert_eq!(
err.to_string(),
"chunk aborted at stream phone:task-123: version mismatch (expected 2, got 5)",
);
}
#[test]
fn import_error_version_overflow_display_is_exact() {
let err: ImportError<std::io::Error> = ImportError::VersionOverflow;
assert_eq!(err.to_string(), "version overflow");
}
#[test]
fn import_error_store_is_transparent_forwarding_display_and_source() {
#[derive(Debug, Error)]
#[error("root cause")]
struct RootCause;
#[derive(Debug, Error)]
#[error("store failed")]
struct DummyStore(#[source] RootCause);
let err: ImportError<DummyStore> = ImportError::Store(DummyStore(RootCause));
assert_eq!(err.to_string(), "store failed");
let source = err
.source()
.expect("transparent Store must forward a source");
assert_eq!(source.to_string(), "root cause");
}
#[test]
fn non_store_variants_expose_no_error_source() {
let overflow: ImportError<std::io::Error> = ImportError::VersionOverflow;
assert!(overflow.source().is_none());
let aborted: ImportError<std::io::Error> = ImportError::Aborted {
stream: StreamKey::from_slice(b"s"),
reason: AbortReason::Corrupt,
};
assert!(aborted.source().is_none());
}
fn sk(s: &str) -> StreamKey {
StreamKey::from_slice(s.as_bytes())
}
fn pending(ver: u64, payload: &[u8]) -> PendingEnvelope {
pending_envelope(v(ver))
.event_type("E")
.payload(Bytes::copy_from_slice(payload))
.build()
.expect("valid envelope")
}
fn planned(target: &str, expected: Option<u64>, versions: &[u64]) -> PlannedAppend {
let (first, rest) = versions.split_first().expect("planned run is non-empty");
PlannedAppend {
target: sk(target),
expected_version: expected.and_then(Version::new),
head: pending(*first, b"p"),
tail: rest.iter().map(|n| pending(*n, b"p")).collect(),
}
}
async fn head_len(store: &mnesis_inmemory::InMemoryStore, id: &StreamKey) -> usize {
store
.read_stream(id, Version::INITIAL)
.await
.expect("read opens")
.filter_map(|r| async move { r.ok() })
.count()
.await
}
#[tokio::test]
async fn atomic_append_many_commits_all_writes() {
let store = mnesis_inmemory::InMemoryStore::new();
let writes = vec![planned("a", None, &[1, 2]), planned("b", None, &[1])];
store.atomic_append_many(&writes).await.expect("commits");
assert_eq!(head_len(&store, &sk("a")).await, 2);
assert_eq!(head_len(&store, &sk("b")).await, 1);
}
#[tokio::test]
async fn store_handle_forwards_atomic_append_position() {
let writes = vec![planned("a", None, &[1, 2]), planned("b", None, &[1])];
let handle = Store::new(mnesis_inmemory::InMemoryStore::new());
let via_handle = handle
.atomic_append_many(&writes)
.await
.expect("handle commits");
let raw = mnesis_inmemory::InMemoryStore::new();
let via_raw = raw.atomic_append_many(&writes).await.expect("raw commits");
assert!(
via_handle.is_some(),
"a non-empty write returns a position, never None"
);
assert_eq!(
via_handle, via_raw,
"the Store handle must forward the raw store's position verbatim"
);
}
#[tokio::test]
async fn atomic_append_many_rolls_back_all_on_one_conflict() {
let store = mnesis_inmemory::InMemoryStore::new();
store
.append(
&sk("b"),
None,
PendingBatch::new(&[pending(1, b"seed")]).expect("non-empty batch"),
)
.await
.expect("seed");
let writes = vec![
planned("a", None, &[1, 2]), planned("b", None, &[1]), ];
let err = store
.atomic_append_many(&writes)
.await
.expect_err("must conflict");
match err {
AtomicAppendError::Conflict { index, actual } => {
assert_eq!(index, 1);
assert_eq!(actual, Version::new(1));
}
other => panic!("expected Conflict, got {other:?}"),
}
assert_eq!(head_len(&store, &sk("a")).await, 0, "rolled back");
assert_eq!(head_len(&store, &sk("b")).await, 1, "unchanged");
}
fn identity_route(origin: &[u8]) -> StreamKey {
StreamKey::from_slice(origin)
}
async fn versions(store: &mnesis_inmemory::InMemoryStore, id: &StreamKey) -> Vec<u64> {
store
.read_stream(id, Version::INITIAL)
.await
.expect("read opens")
.map(|r| r.expect("no read error").version().as_u64())
.collect()
.await
}
#[tokio::test]
async fn per_stream_clean_multi_stream_all_complete() {
let store = mnesis_inmemory::InMemoryStore::new();
let sections = vec![
section("a", vec![evt(1), evt(2), evt(3)]),
section("b", vec![evt(1), evt(2)]),
];
let report = store
.import(§ions, identity_route, Atomicity::PerStream)
.await
.expect("import ok");
assert!(report.all_complete());
assert_eq!(report.streams().len(), 2);
assert_eq!(versions(&store, &sk("a")).await, vec![1, 2, 3]);
assert_eq!(versions(&store, &sk("b")).await, vec![1, 2]);
assert_eq!(
report.streams()[0].outcome,
StreamOutcome::Complete { version: v(3) }
);
}
#[tokio::test]
async fn per_stream_partial_overlap_is_picky_rejected() {
let store = mnesis_inmemory::InMemoryStore::new();
for n in 1..=3 {
store
.append(
&sk("a"),
Version::new(n - 1),
PendingBatch::new(&[pending(n, b"seed")]).expect("non-empty batch"),
)
.await
.expect("seed");
}
let sections = vec![section("a", vec![evt(2), evt(3), evt(4), evt(5)])];
let report = store
.import(§ions, identity_route, Atomicity::PerStream)
.await
.expect("import ok");
assert_eq!(
report.streams()[0].outcome,
StreamOutcome::Mismatch {
reached: None,
got: v(2)
}
);
assert_eq!(versions(&store, &sk("a")).await, vec![1, 2, 3]);
}
#[tokio::test]
async fn per_stream_internal_gap_applies_prefix_then_halts() {
let store = mnesis_inmemory::InMemoryStore::new();
let sections = vec![section("a", vec![evt(1), evt(2), evt(4)])];
let report = store
.import(§ions, identity_route, Atomicity::PerStream)
.await
.expect("import ok");
assert_eq!(
report.streams()[0].outcome,
StreamOutcome::Mismatch {
reached: Some(v(2)),
got: v(4)
}
);
assert_eq!(versions(&store, &sk("a")).await, vec![1, 2], "v4 held back");
}
#[tokio::test]
async fn per_stream_mid_section_corrupt_applies_prefix_then_corrupt() {
let store = mnesis_inmemory::InMemoryStore::new();
let sections = vec![section(
"a",
vec![evt(1), evt(2), ImportBlock::Corrupt, evt(3)],
)];
let report = store
.import(§ions, identity_route, Atomicity::PerStream)
.await
.expect("import ok");
assert_eq!(
report.streams()[0].outcome,
StreamOutcome::Corrupt {
reached: Some(v(2))
}
);
assert_eq!(versions(&store, &sk("a")).await, vec![1, 2]);
}
#[tokio::test]
async fn per_stream_first_block_corrupt_reaches_none() {
let store = mnesis_inmemory::InMemoryStore::new();
let sections = vec![section("a", vec![ImportBlock::Corrupt, evt(1)])];
let report = store
.import(§ions, identity_route, Atomicity::PerStream)
.await
.expect("import ok");
assert_eq!(
report.streams()[0].outcome,
StreamOutcome::Corrupt { reached: None }
);
assert_eq!(versions(&store, &sk("a")).await, Vec::<u64>::new());
}
#[tokio::test]
async fn per_stream_empty_chunk_is_empty_report() {
let store = mnesis_inmemory::InMemoryStore::new();
let report = store
.import(&[], identity_route, Atomicity::PerStream)
.await
.expect("import ok");
assert!(report.streams().is_empty());
assert!(report.all_complete());
}
#[tokio::test]
async fn per_stream_version_overflow_surfaces_error() {
let store = mnesis_inmemory::InMemoryStore::new();
let sections = vec![section("a", vec![evt(u64::MAX), evt(1)])];
let err = store
.import(§ions, identity_route, Atomicity::PerStream)
.await
.expect_err("overflow");
assert!(matches!(err, ImportError::VersionOverflow));
}
struct FailingAppendStore {
inner: mnesis_inmemory::InMemoryStore,
fail_on: String,
}
impl mnesis_store::store::RawEventStore for FailingAppendStore {
type Error = mnesis_inmemory::InMemoryStoreError;
type Stream = mnesis_inmemory::InMemoryStream;
type AllPosition = mnesis_inmemory::InMemoryAllPos;
type AllStream = mnesis_inmemory::InMemoryAllStream;
async fn append(
&self,
id: &StreamKey,
expected_version: Option<Version>,
envelopes: mnesis_store::envelope::PendingBatch<'_>,
) -> Result<Self::AllPosition, AppendError<Self::Error>> {
if id.to_string() == self.fail_on {
return Err(AppendError::Store(
mnesis_inmemory::InMemoryStoreError::VersionOverflow,
));
}
self.inner.append(id, expected_version, envelopes).await
}
async fn read_stream(
&self,
id: &StreamKey,
from: Version,
) -> Result<Self::Stream, Self::Error> {
self.inner.read_stream(id, from).await
}
async fn read_all(
&self,
from: Option<Self::AllPosition>,
) -> Result<Self::AllStream, Self::Error> {
self.inner.read_all(from).await
}
}
impl mnesis_store::import::AtomicAppend for FailingAppendStore {
async fn atomic_append_many(
&self,
writes: &[mnesis_store::import::PlannedAppend],
) -> Result<Option<Self::AllPosition>, mnesis_store::import::AtomicAppendError<Self::Error>>
{
self.inner.atomic_append_many(writes).await
}
}
#[tokio::test]
async fn per_stream_store_error_propagates_and_keeps_prior_commits() {
let store = FailingAppendStore {
inner: mnesis_inmemory::InMemoryStore::new(),
fail_on: "b".to_owned(),
};
let sections = vec![
section("a", vec![evt(1), evt(2)]), section("b", vec![evt(1)]), ];
let err = store
.import(§ions, identity_route, Atomicity::PerStream)
.await
.expect_err("store error must surface");
assert!(matches!(err, ImportError::Store(_)));
assert_eq!(versions(&store.inner, &sk("a")).await, vec![1, 2]);
}
#[tokio::test]
async fn whole_chunk_clean_commits_all_complete() {
let store = mnesis_inmemory::InMemoryStore::new();
let sections = vec![
section("a", vec![evt(1), evt(2)]),
section("b", vec![evt(1)]),
];
let report = store
.import(§ions, identity_route, Atomicity::WholeChunk)
.await
.expect("import ok");
assert!(report.all_complete());
assert_eq!(versions(&store, &sk("a")).await, vec![1, 2]);
assert_eq!(versions(&store, &sk("b")).await, vec![1]);
}
#[tokio::test]
async fn whole_chunk_corrupt_block_aborts_nothing_lands() {
let store = mnesis_inmemory::InMemoryStore::new();
let sections = vec![
section("a", vec![evt(1), evt(2)]), section("b", vec![evt(1), ImportBlock::Corrupt]), ];
let err = store
.import(§ions, identity_route, Atomicity::WholeChunk)
.await
.expect_err("aborts");
match err {
ImportError::Aborted { stream, reason } => {
assert_eq!(stream, sk("b"));
assert_eq!(reason, AbortReason::Corrupt);
}
other => panic!("expected Aborted, got {other:?}"),
}
assert_eq!(versions(&store, &sk("a")).await, Vec::<u64>::new());
assert_eq!(versions(&store, &sk("b")).await, Vec::<u64>::new());
}
#[tokio::test]
async fn whole_chunk_internal_gap_aborts_mismatch() {
let store = mnesis_inmemory::InMemoryStore::new();
let sections = vec![section("a", vec![evt(1), evt(2), evt(4)])];
let err = store
.import(§ions, identity_route, Atomicity::WholeChunk)
.await
.expect_err("aborts");
match err {
ImportError::Aborted { stream, reason } => {
assert_eq!(stream, sk("a"));
assert_eq!(
reason,
AbortReason::Mismatch {
expected: v(3),
got: v(4)
}
);
}
other => panic!("expected Aborted, got {other:?}"),
}
assert_eq!(versions(&store, &sk("a")).await, Vec::<u64>::new());
}
#[tokio::test]
async fn whole_chunk_head_conflict_aborts_mismatch() {
let store = mnesis_inmemory::InMemoryStore::new();
for n in 1..=2 {
store
.append(
&sk("a"),
Version::new(n - 1),
PendingBatch::new(&[pending(n, b"seed")]).expect("non-empty batch"),
)
.await
.expect("seed");
}
let sections = vec![section("a", vec![evt(1), evt(2)])];
let err = store
.import(§ions, identity_route, Atomicity::WholeChunk)
.await
.expect_err("aborts");
match err {
ImportError::Aborted { stream, reason } => {
assert_eq!(stream, sk("a"));
assert_eq!(
reason,
AbortReason::Mismatch {
expected: v(3),
got: v(1)
}
);
}
other => panic!("expected Aborted, got {other:?}"),
}
assert_eq!(versions(&store, &sk("a")).await, vec![1, 2], "unchanged");
}
#[tokio::test]
async fn whole_chunk_first_block_corrupt_aborts() {
let store = mnesis_inmemory::InMemoryStore::new();
let sections = vec![section("a", vec![ImportBlock::Corrupt])];
let err = store
.import(§ions, identity_route, Atomicity::WholeChunk)
.await
.expect_err("aborts");
assert!(matches!(
err,
ImportError::Aborted { ref stream, reason: AbortReason::Corrupt } if *stream == sk("a")
));
}
#[tokio::test]
async fn whole_chunk_conflict_on_later_section_reports_that_section() {
let store = mnesis_inmemory::InMemoryStore::new();
store
.append(
&sk("b"),
None,
PendingBatch::new(&[pending(1, b"seed")]).expect("non-empty batch"),
)
.await
.expect("seed");
let sections = vec![
section("a", vec![evt(1), evt(2)]),
section("b", vec![evt(1)]),
];
let err = store
.import(§ions, identity_route, Atomicity::WholeChunk)
.await
.expect_err("aborts");
match err {
ImportError::Aborted { stream, reason } => {
assert_eq!(
stream,
sk("b"),
"must report the conflicting SECOND section"
);
assert_eq!(
reason,
AbortReason::Mismatch {
expected: v(2),
got: v(1)
}
);
}
other => panic!("expected Aborted, got {other:?}"),
}
assert_eq!(versions(&store, &sk("a")).await, Vec::<u64>::new());
assert_eq!(versions(&store, &sk("b")).await, vec![1], "unchanged");
}
#[tokio::test]
async fn whole_chunk_empty_section_skipped_others_commit() {
let store = mnesis_inmemory::InMemoryStore::new();
let sections = vec![section("empty", vec![]), section("b", vec![evt(1), evt(2)])];
let report = store
.import(§ions, identity_route, Atomicity::WholeChunk)
.await
.expect("import ok");
assert!(report.all_complete());
assert_eq!(report.streams().len(), 1);
assert_eq!(report.streams()[0].stream, sk("b"));
assert_eq!(versions(&store, &sk("b")).await, vec![1, 2]);
}
#[tokio::test]
async fn whole_chunk_non_injective_route_aborts_no_corruption() {
let store = mnesis_inmemory::InMemoryStore::new();
let sections = vec![
section("o1", vec![evt(1), evt(2)]),
section("o2", vec![evt(1), evt(2)]),
];
let to_same = |_origin: &[u8]| sk("T");
let err = store
.import(§ions, to_same, Atomicity::WholeChunk)
.await
.expect_err("non-injective route must abort");
assert!(matches!(err, ImportError::Aborted { .. }));
assert_eq!(versions(&store, &sk("T")).await, Vec::<u64>::new());
}
#[tokio::test]
async fn per_stream_non_injective_route_second_rejected_no_corruption() {
let store = mnesis_inmemory::InMemoryStore::new();
let sections = vec![
section("o1", vec![evt(1), evt(2)]),
section("o2", vec![evt(1), evt(2)]),
];
let to_same = |_origin: &[u8]| sk("T");
let report = store
.import(§ions, to_same, Atomicity::PerStream)
.await
.expect("import ok");
assert_eq!(
versions(&store, &sk("T")).await,
vec![1, 2],
"no [1,2,1,2] corruption — second section rejected, not appended"
);
assert_eq!(
report.streams()[0].outcome,
StreamOutcome::Complete { version: v(2) }
);
assert_eq!(
report.streams()[1].outcome,
StreamOutcome::Mismatch {
reached: None,
got: v(1)
}
);
}
#[tokio::test]
async fn store_handle_imports_without_raw() {
fn assert_importer<I: EventImporter>(_: &I) {}
fn assert_atomic<A: AtomicAppend>(_: &A) {}
let store = Store::new(mnesis_inmemory::InMemoryStore::new());
let sections = vec![section("a", vec![evt(1), evt(2)])];
let report = store
.import(§ions, identity_route, Atomicity::PerStream)
.await
.expect("import via the handle");
assert!(report.all_complete());
assert_eq!(
report.streams()[0].outcome,
StreamOutcome::Complete { version: v(2) }
);
store
.atomic_append_many(&[planned("b", None, &[1])])
.await
.expect("atomic append via the handle");
assert_importer(&store);
assert_atomic(&store);
}
async fn export_section(store: &mnesis_inmemory::InMemoryStore, origin: &str) -> StreamSection {
let blocks = store
.export_stream(&sk(origin), Version::INITIAL)
.await
.expect("export opens")
.map(|r| ImportBlock::Event(r.expect("no read error")))
.collect::<Vec<_>>()
.await;
section(origin, blocks)
}
#[tokio::test]
async fn export_then_import_round_trips_byte_equal_modulo_all_position() {
let source = mnesis_inmemory::InMemoryStore::new();
for n in 1..=3 {
source
.append(
&sk("acct-1"),
Version::new(n - 1),
PendingBatch::new(&[pending(n, b"a")]).expect("non-empty batch"),
)
.await
.expect("seed a");
}
source
.append(
&sk("acct-2"),
None,
PendingBatch::new(&[pending(1, b"b")]).expect("non-empty batch"),
)
.await
.expect("seed b");
let sections = vec![
export_section(&source, "acct-1").await,
export_section(&source, "acct-2").await,
];
let target = mnesis_inmemory::InMemoryStore::new();
target
.append(
&sk("warmup"),
None,
PendingBatch::new(&[pending(1, b"warmup")]).expect("non-empty batch"),
)
.await
.expect("warmup seed");
let report = target
.import(§ions, identity_route, Atomicity::PerStream)
.await
.expect("import ok");
assert!(report.all_complete());
for origin in ["acct-1", "acct-2"] {
let src: Vec<PersistedEnvelope> = source
.read_stream(&sk(origin), Version::INITIAL)
.await
.expect("read")
.map(|r| r.expect("ok"))
.collect()
.await;
let dst: Vec<PersistedEnvelope> = target
.read_stream(&sk(origin), Version::INITIAL)
.await
.expect("read")
.map(|r| r.expect("ok"))
.collect()
.await;
assert_eq!(src.len(), dst.len(), "{origin} length");
for (s, d) in src.iter().zip(dst.iter()) {
assert_eq!(s.version(), d.version(), "{origin} version");
assert_eq!(s.event_type(), d.event_type(), "{origin} type");
assert_eq!(s.payload(), d.payload(), "{origin} payload");
assert_eq!(s.metadata(), d.metadata(), "{origin} metadata");
assert_eq!(s.schema_version(), d.schema_version(), "{origin} schema");
}
}
let target_positions: Vec<u64> = target
.read_all(None)
.await
.expect("read_all")
.map(|r| r.expect("ok").0.as_u64())
.collect()
.await;
assert_eq!(
target_positions,
vec![1, 2, 3, 4, 5],
"import must re-stamp fresh $all positions, not copy the source's",
);
}
#[tokio::test]
async fn reimport_same_chunk_is_idempotent_per_stream() {
let store = mnesis_inmemory::InMemoryStore::new();
let sections = vec![section("a", vec![evt(1), evt(2), evt(3)])];
let first = store
.import(§ions, identity_route, Atomicity::PerStream)
.await
.expect("import ok");
assert!(first.all_complete());
let second = store
.import(§ions, identity_route, Atomicity::PerStream)
.await
.expect("import ok");
assert_eq!(
second.streams()[0].outcome,
StreamOutcome::Mismatch {
reached: None,
got: v(1)
}
);
assert_eq!(versions(&store, &sk("a")).await, vec![1, 2, 3], "no dupes");
}
#[tokio::test]
async fn reimport_same_chunk_whole_chunk_aborts() {
let store = mnesis_inmemory::InMemoryStore::new();
let sections = vec![section("a", vec![evt(1), evt(2)])];
store
.import(§ions, identity_route, Atomicity::WholeChunk)
.await
.expect("first import ok");
let err = store
.import(§ions, identity_route, Atomicity::WholeChunk)
.await
.expect_err("re-import aborts");
match err {
ImportError::Aborted { stream, reason } => {
assert_eq!(stream, sk("a"));
assert_eq!(
reason,
AbortReason::Mismatch {
expected: v(3),
got: v(1)
}
);
}
other => panic!("expected Aborted, got {other:?}"),
}
assert_eq!(versions(&store, &sk("a")).await, vec![1, 2], "no dupes");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn per_stream_import_concurrent_with_writer_surfaces_conflict_no_torn_stream() {
let store = Arc::new(mnesis_inmemory::InMemoryStore::new());
let n = 50u64;
let barrier = Arc::new(Barrier::new(2));
let writer_store = Arc::clone(&store);
let writer_barrier = Arc::clone(&barrier);
let writer = tokio::spawn(async move {
writer_barrier.wait().await;
for vn in 1..=n {
let _ = writer_store
.append(
&sk("race"),
Version::new(vn - 1),
PendingBatch::new(&[pending(vn, b"w")]).expect("non-empty batch"),
)
.await;
}
});
let importer_store = Arc::clone(&store);
let importer_barrier = Arc::clone(&barrier);
let importer = tokio::spawn(async move {
importer_barrier.wait().await;
let blocks: Vec<ImportBlock> = (1..=n).map(evt).collect();
let sections = vec![section("race", blocks)];
importer_store
.import(§ions, identity_route, Atomicity::PerStream)
.await
.expect("import returns Ok (conflicts are per-stream outcomes)")
});
writer.await.expect("writer task");
let report = importer.await.expect("importer task");
let outcome = report.streams().first().map(|s| s.outcome);
match outcome {
Some(
StreamOutcome::Complete { .. } | StreamOutcome::Mismatch { reached: None, .. },
) => {}
other => panic!("unexpected importer outcome: {other:?}"),
}
let final_versions = versions(&store, &sk("race")).await;
for (expected, got) in (1u64..).zip(final_versions.iter()) {
assert_eq!(*got, expected, "stream must be a gapless prefix from 1");
}
assert!(!final_versions.is_empty(), "some events landed");
}
#[derive(Debug, Clone)]
enum Cmd {
Import { first: u64, count: u64 },
}
fn cmd_strategy() -> impl Strategy<Value = Cmd> {
(prop_oneof![Just(1u64), Just(2u64), Just(3u64)], 1u64..=3)
.prop_map(|(first, count)| Cmd::Import { first, count })
}
proptest! {
#![proptest_config(ProptestConfig::with_cases(64))]
#[test]
fn state_machine_per_stream_never_gaps_and_report_matches_model(
cmds in proptest::collection::vec(cmd_strategy(), 1..12),
) {
let rt = tokio::runtime::Builder::new_current_thread()
.build()
.expect("runtime");
rt.block_on(async {
let store = mnesis_inmemory::InMemoryStore::new();
let id = sk("m");
let mut next: u64 = 1;
for cmd in cmds {
let Cmd::Import { first, count } = cmd;
let blocks: Vec<ImportBlock> =
(first..first + count).map(evt).collect();
let sections = vec![section("m", blocks)];
let report = {
store
.import(§ions, identity_route, Atomicity::PerStream)
.await
.expect("import ok")
};
let outcome = report.streams()[0].outcome;
if first == next {
let landed_last = first + count - 1;
prop_assert_eq!(
outcome,
StreamOutcome::Complete { version: v(landed_last) }
);
next = landed_last + 1;
} else {
prop_assert_eq!(
outcome,
StreamOutcome::Mismatch { reached: None, got: v(first) }
);
}
let got = versions(&store, &id).await;
let expected: Vec<u64> = (1..next).collect();
prop_assert_eq!(got, expected);
}
Ok(())
})?;
}
}
}