use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use turnframe_core::case::CaseRef;
use turnframe_core::command::{CommandEnvelope, CommandOrigin, IdempotencyKey};
use turnframe_core::error::{DomainRejection, ExecutionError};
use turnframe_core::ids::{AccountId, AttemptId, CaseRevision, CommandId, EventId, TurnId};
use crate::error::StoreError;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CommandJournalStatus {
Pending,
AwaitingConfirmation,
Executing,
Committed,
Failed,
OutcomeUnknown,
}
impl CommandJournalStatus {
pub const ALL: [Self; 6] = [
Self::Pending,
Self::AwaitingConfirmation,
Self::Executing,
Self::Committed,
Self::Failed,
Self::OutcomeUnknown,
];
#[must_use]
pub fn is_terminal(self) -> bool {
matches!(self, Self::Committed | Self::Failed)
}
#[must_use]
pub fn is_pending(self) -> bool {
matches!(self, Self::Pending | Self::Executing)
}
#[must_use]
pub fn can_transition(from: Self, to: Self) -> bool {
matches!(
(from, to),
(Self::Pending | Self::AwaitingConfirmation, Self::Executing)
| (
Self::Pending | Self::AwaitingConfirmation | Self::Executing,
Self::Committed | Self::Failed | Self::OutcomeUnknown
)
| (Self::OutcomeUnknown, Self::Committed)
| (Self::OutcomeUnknown, Self::Failed)
)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
#[non_exhaustive]
pub enum JournalOutcome {
Committed {
new_revision: CaseRevision,
event_ids: Vec<EventId>,
},
Rejected {
rejection: DomainRejection,
},
RevisionConflict {
current_revision: CaseRevision,
},
Failed {
code: String,
},
OutcomeUnknown {
attempt_id: AttemptId,
#[serde(default, skip_serializing_if = "Option::is_none")]
remote_ref: Option<String>,
reason: String,
},
}
impl JournalOutcome {
#[must_use]
pub fn status(&self) -> CommandJournalStatus {
match self {
Self::Committed { .. } => CommandJournalStatus::Committed,
Self::Rejected { .. } | Self::RevisionConflict { .. } | Self::Failed { .. } => {
CommandJournalStatus::Failed
}
Self::OutcomeUnknown { .. } => CommandJournalStatus::OutcomeUnknown,
}
}
#[must_use]
pub fn from_execution_error(command_id: CommandId, error: &ExecutionError) -> Self {
match error {
ExecutionError::RevisionConflict(conflict) => Self::RevisionConflict {
current_revision: conflict.current_revision,
},
ExecutionError::Rejected(rejection) => Self::Rejected {
rejection: rejection.clone(),
},
ExecutionError::OutcomeUnknown(unknown) => Self::OutcomeUnknown {
attempt_id: unknown.attempt_id.clone(),
remote_ref: unknown.remote_ref.clone(),
reason: unknown.reason.clone(),
},
ExecutionError::Timeout => Self::OutcomeUnknown {
attempt_id: AttemptId::new(command_id.to_string()),
remote_ref: None,
reason: "execution_timeout".to_owned(),
},
ExecutionError::Store(store) => Self::Failed {
code: format!("store.{}", store_code(store)),
},
ExecutionError::IdempotencyMismatch { .. } => Self::Failed {
code: "idempotency_mismatch".to_owned(),
},
ExecutionError::ScopeViolation => Self::Failed {
code: "scope_violation".to_owned(),
},
ExecutionError::Erasure(_) => Self::Failed {
code: "erasure".to_owned(),
},
ExecutionError::Other { code } => Self::Failed { code: code.clone() },
_ => Self::Failed {
code: "other".to_owned(),
},
}
}
}
fn store_code(error: &StoreError) -> &str {
match error {
StoreError::NotFound => "not_found",
StoreError::Conflict => "conflict",
StoreError::Unavailable => "unavailable",
StoreError::Timeout => "timeout",
StoreError::Serialization => "serialization",
StoreError::Corrupt => "corrupt",
StoreError::Other { code } => code,
_ => "other",
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CommandJournalEntry {
pub command_id: CommandId,
pub account_id: AccountId,
pub idempotency_key: IdempotencyKey,
pub turn_id: TurnId,
pub case_ref: CaseRef,
pub command_type: String,
pub command_payload: serde_json::Value,
pub origin: CommandOrigin,
pub status: CommandJournalStatus,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub result: Option<JournalOutcome>,
pub created_at: DateTime<Utc>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub completed_at: Option<DateTime<Utc>>,
}
impl CommandJournalEntry {
pub fn from_envelope<C: Serialize>(
envelope: &CommandEnvelope<C>,
command_type: impl Into<String>,
created_at: DateTime<Utc>,
) -> Result<Self, StoreError> {
let command_payload =
serde_json::to_value(&envelope.command).map_err(|_| StoreError::Serialization)?;
Ok(Self {
command_id: envelope.command_id,
account_id: envelope.actor.account_id.clone(),
idempotency_key: envelope.idempotency_key.clone(),
turn_id: envelope.turn_id,
case_ref: envelope.case_ref.clone(),
command_type: command_type.into(),
command_payload,
origin: envelope.origin.clone(),
status: CommandJournalStatus::Pending,
result: None,
created_at,
completed_at: None,
})
}
#[must_use]
pub fn same_command(&self, other: &Self) -> bool {
self.case_ref == other.case_ref
&& self.command_type == other.command_type
&& self.command_payload == other.command_payload
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum JournalAdmission {
Fresh,
Replay(Box<CommandJournalEntry>),
}
impl JournalAdmission {
#[must_use]
pub fn replay(entry: CommandJournalEntry) -> Self {
Self::Replay(Box::new(entry))
}
#[must_use]
pub fn is_fresh(&self) -> bool {
matches!(self, Self::Fresh)
}
#[must_use]
pub fn replayed(&self) -> Option<&CommandJournalEntry> {
match self {
Self::Fresh => None,
Self::Replay(entry) => Some(entry),
}
}
#[must_use]
pub fn into_replayed(self) -> Option<CommandJournalEntry> {
match self {
Self::Fresh => None,
Self::Replay(entry) => Some(*entry),
}
}
}
#[async_trait]
pub trait CommandJournalReader: Send + Sync {
async fn get(
&self,
account: &AccountId,
command_id: &CommandId,
) -> Result<CommandJournalEntry, StoreError>;
async fn for_turn(
&self,
account: &AccountId,
turn_id: &TurnId,
) -> Result<Vec<CommandJournalEntry>, StoreError>;
async fn pending_for_turn(
&self,
account: &AccountId,
turn_id: &TurnId,
) -> Result<Vec<CommandJournalEntry>, StoreError>;
}
#[async_trait]
pub trait CommandJournalWriter: Send + Sync {
async fn begin(&self, entry: CommandJournalEntry) -> Result<JournalAdmission, StoreError>;
async fn mark_executing(
&self,
account: &AccountId,
command_id: &CommandId,
) -> Result<(), StoreError>;
async fn complete(
&self,
account: &AccountId,
command_id: &CommandId,
outcome: JournalOutcome,
) -> Result<(), StoreError>;
async fn fail(
&self,
account: &AccountId,
command_id: &CommandId,
error: &ExecutionError,
) -> Result<(), StoreError> {
self.complete(
account,
command_id,
JournalOutcome::from_execution_error(*command_id, error),
)
.await
}
}
pub trait CommandJournal: CommandJournalReader + CommandJournalWriter {}
impl<T: CommandJournalReader + CommandJournalWriter + ?Sized> CommandJournal for T {}
#[cfg(test)]
mod tests {
use super::*;
use turnframe_core::error::RevisionConflict;
use turnframe_core::event::UnknownOutcome;
#[test]
fn transition_table() {
use CommandJournalStatus as S;
let allowed: &[(S, S)] = &[
(S::Pending, S::Executing),
(S::Pending, S::Committed),
(S::Pending, S::Failed),
(S::Pending, S::OutcomeUnknown),
(S::AwaitingConfirmation, S::Executing),
(S::AwaitingConfirmation, S::Committed),
(S::AwaitingConfirmation, S::Failed),
(S::AwaitingConfirmation, S::OutcomeUnknown),
(S::Executing, S::Committed),
(S::Executing, S::Failed),
(S::Executing, S::OutcomeUnknown),
(S::OutcomeUnknown, S::Committed),
(S::OutcomeUnknown, S::Failed),
];
for from in S::ALL {
for to in S::ALL {
assert_eq!(
S::can_transition(from, to),
allowed.contains(&(from, to)),
"{from:?} -> {to:?}"
);
}
}
assert!(S::Committed.is_terminal());
assert!(S::Failed.is_terminal());
assert!(!S::OutcomeUnknown.is_terminal());
assert!(S::Pending.is_pending());
assert!(S::Executing.is_pending());
assert!(!S::AwaitingConfirmation.is_pending());
assert!(!S::AwaitingConfirmation.is_terminal());
}
#[test]
fn outcome_status_mapping_and_error_conversion() {
let id = CommandId::nil();
let conflict = ExecutionError::RevisionConflict(RevisionConflict {
expected: CaseRef::new("w", "c", CaseRevision(1)),
current_revision: CaseRevision(2),
});
let out = JournalOutcome::from_execution_error(id, &conflict);
assert_eq!(out.status(), CommandJournalStatus::Failed);
assert_eq!(
out,
JournalOutcome::RevisionConflict {
current_revision: CaseRevision(2)
}
);
let unknown = ExecutionError::OutcomeUnknown(UnknownOutcome {
attempt_id: "a".into(),
remote_ref: Some("r".into()),
reason: "timeout_after_send".into(),
});
assert_eq!(
JournalOutcome::from_execution_error(id, &unknown).status(),
CommandJournalStatus::OutcomeUnknown
);
assert_eq!(
JournalOutcome::from_execution_error(id, &ExecutionError::Timeout).status(),
CommandJournalStatus::OutcomeUnknown
);
assert_eq!(
JournalOutcome::from_execution_error(id, &ExecutionError::Store(StoreError::Timeout)),
JournalOutcome::Failed {
code: "store.timeout".into()
}
);
assert_eq!(
JournalOutcome::from_execution_error(id, &ExecutionError::ScopeViolation).status(),
CommandJournalStatus::Failed
);
}
#[test]
fn admission_helpers() {
assert!(JournalAdmission::Fresh.is_fresh());
assert!(JournalAdmission::Fresh.replayed().is_none());
assert!(JournalAdmission::Fresh.into_replayed().is_none());
}
}