Skip to main content

turnframe_runtime/
execute.rs

1//! Command execution: the journal, the domain, the outbox and the one atomic
2//! write (spec §16, §23 steps M and N).
3//!
4//! Everything the runtime does that a user could notice happens here, and it
5//! happens in a fixed order that the rest of the library depends on:
6//!
7//! 1. **Admission before effect.** Every envelope is admitted to the command
8//!    journal under `UNIQUE (account_id, idempotency_key)` *before* the domain
9//!    is asked to do anything (§16.2). A key that is already there is not a
10//!    second command: [`Admission::Settled`] returns the outcome that was
11//!    recorded the first time, without running the effect again (I14).
12//! 2. **Optimistic concurrency.** The batch carries the revision it was planned
13//!    against and the executor checks it in the same statement that writes
14//!    (§16.1, I13). A mismatch is [`CommandOutcome::RevisionConflict`], never a
15//!    blind overwrite.
16//! 3. **Uncertainty is a state.** A timeout after transmission is
17//!    [`CommandOutcome::OutcomeUnknown`] carrying an [`AttemptId`], because the
18//!    effect may exist. Nothing in this module retries it (I15, §16.5).
19//! 4. **One write.** The journal outcomes, the events, the card resolutions,
20//!    the new cards, the outbox rows, the replay record and the phase marker
21//!    travel in one [`CommitBundle`] and land together or not at all
22//!    (§16.3, §23 step N).
23//!
24//! # What the batch is the unit of
25//!
26//! A [`CommandBatch`] commits as a whole: one revision, one set of events. The
27//! per-command [`CommandOutcomeRecord`]s therefore all report the same
28//! revision, and the batch's event identifiers are recorded once, against its
29//! first envelope, so two receipts can never cite the same events as if they
30//! were two separate outcomes.
31//!
32//! # What is deliberately not atomic
33//!
34//! The domain's own state commit may live in another database, and Turnframe
35//! does not attempt a distributed transaction. Safety across that seam is the
36//! journal: the entry exists before the effect, the executor is idempotent on
37//! the key, and [`crate::recover`] resumes a pending entry by that key
38//! (§23.1).
39
40use std::fmt;
41use std::sync::Arc;
42
43use chrono::{DateTime, Utc};
44use indexmap::IndexMap;
45use turnframe_core::case::{CaseKey, CaseRef};
46use turnframe_core::command::{AtomicityScope, CommandBatch, CommandEnvelope};
47use turnframe_core::error::{
48    DomainRejection, ErrorClassification, ExecutionError, OrchestratorError,
49};
50use turnframe_core::event::{Commit, CommittedEvent, OutboxEntry, OutboxStatus};
51use turnframe_core::flow::WorkflowRegistry;
52use turnframe_core::hash::derive_uuid;
53use turnframe_core::ids::{AccountId, AttemptId, CaseRevision, CommandId, EventId, OutboxId};
54use turnframe_core::reduce::CommandRef;
55use turnframe_core::replay::{CommandOutcome, CommandOutcomeRecord};
56use turnframe_store::commit::{CommitBundle, CommitReceipt, CommitStore};
57use turnframe_store::events::EventBatch;
58use turnframe_store::journal::{
59    CommandJournal, CommandJournalEntry, JournalAdmission, JournalOutcome,
60};
61use turnframe_store::outbox::OutboxStore;
62
63use crate::config::ExecutionConfig;
64
65/// Domain separation of the derived outbox identifiers.
66const OUTBOX_ID_DOMAIN: &str = "turnframe.outbox_id.v1";
67
68/// Domain separation of the derived external attempt identifiers.
69const ATTEMPT_ID_DOMAIN: &str = "turnframe.attempt_id.v1";
70
71/// Derives the outbox row identifier of `command_id`, so replaying a turn
72/// enqueues the same row instead of a second one.
73#[must_use]
74pub fn derive_outbox_id(command_id: &CommandId) -> OutboxId {
75    OutboxId::from(derive_uuid(OUTBOX_ID_DOMAIN, &[&command_id.to_string()]))
76}
77
78/// Derives the attempt identifier a command's unknown outcome is reconciled by.
79///
80/// It names one attempt at one external effect, which is what a reconciler
81/// quotes back to the remote system (§16.5, I15).
82#[must_use]
83pub fn derive_attempt_id(command_id: &CommandId) -> AttemptId {
84    AttemptId::new(
85        derive_uuid(ATTEMPT_ID_DOMAIN, &[&command_id.to_string()])
86            .simple()
87            .to_string(),
88    )
89}
90
91/// A stable label for an erased command, for the journal's `command_type`.
92///
93/// Erased commands are the JSON a domain's own enum serializes to, so the
94/// externally tagged single-key object and the unit-variant string both name
95/// their variant; anything else is labelled by the workflow alone rather than
96/// by a guess.
97#[must_use]
98pub fn command_type(case_ref: &CaseRef, command: &serde_json::Value) -> String {
99    let workflow = case_ref.workflow.as_str();
100    match command {
101        serde_json::Value::String(variant) => format!("{workflow}.{variant}"),
102        serde_json::Value::Object(map) if map.len() == 1 => match map.keys().next() {
103            Some(variant) => format!("{workflow}.{variant}"),
104            None => format!("{workflow}.command"),
105        },
106        _ => format!("{workflow}.command"),
107    }
108}
109
110/// What the journal said about one envelope before it ran.
111#[derive(Debug, Clone, PartialEq, Eq)]
112#[non_exhaustive]
113pub enum Admission {
114    /// The key is new; the entry is now persisted as `Pending`.
115    Fresh,
116    /// The key exists in `Pending` or `Executing`: a previous attempt was
117    /// interrupted and this one resumes it by key (§23.1).
118    Resume,
119    /// The key exists with a recorded outcome. The original outcome is
120    /// returned and nothing runs again (I14).
121    Settled(Box<JournalOutcome>),
122}
123
124impl Admission {
125    /// Returns `true` when the domain still has to be called.
126    #[must_use]
127    pub const fn needs_execution(&self) -> bool {
128        matches!(self, Self::Fresh | Self::Resume)
129    }
130}
131
132/// What executing one turn's batches produced.
133#[derive(Debug, Clone, Default)]
134#[non_exhaustive]
135pub struct ExecutionReport {
136    /// One record per envelope, in batch then envelope order.
137    pub outcomes: Vec<CommandOutcomeRecord>,
138    /// Event batches to append, one per committed command batch.
139    pub events: Vec<EventBatch>,
140    /// Every committed event, in commit order, for receipts (§17.3).
141    pub committed: Vec<CommittedEvent<serde_json::Value>>,
142    /// Cases whose revision moved, with the revision they moved to.
143    pub changed: IndexMap<CaseKey, CaseRevision>,
144    /// Outbox rows for the external effects that were accepted locally (§16.4).
145    pub outbox: Vec<OutboxEntry>,
146    /// Journal completions for the commit bundle.
147    pub completions: Vec<(CommandId, JournalOutcome)>,
148    /// Domain rejections this stage decided, with the case each was aimed at. Some can only be
149    /// decided here, such as a name no registry entry matches; each reaches the writing stage
150    /// with the domain's explanation, as a reducer refusal does, so the reply cannot claim it.
151    pub rejections: Vec<(CaseRef, DomainRejection)>,
152}
153
154impl ExecutionReport {
155    /// Folds an earlier report of the same turn into this one.
156    ///
157    /// A turn that commits twice executed twice, and everything downstream of
158    /// the commit reads one report: the receipts the writing stage may rest on,
159    /// the refusals it must not contradict, and whether anything failed. So the
160    /// halves are joined **after** the second commit and never before — a bundle
161    /// built from a joined report would append the first half's events a second
162    /// time.
163    ///
164    /// `earlier` goes first in every list, because commit order is the order
165    /// receipts are read in. A case that moved in both halves keeps the later
166    /// revision, which is the one it is at.
167    pub fn absorb(&mut self, earlier: Self) {
168        let mut merged = earlier;
169        merged.outcomes.append(&mut self.outcomes);
170        merged.events.append(&mut self.events);
171        merged.committed.append(&mut self.committed);
172        merged.outbox.append(&mut self.outbox);
173        merged.completions.append(&mut self.completions);
174        merged.rejections.append(&mut self.rejections);
175        for (key, revision) in std::mem::take(&mut self.changed) {
176            merged.changed.insert(key, revision);
177        }
178        *self = merged;
179    }
180
181    /// Every event identifier the turn committed, in commit order.
182    #[must_use]
183    pub fn event_ids(&self) -> Vec<EventId> {
184        self.committed.iter().map(|event| event.event_id).collect()
185    }
186
187    /// Returns `true` when every command committed.
188    #[must_use]
189    pub fn all_committed(&self) -> bool {
190        !self.outcomes.is_empty()
191            && self
192                .outcomes
193                .iter()
194                .all(|record| matches!(record.outcome, CommandOutcome::Committed { .. }))
195    }
196
197    /// Returns `true` when at least one command committed.
198    #[must_use]
199    pub fn any_committed(&self) -> bool {
200        self.outcomes
201            .iter()
202            .any(|record| matches!(record.outcome, CommandOutcome::Committed { .. }))
203    }
204
205    /// Returns `true` when a command ended with an effect that may or may not
206    /// have happened, so the turn must say "verification in progress" and a
207    /// reconciler must settle it (I15).
208    #[must_use]
209    pub fn has_unknown_outcome(&self) -> bool {
210        self.outcomes
211            .iter()
212            .any(|record| matches!(record.outcome, CommandOutcome::OutcomeUnknown { .. }))
213    }
214
215    /// The attempts a reconciler has to settle.
216    #[must_use]
217    pub fn pending_attempts(&self) -> Vec<AttemptId> {
218        self.outcomes
219            .iter()
220            .filter_map(|record| match &record.outcome {
221                CommandOutcome::OutcomeUnknown { attempt_id } => Some(attempt_id.clone()),
222                _ => None,
223            })
224            .collect()
225    }
226
227    /// Returns `true` when no command committed and at least one failed, which
228    /// is what makes a receipt impossible.
229    #[must_use]
230    pub fn failed_outright(&self) -> bool {
231        !self.outcomes.is_empty() && !self.any_committed()
232    }
233
234    /// Returns `true` when any command did **not** commit — a rejection, a
235    /// conflict, a failure, or an outcome nobody knows yet.
236    ///
237    /// This, rather than [`Self::failed_outright`], is what makes a notice
238    /// mandatory: a turn where two commands committed and a third did not is
239    /// still a turn that has to say so, and receipts alone would let the user
240    /// read the silence as success.
241    #[must_use]
242    pub fn any_uncommitted(&self) -> bool {
243        self.outcomes
244            .iter()
245            .any(|record| !matches!(record.outcome, CommandOutcome::Committed { .. }))
246    }
247
248    /// The bundle items this execution produced: journal completions, events
249    /// and outbox rows.
250    ///
251    /// The caller adds the card resolutions, the new cards, the replay record
252    /// and the phase marker, then writes everything at once (§16.3).
253    #[must_use]
254    pub fn bundle(&self) -> CommitBundle {
255        let mut bundle = CommitBundle::new();
256        for (command_id, outcome) in &self.completions {
257            bundle = bundle.with_journal_completion(*command_id, outcome.clone());
258        }
259        for batch in &self.events {
260            bundle = bundle.with_events(batch.clone());
261        }
262        for entry in &self.outbox {
263            bundle = bundle.with_outbox_entry(entry.clone());
264        }
265        bundle
266    }
267}
268
269/// Runs command batches against the domain, the journal and the outbox.
270#[derive(Clone)]
271pub struct CommandExecutor {
272    workflows: Arc<WorkflowRegistry>,
273    journal: Arc<dyn CommandJournal>,
274    commit: Arc<dyn CommitStore>,
275    outbox: Arc<dyn OutboxStore>,
276    config: ExecutionConfig,
277}
278
279impl fmt::Debug for CommandExecutor {
280    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
281        f.debug_struct("CommandExecutor")
282            .field("workflows", &self.workflows.len())
283            .field("config", &self.config)
284            .finish_non_exhaustive()
285    }
286}
287
288impl CommandExecutor {
289    /// Builds an executor.
290    #[must_use]
291    pub fn new(
292        workflows: Arc<WorkflowRegistry>,
293        journal: Arc<dyn CommandJournal>,
294        commit: Arc<dyn CommitStore>,
295        outbox: Arc<dyn OutboxStore>,
296        config: ExecutionConfig,
297    ) -> Self {
298        Self {
299            workflows,
300            journal,
301            commit,
302            outbox,
303            config,
304        }
305    }
306
307    /// The execution configuration in force.
308    #[must_use]
309    pub const fn config(&self) -> &ExecutionConfig {
310        &self.config
311    }
312
313    /// Admits the commands a confirmation card will authorize, in `Pending`
314    /// (spec §15.3).
315    ///
316    /// They do not run: the card names them, and a later click resumes exactly
317    /// these entries by idempotency key. Re-journaling the same entry (a
318    /// replayed turn) is accepted and changes nothing.
319    ///
320    /// # Errors
321    ///
322    /// [`OrchestratorError::Store`] when the journal could not be written.
323    pub async fn journal_pending(
324        &self,
325        batches: &[CommandBatch<serde_json::Value>],
326        now: DateTime<Utc>,
327    ) -> Result<Vec<CommandId>, OrchestratorError> {
328        let mut admitted = Vec::new();
329        for batch in batches {
330            for envelope in &batch.envelopes {
331                let mut entry = self.entry_for(envelope, now)?;
332                entry.status = turnframe_store::journal::CommandJournalStatus::AwaitingConfirmation;
333                match self.journal.begin(entry).await {
334                    Ok(_) => admitted.push(envelope.command_id),
335                    Err(error) => return Err(OrchestratorError::Store(error)),
336                }
337            }
338        }
339        Ok(admitted)
340    }
341
342    /// Rebuilds the batches a confirmed card authorized, from the journal
343    /// entries the card names (spec §15.3).
344    ///
345    /// The commands come back exactly as they were reviewed, with the
346    /// idempotency key they were admitted under, so answering the card twice
347    /// cannot execute twice. The `origin` replaces the one recorded at
348    /// admission: the authority is now the click, and policy is re-checked
349    /// against it.
350    ///
351    /// The rebuilt batches carry [`AtomicityScope::PerCase`], which is what a
352    /// reviewed set has to be: the user confirmed one change to one case, so
353    /// its commands commit together or not at all.
354    ///
355    /// Every rebuilt envelope is policed again before it is returned: the
356    /// domain's [`CommandPolicy`](turnframe_core::command::CommandPolicy) for
357    /// the command, against the origin the click minted
358    /// ([`origin_satisfies`](turnframe_core::command::origin_satisfies)). The
359    /// card was built for exactly these commands under exactly that policy, so
360    /// the check should never fire — which is the point of running it. A
361    /// command it refuses is dropped rather than executed.
362    ///
363    /// # Errors
364    ///
365    /// [`OrchestratorError::Store`] when an entry named by the card is not in
366    /// the journal for this account.
367    pub async fn resume_confirmed(
368        &self,
369        account: &AccountId,
370        command_refs: &[CommandRef],
371        origin: &turnframe_core::command::CommandOrigin,
372        actor: &turnframe_core::turn::ActorContext,
373        turn_id: turnframe_core::ids::TurnId,
374    ) -> Result<Vec<CommandBatch<serde_json::Value>>, OrchestratorError> {
375        let mut batches: IndexMap<turnframe_core::ids::BatchId, CommandBatch<serde_json::Value>> =
376            IndexMap::new();
377        for command_ref in command_refs {
378            let entry = self
379                .journal
380                .get(account, &command_ref.command_id)
381                .await
382                .map_err(OrchestratorError::Store)?;
383            if entry.status.is_terminal() {
384                // The card is being answered a second time; the original
385                // outcome stands and nothing is rebuilt for it (I14).
386                continue;
387            }
388            if !self.authorizes(account, &entry, origin).await {
389                tracing::warn!(
390                    target: "turnframe.execute",
391                    "a confirmed command's policy refuses the origin that confirmed it; dropped"
392                );
393                continue;
394            }
395            let envelope = CommandEnvelope {
396                command_id: entry.command_id,
397                turn_id,
398                actor: actor.clone(),
399                case_ref: entry.case_ref.clone(),
400                idempotency_key: entry.idempotency_key.clone(),
401                origin: origin.clone(),
402                command: entry.command_payload.clone(),
403            };
404            batches
405                .entry(command_ref.batch_id)
406                .or_insert_with(|| CommandBatch {
407                    batch_id: command_ref.batch_id,
408                    scope: AtomicityScope::PerCase,
409                    envelopes: Vec::new(),
410                })
411                .envelopes
412                .push(envelope);
413        }
414        Ok(batches.into_values().collect())
415    }
416
417    /// Whether the click that answered a card actually satisfies the policy of
418    /// the command it is about to run (I12).
419    ///
420    /// The policy is asked for against the case's **current** state, because
421    /// that is the state the command will meet. A workflow that cannot be read
422    /// answers `false`: a policy nobody could consult is not a policy that was
423    /// satisfied (I19).
424    async fn authorizes(
425        &self,
426        account: &AccountId,
427        entry: &CommandJournalEntry,
428        origin: &turnframe_core::command::CommandOrigin,
429    ) -> bool {
430        let Ok(registered) = self.workflows.require(&entry.case_ref.workflow) else {
431            return false;
432        };
433        let Ok(loaded) = registered
434            .executor
435            .load(account, &entry.case_ref.case_id)
436            .await
437        else {
438            return false;
439        };
440        let Ok(policy) = registered
441            .definition
442            .command_policy(loaded.value.as_ref(), &entry.command_payload)
443        else {
444            return false;
445        };
446        turnframe_core::command::origin_satisfies(origin, &policy)
447    }
448
449    /// Admits one envelope to the journal, before anything runs (§16.2, I14).
450    ///
451    /// This is the single door every effect goes through, and it is public
452    /// because an adopter driving execution themselves has to go through it
453    /// too. The three answers are the whole contract: the key is new, the key
454    /// is there and unsettled (resume it — never re-plan it), or the key is
455    /// there with an outcome (return that outcome and run nothing).
456    ///
457    /// # Errors
458    ///
459    /// * [`OrchestratorError::Store`] when the journal could not be written;
460    /// * [`OrchestratorError::Execution`] with
461    ///   [`ExecutionError::IdempotencyMismatch`] when the key names a
462    ///   *different* command, which is a defect nobody may guess past.
463    pub async fn admit(
464        &self,
465        envelope: &CommandEnvelope<serde_json::Value>,
466        now: DateTime<Utc>,
467    ) -> Result<Admission, OrchestratorError> {
468        let entry = self.entry_for(envelope, now)?;
469        match self.journal.begin(entry.clone()).await {
470            Ok(JournalAdmission::Fresh) => Ok(Admission::Fresh),
471            Ok(JournalAdmission::Replay(existing)) => {
472                if !existing.same_command(&entry) {
473                    return Err(OrchestratorError::Execution(
474                        ExecutionError::IdempotencyMismatch {
475                            command_id: envelope.command_id,
476                        },
477                    ));
478                }
479                Ok(match (existing.status, existing.result.clone()) {
480                    (status, Some(outcome)) if !status.is_pending() => {
481                        Admission::Settled(Box::new(outcome))
482                    }
483                    _ => Admission::Resume,
484                })
485            }
486            Err(error) => Err(OrchestratorError::Store(error)),
487        }
488    }
489
490    /// Executes every batch (spec §23 step M).
491    ///
492    /// Batches run in order. When
493    /// [`ExecutionConfig::allow_cross_case_partial_success`] is off, the first
494    /// batch that does not commit stops the rest: a turn that half-happened
495    /// across two cases is harder to explain than one that did not start.
496    ///
497    /// # Errors
498    ///
499    /// [`OrchestratorError::Store`] when the journal itself could not be
500    /// reached. A command that merely *failed* is not an error here: it is an
501    /// outcome, and it is reported as one.
502    pub async fn execute(
503        &self,
504        account: &AccountId,
505        batches: &[CommandBatch<serde_json::Value>],
506        now: DateTime<Utc>,
507    ) -> Result<ExecutionReport, OrchestratorError> {
508        let mut report = ExecutionReport::default();
509        for batch in batches {
510            if batch.is_empty() {
511                continue;
512            }
513            self.execute_batch(account, batch, now, &mut report).await?;
514            if !self.config.allow_cross_case_partial_success
515                && report.outcomes.last().is_some_and(|record| {
516                    !matches!(record.outcome, CommandOutcome::Committed { .. })
517                })
518            {
519                break;
520            }
521        }
522        Ok(report)
523    }
524
525    async fn execute_batch(
526        &self,
527        account: &AccountId,
528        batch: &CommandBatch<serde_json::Value>,
529        now: DateTime<Utc>,
530        report: &mut ExecutionReport,
531    ) -> Result<(), OrchestratorError> {
532        let Some(first) = batch.envelopes.first() else {
533            return Ok(());
534        };
535        let case_ref = first.case_ref.clone();
536        let Ok(registered) = self.workflows.require(&case_ref.workflow) else {
537            self.record_all(
538                batch,
539                report,
540                &CommandOutcome::Failed {
541                    code: "unknown_workflow".to_owned(),
542                },
543                None,
544            );
545            return Ok(());
546        };
547
548        // 1. Admission, before the domain hears about any of it (§16.2).
549        let mut admissions = Vec::with_capacity(batch.envelopes.len());
550        for envelope in &batch.envelopes {
551            match self.admit(envelope, now).await {
552                Ok(admission) => admissions.push(admission),
553                // The same key names a different command: executing either one
554                // would be guessing which the user meant.
555                Err(OrchestratorError::Execution(ExecutionError::IdempotencyMismatch {
556                    ..
557                })) => {
558                    self.record_all(
559                        batch,
560                        report,
561                        &CommandOutcome::Failed {
562                            code: "idempotency_mismatch".to_owned(),
563                        },
564                        None,
565                    );
566                    return Ok(());
567                }
568                Err(error) => return Err(error),
569            }
570        }
571
572        // 2. Everything already settled: return the original outcomes.
573        if admissions
574            .iter()
575            .all(|admission| !admission.needs_execution())
576        {
577            for (envelope, admission) in batch.envelopes.iter().zip(&admissions) {
578                let Admission::Settled(outcome) = admission else {
579                    continue;
580                };
581                // A settled rejection read back from the journal says the same
582                // thing it said the first time, so the turn can too.
583                if let JournalOutcome::Rejected { rejection } = &**outcome {
584                    report
585                        .rejections
586                        .push((envelope.case_ref.clone(), rejection.clone()));
587                }
588                report.outcomes.push(CommandOutcomeRecord {
589                    command_ref: CommandRef {
590                        batch_id: batch.batch_id,
591                        command_id: envelope.command_id,
592                    },
593                    idempotency_key: envelope.idempotency_key.clone(),
594                    case_ref: envelope.case_ref.clone(),
595                    origin: Some(envelope.origin.clone()),
596                    outcome: replayed_outcome(outcome),
597                });
598            }
599            return Ok(());
600        }
601
602        // 3. Hand the batch to the domain, under its expected revision (I13).
603        for envelope in &batch.envelopes {
604            if let Err(error) = self
605                .journal
606                .mark_executing(account, &envelope.command_id)
607                .await
608            {
609                return Err(OrchestratorError::Store(error));
610            }
611        }
612        let executed = registered.executor.execute(batch.clone()).await;
613        match executed {
614            Ok(commit) => self.record_commit(account, batch, &case_ref, commit, now, report),
615            Err(error) => {
616                let outcome = failure_outcome(first.command_id, &error);
617                self.record_all(batch, report, &outcome, Some(&error));
618            }
619        }
620        Ok(())
621    }
622
623    fn record_commit(
624        &self,
625        account: &AccountId,
626        batch: &CommandBatch<serde_json::Value>,
627        case_ref: &CaseRef,
628        commit: Commit<serde_json::Value, serde_json::Value>,
629        now: DateTime<Utc>,
630        report: &mut ExecutionReport,
631    ) {
632        let Some(first) = batch.envelopes.first() else {
633            return;
634        };
635        let event_ids: Vec<EventId> = commit.event_ids();
636        if !commit.events.is_empty() {
637            report.events.push(EventBatch::new(
638                account.clone(),
639                case_ref.key(),
640                first.command_id,
641                commit.new_revision,
642                commit.events.clone(),
643            ));
644            report.committed.extend(commit.events.iter().cloned());
645        }
646        report.changed.insert(case_ref.key(), commit.new_revision);
647
648        for (position, envelope) in batch.envelopes.iter().enumerate() {
649            // The batch is the unit of commit, so its events are recorded once:
650            // against the first envelope. The rest report the revision they
651            // reached and cite nothing, which is what keeps two receipts from
652            // claiming the same events twice.
653            let outcome = CommandOutcome::Committed {
654                new_revision: commit.new_revision,
655                event_ids: if position == 0 {
656                    event_ids.clone()
657                } else {
658                    Vec::new()
659                },
660            };
661            report.outcomes.push(CommandOutcomeRecord {
662                command_ref: CommandRef {
663                    batch_id: batch.batch_id,
664                    command_id: envelope.command_id,
665                },
666                idempotency_key: envelope.idempotency_key.clone(),
667                case_ref: envelope.case_ref.clone(),
668                origin: Some(envelope.origin.clone()),
669                outcome,
670            });
671            report.completions.push((
672                envelope.command_id,
673                JournalOutcome::Committed {
674                    new_revision: commit.new_revision,
675                    event_ids: if position == 0 {
676                        event_ids.clone()
677                    } else {
678                        Vec::new()
679                    },
680                },
681            ));
682            if let AtomicityScope::ExternalSaga { saga } = &batch.scope {
683                report.outbox.push(outbox_row(envelope, saga, now));
684            }
685        }
686    }
687
688    /// Records the same outcome for every envelope of a batch that did not
689    /// commit, and the matching journal completion.
690    fn record_all(
691        &self,
692        batch: &CommandBatch<serde_json::Value>,
693        report: &mut ExecutionReport,
694        outcome: &CommandOutcome,
695        error: Option<&ExecutionError>,
696    ) {
697        // The domain's own words, once for the batch that carried them. Without
698        // this the explanation reaches nobody and the writing stage has nothing
699        // either way about the write it is about to claim.
700        if let Some(ExecutionError::Rejected(rejection)) = error
701            && let Some(first) = batch.envelopes.first()
702        {
703            report
704                .rejections
705                .push((first.case_ref.clone(), rejection.clone()));
706        }
707        for envelope in &batch.envelopes {
708            report.outcomes.push(CommandOutcomeRecord {
709                command_ref: CommandRef {
710                    batch_id: batch.batch_id,
711                    command_id: envelope.command_id,
712                },
713                idempotency_key: envelope.idempotency_key.clone(),
714                case_ref: envelope.case_ref.clone(),
715                origin: Some(envelope.origin.clone()),
716                outcome: outcome.clone(),
717            });
718            let completion = match error {
719                Some(error) => JournalOutcome::from_execution_error(envelope.command_id, error),
720                None => JournalOutcome::Failed {
721                    code: outcome_code(outcome),
722                },
723            };
724            report.completions.push((envelope.command_id, completion));
725        }
726    }
727
728    fn entry_for(
729        &self,
730        envelope: &CommandEnvelope<serde_json::Value>,
731        now: DateTime<Utc>,
732    ) -> Result<CommandJournalEntry, OrchestratorError> {
733        let label = command_type(&envelope.case_ref, &envelope.command);
734        CommandJournalEntry::from_envelope(envelope, label, now).map_err(OrchestratorError::Store)
735    }
736
737    /// Writes the bundle, all or nothing (spec §16.3, §23 step N).
738    ///
739    /// # Errors
740    ///
741    /// [`OrchestratorError::Store`]. A [`StoreError::Timeout`](turnframe_core::error::StoreError::Timeout) means the write
742    /// *may* have landed: the caller must re-read rather than retry, and the
743    /// bundle's own atomicity guarantees that whatever it finds is either all
744    /// of it or none of it (§16.5).
745    pub async fn commit(
746        &self,
747        account: &AccountId,
748        bundle: CommitBundle,
749    ) -> Result<CommitReceipt, OrchestratorError> {
750        self.commit.commit(account, bundle).await.map_err(|error| {
751            if error.reconciliation_required() {
752                tracing::error!(
753                    target: "turnframe.execute",
754                    "commit bundle did not confirm; the turn must re-read rather than retry"
755                );
756            }
757            OrchestratorError::Store(error)
758        })
759    }
760
761    /// The outbox the dispatcher reads, for callers that reconcile from the
762    /// same executor.
763    #[must_use]
764    pub fn outbox(&self) -> &Arc<dyn OutboxStore> {
765        &self.outbox
766    }
767
768    /// The command journal.
769    #[must_use]
770    pub fn journal(&self) -> &Arc<dyn CommandJournal> {
771        &self.journal
772    }
773}
774
775/// One outbox row for an external effect that was accepted locally (§16.4).
776fn outbox_row(
777    envelope: &CommandEnvelope<serde_json::Value>,
778    saga: &str,
779    now: DateTime<Utc>,
780) -> OutboxEntry {
781    OutboxEntry {
782        outbox_id: derive_outbox_id(&envelope.command_id),
783        command_id: envelope.command_id,
784        destination: saga.to_owned(),
785        payload: envelope.command.clone(),
786        idempotency_key: envelope.idempotency_key.clone(),
787        status: OutboxStatus::Pending,
788        attempt_count: 0,
789        next_attempt_at: None,
790        created_at: now,
791        completed_at: None,
792    }
793}
794
795/// The outcome a persisted journal result stands for on a repeat.
796fn replayed_outcome(outcome: &JournalOutcome) -> CommandOutcome {
797    match outcome {
798        JournalOutcome::Committed {
799            new_revision,
800            event_ids,
801        } => CommandOutcome::Committed {
802            new_revision: *new_revision,
803            event_ids: event_ids.clone(),
804        },
805        JournalOutcome::Rejected { rejection } => CommandOutcome::Rejected {
806            code: rejection.code.clone(),
807        },
808        JournalOutcome::RevisionConflict { current_revision } => CommandOutcome::RevisionConflict {
809            current_revision: *current_revision,
810        },
811        JournalOutcome::Failed { code } => CommandOutcome::Failed { code: code.clone() },
812        JournalOutcome::OutcomeUnknown { attempt_id, .. } => CommandOutcome::OutcomeUnknown {
813            attempt_id: attempt_id.clone(),
814        },
815        // A persisted outcome this version does not know is still a settled
816        // one: it is reported as a failure rather than re-executed.
817        _ => CommandOutcome::Failed {
818            code: "unknown_recorded_outcome".to_owned(),
819        },
820    }
821}
822
823/// Maps an execution failure onto the outcome the turn records.
824fn failure_outcome(command_id: CommandId, error: &ExecutionError) -> CommandOutcome {
825    match error {
826        ExecutionError::RevisionConflict(conflict) => CommandOutcome::RevisionConflict {
827            current_revision: conflict.current_revision,
828        },
829        ExecutionError::Rejected(rejection) => CommandOutcome::Rejected {
830            code: rejection.code.clone(),
831        },
832        ExecutionError::OutcomeUnknown(unknown) => CommandOutcome::OutcomeUnknown {
833            attempt_id: unknown.attempt_id.clone(),
834        },
835        // A timeout is not a failure: the effect may exist, so it becomes an
836        // attempt somebody has to settle rather than one to repeat (I15).
837        ExecutionError::Timeout => CommandOutcome::OutcomeUnknown {
838            attempt_id: derive_attempt_id(&command_id),
839        },
840        ExecutionError::Store(store) if store.effect_may_have_happened() => {
841            CommandOutcome::OutcomeUnknown {
842                attempt_id: derive_attempt_id(&command_id),
843            }
844        }
845        ExecutionError::Store(_) => CommandOutcome::Failed {
846            code: "store".to_owned(),
847        },
848        ExecutionError::IdempotencyMismatch { .. } => CommandOutcome::Failed {
849            code: "idempotency_mismatch".to_owned(),
850        },
851        ExecutionError::ScopeViolation => CommandOutcome::Failed {
852            code: "scope_violation".to_owned(),
853        },
854        ExecutionError::Erasure(_) => CommandOutcome::Failed {
855            code: "erasure".to_owned(),
856        },
857        ExecutionError::Other { code } => CommandOutcome::Failed { code: code.clone() },
858        _ => CommandOutcome::Failed {
859            code: "other".to_owned(),
860        },
861    }
862}
863
864/// The stable code of an outcome, for a journal completion that has no error.
865fn outcome_code(outcome: &CommandOutcome) -> String {
866    match outcome {
867        CommandOutcome::Failed { code } => code.clone(),
868        CommandOutcome::Rejected { code } => code.as_str().to_owned(),
869        CommandOutcome::RevisionConflict { .. } => "revision_conflict".to_owned(),
870        _ => "other".to_owned(),
871    }
872}
873
874#[cfg(test)]
875mod tests {
876    use turnframe_core::ids::CaseRevision;
877
878    use super::*;
879
880    fn case() -> CaseRef {
881        CaseRef::new("trip", "i1", CaseRevision(3))
882    }
883
884    #[test]
885    fn a_command_type_names_its_variant() {
886        assert_eq!(
887            command_type(&case(), &serde_json::json!({"set_name": {"value": "x"}})),
888            "trip.set_name"
889        );
890        assert_eq!(
891            command_type(&case(), &serde_json::json!("rebook")),
892            "trip.rebook"
893        );
894        assert_eq!(
895            command_type(&case(), &serde_json::json!({"a": 1, "b": 2})),
896            "trip.command"
897        );
898    }
899
900    #[test]
901    fn a_timeout_is_an_unknown_outcome_and_not_a_failure() {
902        let command_id = CommandId::nil();
903        let outcome = failure_outcome(command_id, &ExecutionError::Timeout);
904        assert!(matches!(outcome, CommandOutcome::OutcomeUnknown { .. }));
905        assert_eq!(
906            derive_attempt_id(&command_id),
907            derive_attempt_id(&command_id),
908            "a reconciler must be able to name the same attempt twice"
909        );
910    }
911
912    #[test]
913    fn outbox_identifiers_are_derived_from_the_command() {
914        let command_id = CommandId::nil();
915        assert_eq!(derive_outbox_id(&command_id), derive_outbox_id(&command_id));
916    }
917}