use gwk_domain::envelope::{Actor, EventEnvelope, Origin};
use gwk_domain::ids::{AggregateId, EventId, ProjectId, Seq, Timestamp};
use gwk_domain::port::{AppendError, EventStore, MAX_READ_LIMIT};
pub fn fixture_event(aggregate_id: &str, aggregate_version: u32) -> EventEnvelope {
EventEnvelope {
event_id: EventId::new(format!("evt-{aggregate_id}-{aggregate_version}")),
project_id: ProjectId::new("conformance"),
aggregate_type: "task".into(),
aggregate_id: AggregateId::new(aggregate_id),
aggregate_version,
event_type: "conformance_tick".into(),
schema_version: gwk_domain::envelope::ENVELOPE_SCHEMA_VERSION,
global_sequence: Seq::new(u64::MAX), occurred_at: Timestamp::new("2026-01-01T00:00:00Z"),
appended_at: Timestamp::new("9999-01-01T00:00:00Z"), actor: Actor {
kind: "conformance".into(),
id: None,
},
origin: Origin {
system: "gwk-cert".into(),
r#ref: None,
},
causation_id: None,
correlation_id: None,
idempotency_key: None,
payload: serde_json::json!({}),
payload_ref: None,
}
}
pub async fn check_append_assigns_commit_order<S: EventStore>(store: &S) {
let a = store
.append(0, None, vec![fixture_event("agg-a", 1)])
.await
.expect("first append");
let b = store
.append(0, None, vec![fixture_event("agg-b", 1)])
.await
.expect("second append");
let seq_a = a[0].global_sequence;
let seq_b = b[0].global_sequence;
assert!(seq_a.value() < seq_b.value(), "sequence not increasing");
assert_ne!(
seq_a.value(),
u64::MAX,
"store kept the client sentinel seq"
);
assert_ne!(
a[0].appended_at.as_str(),
"9999-01-01T00:00:00Z",
"store kept the client sentinel appended_at"
);
}
pub async fn check_expected_version_conflict<S: EventStore>(store: &S) {
store
.append(0, None, vec![fixture_event("agg", 1)])
.await
.expect("seed append");
let err = store
.append(0, None, vec![fixture_event("agg", 1)])
.await
.expect_err("stale append must conflict");
assert_eq!(
err,
AppendError::VersionConflict {
actual: 1,
expected: 0
}
);
}
pub async fn check_cas_refusal_and_recovery<S: EventStore>(store: &S) {
store
.append(0, None, vec![fixture_event("agg", 1)])
.await
.expect("writer A wins round 1");
let loser = store
.append(0, None, vec![fixture_event("agg", 1)])
.await
.expect_err("writer B must lose round 1");
let AppendError::VersionConflict { actual, .. } = loser else {
panic!("expected VersionConflict, got {loser:?}");
};
store
.append(actual, None, vec![fixture_event("agg", actual + 1)])
.await
.expect("writer B wins after re-reading the actual version");
}
pub async fn check_fencing<S: EventStore>(store: &S) {
let old = store.grant_fence().await.expect("first grant");
store
.append(0, Some(old), vec![fixture_event("agg", 1)])
.await
.expect("current fence accepted");
let new = store.grant_fence().await.expect("second grant");
assert!(old.value() < new.value(), "fence tokens must increase");
let err = store
.append(1, Some(old), vec![fixture_event("agg", 2)])
.await
.expect_err("stale fence must be refused");
assert!(matches!(err, AppendError::Fenced { .. }), "got {err:?}");
let err = store
.append(1, None, vec![fixture_event("agg", 2)])
.await
.expect_err("omitted fence must be refused once one is granted");
assert!(matches!(err, AppendError::Fenced { .. }), "got {err:?}");
store
.append(1, Some(new), vec![fixture_event("agg", 2)])
.await
.expect("current fence still works");
}
pub async fn check_cursor_recovery<S: EventStore>(store: &S) {
for i in 1..=5u32 {
store
.append(i - 1, None, vec![fixture_event("agg", i)])
.await
.expect("append");
}
let first_two = store.read_from(None, 2).await.expect("read page 1");
assert_eq!(first_two.len(), 2);
let cursor = first_two[1].global_sequence;
let rest = store.read_from(Some(cursor), 100).await.expect("read rest");
assert_eq!(rest.len(), 3, "cursor recovery lost events");
let mut all: Vec<u64> = first_two
.iter()
.chain(rest.iter())
.map(|e| e.global_sequence.value())
.collect();
let sorted = {
let mut s = all.clone();
s.sort_unstable();
s
};
assert_eq!(all.len(), 5);
assert_eq!(all, sorted, "read order must be ascending");
all.dedup();
assert_eq!(all.len(), 5, "duplicate sequences across pages");
}
pub async fn check_deterministic_rebuild<S: EventStore>(store: &S) {
for i in 1..=4u32 {
store
.append(i - 1, None, vec![fixture_event("agg", i)])
.await
.expect("append");
}
let once: Vec<(String, u32, u64)> = read_all(store).await;
let twice: Vec<(String, u32, u64)> = read_all(store).await;
assert_eq!(once, twice, "rebuild is not deterministic");
}
pub async fn check_watermark<S: EventStore>(store: &S) {
assert_eq!(
store.watermark().await.expect("empty watermark"),
None,
"empty store must have no watermark"
);
store
.append(0, None, vec![fixture_event("agg", 1)])
.await
.expect("append");
let last = store
.append(1, None, vec![fixture_event("agg", 2)])
.await
.expect("append")[0]
.global_sequence;
assert_eq!(
store.watermark().await.expect("watermark"),
Some(last),
"watermark must equal the last committed sequence"
);
}
pub async fn check_read_limit_is_clamped<S: EventStore>(store: &S) {
for i in 1..=8u32 {
store
.append(i - 1, None, vec![fixture_event("agg", i)])
.await
.expect("append");
}
let short = store.read_from(None, 3).await.expect("read");
assert_eq!(
short.len(),
3,
"read_from must honour a limit below the page size"
);
let huge = store
.read_from(None, usize::MAX)
.await
.expect("a huge limit must not error");
assert_eq!(huge.len(), 8, "a huge limit returns all available events");
}
pub async fn check_read_limit_ceiling<S: EventStore>(store: &S) {
let page = store
.read_from(None, usize::MAX)
.await
.expect("an unbounded read must be clamped, not refused");
assert_eq!(
page.len(),
MAX_READ_LIMIT,
"a log longer than a page must read back as exactly one page"
);
let sequences: Vec<u64> = page.iter().map(|e| e.global_sequence.value()).collect();
let mut ascending = sequences.clone();
ascending.sort_unstable();
ascending.dedup();
assert_eq!(
sequences, ascending,
"a clamped page is not in sequence order"
);
let after = store
.read_from(page.last().map(|e| e.global_sequence), MAX_READ_LIMIT)
.await
.expect("read past the first page");
assert!(
!after.is_empty(),
"the log is not longer than one page: this check proved nothing"
);
}
async fn read_all<S: EventStore>(store: &S) -> Vec<(String, u32, u64)> {
store
.read_from(None, usize::MAX)
.await
.expect("read_all")
.into_iter()
.map(|e| {
(
e.aggregate_id.as_str().to_string(),
e.aggregate_version,
e.global_sequence.value(),
)
})
.collect()
}
pub async fn run_all<S: EventStore, F: Fn() -> S>(fresh: F) {
check_append_assigns_commit_order(&fresh()).await;
check_expected_version_conflict(&fresh()).await;
check_cas_refusal_and_recovery(&fresh()).await;
check_fencing(&fresh()).await;
check_cursor_recovery(&fresh()).await;
check_deterministic_rebuild(&fresh()).await;
check_watermark(&fresh()).await;
check_read_limit_is_clamped(&fresh()).await;
}
pub mod memory {
use std::sync::Mutex;
use gwk_domain::envelope::EventEnvelope;
use gwk_domain::ids::{FenceToken, Seq, Timestamp};
use gwk_domain::port::{AppendError, EventStore, MAX_READ_LIMIT, StorageError};
#[derive(Default)]
struct Inner {
events: Vec<EventEnvelope>,
next_seq: u64,
fence: u64,
}
#[derive(Default)]
pub struct InMemoryStore {
inner: Mutex<Inner>,
}
impl InMemoryStore {
pub fn new() -> Self {
Self::default()
}
pub fn seed(&self, count: usize) {
let mut inner = self.inner.lock().expect("seed a store nothing else holds");
for _ in 0..count {
inner.next_seq += 1;
let seq = inner.next_seq;
let mut event = super::fixture_event(&format!("agg-seed-{seq}"), 1);
event.global_sequence = Seq::new(seq);
event.appended_at = Timestamp::new(format!("seq:{seq}"));
inner.events.push(event);
}
}
}
impl EventStore for InMemoryStore {
async fn append(
&self,
expected_version: u32,
fence: Option<FenceToken>,
mut events: Vec<EventEnvelope>,
) -> Result<Vec<EventEnvelope>, AppendError> {
let mut inner = self
.inner
.lock()
.map_err(|e| AppendError::Storage(e.to_string()))?;
let presented = fence.unwrap_or(FenceToken::new(0));
if presented.value() != inner.fence {
return Err(AppendError::Fenced {
presented,
current: FenceToken::new(inner.fence),
});
}
let Some(first) = events.first() else {
return Err(AppendError::MalformedBatch("empty batch".into()));
};
let agg = (first.aggregate_type.clone(), first.aggregate_id.clone());
for (i, event) in events.iter().enumerate() {
if (event.aggregate_type.clone(), event.aggregate_id.clone()) != agg {
return Err(AppendError::MalformedBatch("mixed aggregates".into()));
}
let wanted = expected_version + 1 + i as u32;
if event.aggregate_version != wanted {
return Err(AppendError::MalformedBatch(format!(
"non-contiguous version: got {}, want {wanted}",
event.aggregate_version
)));
}
}
let actual = inner
.events
.iter()
.filter(|e| (e.aggregate_type.clone(), e.aggregate_id.clone()) == agg)
.map(|e| e.aggregate_version)
.max()
.unwrap_or(0);
if actual != expected_version {
return Err(AppendError::VersionConflict {
actual,
expected: expected_version,
});
}
for event in &mut events {
inner.next_seq += 1;
event.global_sequence = Seq::new(inner.next_seq);
event.appended_at = Timestamp::new(format!("seq:{}", inner.next_seq));
}
inner.events.extend(events.iter().cloned());
Ok(events)
}
async fn read_from(
&self,
cursor: Option<Seq>,
limit: usize,
) -> Result<Vec<EventEnvelope>, StorageError> {
let inner = self.inner.lock().map_err(|e| StorageError(e.to_string()))?;
let after = cursor.map(|s| s.value()).unwrap_or(0);
Ok(inner
.events
.iter()
.filter(|e| e.global_sequence.value() > after)
.take(limit.min(MAX_READ_LIMIT))
.cloned()
.collect())
}
async fn watermark(&self) -> Result<Option<Seq>, StorageError> {
let inner = self.inner.lock().map_err(|e| StorageError(e.to_string()))?;
Ok(inner.events.last().map(|e| e.global_sequence))
}
async fn grant_fence(&self) -> Result<FenceToken, StorageError> {
let mut inner = self.inner.lock().map_err(|e| StorageError(e.to_string()))?;
inner.fence += 1;
Ok(FenceToken::new(inner.fence))
}
}
}
#[cfg(test)]
mod tests {
use super::memory::InMemoryStore;
fn block_on<F: Future>(fut: F) -> F::Output {
let waker = std::task::Waker::noop();
let mut cx = std::task::Context::from_waker(waker);
let mut fut = std::pin::pin!(fut);
loop {
match fut.as_mut().poll(&mut cx) {
std::task::Poll::Ready(v) => return v,
std::task::Poll::Pending => std::thread::yield_now(),
}
}
}
#[test]
fn reference_store_passes_the_full_suite() {
block_on(super::run_all(InMemoryStore::new));
}
#[test]
fn the_reference_store_clamps_a_log_longer_than_a_page() {
let store = InMemoryStore::new();
store.seed(gwk_domain::port::MAX_READ_LIMIT + 1);
block_on(super::check_read_limit_ceiling(&store));
}
}