#![allow(
clippy::unwrap_used,
reason = "tests use unwrap for clarity and brevity"
)]
#![allow(
clippy::expect_used,
reason = "tests use expect for clarity and brevity"
)]
#![allow(clippy::str_to_string, reason = "tests use to_string/to_owned freely")]
#![allow(
clippy::shadow_reuse,
reason = "tests shadow variables for readability"
)]
#![allow(
clippy::shadow_unrelated,
reason = "tests shadow variables for readability"
)]
#![allow(clippy::panic, reason = "tests use panic for error propagation")]
use futures::StreamExt;
use mnesis::Version;
use mnesis_inmemory::{InMemoryStore, InMemoryStream};
use mnesis_store::PendingBatch;
use mnesis_store::pending_envelope;
use mnesis_store::store::RawEventStore;
fn make_envelopes(start: u64, count: u64) -> Vec<mnesis_store::envelope::PendingEnvelope> {
(start..start + count)
.map(|v| {
pending_envelope(Version::new(v).unwrap())
.event_type("Event")
.payload(v.to_le_bytes().to_vec())
.build()
.expect("valid envelope")
})
.collect()
}
async fn collect_stream(stream: &mut InMemoryStream) -> Vec<(u64, Vec<u8>)> {
let mut result = Vec::new();
loop {
let item = stream.next().await;
match item {
None => break,
Some(Ok(env)) => {
result.push((env.version().as_u64(), env.payload().to_vec()));
}
Some(Err(e)) => panic!("unexpected stream error: {e}"),
}
}
result
}
#[tokio::test]
async fn multi_batch_append_then_read_all() {
let store = InMemoryStore::new();
let batch1 = make_envelopes(1, 3);
store
.append(
&mnesis_store::StreamKey::from_slice(b"multi-batch-stream"),
None,
PendingBatch::new(&batch1).expect("non-empty batch"),
)
.await
.unwrap();
let batch2 = make_envelopes(4, 3);
store
.append(
&mnesis_store::StreamKey::from_slice(b"multi-batch-stream"),
Version::new(3),
PendingBatch::new(&batch2).expect("non-empty batch"),
)
.await
.unwrap();
let mut cursor = store
.read_stream(
&mnesis_store::StreamKey::from_slice(b"multi-batch-stream"),
Version::INITIAL,
)
.await
.unwrap();
let events = collect_stream(&mut cursor).await;
let versions: Vec<u64> = events.iter().map(|(v, _)| *v).collect();
assert_eq!(versions, vec![1, 2, 3, 4, 5, 6]);
}
#[tokio::test]
async fn read_stream_from_version_filters_earlier() {
let store = InMemoryStore::new();
let envelopes = make_envelopes(1, 5);
store
.append(
&mnesis_store::StreamKey::from_slice(b"filter-stream"),
None,
PendingBatch::new(&envelopes).expect("non-empty batch"),
)
.await
.unwrap();
let mut cursor = store
.read_stream(
&mnesis_store::StreamKey::from_slice(b"filter-stream"),
Version::new(3).unwrap(),
)
.await
.unwrap();
let events = collect_stream(&mut cursor).await;
let versions: Vec<u64> = events.iter().map(|(v, _)| *v).collect();
assert_eq!(versions, vec![3, 4, 5]);
}
#[tokio::test]
async fn concurrent_append_detects_conflict() {
let store = InMemoryStore::new();
let seed = make_envelopes(1, 1);
store
.append(
&mnesis_store::StreamKey::from_slice(b"conflict-stream"),
None,
PendingBatch::new(&seed).expect("non-empty batch"),
)
.await
.unwrap();
let writer_a = make_envelopes(2, 1);
let result_a = store
.append(
&mnesis_store::StreamKey::from_slice(b"conflict-stream"),
Version::new(1),
PendingBatch::new(&writer_a).expect("non-empty batch"),
)
.await;
assert!(result_a.is_ok(), "writer A should succeed");
let writer_b = make_envelopes(2, 1);
let result_b = store
.append(
&mnesis_store::StreamKey::from_slice(b"conflict-stream"),
Version::new(1),
PendingBatch::new(&writer_b).expect("non-empty batch"),
)
.await;
assert!(result_b.is_err(), "writer B should get a conflict error");
let err = result_b.unwrap_err();
let err_msg = format!("{err}");
assert!(
err_msg.contains("conflict"),
"error should mention conflict, got: {err_msg}"
);
let mut cursor = store
.read_stream(
&mnesis_store::StreamKey::from_slice(b"conflict-stream"),
Version::INITIAL,
)
.await
.unwrap();
let events = collect_stream(&mut cursor).await;
assert_eq!(events.len(), 2);
}
#[tokio::test]
async fn large_batch_append_and_sequential_readback() {
let store = InMemoryStore::new();
let count: u64 = 1000;
let envelopes = make_envelopes(1, count);
store
.append(
&mnesis_store::StreamKey::from_slice(b"large-batch-stream"),
None,
PendingBatch::new(&envelopes).expect("non-empty batch"),
)
.await
.unwrap();
let mut cursor = store
.read_stream(
&mnesis_store::StreamKey::from_slice(b"large-batch-stream"),
Version::INITIAL,
)
.await
.unwrap();
let events = collect_stream(&mut cursor).await;
assert_eq!(
u64::try_from(events.len()).unwrap_or(0),
count,
"should have exactly {count} events"
);
for (version, payload) in &events {
let expected_payload = version.to_le_bytes().to_vec();
assert_eq!(
payload, &expected_payload,
"payload mismatch at version {version}"
);
}
let versions: Vec<u64> = events.iter().map(|(v, _)| *v).collect();
let expected_versions: Vec<u64> = (1..=count).collect();
assert_eq!(versions, expected_versions, "versions should be 1..=1000");
}
#[tokio::test]
async fn read_from_future_version_returns_empty() {
let store = InMemoryStore::new();
let envelopes = make_envelopes(1, 3);
store
.append(
&mnesis_store::StreamKey::from_slice(b"future-version-stream"),
None,
PendingBatch::new(&envelopes).expect("non-empty batch"),
)
.await
.unwrap();
let mut cursor = store
.read_stream(
&mnesis_store::StreamKey::from_slice(b"future-version-stream"),
Version::new(100).unwrap(),
)
.await
.unwrap();
assert!(
cursor.next().await.is_none(),
"should return empty for future version"
);
}