use chrono::{DateTime, Utc};
use turnframe_core::ids::{CaseRevision, CommandId, EventId};
use super::ConformanceFailure;
use super::fixtures::{
PERSONAL_DATA, account, at, case_key, ensure, ensure_code, ensure_eq, ensure_error, ensure_ok,
epoch, event_batch, event_batch_for, event_batch_with_personal_data, other_account,
other_case_key, other_redaction_authority, redaction_authority,
};
use crate::error::{StoreError, codes};
use crate::events::{EventCursor, StoredEvent};
use crate::stores::Stores;
pub async fn check_event_append_ordering_and_get_by_ids(
stores: &Stores,
) -> Result<(), ConformanceFailure> {
const CHECK: &str = "check_event_append_ordering_and_get_by_ids";
let events = stores.events();
let account = account();
let command = CommandId::new();
let first_ids = vec![EventId::new(), EventId::new()];
let appended = ensure_ok(
CHECK,
"appending the first batch",
events
.append(event_batch(&account, command, 1, &first_ids, epoch()))
.await,
)?;
ensure_eq(
CHECK,
"append returns the identifiers in batch order",
&appended,
&first_ids,
)?;
let second_ids = vec![EventId::new()];
ensure_ok(
CHECK,
"appending the second batch",
events
.append(event_batch(
&account,
CommandId::new(),
2,
&second_ids,
at(1),
))
.await,
)?;
let history = ensure_ok(
CHECK,
"listing the whole history",
events
.list_since(&account, &case_key(), CaseRevision::ZERO, 10)
.await,
)?;
ensure_eq(
CHECK,
"history in append order",
&history.iter().map(|e| e.event_id).collect::<Vec<_>>(),
&vec![first_ids[0], first_ids[1], second_ids[0]],
)?;
ensure(
CHECK,
history
.windows(2)
.all(|pair| pair[0].sequence < pair[1].sequence),
"the store-assigned sequence must strictly increase in append order",
)?;
ensure_eq(
CHECK,
"the revision stamped on an event",
&history[2].case_revision,
&CaseRevision(2),
)?;
let since = ensure_ok(
CHECK,
"listing from revision 1",
events
.list_since(&account, &case_key(), CaseRevision(1), 10)
.await,
)?;
ensure_eq(
CHECK,
"list_since is exclusive on the cursor",
&since.iter().map(|e| e.event_id).collect::<Vec<_>>(),
&second_ids,
)?;
let limited = ensure_ok(
CHECK,
"listing with a limit",
events
.list_since(&account, &case_key(), CaseRevision::ZERO, 2)
.await,
)?;
ensure_eq(CHECK, "the limit is honoured", &limited.len(), &2)?;
let requested = vec![second_ids[0], EventId::new(), first_ids[0]];
let by_ids = ensure_ok(
CHECK,
"reading events by identifier",
events.get_by_ids(&account, &requested).await,
)?;
ensure_eq(
CHECK,
"events come back in the order requested, unknown ones omitted",
&by_ids.iter().map(|e| e.event_id).collect::<Vec<_>>(),
&vec![second_ids[0], first_ids[0]],
)?;
let foreign = ensure_ok(
CHECK,
"reading this tenant's events as another tenant",
events.get_by_ids(&other_account(), &requested).await,
)?;
ensure(
CHECK,
foreign.is_empty(),
"another tenant must not read these events by identifier",
)?;
let count = ensure_ok(
CHECK,
"counting the events of the case",
events.count(&account, &case_key()).await,
)?;
ensure_eq(CHECK, "events on the case", &count, &3)?;
let fresh = EventId::new();
ensure_error(
CHECK,
"appending a batch containing an event that already exists",
events
.append(event_batch(
&account,
command,
3,
&[fresh, first_ids[0]],
at(2),
))
.await,
&StoreError::Conflict,
)?;
let after_conflict = ensure_ok(
CHECK,
"counting after the rejected batch",
events.count(&account, &case_key()).await,
)?;
ensure_eq(
CHECK,
"a rejected batch must leave nothing behind",
&after_conflict,
&3,
)?;
let orphan = ensure_ok(
CHECK,
"reading the event of the rejected batch",
events.get_by_ids(&account, &[fresh]).await,
)?;
ensure(
CHECK,
orphan.is_empty(),
"no event of a rejected batch may be visible",
)?;
ensure_code(
CHECK,
"appending an empty batch",
events
.append(event_batch(&account, command, 4, &[], at(3)))
.await,
codes::INVALID_RECORD,
)
}
pub async fn check_event_stream_cursor_pages_exactly_once(
stores: &Stores,
) -> Result<(), ConformanceFailure> {
const CHECK: &str = "check_event_stream_cursor_pages_exactly_once";
let events = stores.events();
let account = account();
let first: Vec<EventId> = (0..3).map(|_| EventId::new()).collect();
ensure_ok(
CHECK,
"appending three events at one revision",
events
.append(event_batch(&account, CommandId::new(), 1, &first, epoch()))
.await,
)?;
let page = ensure_ok(
CHECK,
"reading the first page",
events.read_from(&account, EventCursor::START, 2).await,
)?;
ensure_eq(
CHECK,
"the first page honours the limit",
&page.events.iter().map(|e| e.event_id).collect::<Vec<_>>(),
&first[..2].to_vec(),
)?;
ensure(
CHECK,
page.next_cursor == EventCursor::after(page.events[1].sequence),
"the page must resume after the last event it carried",
)?;
let second: Vec<EventId> = (0..2).map(|_| EventId::new()).collect();
ensure_ok(
CHECK,
"appending to another case of the same account",
events
.append(event_batch_for(
&account,
&other_case_key(),
CommandId::new(),
1,
&second,
at(1),
))
.await,
)?;
let foreign = vec![EventId::new()];
ensure_ok(
CHECK,
"appending an event of another tenant",
events
.append(event_batch(
&other_account(),
CommandId::new(),
1,
&foreign,
at(2),
))
.await,
)?;
let third: Vec<EventId> = vec![EventId::new()];
ensure_ok(
CHECK,
"appending one more event of the first case",
events
.append(event_batch(&account, CommandId::new(), 2, &third, at(3)))
.await,
)?;
let mut cursor = page.next_cursor;
let mut seen: Vec<EventId> = page.events.iter().map(|e| e.event_id).collect();
let mut sequences: Vec<u64> = page.events.iter().map(|e| e.sequence).collect();
for _ in 0..10 {
let next = ensure_ok(
CHECK,
"reading the next page",
events.read_from(&account, cursor, 2).await,
)?;
if next.is_empty() {
ensure(
CHECK,
next.next_cursor == cursor,
"an empty page must leave the cursor where it was",
)?;
break;
}
ensure(
CHECK,
next.next_cursor > cursor,
"a non-empty page must move the cursor forward",
)?;
cursor = next.next_cursor;
seen.extend(next.events.iter().map(|e| e.event_id));
sequences.extend(next.events.iter().map(|e| e.sequence));
}
let expected: Vec<EventId> = first
.iter()
.chain(second.iter())
.chain(third.iter())
.copied()
.collect();
ensure_eq(
CHECK,
"every event of the account, once each, in commit order",
&seen,
&expected,
)?;
ensure(
CHECK,
sequences.windows(2).all(|pair| pair[0] < pair[1]),
"the pages must not repeat or reorder a sequence",
)?;
ensure(
CHECK,
!seen.contains(&foreign[0]),
"another tenant's event must never appear in this account's stream",
)?;
let caught_up = ensure_ok(
CHECK,
"reading after the last event",
events.read_from(&account, cursor, 10).await,
)?;
ensure(
CHECK,
caught_up.is_empty() && caught_up.next_cursor == cursor,
"a consumer that has caught up gets an empty page and keeps its cursor",
)?;
let late = vec![EventId::new()];
ensure_ok(
CHECK,
"appending after the consumer caught up",
events
.append(event_batch(&account, CommandId::new(), 3, &late, at(4)))
.await,
)?;
let resumed = ensure_ok(
CHECK,
"reading again from the stored cursor",
events.read_from(&account, cursor, 10).await,
)?;
ensure_eq(
CHECK,
"only the event appended after the consumer caught up",
&resumed
.events
.iter()
.map(|e| e.event_id)
.collect::<Vec<_>>(),
&late,
)?;
let mut replayed: Vec<EventId> = Vec::new();
let mut from_start = EventCursor::START;
for _ in 0..10 {
let page = ensure_ok(
CHECK,
"replaying the whole stream one event at a time",
events.read_from(&account, from_start, 1).await,
)?;
if page.is_empty() {
break;
}
from_start = page.next_cursor;
replayed.extend(page.events.iter().map(|e| e.event_id));
}
let mut all = expected;
all.extend(late);
ensure_eq(
CHECK,
"a full replay sees the same events in the same order",
&replayed,
&all,
)?;
let foreign_view = ensure_ok(
CHECK,
"reading this account's stream as another tenant",
events
.read_from(&other_account(), EventCursor::START, 10)
.await,
)?;
ensure_eq(
CHECK,
"the other tenant sees only its own event",
&foreign_view
.events
.iter()
.map(|e| e.event_id)
.collect::<Vec<_>>(),
&foreign,
)
}
pub async fn check_revision_read_truncates_inside_a_revision(
stores: &Stores,
) -> Result<(), ConformanceFailure> {
const CHECK: &str = "check_revision_read_truncates_inside_a_revision";
let events = stores.events();
let account = account();
let ids: Vec<EventId> = (0..3).map(|_| EventId::new()).collect();
ensure_ok(
CHECK,
"appending three events at one revision",
events
.append(event_batch(&account, CommandId::new(), 1, &ids, epoch()))
.await,
)?;
let truncated = ensure_ok(
CHECK,
"listing the case history with a limit of two",
events
.list_since(&account, &case_key(), CaseRevision::ZERO, 2)
.await,
)?;
ensure_eq(
CHECK,
"the limit cuts the revision in half",
&truncated.iter().map(|e| e.event_id).collect::<Vec<_>>(),
&ids[..2].to_vec(),
)?;
let resumed_by_revision = ensure_ok(
CHECK,
"trying to continue by revision",
events
.list_since(&account, &case_key(), CaseRevision(1), 2)
.await,
)?;
ensure(
CHECK,
resumed_by_revision.is_empty(),
"continuing past revision 1 skips the rest of revision 1: that is why \
the revision is not a page boundary",
)?;
let by_sequence = ensure_ok(
CHECK,
"continuing by sequence instead",
events
.read_from(
&account,
EventCursor::after(truncated[1].sequence),
usize::MAX,
)
.await,
)?;
ensure_eq(
CHECK,
"the sequence cursor recovers the event the revision cursor lost",
&by_sequence
.events
.iter()
.map(|e| e.event_id)
.collect::<Vec<_>>(),
&ids[2..].to_vec(),
)
}
pub async fn check_event_redaction_preserves_identity_and_order(
stores: &Stores,
) -> Result<(), ConformanceFailure> {
const CHECK: &str = "check_event_redaction_preserves_identity_and_order";
let events = stores.events();
let account = account();
let first: Vec<EventId> = (0..2).map(|_| EventId::new()).collect();
ensure_ok(
CHECK,
"appending the first batch",
events
.append(event_batch_with_personal_data(
&account,
CommandId::new(),
1,
&first,
epoch(),
))
.await,
)?;
let second = vec![EventId::new()];
ensure_ok(
CHECK,
"appending the second batch",
events
.append(event_batch_with_personal_data(
&account,
CommandId::new(),
2,
&second,
at(1),
))
.await,
)?;
let foreign = vec![EventId::new()];
ensure_ok(
CHECK,
"appending an event of another tenant",
events
.append(event_batch_with_personal_data(
&other_account(),
CommandId::new(),
1,
&foreign,
at(2),
))
.await,
)?;
let before = ensure_ok(
CHECK,
"reading the whole stream before the erasure",
events.read_from(&account, EventCursor::START, 10).await,
)?;
ensure_eq(
CHECK,
"the stream before the erasure",
&before.events.len(),
&3,
)?;
let target = first[1];
let record = ensure_ok(
CHECK,
"redacting the payload of one event",
events
.redact_payload(&account, &target, &redaction_authority())
.await,
)?;
ensure_eq(
CHECK,
"the erasure is recorded under the authority it was asked for",
&record.authority,
&redaction_authority(),
)?;
let after = ensure_ok(
CHECK,
"reading the whole stream after the erasure",
events.read_from(&account, EventCursor::START, 10).await,
)?;
ensure_eq(
CHECK,
"an erasure must not add, remove or move an event",
&after.events.iter().map(identity_of).collect::<Vec<_>>(),
&before.events.iter().map(identity_of).collect::<Vec<_>>(),
)?;
ensure_eq(
CHECK,
"the cursor at the end of the stream must not move",
&after.next_cursor,
&before.next_cursor,
)?;
ensure_eq(
CHECK,
"the events of the case are still counted",
&ensure_ok(
CHECK,
"counting after the erasure",
events.count(&account, &case_key()).await,
)?,
&3,
)?;
for (was, is) in before.events.iter().zip(after.events.iter()) {
if is.event_id == target {
ensure(
CHECK,
is.is_redacted(),
"the redacted event must be marked as redacted",
)?;
ensure_eq(
CHECK,
"the payload of a redacted event",
&is.payload,
&serde_json::Value::Null,
)?;
ensure(
CHECK,
!serde_json::to_string(is)
.unwrap_or_default()
.contains(PERSONAL_DATA),
"the erased value must not survive anywhere on the stored event, \
including in the record of the erasure",
)?;
ensure(
CHECK,
is.to_receipt_event().is_redacted(),
"a redacted event must reach a receipt renderer as redacted",
)?;
} else {
ensure_eq(CHECK, "an event nobody asked to erase", is, was)?;
ensure(
CHECK,
!is.is_redacted() && !is.to_receipt_event().is_redacted(),
"erasing one payload must not mark another event as redacted",
)?;
}
}
let by_ids = ensure_ok(
CHECK,
"reading the redacted event by identifier",
events.get_by_ids(&account, &[target]).await,
)?;
ensure_eq(
CHECK,
"a redacted event is still readable by identifier",
&by_ids.iter().map(|e| e.event_id).collect::<Vec<_>>(),
&vec![target],
)?;
let history = ensure_ok(
CHECK,
"listing the case history after the erasure",
events
.list_since(&account, &case_key(), CaseRevision::ZERO, 10)
.await,
)?;
ensure_eq(
CHECK,
"the case history still carries every event, in order",
&history.iter().map(|e| e.event_id).collect::<Vec<_>>(),
&vec![first[0], first[1], second[0]],
)?;
let mut walked: Vec<EventId> = Vec::new();
let mut cursor = EventCursor::START;
for _ in 0..10 {
let page = ensure_ok(
CHECK,
"paging the stream after the erasure",
events.read_from(&account, cursor, 1).await,
)?;
if page.is_empty() {
break;
}
cursor = page.next_cursor;
walked.extend(page.events.iter().map(|e| e.event_id));
}
ensure_eq(
CHECK,
"every position of the journal is still delivered exactly once",
&walked,
&vec![first[0], first[1], second[0]],
)?;
let elsewhere = ensure_ok(
CHECK,
"reading the other tenant's stream",
events
.read_from(&other_account(), EventCursor::START, 10)
.await,
)?;
ensure(
CHECK,
elsewhere.events.len() == 1 && !elsewhere.events[0].is_redacted(),
"an erasure in one tenant must not touch another tenant's events",
)
}
pub async fn check_event_redaction_is_audited_and_idempotent(
stores: &Stores,
) -> Result<(), ConformanceFailure> {
const CHECK: &str = "check_event_redaction_is_audited_and_idempotent";
let events = stores.events();
let account = account();
let ids: Vec<EventId> = (0..2).map(|_| EventId::new()).collect();
ensure_ok(
CHECK,
"appending two events",
events
.append(event_batch_with_personal_data(
&account,
CommandId::new(),
1,
&ids,
epoch(),
))
.await,
)?;
let record = ensure_ok(
CHECK,
"redacting a payload",
events
.redact_payload(&account, &ids[0], &redaction_authority())
.await,
)?;
let stored = ensure_ok(
CHECK,
"reading the redacted event back",
events.get_by_ids(&account, &[ids[0]]).await,
)?;
let stored = stored
.first()
.ok_or_else(|| ConformanceFailure::new(CHECK, "the redacted event is gone"))?;
ensure_eq(
CHECK,
"the erasure the store reported is the erasure it persisted",
&stored.redaction,
&Some(record.clone()),
)?;
let repeated = ensure_ok(
CHECK,
"redacting the same payload again under another authority",
events
.redact_payload(&account, &ids[0], &other_redaction_authority())
.await,
)?;
ensure_eq(
CHECK,
"a repeated erasure keeps the record of the first one",
&repeated,
&record,
)?;
let after_repeat = ensure_ok(
CHECK,
"reading the event after the repeated erasure",
events.get_by_ids(&account, &[ids[0]]).await,
)?;
ensure_eq(
CHECK,
"a repeated erasure rewrites nothing",
&after_repeat.first(),
&Some(stored),
)?;
ensure_eq(
CHECK,
"the events of the case after two erasure requests",
&ensure_ok(
CHECK,
"counting the events of the case",
events.count(&account, &case_key()).await,
)?,
&2,
)?;
ensure_error(
CHECK,
"redacting an event that does not exist",
events
.redact_payload(&account, &EventId::new(), &redaction_authority())
.await,
&StoreError::NotFound,
)?;
ensure_error(
CHECK,
"redacting this tenant's event as another tenant",
events
.redact_payload(&other_account(), &ids[1], &redaction_authority())
.await,
&StoreError::NotFound,
)?;
let untouched = ensure_ok(
CHECK,
"reading the event another tenant tried to redact",
events.get_by_ids(&account, &[ids[1]]).await,
)?;
ensure(
CHECK,
untouched
.first()
.is_some_and(|event| !event.is_redacted() && event.payload.get("full_name").is_some()),
"a refused erasure must leave the event exactly as it was",
)
}
fn identity_of(event: &StoredEvent) -> (u64, EventId, CaseRevision, &str, DateTime<Utc>) {
(
event.sequence,
event.event_id,
event.case_revision,
event.event_type.as_str(),
event.occurred_at,
)
}