use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use turnframe_core::case::CaseKey;
use turnframe_core::event::{Commit, CommittedEvent, EventRedaction, ReceiptEvent, RedactedEvent};
use turnframe_core::ids::{AccountId, CaseRevision, CommandId, EventId, RedactionAuthority};
use crate::error::StoreError;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct StoredEvent {
pub sequence: u64,
pub event_id: EventId,
pub account_id: AccountId,
pub case_key: CaseKey,
pub case_revision: CaseRevision,
pub command_id: CommandId,
pub event_type: String,
pub payload: serde_json::Value,
pub occurred_at: DateTime<Utc>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub redaction: Option<EventRedaction>,
}
impl StoredEvent {
#[must_use]
pub fn to_committed(&self) -> CommittedEvent<serde_json::Value> {
CommittedEvent {
event_id: self.event_id,
event_type: self.event_type.clone(),
occurred_at: self.occurred_at,
payload: self.payload.clone(),
}
}
#[must_use]
pub fn is_redacted(&self) -> bool {
self.redaction.is_some()
}
#[must_use]
pub fn to_receipt_event(&self) -> ReceiptEvent<serde_json::Value> {
match &self.redaction {
None => ReceiptEvent::Committed(self.to_committed()),
Some(redaction) => ReceiptEvent::Redacted(RedactedEvent {
event_id: self.event_id,
event_type: self.event_type.clone(),
occurred_at: self.occurred_at,
redaction: redaction.clone(),
}),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct EventBatch {
pub account_id: AccountId,
pub case_key: CaseKey,
pub command_id: CommandId,
pub revision: CaseRevision,
pub events: Vec<CommittedEvent<serde_json::Value>>,
}
impl EventBatch {
#[must_use]
pub fn new(
account_id: AccountId,
case_key: CaseKey,
command_id: CommandId,
revision: CaseRevision,
events: Vec<CommittedEvent<serde_json::Value>>,
) -> Self {
Self {
account_id,
case_key,
command_id,
revision,
events,
}
}
pub fn from_commit<S, E: Serialize>(
account_id: AccountId,
case_key: CaseKey,
command_id: CommandId,
commit: &Commit<S, E>,
) -> Result<Self, StoreError> {
let mut events = Vec::with_capacity(commit.events.len());
for event in &commit.events {
let payload =
serde_json::to_value(&event.payload).map_err(|_| StoreError::Serialization)?;
events.push(CommittedEvent {
event_id: event.event_id,
event_type: event.event_type.clone(),
occurred_at: event.occurred_at,
payload,
});
}
Ok(Self::new(
account_id,
case_key,
command_id,
commit.new_revision,
events,
))
}
#[must_use]
pub fn event_ids(&self) -> Vec<EventId> {
self.events.iter().map(|e| e.event_id).collect()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.events.is_empty()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LedgerReceiptGroup {
pub case_key: CaseKey,
pub command_id: CommandId,
pub revision: CaseRevision,
pub events: Vec<ReceiptEvent<serde_json::Value>>,
}
#[must_use]
pub fn group_for_receipts(events: &[StoredEvent]) -> Vec<LedgerReceiptGroup> {
let mut groups: Vec<LedgerReceiptGroup> = Vec::new();
for event in events {
let existing = groups.iter().position(|group| {
group.case_key == event.case_key && group.command_id == event.command_id
});
match existing.and_then(|index| groups.get_mut(index)) {
Some(group) => group.events.push(event.to_receipt_event()),
None => groups.push(LedgerReceiptGroup {
case_key: event.case_key.clone(),
command_id: event.command_id,
revision: event.case_revision,
events: vec![event.to_receipt_event()],
}),
}
}
groups
}
#[derive(
Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Default, Serialize, Deserialize,
)]
#[serde(transparent)]
pub struct EventCursor(pub u64);
impl EventCursor {
pub const START: Self = Self(0);
#[must_use]
pub const fn after(sequence: u64) -> Self {
Self(sequence)
}
#[must_use]
pub const fn value(self) -> u64 {
self.0
}
}
impl std::fmt::Display for EventCursor {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.0)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct EventPage {
pub events: Vec<StoredEvent>,
pub next_cursor: EventCursor,
}
impl EventPage {
#[must_use]
pub fn new(events: Vec<StoredEvent>, asked: EventCursor) -> Self {
let next_cursor = events
.last()
.map_or(asked, |event| EventCursor::after(event.sequence));
Self {
events,
next_cursor,
}
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.events.is_empty()
}
}
#[async_trait]
pub trait EventJournalReader: Send + Sync {
async fn list_since(
&self,
account: &AccountId,
case_key: &CaseKey,
since: CaseRevision,
limit: usize,
) -> Result<Vec<StoredEvent>, StoreError>;
async fn read_from(
&self,
account: &AccountId,
after: EventCursor,
limit: usize,
) -> Result<EventPage, StoreError>;
async fn get_by_ids(
&self,
account: &AccountId,
ids: &[EventId],
) -> Result<Vec<StoredEvent>, StoreError>;
async fn count(&self, account: &AccountId, case_key: &CaseKey) -> Result<u64, StoreError>;
}
#[async_trait]
pub trait EventJournalWriter: Send + Sync {
async fn append(&self, batch: EventBatch) -> Result<Vec<EventId>, StoreError>;
async fn redact_payload(
&self,
account: &AccountId,
event_id: &EventId,
authority: &RedactionAuthority,
) -> Result<EventRedaction, StoreError>;
}
pub trait EventJournal: EventJournalReader + EventJournalWriter {}
impl<T: EventJournalReader + EventJournalWriter + ?Sized> EventJournal for T {}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn batch_from_commit_erases_payloads() {
#[derive(Serialize)]
struct Ev {
n: u32,
}
let commit: Commit<(), Ev> = Commit {
state: None,
new_revision: CaseRevision(3),
events: vec![CommittedEvent {
event_id: EventId::nil(),
event_type: "t".into(),
occurred_at: DateTime::<Utc>::UNIX_EPOCH,
payload: Ev { n: 7 },
}],
idempotency_replay: false,
};
let batch = EventBatch::from_commit(
AccountId::from("a"),
CaseKey::new("w", "c"),
CommandId::nil(),
&commit,
)
.unwrap();
assert_eq!(batch.revision, CaseRevision(3));
assert_eq!(batch.events[0].payload, serde_json::json!({"n": 7}));
assert_eq!(batch.event_ids(), vec![EventId::nil()]);
assert!(!batch.is_empty());
}
#[test]
fn cursor_is_exclusive_and_pages_carry_where_to_resume() {
let event = |sequence: u64| StoredEvent {
sequence,
event_id: EventId::new(),
account_id: AccountId::from("a"),
case_key: CaseKey::new("w", "c"),
case_revision: CaseRevision(1),
command_id: CommandId::nil(),
event_type: "t".into(),
payload: serde_json::Value::Null,
occurred_at: DateTime::<Utc>::UNIX_EPOCH,
redaction: None,
};
assert_eq!(EventCursor::START, EventCursor(0));
assert_eq!(EventCursor::default(), EventCursor::START);
assert_eq!(EventCursor::after(7).value(), 7);
assert_eq!(EventCursor::after(7).to_string(), "7");
let page = EventPage::new(vec![event(4), event(5)], EventCursor::after(3));
assert!(!page.is_empty());
assert_eq!(
page.next_cursor,
EventCursor::after(5),
"resume after the last event of the page"
);
let caught_up = EventPage::new(Vec::new(), page.next_cursor);
assert!(caught_up.is_empty());
assert_eq!(
caught_up.next_cursor, page.next_cursor,
"an empty page must not move the cursor"
);
let json = serde_json::to_string(&page).unwrap();
assert!(json.contains("\"next_cursor\":5"), "{json}");
assert_eq!(serde_json::from_str::<EventPage>(&json).unwrap(), page);
}
fn stored(sequence: u64, case: &str, command: CommandId, redacted: bool) -> StoredEvent {
StoredEvent {
sequence,
event_id: EventId::new(),
account_id: AccountId::from("a"),
case_key: CaseKey::new("w", case),
case_revision: CaseRevision(1),
command_id: command,
event_type: "t".into(),
payload: if redacted {
serde_json::Value::Null
} else {
serde_json::json!({ "full_name": "Marta Bianchi" })
},
occurred_at: DateTime::<Utc>::UNIX_EPOCH,
redaction: redacted.then(|| EventRedaction {
redacted_at: DateTime::<Utc>::UNIX_EPOCH,
authority: RedactionAuthority::from("erasure-request-1"),
}),
}
}
#[test]
fn a_redacted_event_reaches_a_receipt_renderer_as_redacted() {
let intact = stored(1, "c", CommandId::nil(), false);
assert!(!intact.is_redacted());
assert!(!intact.to_receipt_event().is_redacted());
assert_eq!(
intact.to_receipt_event().payload(),
Some(&serde_json::json!({ "full_name": "Marta Bianchi" }))
);
let erased = stored(2, "c", CommandId::nil(), true);
assert!(erased.is_redacted());
let event = erased.to_receipt_event();
assert!(event.is_redacted(), "the payload is gone and it says so");
assert_eq!(event.event_id(), erased.event_id, "identity survives");
assert_eq!(event.event_type(), "t", "the type survives");
assert_eq!(
event.occurred_at(),
erased.occurred_at,
"the instant survives"
);
assert_eq!(
event.redaction().map(|r| r.authority.as_str()),
Some("erasure-request-1")
);
let json = serde_json::to_string(&erased).expect("a stored event serializes");
assert_eq!(
serde_json::from_str::<StoredEvent>(&json).expect("and deserializes"),
erased
);
let older = serde_json::json!({
"sequence": 1,
"event_id": EventId::nil(),
"account_id": "a",
"case_key": { "workflow": "w", "case_id": "c" },
"case_revision": 1,
"command_id": CommandId::nil(),
"event_type": "t",
"payload": {},
"occurred_at": "1970-01-01T00:00:00Z",
});
let older: StoredEvent = serde_json::from_value(older).expect("a record without the field");
assert!(!older.is_redacted(), "no record means nothing was erased");
}
#[test]
fn grouping_a_ledger_read_keeps_order_and_carries_erasures() {
let (first, second) = (CommandId::new(), CommandId::new());
let events = vec![
stored(1, "c1", first, false),
stored(2, "c1", first, true),
stored(3, "c2", second, false),
stored(4, "c1", first, false),
];
let groups = group_for_receipts(&events);
assert_eq!(groups.len(), 2, "one group per command");
assert_eq!(groups[0].case_key, CaseKey::new("w", "c1"));
assert_eq!(groups[0].command_id, first);
assert_eq!(groups[0].events.len(), 3);
assert_eq!(groups[1].case_key, CaseKey::new("w", "c2"));
assert_eq!(groups[1].events.len(), 1);
assert_eq!(
groups[0]
.events
.iter()
.map(ReceiptEvent::is_redacted)
.collect::<Vec<_>>(),
vec![false, true, false],
"the erasure of the middle event survives the grouping"
);
assert_eq!(
groups[0]
.events
.iter()
.map(ReceiptEvent::event_id)
.collect::<Vec<_>>(),
vec![events[0].event_id, events[1].event_id, events[3].event_id],
"append order inside a group"
);
assert!(group_for_receipts(&[]).is_empty());
}
#[test]
fn stored_event_to_committed() {
let stored = StoredEvent {
sequence: 1,
event_id: EventId::nil(),
account_id: AccountId::from("a"),
case_key: CaseKey::new("w", "c"),
case_revision: CaseRevision(1),
command_id: CommandId::nil(),
event_type: "t".into(),
payload: serde_json::Value::Null,
occurred_at: DateTime::<Utc>::UNIX_EPOCH,
redaction: None,
};
let committed = stored.to_committed();
assert_eq!(committed.event_id, EventId::nil());
assert_eq!(committed.event_type, "t");
}
}