use std::fmt;
use std::sync::Arc;
use chrono::{DateTime, Utc};
use indexmap::IndexMap;
use turnframe_core::case::{CaseKey, CaseRef};
use turnframe_core::command::{AtomicityScope, CommandBatch, CommandEnvelope};
use turnframe_core::error::{
DomainRejection, ErrorClassification, ExecutionError, OrchestratorError,
};
use turnframe_core::event::{Commit, CommittedEvent, OutboxEntry, OutboxStatus};
use turnframe_core::flow::WorkflowRegistry;
use turnframe_core::hash::derive_uuid;
use turnframe_core::ids::{AccountId, AttemptId, CaseRevision, CommandId, EventId, OutboxId};
use turnframe_core::reduce::CommandRef;
use turnframe_core::replay::{CommandOutcome, CommandOutcomeRecord};
use turnframe_store::commit::{CommitBundle, CommitReceipt, CommitStore};
use turnframe_store::events::EventBatch;
use turnframe_store::journal::{
CommandJournal, CommandJournalEntry, JournalAdmission, JournalOutcome,
};
use turnframe_store::outbox::OutboxStore;
use crate::config::ExecutionConfig;
const OUTBOX_ID_DOMAIN: &str = "turnframe.outbox_id.v1";
const ATTEMPT_ID_DOMAIN: &str = "turnframe.attempt_id.v1";
#[must_use]
pub fn derive_outbox_id(command_id: &CommandId) -> OutboxId {
OutboxId::from(derive_uuid(OUTBOX_ID_DOMAIN, &[&command_id.to_string()]))
}
#[must_use]
pub fn derive_attempt_id(command_id: &CommandId) -> AttemptId {
AttemptId::new(
derive_uuid(ATTEMPT_ID_DOMAIN, &[&command_id.to_string()])
.simple()
.to_string(),
)
}
#[must_use]
pub fn command_type(case_ref: &CaseRef, command: &serde_json::Value) -> String {
let workflow = case_ref.workflow.as_str();
match command {
serde_json::Value::String(variant) => format!("{workflow}.{variant}"),
serde_json::Value::Object(map) if map.len() == 1 => match map.keys().next() {
Some(variant) => format!("{workflow}.{variant}"),
None => format!("{workflow}.command"),
},
_ => format!("{workflow}.command"),
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum Admission {
Fresh,
Resume,
Settled(Box<JournalOutcome>),
}
impl Admission {
#[must_use]
pub const fn needs_execution(&self) -> bool {
matches!(self, Self::Fresh | Self::Resume)
}
}
#[derive(Debug, Clone, Default)]
#[non_exhaustive]
pub struct ExecutionReport {
pub outcomes: Vec<CommandOutcomeRecord>,
pub events: Vec<EventBatch>,
pub committed: Vec<CommittedEvent<serde_json::Value>>,
pub changed: IndexMap<CaseKey, CaseRevision>,
pub outbox: Vec<OutboxEntry>,
pub completions: Vec<(CommandId, JournalOutcome)>,
pub rejections: Vec<(CaseRef, DomainRejection)>,
}
impl ExecutionReport {
pub fn absorb(&mut self, earlier: Self) {
let mut merged = earlier;
merged.outcomes.append(&mut self.outcomes);
merged.events.append(&mut self.events);
merged.committed.append(&mut self.committed);
merged.outbox.append(&mut self.outbox);
merged.completions.append(&mut self.completions);
merged.rejections.append(&mut self.rejections);
for (key, revision) in std::mem::take(&mut self.changed) {
merged.changed.insert(key, revision);
}
*self = merged;
}
#[must_use]
pub fn event_ids(&self) -> Vec<EventId> {
self.committed.iter().map(|event| event.event_id).collect()
}
#[must_use]
pub fn all_committed(&self) -> bool {
!self.outcomes.is_empty()
&& self
.outcomes
.iter()
.all(|record| matches!(record.outcome, CommandOutcome::Committed { .. }))
}
#[must_use]
pub fn any_committed(&self) -> bool {
self.outcomes
.iter()
.any(|record| matches!(record.outcome, CommandOutcome::Committed { .. }))
}
#[must_use]
pub fn has_unknown_outcome(&self) -> bool {
self.outcomes
.iter()
.any(|record| matches!(record.outcome, CommandOutcome::OutcomeUnknown { .. }))
}
#[must_use]
pub fn pending_attempts(&self) -> Vec<AttemptId> {
self.outcomes
.iter()
.filter_map(|record| match &record.outcome {
CommandOutcome::OutcomeUnknown { attempt_id } => Some(attempt_id.clone()),
_ => None,
})
.collect()
}
#[must_use]
pub fn failed_outright(&self) -> bool {
!self.outcomes.is_empty() && !self.any_committed()
}
#[must_use]
pub fn any_uncommitted(&self) -> bool {
self.outcomes
.iter()
.any(|record| !matches!(record.outcome, CommandOutcome::Committed { .. }))
}
#[must_use]
pub fn bundle(&self) -> CommitBundle {
let mut bundle = CommitBundle::new();
for (command_id, outcome) in &self.completions {
bundle = bundle.with_journal_completion(*command_id, outcome.clone());
}
for batch in &self.events {
bundle = bundle.with_events(batch.clone());
}
for entry in &self.outbox {
bundle = bundle.with_outbox_entry(entry.clone());
}
bundle
}
}
#[derive(Clone)]
pub struct CommandExecutor {
workflows: Arc<WorkflowRegistry>,
journal: Arc<dyn CommandJournal>,
commit: Arc<dyn CommitStore>,
outbox: Arc<dyn OutboxStore>,
config: ExecutionConfig,
}
impl fmt::Debug for CommandExecutor {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("CommandExecutor")
.field("workflows", &self.workflows.len())
.field("config", &self.config)
.finish_non_exhaustive()
}
}
impl CommandExecutor {
#[must_use]
pub fn new(
workflows: Arc<WorkflowRegistry>,
journal: Arc<dyn CommandJournal>,
commit: Arc<dyn CommitStore>,
outbox: Arc<dyn OutboxStore>,
config: ExecutionConfig,
) -> Self {
Self {
workflows,
journal,
commit,
outbox,
config,
}
}
#[must_use]
pub const fn config(&self) -> &ExecutionConfig {
&self.config
}
pub async fn journal_pending(
&self,
batches: &[CommandBatch<serde_json::Value>],
now: DateTime<Utc>,
) -> Result<Vec<CommandId>, OrchestratorError> {
let mut admitted = Vec::new();
for batch in batches {
for envelope in &batch.envelopes {
let mut entry = self.entry_for(envelope, now)?;
entry.status = turnframe_store::journal::CommandJournalStatus::AwaitingConfirmation;
match self.journal.begin(entry).await {
Ok(_) => admitted.push(envelope.command_id),
Err(error) => return Err(OrchestratorError::Store(error)),
}
}
}
Ok(admitted)
}
pub async fn resume_confirmed(
&self,
account: &AccountId,
command_refs: &[CommandRef],
origin: &turnframe_core::command::CommandOrigin,
actor: &turnframe_core::turn::ActorContext,
turn_id: turnframe_core::ids::TurnId,
) -> Result<Vec<CommandBatch<serde_json::Value>>, OrchestratorError> {
let mut batches: IndexMap<turnframe_core::ids::BatchId, CommandBatch<serde_json::Value>> =
IndexMap::new();
for command_ref in command_refs {
let entry = self
.journal
.get(account, &command_ref.command_id)
.await
.map_err(OrchestratorError::Store)?;
if entry.status.is_terminal() {
continue;
}
if !self.authorizes(account, &entry, origin).await {
tracing::warn!(
target: "turnframe.execute",
"a confirmed command's policy refuses the origin that confirmed it; dropped"
);
continue;
}
let envelope = CommandEnvelope {
command_id: entry.command_id,
turn_id,
actor: actor.clone(),
case_ref: entry.case_ref.clone(),
idempotency_key: entry.idempotency_key.clone(),
origin: origin.clone(),
command: entry.command_payload.clone(),
};
batches
.entry(command_ref.batch_id)
.or_insert_with(|| CommandBatch {
batch_id: command_ref.batch_id,
scope: AtomicityScope::PerCase,
envelopes: Vec::new(),
})
.envelopes
.push(envelope);
}
Ok(batches.into_values().collect())
}
async fn authorizes(
&self,
account: &AccountId,
entry: &CommandJournalEntry,
origin: &turnframe_core::command::CommandOrigin,
) -> bool {
let Ok(registered) = self.workflows.require(&entry.case_ref.workflow) else {
return false;
};
let Ok(loaded) = registered
.executor
.load(account, &entry.case_ref.case_id)
.await
else {
return false;
};
let Ok(policy) = registered
.definition
.command_policy(loaded.value.as_ref(), &entry.command_payload)
else {
return false;
};
turnframe_core::command::origin_satisfies(origin, &policy)
}
pub async fn admit(
&self,
envelope: &CommandEnvelope<serde_json::Value>,
now: DateTime<Utc>,
) -> Result<Admission, OrchestratorError> {
let entry = self.entry_for(envelope, now)?;
match self.journal.begin(entry.clone()).await {
Ok(JournalAdmission::Fresh) => Ok(Admission::Fresh),
Ok(JournalAdmission::Replay(existing)) => {
if !existing.same_command(&entry) {
return Err(OrchestratorError::Execution(
ExecutionError::IdempotencyMismatch {
command_id: envelope.command_id,
},
));
}
Ok(match (existing.status, existing.result.clone()) {
(status, Some(outcome)) if !status.is_pending() => {
Admission::Settled(Box::new(outcome))
}
_ => Admission::Resume,
})
}
Err(error) => Err(OrchestratorError::Store(error)),
}
}
pub async fn execute(
&self,
account: &AccountId,
batches: &[CommandBatch<serde_json::Value>],
now: DateTime<Utc>,
) -> Result<ExecutionReport, OrchestratorError> {
let mut report = ExecutionReport::default();
for batch in batches {
if batch.is_empty() {
continue;
}
self.execute_batch(account, batch, now, &mut report).await?;
if !self.config.allow_cross_case_partial_success
&& report.outcomes.last().is_some_and(|record| {
!matches!(record.outcome, CommandOutcome::Committed { .. })
})
{
break;
}
}
Ok(report)
}
async fn execute_batch(
&self,
account: &AccountId,
batch: &CommandBatch<serde_json::Value>,
now: DateTime<Utc>,
report: &mut ExecutionReport,
) -> Result<(), OrchestratorError> {
let Some(first) = batch.envelopes.first() else {
return Ok(());
};
let case_ref = first.case_ref.clone();
let Ok(registered) = self.workflows.require(&case_ref.workflow) else {
self.record_all(
batch,
report,
&CommandOutcome::Failed {
code: "unknown_workflow".to_owned(),
},
None,
);
return Ok(());
};
let mut admissions = Vec::with_capacity(batch.envelopes.len());
for envelope in &batch.envelopes {
match self.admit(envelope, now).await {
Ok(admission) => admissions.push(admission),
Err(OrchestratorError::Execution(ExecutionError::IdempotencyMismatch {
..
})) => {
self.record_all(
batch,
report,
&CommandOutcome::Failed {
code: "idempotency_mismatch".to_owned(),
},
None,
);
return Ok(());
}
Err(error) => return Err(error),
}
}
if admissions
.iter()
.all(|admission| !admission.needs_execution())
{
for (envelope, admission) in batch.envelopes.iter().zip(&admissions) {
let Admission::Settled(outcome) = admission else {
continue;
};
if let JournalOutcome::Rejected { rejection } = &**outcome {
report
.rejections
.push((envelope.case_ref.clone(), rejection.clone()));
}
report.outcomes.push(CommandOutcomeRecord {
command_ref: CommandRef {
batch_id: batch.batch_id,
command_id: envelope.command_id,
},
idempotency_key: envelope.idempotency_key.clone(),
case_ref: envelope.case_ref.clone(),
origin: Some(envelope.origin.clone()),
outcome: replayed_outcome(outcome),
});
}
return Ok(());
}
for envelope in &batch.envelopes {
if let Err(error) = self
.journal
.mark_executing(account, &envelope.command_id)
.await
{
return Err(OrchestratorError::Store(error));
}
}
let executed = registered.executor.execute(batch.clone()).await;
match executed {
Ok(commit) => self.record_commit(account, batch, &case_ref, commit, now, report),
Err(error) => {
let outcome = failure_outcome(first.command_id, &error);
self.record_all(batch, report, &outcome, Some(&error));
}
}
Ok(())
}
fn record_commit(
&self,
account: &AccountId,
batch: &CommandBatch<serde_json::Value>,
case_ref: &CaseRef,
commit: Commit<serde_json::Value, serde_json::Value>,
now: DateTime<Utc>,
report: &mut ExecutionReport,
) {
let Some(first) = batch.envelopes.first() else {
return;
};
let event_ids: Vec<EventId> = commit.event_ids();
if !commit.events.is_empty() {
report.events.push(EventBatch::new(
account.clone(),
case_ref.key(),
first.command_id,
commit.new_revision,
commit.events.clone(),
));
report.committed.extend(commit.events.iter().cloned());
}
report.changed.insert(case_ref.key(), commit.new_revision);
for (position, envelope) in batch.envelopes.iter().enumerate() {
let outcome = CommandOutcome::Committed {
new_revision: commit.new_revision,
event_ids: if position == 0 {
event_ids.clone()
} else {
Vec::new()
},
};
report.outcomes.push(CommandOutcomeRecord {
command_ref: CommandRef {
batch_id: batch.batch_id,
command_id: envelope.command_id,
},
idempotency_key: envelope.idempotency_key.clone(),
case_ref: envelope.case_ref.clone(),
origin: Some(envelope.origin.clone()),
outcome,
});
report.completions.push((
envelope.command_id,
JournalOutcome::Committed {
new_revision: commit.new_revision,
event_ids: if position == 0 {
event_ids.clone()
} else {
Vec::new()
},
},
));
if let AtomicityScope::ExternalSaga { saga } = &batch.scope {
report.outbox.push(outbox_row(envelope, saga, now));
}
}
}
fn record_all(
&self,
batch: &CommandBatch<serde_json::Value>,
report: &mut ExecutionReport,
outcome: &CommandOutcome,
error: Option<&ExecutionError>,
) {
if let Some(ExecutionError::Rejected(rejection)) = error
&& let Some(first) = batch.envelopes.first()
{
report
.rejections
.push((first.case_ref.clone(), rejection.clone()));
}
for envelope in &batch.envelopes {
report.outcomes.push(CommandOutcomeRecord {
command_ref: CommandRef {
batch_id: batch.batch_id,
command_id: envelope.command_id,
},
idempotency_key: envelope.idempotency_key.clone(),
case_ref: envelope.case_ref.clone(),
origin: Some(envelope.origin.clone()),
outcome: outcome.clone(),
});
let completion = match error {
Some(error) => JournalOutcome::from_execution_error(envelope.command_id, error),
None => JournalOutcome::Failed {
code: outcome_code(outcome),
},
};
report.completions.push((envelope.command_id, completion));
}
}
fn entry_for(
&self,
envelope: &CommandEnvelope<serde_json::Value>,
now: DateTime<Utc>,
) -> Result<CommandJournalEntry, OrchestratorError> {
let label = command_type(&envelope.case_ref, &envelope.command);
CommandJournalEntry::from_envelope(envelope, label, now).map_err(OrchestratorError::Store)
}
pub async fn commit(
&self,
account: &AccountId,
bundle: CommitBundle,
) -> Result<CommitReceipt, OrchestratorError> {
self.commit.commit(account, bundle).await.map_err(|error| {
if error.reconciliation_required() {
tracing::error!(
target: "turnframe.execute",
"commit bundle did not confirm; the turn must re-read rather than retry"
);
}
OrchestratorError::Store(error)
})
}
#[must_use]
pub fn outbox(&self) -> &Arc<dyn OutboxStore> {
&self.outbox
}
#[must_use]
pub fn journal(&self) -> &Arc<dyn CommandJournal> {
&self.journal
}
}
fn outbox_row(
envelope: &CommandEnvelope<serde_json::Value>,
saga: &str,
now: DateTime<Utc>,
) -> OutboxEntry {
OutboxEntry {
outbox_id: derive_outbox_id(&envelope.command_id),
command_id: envelope.command_id,
destination: saga.to_owned(),
payload: envelope.command.clone(),
idempotency_key: envelope.idempotency_key.clone(),
status: OutboxStatus::Pending,
attempt_count: 0,
next_attempt_at: None,
created_at: now,
completed_at: None,
}
}
fn replayed_outcome(outcome: &JournalOutcome) -> CommandOutcome {
match outcome {
JournalOutcome::Committed {
new_revision,
event_ids,
} => CommandOutcome::Committed {
new_revision: *new_revision,
event_ids: event_ids.clone(),
},
JournalOutcome::Rejected { rejection } => CommandOutcome::Rejected {
code: rejection.code.clone(),
},
JournalOutcome::RevisionConflict { current_revision } => CommandOutcome::RevisionConflict {
current_revision: *current_revision,
},
JournalOutcome::Failed { code } => CommandOutcome::Failed { code: code.clone() },
JournalOutcome::OutcomeUnknown { attempt_id, .. } => CommandOutcome::OutcomeUnknown {
attempt_id: attempt_id.clone(),
},
_ => CommandOutcome::Failed {
code: "unknown_recorded_outcome".to_owned(),
},
}
}
fn failure_outcome(command_id: CommandId, error: &ExecutionError) -> CommandOutcome {
match error {
ExecutionError::RevisionConflict(conflict) => CommandOutcome::RevisionConflict {
current_revision: conflict.current_revision,
},
ExecutionError::Rejected(rejection) => CommandOutcome::Rejected {
code: rejection.code.clone(),
},
ExecutionError::OutcomeUnknown(unknown) => CommandOutcome::OutcomeUnknown {
attempt_id: unknown.attempt_id.clone(),
},
ExecutionError::Timeout => CommandOutcome::OutcomeUnknown {
attempt_id: derive_attempt_id(&command_id),
},
ExecutionError::Store(store) if store.effect_may_have_happened() => {
CommandOutcome::OutcomeUnknown {
attempt_id: derive_attempt_id(&command_id),
}
}
ExecutionError::Store(_) => CommandOutcome::Failed {
code: "store".to_owned(),
},
ExecutionError::IdempotencyMismatch { .. } => CommandOutcome::Failed {
code: "idempotency_mismatch".to_owned(),
},
ExecutionError::ScopeViolation => CommandOutcome::Failed {
code: "scope_violation".to_owned(),
},
ExecutionError::Erasure(_) => CommandOutcome::Failed {
code: "erasure".to_owned(),
},
ExecutionError::Other { code } => CommandOutcome::Failed { code: code.clone() },
_ => CommandOutcome::Failed {
code: "other".to_owned(),
},
}
}
fn outcome_code(outcome: &CommandOutcome) -> String {
match outcome {
CommandOutcome::Failed { code } => code.clone(),
CommandOutcome::Rejected { code } => code.as_str().to_owned(),
CommandOutcome::RevisionConflict { .. } => "revision_conflict".to_owned(),
_ => "other".to_owned(),
}
}
#[cfg(test)]
mod tests {
use turnframe_core::ids::CaseRevision;
use super::*;
fn case() -> CaseRef {
CaseRef::new("trip", "i1", CaseRevision(3))
}
#[test]
fn a_command_type_names_its_variant() {
assert_eq!(
command_type(&case(), &serde_json::json!({"set_name": {"value": "x"}})),
"trip.set_name"
);
assert_eq!(
command_type(&case(), &serde_json::json!("rebook")),
"trip.rebook"
);
assert_eq!(
command_type(&case(), &serde_json::json!({"a": 1, "b": 2})),
"trip.command"
);
}
#[test]
fn a_timeout_is_an_unknown_outcome_and_not_a_failure() {
let command_id = CommandId::nil();
let outcome = failure_outcome(command_id, &ExecutionError::Timeout);
assert!(matches!(outcome, CommandOutcome::OutcomeUnknown { .. }));
assert_eq!(
derive_attempt_id(&command_id),
derive_attempt_id(&command_id),
"a reconciler must be able to name the same attempt twice"
);
}
#[test]
fn outbox_identifiers_are_derived_from_the_command() {
let command_id = CommandId::nil();
assert_eq!(derive_outbox_id(&command_id), derive_outbox_id(&command_id));
}
}