1use 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
65const OUTBOX_ID_DOMAIN: &str = "turnframe.outbox_id.v1";
67
68const ATTEMPT_ID_DOMAIN: &str = "turnframe.attempt_id.v1";
70
71#[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#[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#[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#[derive(Debug, Clone, PartialEq, Eq)]
112#[non_exhaustive]
113pub enum Admission {
114 Fresh,
116 Resume,
119 Settled(Box<JournalOutcome>),
122}
123
124impl Admission {
125 #[must_use]
127 pub const fn needs_execution(&self) -> bool {
128 matches!(self, Self::Fresh | Self::Resume)
129 }
130}
131
132#[derive(Debug, Clone, Default)]
134#[non_exhaustive]
135pub struct ExecutionReport {
136 pub outcomes: Vec<CommandOutcomeRecord>,
138 pub events: Vec<EventBatch>,
140 pub committed: Vec<CommittedEvent<serde_json::Value>>,
142 pub changed: IndexMap<CaseKey, CaseRevision>,
144 pub outbox: Vec<OutboxEntry>,
146 pub completions: Vec<(CommandId, JournalOutcome)>,
148 pub rejections: Vec<(CaseRef, DomainRejection)>,
152}
153
154impl ExecutionReport {
155 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 #[must_use]
183 pub fn event_ids(&self) -> Vec<EventId> {
184 self.committed.iter().map(|event| event.event_id).collect()
185 }
186
187 #[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 #[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 #[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 #[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 #[must_use]
230 pub fn failed_outright(&self) -> bool {
231 !self.outcomes.is_empty() && !self.any_committed()
232 }
233
234 #[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 #[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#[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 #[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 #[must_use]
309 pub const fn config(&self) -> &ExecutionConfig {
310 &self.config
311 }
312
313 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 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 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 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 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 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 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 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 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 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 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 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 fn record_all(
691 &self,
692 batch: &CommandBatch<serde_json::Value>,
693 report: &mut ExecutionReport,
694 outcome: &CommandOutcome,
695 error: Option<&ExecutionError>,
696 ) {
697 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 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 #[must_use]
764 pub fn outbox(&self) -> &Arc<dyn OutboxStore> {
765 &self.outbox
766 }
767
768 #[must_use]
770 pub fn journal(&self) -> &Arc<dyn CommandJournal> {
771 &self.journal
772 }
773}
774
775fn 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
795fn 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 _ => CommandOutcome::Failed {
818 code: "unknown_recorded_outcome".to_owned(),
819 },
820 }
821}
822
823fn 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 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
864fn 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}