1use async_trait::async_trait;
27use chrono::{DateTime, Utc};
28use serde::{Deserialize, Serialize};
29use turnframe_core::case::CaseRef;
30use turnframe_core::command::{CommandEnvelope, CommandOrigin, IdempotencyKey};
31use turnframe_core::error::{DomainRejection, ExecutionError};
32use turnframe_core::ids::{AccountId, AttemptId, CaseRevision, CommandId, EventId, TurnId};
33
34use crate::error::StoreError;
35
36#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
38#[serde(rename_all = "snake_case")]
39pub enum CommandJournalStatus {
40 Pending,
42 AwaitingConfirmation,
45 Executing,
47 Committed,
49 Failed,
51 OutcomeUnknown,
54}
55
56impl CommandJournalStatus {
57 pub const ALL: [Self; 6] = [
59 Self::Pending,
60 Self::AwaitingConfirmation,
61 Self::Executing,
62 Self::Committed,
63 Self::Failed,
64 Self::OutcomeUnknown,
65 ];
66
67 #[must_use]
69 pub fn is_terminal(self) -> bool {
70 matches!(self, Self::Committed | Self::Failed)
71 }
72
73 #[must_use]
76 pub fn is_pending(self) -> bool {
77 matches!(self, Self::Pending | Self::Executing)
78 }
79
80 #[must_use]
82 pub fn can_transition(from: Self, to: Self) -> bool {
83 matches!(
84 (from, to),
85 (Self::Pending | Self::AwaitingConfirmation, Self::Executing)
86 | (
87 Self::Pending | Self::AwaitingConfirmation | Self::Executing,
88 Self::Committed | Self::Failed | Self::OutcomeUnknown
89 )
90 | (Self::OutcomeUnknown, Self::Committed)
91 | (Self::OutcomeUnknown, Self::Failed)
92 )
93 }
94}
95
96#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
98#[serde(tag = "kind", rename_all = "snake_case")]
99#[non_exhaustive]
100pub enum JournalOutcome {
101 Committed {
103 new_revision: CaseRevision,
105 event_ids: Vec<EventId>,
107 },
108 Rejected {
110 rejection: DomainRejection,
112 },
113 RevisionConflict {
115 current_revision: CaseRevision,
117 },
118 Failed {
120 code: String,
122 },
123 OutcomeUnknown {
125 attempt_id: AttemptId,
127 #[serde(default, skip_serializing_if = "Option::is_none")]
129 remote_ref: Option<String>,
130 reason: String,
132 },
133}
134
135impl JournalOutcome {
136 #[must_use]
138 pub fn status(&self) -> CommandJournalStatus {
139 match self {
140 Self::Committed { .. } => CommandJournalStatus::Committed,
141 Self::Rejected { .. } | Self::RevisionConflict { .. } | Self::Failed { .. } => {
142 CommandJournalStatus::Failed
143 }
144 Self::OutcomeUnknown { .. } => CommandJournalStatus::OutcomeUnknown,
145 }
146 }
147
148 #[must_use]
154 pub fn from_execution_error(command_id: CommandId, error: &ExecutionError) -> Self {
155 match error {
156 ExecutionError::RevisionConflict(conflict) => Self::RevisionConflict {
157 current_revision: conflict.current_revision,
158 },
159 ExecutionError::Rejected(rejection) => Self::Rejected {
160 rejection: rejection.clone(),
161 },
162 ExecutionError::OutcomeUnknown(unknown) => Self::OutcomeUnknown {
163 attempt_id: unknown.attempt_id.clone(),
164 remote_ref: unknown.remote_ref.clone(),
165 reason: unknown.reason.clone(),
166 },
167 ExecutionError::Timeout => Self::OutcomeUnknown {
168 attempt_id: AttemptId::new(command_id.to_string()),
169 remote_ref: None,
170 reason: "execution_timeout".to_owned(),
171 },
172 ExecutionError::Store(store) => Self::Failed {
173 code: format!("store.{}", store_code(store)),
174 },
175 ExecutionError::IdempotencyMismatch { .. } => Self::Failed {
176 code: "idempotency_mismatch".to_owned(),
177 },
178 ExecutionError::ScopeViolation => Self::Failed {
179 code: "scope_violation".to_owned(),
180 },
181 ExecutionError::Erasure(_) => Self::Failed {
182 code: "erasure".to_owned(),
183 },
184 ExecutionError::Other { code } => Self::Failed { code: code.clone() },
185 _ => Self::Failed {
186 code: "other".to_owned(),
187 },
188 }
189 }
190}
191
192fn store_code(error: &StoreError) -> &str {
193 match error {
194 StoreError::NotFound => "not_found",
195 StoreError::Conflict => "conflict",
196 StoreError::Unavailable => "unavailable",
197 StoreError::Timeout => "timeout",
198 StoreError::Serialization => "serialization",
199 StoreError::Corrupt => "corrupt",
200 StoreError::Other { code } => code,
201 _ => "other",
202 }
203}
204
205#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
207pub struct CommandJournalEntry {
208 pub command_id: CommandId,
210 pub account_id: AccountId,
212 pub idempotency_key: IdempotencyKey,
214 pub turn_id: TurnId,
216 pub case_ref: CaseRef,
218 pub command_type: String,
220 pub command_payload: serde_json::Value,
222 pub origin: CommandOrigin,
224 pub status: CommandJournalStatus,
226 #[serde(default, skip_serializing_if = "Option::is_none")]
228 pub result: Option<JournalOutcome>,
229 pub created_at: DateTime<Utc>,
231 #[serde(default, skip_serializing_if = "Option::is_none")]
233 pub completed_at: Option<DateTime<Utc>>,
234}
235
236impl CommandJournalEntry {
237 pub fn from_envelope<C: Serialize>(
242 envelope: &CommandEnvelope<C>,
243 command_type: impl Into<String>,
244 created_at: DateTime<Utc>,
245 ) -> Result<Self, StoreError> {
246 let command_payload =
247 serde_json::to_value(&envelope.command).map_err(|_| StoreError::Serialization)?;
248 Ok(Self {
249 command_id: envelope.command_id,
250 account_id: envelope.actor.account_id.clone(),
251 idempotency_key: envelope.idempotency_key.clone(),
252 turn_id: envelope.turn_id,
253 case_ref: envelope.case_ref.clone(),
254 command_type: command_type.into(),
255 command_payload,
256 origin: envelope.origin.clone(),
257 status: CommandJournalStatus::Pending,
258 result: None,
259 created_at,
260 completed_at: None,
261 })
262 }
263
264 #[must_use]
268 pub fn same_command(&self, other: &Self) -> bool {
269 self.case_ref == other.case_ref
270 && self.command_type == other.command_type
271 && self.command_payload == other.command_payload
272 }
273}
274
275#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
282#[serde(tag = "kind", rename_all = "snake_case")]
283pub enum JournalAdmission {
284 Fresh,
286 Replay(Box<CommandJournalEntry>),
289}
290
291impl JournalAdmission {
292 #[must_use]
294 pub fn replay(entry: CommandJournalEntry) -> Self {
295 Self::Replay(Box::new(entry))
296 }
297
298 #[must_use]
300 pub fn is_fresh(&self) -> bool {
301 matches!(self, Self::Fresh)
302 }
303
304 #[must_use]
306 pub fn replayed(&self) -> Option<&CommandJournalEntry> {
307 match self {
308 Self::Fresh => None,
309 Self::Replay(entry) => Some(entry),
310 }
311 }
312
313 #[must_use]
315 pub fn into_replayed(self) -> Option<CommandJournalEntry> {
316 match self {
317 Self::Fresh => None,
318 Self::Replay(entry) => Some(*entry),
319 }
320 }
321}
322
323#[async_trait]
328pub trait CommandJournalReader: Send + Sync {
329 async fn get(
334 &self,
335 account: &AccountId,
336 command_id: &CommandId,
337 ) -> Result<CommandJournalEntry, StoreError>;
338
339 async fn for_turn(
344 &self,
345 account: &AccountId,
346 turn_id: &TurnId,
347 ) -> Result<Vec<CommandJournalEntry>, StoreError>;
348
349 async fn pending_for_turn(
355 &self,
356 account: &AccountId,
357 turn_id: &TurnId,
358 ) -> Result<Vec<CommandJournalEntry>, StoreError>;
359}
360
361#[async_trait]
363pub trait CommandJournalWriter: Send + Sync {
364 async fn begin(&self, entry: CommandJournalEntry) -> Result<JournalAdmission, StoreError>;
374
375 async fn mark_executing(
382 &self,
383 account: &AccountId,
384 command_id: &CommandId,
385 ) -> Result<(), StoreError>;
386
387 async fn complete(
396 &self,
397 account: &AccountId,
398 command_id: &CommandId,
399 outcome: JournalOutcome,
400 ) -> Result<(), StoreError>;
401
402 async fn fail(
412 &self,
413 account: &AccountId,
414 command_id: &CommandId,
415 error: &ExecutionError,
416 ) -> Result<(), StoreError> {
417 self.complete(
418 account,
419 command_id,
420 JournalOutcome::from_execution_error(*command_id, error),
421 )
422 .await
423 }
424}
425
426pub trait CommandJournal: CommandJournalReader + CommandJournalWriter {}
432
433impl<T: CommandJournalReader + CommandJournalWriter + ?Sized> CommandJournal for T {}
434
435#[cfg(test)]
436mod tests {
437 use super::*;
438 use turnframe_core::error::RevisionConflict;
439 use turnframe_core::event::UnknownOutcome;
440
441 #[test]
442 fn transition_table() {
443 use CommandJournalStatus as S;
444 let allowed: &[(S, S)] = &[
445 (S::Pending, S::Executing),
446 (S::Pending, S::Committed),
447 (S::Pending, S::Failed),
448 (S::Pending, S::OutcomeUnknown),
449 (S::AwaitingConfirmation, S::Executing),
450 (S::AwaitingConfirmation, S::Committed),
451 (S::AwaitingConfirmation, S::Failed),
452 (S::AwaitingConfirmation, S::OutcomeUnknown),
453 (S::Executing, S::Committed),
454 (S::Executing, S::Failed),
455 (S::Executing, S::OutcomeUnknown),
456 (S::OutcomeUnknown, S::Committed),
457 (S::OutcomeUnknown, S::Failed),
458 ];
459 for from in S::ALL {
460 for to in S::ALL {
461 assert_eq!(
462 S::can_transition(from, to),
463 allowed.contains(&(from, to)),
464 "{from:?} -> {to:?}"
465 );
466 }
467 }
468 assert!(S::Committed.is_terminal());
469 assert!(S::Failed.is_terminal());
470 assert!(!S::OutcomeUnknown.is_terminal());
471 assert!(S::Pending.is_pending());
472 assert!(S::Executing.is_pending());
473 assert!(!S::AwaitingConfirmation.is_pending());
474 assert!(!S::AwaitingConfirmation.is_terminal());
475 }
476
477 #[test]
478 fn outcome_status_mapping_and_error_conversion() {
479 let id = CommandId::nil();
480 let conflict = ExecutionError::RevisionConflict(RevisionConflict {
481 expected: CaseRef::new("w", "c", CaseRevision(1)),
482 current_revision: CaseRevision(2),
483 });
484 let out = JournalOutcome::from_execution_error(id, &conflict);
485 assert_eq!(out.status(), CommandJournalStatus::Failed);
486 assert_eq!(
487 out,
488 JournalOutcome::RevisionConflict {
489 current_revision: CaseRevision(2)
490 }
491 );
492 let unknown = ExecutionError::OutcomeUnknown(UnknownOutcome {
493 attempt_id: "a".into(),
494 remote_ref: Some("r".into()),
495 reason: "timeout_after_send".into(),
496 });
497 assert_eq!(
498 JournalOutcome::from_execution_error(id, &unknown).status(),
499 CommandJournalStatus::OutcomeUnknown
500 );
501 assert_eq!(
502 JournalOutcome::from_execution_error(id, &ExecutionError::Timeout).status(),
503 CommandJournalStatus::OutcomeUnknown
504 );
505 assert_eq!(
506 JournalOutcome::from_execution_error(id, &ExecutionError::Store(StoreError::Timeout)),
507 JournalOutcome::Failed {
508 code: "store.timeout".into()
509 }
510 );
511 assert_eq!(
512 JournalOutcome::from_execution_error(id, &ExecutionError::ScopeViolation).status(),
513 CommandJournalStatus::Failed
514 );
515 }
516
517 #[test]
518 fn admission_helpers() {
519 assert!(JournalAdmission::Fresh.is_fresh());
520 assert!(JournalAdmission::Fresh.replayed().is_none());
521 assert!(JournalAdmission::Fresh.into_replayed().is_none());
522 }
523}