use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use crate::command::IdempotencyKey;
use crate::hash::derive_uuid;
use crate::ids::{
AttemptId, CaseRevision, CommandId, EventId, OutboxId, ReceiptId, RedactionAuthority,
};
use crate::locale::LocalizedText;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Commit<S, E> {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub state: Option<S>,
pub new_revision: CaseRevision,
pub events: Vec<CommittedEvent<E>>,
pub idempotency_replay: bool,
}
impl<S, E> Commit<S, E> {
#[must_use]
pub fn event_ids(&self) -> Vec<EventId> {
self.events.iter().map(|e| e.event_id).collect()
}
#[must_use]
pub fn event_refs(&self) -> Vec<EventRef> {
self.events.iter().map(CommittedEvent::event_ref).collect()
}
}
impl<S, E: Clone> Commit<S, E> {
#[must_use]
pub fn receipt_events(&self) -> Vec<ReceiptEvent<E>> {
self.events
.iter()
.cloned()
.map(ReceiptEvent::Committed)
.collect()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CommittedEvent<E> {
pub event_id: EventId,
pub event_type: String,
pub occurred_at: DateTime<Utc>,
pub payload: E,
}
impl<E> CommittedEvent<E> {
#[must_use]
pub fn event_ref(&self) -> EventRef {
EventRef {
event_id: self.event_id,
event_type: self.event_type.clone(),
}
}
pub fn try_map_payload<F, Err>(
self,
f: impl FnOnce(E) -> Result<F, Err>,
) -> Result<CommittedEvent<F>, Err> {
Ok(CommittedEvent {
event_id: self.event_id,
event_type: self.event_type,
occurred_at: self.occurred_at,
payload: f(self.payload)?,
})
}
pub fn try_map_payload_ref<F, Err>(
&self,
f: impl FnOnce(&E) -> Result<F, Err>,
) -> Result<CommittedEvent<F>, Err> {
Ok(CommittedEvent {
event_id: self.event_id,
event_type: self.event_type.clone(),
occurred_at: self.occurred_at,
payload: f(&self.payload)?,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct EventRef {
pub event_id: EventId,
pub event_type: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct EventRedaction {
pub redacted_at: DateTime<Utc>,
pub authority: RedactionAuthority,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RedactedEvent {
pub event_id: EventId,
pub event_type: String,
pub occurred_at: DateTime<Utc>,
pub redaction: EventRedaction,
}
impl RedactedEvent {
#[must_use]
pub fn event_ref(&self) -> EventRef {
EventRef {
event_id: self.event_id,
event_type: self.event_type.clone(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ReceiptEvent<E> {
Committed(CommittedEvent<E>),
Redacted(RedactedEvent),
}
impl<E> ReceiptEvent<E> {
#[must_use]
pub fn event_id(&self) -> EventId {
match self {
Self::Committed(event) => event.event_id,
Self::Redacted(event) => event.event_id,
}
}
#[must_use]
pub fn event_type(&self) -> &str {
match self {
Self::Committed(event) => &event.event_type,
Self::Redacted(event) => &event.event_type,
}
}
#[must_use]
pub fn occurred_at(&self) -> DateTime<Utc> {
match self {
Self::Committed(event) => event.occurred_at,
Self::Redacted(event) => event.occurred_at,
}
}
#[must_use]
pub fn event_ref(&self) -> EventRef {
match self {
Self::Committed(event) => event.event_ref(),
Self::Redacted(event) => event.event_ref(),
}
}
#[must_use]
pub fn payload(&self) -> Option<&E> {
match self {
Self::Committed(event) => Some(&event.payload),
Self::Redacted(_) => None,
}
}
#[must_use]
pub fn redaction(&self) -> Option<&EventRedaction> {
match self {
Self::Committed(_) => None,
Self::Redacted(event) => Some(&event.redaction),
}
}
#[must_use]
pub fn is_redacted(&self) -> bool {
matches!(self, Self::Redacted(_))
}
pub fn try_map_payload_ref<F, Err>(
&self,
f: impl FnOnce(&E) -> Result<F, Err>,
) -> Result<ReceiptEvent<F>, Err> {
match self {
Self::Committed(event) => event.try_map_payload_ref(f).map(ReceiptEvent::Committed),
Self::Redacted(event) => Ok(ReceiptEvent::Redacted(event.clone())),
}
}
}
impl<E> From<CommittedEvent<E>> for ReceiptEvent<E> {
fn from(event: CommittedEvent<E>) -> Self {
Self::Committed(event)
}
}
impl<E> From<RedactedEvent> for ReceiptEvent<E> {
fn from(event: RedactedEvent) -> Self {
Self::Redacted(event)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ReceiptSeverity {
Info,
Success,
Warning,
Error,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ArtifactRef {
pub artifact_id: String,
pub kind: String,
pub label: LocalizedText,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub uri: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub media_type: Option<String>,
}
const RECEIPT_ID_DOMAIN: &str = "turnframe.receipt_id.v1";
impl ReceiptId {
#[must_use]
pub fn derive(event_ids: &[EventId], status_code: &str) -> Self {
let mut ids: Vec<String> = event_ids.iter().map(EventId::to_string).collect();
ids.sort_unstable();
ids.dedup();
let mut parts: Vec<&str> = vec![status_code];
parts.extend(ids.iter().map(String::as_str));
Self(derive_uuid(RECEIPT_ID_DOMAIN, &parts))
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct OperationalReceipt {
pub receipt_id: ReceiptId,
pub event_ids: Vec<EventId>,
pub severity: ReceiptSeverity,
pub title: LocalizedText,
pub body: LocalizedText,
pub status_code: String,
#[serde(default)]
pub artifact_refs: Vec<ArtifactRef>,
}
impl OperationalReceipt {
#[must_use]
pub fn is_event_backed(&self) -> bool {
!self.event_ids.is_empty()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
#[non_exhaustive]
pub enum ExternalStatus {
Prepared,
Validated,
AwaitingConfirmation,
Submitted,
ReceivedByIntermediary,
ReceivedByAuthority,
Accepted,
Rejected,
Issued,
NotDelivered,
Delivered,
Completed,
}
impl ExternalStatus {
#[must_use]
pub fn is_final(self) -> bool {
matches!(self, Self::Rejected | Self::NotDelivered | Self::Completed)
}
#[must_use]
pub fn is_transmitted(self) -> bool {
self >= Self::Submitted
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, thiserror::Error)]
#[error("external outcome unknown for attempt {attempt_id} ({reason})")]
pub struct UnknownOutcome {
pub attempt_id: AttemptId,
pub remote_ref: Option<String>,
pub reason: String,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
#[non_exhaustive]
pub enum OutboxStatus {
Pending,
Dispatching,
OutcomeUnknown,
Completed,
Failed,
}
impl OutboxStatus {
#[must_use]
pub fn is_terminal(self) -> bool {
matches!(self, Self::Completed | Self::Failed)
}
#[must_use]
pub fn can_transition(from: Self, to: Self) -> bool {
matches!(
(from, to),
(Self::Pending, Self::Dispatching)
| (Self::Dispatching, Self::Pending)
| (Self::Dispatching, Self::OutcomeUnknown)
| (Self::Dispatching, Self::Completed)
| (Self::Dispatching, Self::Failed)
| (Self::OutcomeUnknown, Self::Completed)
| (Self::OutcomeUnknown, Self::Failed)
| (Self::OutcomeUnknown, Self::Pending)
)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct OutboxEntry {
pub outbox_id: OutboxId,
pub command_id: CommandId,
pub destination: String,
pub payload: serde_json::Value,
pub idempotency_key: IdempotencyKey,
pub status: OutboxStatus,
pub attempt_count: u32,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub next_attempt_at: Option<DateTime<Utc>>,
pub created_at: DateTime<Utc>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub completed_at: Option<DateTime<Utc>>,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn external_status_finality() {
assert!(ExternalStatus::Completed.is_final());
assert!(ExternalStatus::Rejected.is_final());
assert!(ExternalStatus::NotDelivered.is_final());
assert!(!ExternalStatus::Accepted.is_final());
assert!(ExternalStatus::Submitted.is_transmitted());
assert!(!ExternalStatus::Validated.is_transmitted());
}
#[test]
fn receipt_ids_are_derived_from_events_and_status() {
let a = EventId::nil();
let b = EventId::from(uuid::Uuid::from_u128(7));
let id = ReceiptId::derive(&[a, b], "trip.rebooking_sent");
assert_eq!(id, ReceiptId::derive(&[a, b], "trip.rebooking_sent"));
assert_eq!(
id,
ReceiptId::derive(&[b, a], "trip.rebooking_sent"),
"order free"
);
assert_ne!(id, ReceiptId::derive(&[a], "trip.rebooking_sent"));
assert_ne!(id, ReceiptId::derive(&[a, b], "trip.refused"));
}
#[test]
fn a_redacted_event_keeps_everything_a_claim_rests_on() {
let committed = CommittedEvent {
event_id: EventId::from(uuid::Uuid::from_u128(11)),
event_type: "trip.extra_added".to_owned(),
occurred_at: DateTime::<Utc>::UNIX_EPOCH,
payload: serde_json::json!({ "description": "a name that must go" }),
};
let redacted = RedactedEvent {
event_id: committed.event_id,
event_type: committed.event_type.clone(),
occurred_at: committed.occurred_at,
redaction: EventRedaction {
redacted_at: DateTime::<Utc>::UNIX_EPOCH,
authority: RedactionAuthority::from("erasure-request-8842"),
},
};
let present = ReceiptEvent::from(committed.clone());
let gone = ReceiptEvent::<serde_json::Value>::from(redacted);
assert_eq!(present.event_id(), gone.event_id());
assert_eq!(present.event_type(), gone.event_type());
assert_eq!(present.occurred_at(), gone.occurred_at());
assert_eq!(present.event_ref(), gone.event_ref());
assert!(!present.is_redacted());
assert!(gone.is_redacted());
assert!(present.payload().is_some());
assert!(gone.payload().is_none());
assert!(present.redaction().is_none());
assert_eq!(
gone.redaction().map(|r| r.authority.as_str()),
Some("erasure-request-8842")
);
let mapped = gone.try_map_payload_ref(|_: &serde_json::Value| Err::<(), &str>("never run"));
assert!(mapped.is_ok(), "a redacted event maps without calling f");
assert!(mapped.unwrap_or_else(|_| unreachable!()).is_redacted());
assert!(
present
.try_map_payload_ref(|_: &serde_json::Value| Err::<(), &str>("boom"))
.is_err(),
"a present payload still goes through f"
);
}
#[test]
fn a_commit_offers_its_events_for_receipt_rendering() {
let commit: Commit<(), u32> = Commit {
state: None,
new_revision: CaseRevision(1),
events: vec![CommittedEvent {
event_id: EventId::nil(),
event_type: "t".to_owned(),
occurred_at: DateTime::<Utc>::UNIX_EPOCH,
payload: 7,
}],
idempotency_replay: false,
};
let events = commit.receipt_events();
assert_eq!(events.len(), 1);
assert_eq!(events[0].payload(), Some(&7));
assert!(
!events[0].is_redacted(),
"a fresh commit never carries an erased payload"
);
}
#[test]
fn outbox_transitions() {
assert!(OutboxStatus::can_transition(
OutboxStatus::Pending,
OutboxStatus::Dispatching
));
assert!(OutboxStatus::can_transition(
OutboxStatus::Dispatching,
OutboxStatus::OutcomeUnknown
));
assert!(!OutboxStatus::can_transition(
OutboxStatus::Completed,
OutboxStatus::Pending
));
assert!(!OutboxStatus::can_transition(
OutboxStatus::Pending,
OutboxStatus::Completed
));
}
}