#![cfg(feature = "subscription")]
#![allow(clippy::unwrap_used, reason = "tests")]
#![allow(clippy::expect_used, reason = "tests")]
#![allow(clippy::panic, reason = "tests")]
use std::time::Duration;
use futures::StreamExt;
use mnesis::Version;
use mnesis_inmemory::InMemoryStore;
use mnesis_store::PendingBatch;
use mnesis_store::store::RawEventStore;
use mnesis_store::{StepStreamExt, Store, Subscription, pending_envelope};
use tokio::time::timeout;
#[derive(Debug, Clone, Hash, PartialEq, Eq)]
struct TestId(String);
impl TestId {
fn new(s: &str) -> Self {
Self(s.to_owned())
}
}
impl std::fmt::Display for TestId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.0)
}
}
impl AsRef<[u8]> for TestId {
fn as_ref(&self) -> &[u8] {
self.0.as_bytes()
}
}
fn make_envelope(version: u64, event_type: &'static str) -> mnesis_store::PendingEnvelope {
pending_envelope(Version::new(version).unwrap())
.event_type(event_type)
.payload(format!("payload-{version}").into_bytes())
.build()
.expect("valid envelope")
}
async fn append_one(
store: &Store<InMemoryStore>,
id: &TestId,
version: u64,
expected: Option<Version>,
event_type: &'static str,
) {
let envelope = make_envelope(version, event_type);
store
.append(
&mnesis_store::StreamKey::from_slice(id.as_ref()),
expected,
PendingBatch::new(&[envelope]).expect("non-empty batch"),
)
.await
.unwrap();
}
const TIMEOUT: Duration = Duration::from_secs(2);
#[tokio::test]
async fn subscribe_returns_a_reexported_stream() {
fn assert_stream<T: mnesis_store::Stream>(_: &T) {}
let store = Store::new(InMemoryStore::new());
let id = TestId::new("stream-static");
let per_stream = Subscription::new(&store)
.subscribe(&id, None)
.unwrap()
.events();
assert_stream(&per_stream);
let all = Subscription::new(&store)
.subscribe_all(None)
.unwrap()
.events();
assert_stream(&all);
}
#[tokio::test]
async fn subscribe_catchup_then_live() {
let store = Store::new(InMemoryStore::new());
let id = TestId::new("stream-1");
append_one(&store, &id, 1, None, "E1").await;
append_one(&store, &id, 2, Version::new(1), "E2").await;
let stream = Subscription::new(&store)
.subscribe(&id, None)
.unwrap()
.events();
futures::pin_mut!(stream);
let env1 = timeout(TIMEOUT, stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(env1.event_type(), "E1");
assert_eq!(env1.version(), Version::new(1).unwrap());
let env2 = timeout(TIMEOUT, stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(env2.event_type(), "E2");
assert_eq!(env2.version(), Version::new(2).unwrap());
append_one(&store, &id, 3, Version::new(2), "E3").await;
let env3 = timeout(TIMEOUT, stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(env3.event_type(), "E3");
assert_eq!(env3.version(), Version::new(3).unwrap());
}
#[tokio::test]
async fn subscribe_from_checkpoint() {
let store = Store::new(InMemoryStore::new());
let id = TestId::new("stream-1");
append_one(&store, &id, 1, None, "E1").await;
append_one(&store, &id, 2, Version::new(1), "E2").await;
append_one(&store, &id, 3, Version::new(2), "E3").await;
let stream = Subscription::new(&store)
.subscribe(&id, Some(Version::new(2).unwrap()))
.unwrap()
.events();
futures::pin_mut!(stream);
let env = timeout(TIMEOUT, stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(env.event_type(), "E3");
assert_eq!(env.version(), Version::new(3).unwrap());
}
#[tokio::test]
async fn drop_and_resubscribe_from_position() {
let store = Store::new(InMemoryStore::new());
let id = TestId::new("stream-1");
append_one(&store, &id, 1, None, "E1").await;
let position = {
let sub_stream = Subscription::new(&store)
.subscribe(&id, None)
.unwrap()
.events();
futures::pin_mut!(sub_stream);
let first_env = timeout(TIMEOUT, sub_stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(first_env.version(), Version::new(1).unwrap());
first_env.version()
};
assert_eq!(position, Version::new(1).unwrap());
append_one(&store, &id, 2, Version::new(1), "E2").await;
append_one(&store, &id, 3, Version::new(2), "E3").await;
let stream = Subscription::new(&store)
.subscribe(&id, Some(position))
.unwrap()
.events();
futures::pin_mut!(stream);
let env2 = timeout(TIMEOUT, stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(env2.event_type(), "E2");
assert_eq!(env2.version(), Version::new(2).unwrap());
let env3 = timeout(TIMEOUT, stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(env3.event_type(), "E3");
assert_eq!(env3.version(), Version::new(3).unwrap());
}
#[tokio::test]
async fn catchup_events_appended_before_subscribe() {
let store = Store::new(InMemoryStore::new());
let id = TestId::new("stream-1");
append_one(&store, &id, 1, None, "E1").await;
append_one(&store, &id, 2, Version::new(1), "E2").await;
let stream = Subscription::new(&store)
.subscribe(&id, None)
.unwrap()
.events();
futures::pin_mut!(stream);
let env1 = timeout(TIMEOUT, stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(env1.event_type(), "E1");
assert_eq!(env1.version(), Version::new(1).unwrap());
let env2 = timeout(TIMEOUT, stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(env2.event_type(), "E2");
assert_eq!(env2.version(), Version::new(2).unwrap());
}
#[tokio::test]
async fn subscribe_to_nonexistent_stream_waits() {
let store = Store::new(InMemoryStore::new());
let id = TestId::new("ghost-stream");
let stream = Subscription::new(&store)
.subscribe(&id, None)
.unwrap()
.events();
futures::pin_mut!(stream);
let result = tokio::time::timeout(Duration::from_millis(50), stream.next()).await;
assert!(result.is_err(), "expected timeout, but got an event");
append_one(&store, &id, 1, None, "E1").await;
let env = timeout(TIMEOUT, stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(env.event_type(), "E1");
assert_eq!(env.version(), Version::new(1).unwrap());
}
#[tokio::test]
async fn subscribe_from_beyond_head() {
let store = Store::new(InMemoryStore::new());
let id = TestId::new("stream-1");
append_one(&store, &id, 1, None, "E1").await;
append_one(&store, &id, 2, Version::new(1), "E2").await;
let stream = Subscription::new(&store)
.subscribe(&id, Some(Version::new(5).unwrap()))
.unwrap()
.events();
futures::pin_mut!(stream);
let result = tokio::time::timeout(Duration::from_millis(50), stream.next()).await;
assert!(result.is_err(), "expected timeout, but got an event");
append_one(&store, &id, 3, Version::new(2), "E3").await;
append_one(&store, &id, 4, Version::new(3), "E4").await;
append_one(&store, &id, 5, Version::new(4), "E5").await;
append_one(&store, &id, 6, Version::new(5), "E6").await;
let env = timeout(TIMEOUT, stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(env.event_type(), "E6");
assert_eq!(env.version(), Version::new(6).unwrap());
}
#[tokio::test]
async fn concurrent_append_and_subscribe() {
let store = Store::new(InMemoryStore::new());
let id = TestId::new("concurrent-stream");
let event_count: u64 = 50;
let stream = Subscription::new(&store)
.subscribe(&id, None)
.unwrap()
.events();
futures::pin_mut!(stream);
let writer_store = store.clone();
let writer_id = id.clone();
let writer = tokio::spawn(async move {
for i in 1..=event_count {
let expected = if i == 1 {
None
} else {
Version::new(i.checked_sub(1).unwrap())
};
append_one(&writer_store, &writer_id, i, expected, "ConcurrentEvent").await;
tokio::task::yield_now().await;
}
});
for expected_version in 1..=event_count {
let env = timeout(TIMEOUT, stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(
env.version(),
Version::new(expected_version).unwrap(),
"expected version {expected_version}, got {}",
env.version()
);
assert_eq!(env.event_type(), "ConcurrentEvent");
}
writer.await.unwrap();
}
#[tokio::test]
async fn append_during_catchup_no_loss() {
let store = Store::new(InMemoryStore::new());
let id = TestId::new("stream-1");
for i in 1..=20u64 {
let expected = if i == 1 { None } else { Version::new(i - 1) };
append_one(&store, &id, i, expected, "Prepop").await;
}
let stream = Subscription::new(&store)
.subscribe(&id, None)
.unwrap()
.events();
futures::pin_mut!(stream);
for expected_v in 1..=5u64 {
let env = timeout(TIMEOUT, stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(env.version(), Version::new(expected_v).unwrap());
}
append_one(&store, &id, 21, Version::new(20), "Live").await;
for expected_v in 6..=21u64 {
let env = timeout(TIMEOUT, stream.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(env.version(), Version::new(expected_v).unwrap());
}
}
#[tokio::test]
async fn multiple_subscribers_same_stream() {
let store = Store::new(InMemoryStore::new());
let id = TestId::new("shared-stream");
let sub = Subscription::new(&store);
let sub1 = sub.subscribe(&id, None).unwrap().events();
let sub2 = sub.subscribe(&id, None).unwrap().events();
futures::pin_mut!(sub1, sub2);
append_one(&store, &id, 1, None, "SharedEvent").await;
let env1 = timeout(TIMEOUT, sub1.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(env1.event_type(), "SharedEvent");
assert_eq!(env1.version(), Version::new(1).unwrap());
let env2 = timeout(TIMEOUT, sub2.next())
.await
.unwrap()
.unwrap()
.unwrap();
assert_eq!(env2.event_type(), "SharedEvent");
assert_eq!(env2.version(), Version::new(1).unwrap());
}
#[tokio::test]
async fn subscription_cursor_is_static() {
fn assert_static<T: 'static>(_: &T) {}
let store = Store::new(InMemoryStore::new());
let id = TestId::new("s-1");
let sub = Subscription::new(&store)
.subscribe(&id, None)
.unwrap()
.events();
assert_static(&sub);
}