Skip to main content

turnframe_store/
journal.rs

1//! The command journal: idempotency admission and persisted outcomes
2//! (spec §16.2, §23.1, I14).
3//!
4//! # Contract
5//!
6//! * `(account_id, idempotency_key)` is unique. [`CommandJournalWriter::begin`] is the
7//!   single admission point: the first call for a key persists the entry and
8//!   answers [`JournalAdmission::Fresh`]; every later call for the same key, in
9//!   any status, answers [`JournalAdmission::Replay`] carrying the persisted
10//!   entry and therefore the persisted outcome. There is never a second
11//!   `Fresh` for one key, whatever the interleaving.
12//! * A `Replay` may carry an entry whose command differs from the caller's
13//!   (the key was reused for another command). The store does not judge that;
14//!   the runtime compares with [`CommandJournalEntry::same_command`] and raises
15//!   `ExecutionError::IdempotencyMismatch`.
16//! * Statuses move along [`CommandJournalStatus::can_transition`]:
17//!   `Pending → Executing`, then to `Committed`, `Failed` or `OutcomeUnknown`;
18//!   `OutcomeUnknown` is settled by reconciliation to `Committed` or `Failed`.
19//!   Re-recording the same terminal outcome is accepted; a different one is a
20//!   `Conflict`.
21//! * Crash recovery reads [`CommandJournalReader::for_turn`] and
22//!   [`CommandJournalReader::pending_for_turn`]: entries in `Pending` or `Executing`
23//!   are resumed by idempotency key, `Committed` ones are not re-executed,
24//!   `OutcomeUnknown` ones are reconciled (spec §23.1).
25
26use 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/// Lifecycle status of a journal entry.
37#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
38#[serde(rename_all = "snake_case")]
39pub enum CommandJournalStatus {
40    /// Admitted, not yet handed to the executor.
41    Pending,
42    /// Journaled for a confirmation card; runs only once that card is
43    /// confirmed, and crash recovery never resumes it.
44    AwaitingConfirmation,
45    /// Handed to the executor.
46    Executing,
47    /// The executor committed; the outcome is `Committed`.
48    Committed,
49    /// The command ended without effect or with a definite failure.
50    Failed,
51    /// An effect may exist and its result is unknown; reconcile, never retry
52    /// blindly (I15).
53    OutcomeUnknown,
54}
55
56impl CommandJournalStatus {
57    /// Every status, in declaration order.
58    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    /// Returns `true` for `Committed` and `Failed`.
68    #[must_use]
69    pub fn is_terminal(self) -> bool {
70        matches!(self, Self::Committed | Self::Failed)
71    }
72
73    /// Returns `true` for `Pending` and `Executing`: the entry must be resumed
74    /// after a crash.
75    #[must_use]
76    pub fn is_pending(self) -> bool {
77        matches!(self, Self::Pending | Self::Executing)
78    }
79
80    /// Legal transitions.
81    #[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/// The persisted outcome of a command.
97#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
98#[serde(tag = "kind", rename_all = "snake_case")]
99#[non_exhaustive]
100pub enum JournalOutcome {
101    /// The executor committed.
102    Committed {
103        /// Revision after the commit.
104        new_revision: CaseRevision,
105        /// Events produced.
106        event_ids: Vec<EventId>,
107    },
108    /// The domain refused.
109    Rejected {
110        /// The rejection.
111        rejection: DomainRejection,
112    },
113    /// The expected revision was stale.
114    RevisionConflict {
115        /// Revision found.
116        current_revision: CaseRevision,
117    },
118    /// Execution failed with a stable code and no effect.
119    Failed {
120        /// Stable code.
121        code: String,
122    },
123    /// An effect may exist; its result is unknown (I15).
124    OutcomeUnknown {
125        /// Attempt identifier for reconciliation.
126        attempt_id: AttemptId,
127        /// Remote reference, when the remote returned one.
128        #[serde(default, skip_serializing_if = "Option::is_none")]
129        remote_ref: Option<String>,
130        /// Stable reason code.
131        reason: String,
132    },
133}
134
135impl JournalOutcome {
136    /// The status an entry takes when this outcome is recorded.
137    #[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    /// Maps an execution error to the outcome the journal records for it.
149    ///
150    /// `Timeout` and `OutcomeUnknown` become [`JournalOutcome::OutcomeUnknown`]
151    /// because an effect may exist; every other variant becomes a definite
152    /// outcome without effect.
153    #[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/// One row of the command journal (spec §22.2, `tf_command_journal`).
206#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
207pub struct CommandJournalEntry {
208    /// The command.
209    pub command_id: CommandId,
210    /// Owning tenant.
211    pub account_id: AccountId,
212    /// Idempotency key, unique per account (I14).
213    pub idempotency_key: IdempotencyKey,
214    /// Turn that produced the command.
215    pub turn_id: TurnId,
216    /// Target case and expected revision.
217    pub case_ref: CaseRef,
218    /// Stable command type label (e.g. `"trip.add_extra"`).
219    pub command_type: String,
220    /// The command payload as JSON.
221    pub command_payload: serde_json::Value,
222    /// Authorizing origin.
223    pub origin: CommandOrigin,
224    /// Current status.
225    pub status: CommandJournalStatus,
226    /// Persisted outcome, once recorded.
227    #[serde(default, skip_serializing_if = "Option::is_none")]
228    pub result: Option<JournalOutcome>,
229    /// Admission time, supplied by the caller.
230    pub created_at: DateTime<Utc>,
231    /// When the outcome was recorded, stamped by the store's clock.
232    #[serde(default, skip_serializing_if = "Option::is_none")]
233    pub completed_at: Option<DateTime<Utc>>,
234}
235
236impl CommandJournalEntry {
237    /// Builds a `Pending` entry from a typed envelope, serializing the command.
238    ///
239    /// # Errors
240    /// * `Serialization` when the command cannot be rendered as JSON.
241    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    /// Returns `true` when both entries describe the same command: same case
265    /// reference, type and payload. Used to detect an idempotency key reused
266    /// for a different command.
267    #[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/// Result of [`CommandJournalWriter::begin`].
276///
277/// The entry is boxed because the two answers are wildly different in size and
278/// this value is returned from every admission, including the overwhelmingly
279/// common `Fresh` one; build it with [`JournalAdmission::replay`] rather than
280/// writing the `Box` yourself.
281#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
282#[serde(tag = "kind", rename_all = "snake_case")]
283pub enum JournalAdmission {
284    /// The key was never seen for this account; the entry is now persisted.
285    Fresh,
286    /// The key exists; the persisted entry (and its outcome, if any) is
287    /// returned instead of admitting a second command (spec §16.2).
288    Replay(Box<CommandJournalEntry>),
289}
290
291impl JournalAdmission {
292    /// Builds a [`Self::Replay`] from the persisted entry.
293    #[must_use]
294    pub fn replay(entry: CommandJournalEntry) -> Self {
295        Self::Replay(Box::new(entry))
296    }
297
298    /// Returns `true` for [`Self::Fresh`].
299    #[must_use]
300    pub fn is_fresh(&self) -> bool {
301        matches!(self, Self::Fresh)
302    }
303
304    /// The replayed entry, if any.
305    #[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    /// Takes ownership of the replayed entry, if any.
314    #[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/// The read half of the command journal (spec §22.1).
324///
325/// Reading the journal says which commands were admitted and how they ended;
326/// it admits nothing and settles nothing.
327#[async_trait]
328pub trait CommandJournalReader: Send + Sync {
329    /// Loads one entry.
330    ///
331    /// # Errors
332    /// * `NotFound` when it does not exist for `account`.
333    async fn get(
334        &self,
335        account: &AccountId,
336        command_id: &CommandId,
337    ) -> Result<CommandJournalEntry, StoreError>;
338
339    /// Every entry of a turn, ordered by `created_at` then `command_id`.
340    ///
341    /// # Errors
342    /// * [`StoreError`] when the listing could not be read.
343    async fn for_turn(
344        &self,
345        account: &AccountId,
346        turn_id: &TurnId,
347    ) -> Result<Vec<CommandJournalEntry>, StoreError>;
348
349    /// Entries of a turn whose status is `Pending` or `Executing`, ordered by
350    /// `created_at` then `command_id` (spec §23.1: resume by idempotency key).
351    ///
352    /// # Errors
353    /// * [`StoreError`] when the listing could not be read.
354    async fn pending_for_turn(
355        &self,
356        account: &AccountId,
357        turn_id: &TurnId,
358    ) -> Result<Vec<CommandJournalEntry>, StoreError>;
359}
360
361/// The write half of the command journal (spec §22.1).
362#[async_trait]
363pub trait CommandJournalWriter: Send + Sync {
364    /// Admits a command under `UNIQUE (account_id, idempotency_key)`.
365    ///
366    /// The entry is persisted as given (normally `Pending`) when the key is
367    /// new. When the key exists, nothing is written and the persisted entry is
368    /// returned as [`JournalAdmission::Replay`], whatever its status.
369    ///
370    /// # Errors
371    /// * `Conflict` when `command_id` already exists for the account under
372    ///   another idempotency key.
373    async fn begin(&self, entry: CommandJournalEntry) -> Result<JournalAdmission, StoreError>;
374
375    /// Moves `Pending → Executing`. Calling it on an `Executing` entry is
376    /// accepted without change.
377    ///
378    /// # Errors
379    /// * `NotFound` when the entry does not exist for `account`.
380    /// * `Conflict` when the status is neither `Pending` nor `Executing`.
381    async fn mark_executing(
382        &self,
383        account: &AccountId,
384        command_id: &CommandId,
385    ) -> Result<(), StoreError>;
386
387    /// Records the outcome and moves the entry to
388    /// [`JournalOutcome::status`]. Recording the same outcome again on an
389    /// entry already in that status is accepted without change.
390    ///
391    /// # Errors
392    /// * `NotFound` when the entry does not exist for `account`.
393    /// * `Conflict` when the transition is illegal or a different outcome is
394    ///   already recorded.
395    async fn complete(
396        &self,
397        account: &AccountId,
398        command_id: &CommandId,
399        outcome: JournalOutcome,
400    ) -> Result<(), StoreError>;
401
402    /// Records an execution error as the outcome, through
403    /// [`JournalOutcome::from_execution_error`]. Same rules as
404    /// [`Self::complete`].
405    ///
406    /// The default body is the whole contract; override it only to save a round
407    /// trip, and keep the mapping identical.
408    ///
409    /// # Errors
410    /// * Whatever [`Self::complete`] returns.
411    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
426/// The command journal (spec §22.1): both halves.
427///
428/// There is nothing to implement here: write [`CommandJournalReader`] and
429/// [`CommandJournalWriter`] and the blanket implementation below supplies this
430/// trait.
431pub 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}