Skip to main content

ironflow_engine/
engine.rs

1//! The core [`Engine`] -- orchestrates workflow execution and persistence.
2//!
3//! The engine ties together a `RunStore` for persistence, an [`AgentProvider`]
4//! for AI operations, and a registry of [`WorkflowHandler`]s.
5//!
6//! Handlers are Rust-native: steps can reference previous outputs, use native
7//! `if`/`else`/`match` for conditional branching, and execute in parallel.
8
9use std::collections::{HashMap, HashSet};
10use std::fmt;
11use std::sync::{Arc, Mutex};
12use std::time::Instant;
13
14use chrono::{DateTime, TimeDelta, Utc};
15use rust_decimal::Decimal;
16use serde_json::{Value, to_value};
17use tokio::spawn;
18use tracing::{error, info, warn};
19use uuid::Uuid;
20
21use ironflow_core::error::OperationError;
22#[cfg(feature = "prometheus")]
23use ironflow_core::metric_names::{
24    RUN_BUDGET_EXCEEDED_TOTAL, RUN_COST_USD, RUN_DURATION_SECONDS, RUNS_ACTIVE, RUNS_TOTAL,
25};
26use ironflow_core::provider::{AgentProvider, LABEL_ROOT_RUN_ID};
27use ironflow_store::error::StoreError;
28use ironflow_store::models::{
29    NewRun, NewSignal, Run, RunActor, RunCreation, RunFilter, RunStatus, RunUpdate, SignalInsert,
30    SignalStepResolution, StepStatus, StepUpdate, TriggerKind,
31};
32use ironflow_store::store::Store;
33#[cfg(feature = "prometheus")]
34use metrics::{counter, gauge, histogram};
35
36use crate::artifact::ArtifactSink;
37use crate::budget::{BudgetConfig, month_start};
38use crate::context::{PARENT_RUN_ID_LABEL, WorkflowContext};
39use crate::error::EngineError;
40use crate::executor::{StepInterceptor, StepResult};
41use crate::guard::{WorkflowGuardConfig, new_shared_guard_state};
42use crate::handler::{WorkflowHandler, WorkflowInfo};
43use crate::log_sender::LogSender;
44use crate::notify::{
45    ApprovalRequestedEvent, Event, EventPublisher, EventSubscriber, RunBudgetExceededEvent,
46    RunFailedEvent, RunStatusChangedEvent, SignalAwaitedEvent, SignalReceivedEvent,
47    WorkflowEventBus,
48};
49use crate::plan::{
50    ExecutionPlan, PlanOptions, PlanRecorder, SharedPlanRecorder, estimate_durations, lock_plan,
51};
52use crate::retry_policy::{backoff_for_retry, is_run_retryable};
53use crate::schedule::CronSchedule;
54use crate::signal::{
55    Signal, SignalDelivery, SignalRejected, SignalResumed, received_output, validate_step_payload,
56};
57use ironflow_core::decision::DecisionProvider;
58
59/// Result of a workflow execution, carrying the final [`Run`] and per-step
60/// metrics collected during execution.
61///
62/// Returned by [`Engine::run_handler`], [`Engine::execute_handler_run`],
63/// [`Engine::execute_run`], and [`Engine::resume_run`].
64///
65/// # Examples
66///
67/// ```no_run
68/// use ironflow_engine::engine::WorkflowResult;
69///
70/// # fn example(result: WorkflowResult) {
71/// println!("run {} finished with {} steps", result.run.id, result.steps.len());
72/// for step in &result.steps {
73///     println!("  {} ({:?}): {}ms", step.name, step.status, step.duration_ms);
74/// }
75/// # }
76/// ```
77#[derive(Debug, Clone)]
78pub struct WorkflowResult {
79    /// The finalized run record.
80    pub run: Run,
81    /// Per-step results in execution order.
82    pub steps: Vec<StepResult>,
83}
84
85/// Optional settings for [`Engine::enqueue_handler_with_options`].
86///
87/// All fields fall back to handler or server defaults when left at their
88/// [`Default`] value.
89///
90/// # Examples
91///
92/// ```
93/// use ironflow_engine::engine::EnqueueOptions;
94/// use rust_decimal::Decimal;
95///
96/// let options = EnqueueOptions {
97///     max_retries: 3,
98///     max_cost_usd: Some(Decimal::new(50, 2)),
99///     ..Default::default()
100/// };
101/// assert_eq!(options.max_retries, 3);
102/// ```
103#[derive(Debug, Clone, Default)]
104pub struct EnqueueOptions {
105    /// Number of automatic retries granted to the run.
106    pub max_retries: u32,
107    /// Labels merged on top of the handler's default labels.
108    pub labels: HashMap<String, String>,
109    /// Defer execution until this instant instead of running as soon as a
110    /// worker picks the run up.
111    pub scheduled_at: Option<DateTime<Utc>>,
112    /// Cost cap for the run. Overrides both the handler default and the server
113    /// default. `None` falls back to
114    /// [`BudgetConfig::resolve_run_cap`](crate::budget::BudgetConfig::resolve_run_cap).
115    pub max_cost_usd: Option<Decimal>,
116    /// Authenticated principal that triggered the run. `None` for cron,
117    /// webhook, and programmatic triggers.
118    pub created_by: Option<RunActor>,
119    /// Idempotency key binding this enqueue to a single run.
120    ///
121    /// When set and already bound to a run created within
122    /// [`IDEMPOTENCY_WINDOW`](ironflow_store::entities::IDEMPOTENCY_WINDOW),
123    /// nothing is enqueued and the original run is replayed.
124    pub idempotency_key: Option<String>,
125    /// Concurrency key making the run exclusive.
126    ///
127    /// When set, the run is refused with [`EngineError::ConcurrencyConflict`]
128    /// while another non-terminal run holds the same key. The key is released
129    /// once the holder reaches a terminal state.
130    pub concurrency_key: Option<String>,
131}
132
133/// Where a run resumes once an approval, a human input or an escalation
134/// resolves the gate it was suspended on.
135///
136/// # Examples
137///
138/// ```
139/// use ironflow_engine::engine::ExecutionMode;
140///
141/// assert_eq!(ExecutionMode::default(), ExecutionMode::Local);
142/// ```
143#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
144pub enum ExecutionMode {
145    /// Resume in-process via [`Engine::resume_run`]. Used by
146    /// [`crate::testing::TestEngine`] and single-process deployments where
147    /// the API also registers the handlers.
148    #[default]
149    Local,
150    /// Requeue the run to `Pending` instead. A worker's `pick_next_pending`
151    /// claims it and finishes it via [`Engine::execute_handler_run`].
152    Workers,
153}
154
155/// The workflow orchestration engine.
156///
157/// Holds references to the store, agent provider, and a registry of
158/// [`WorkflowHandler`]s.
159///
160/// # Examples
161///
162/// ```no_run
163/// use std::sync::Arc;
164/// use ironflow_engine::engine::Engine;
165/// use ironflow_engine::config::ShellConfig;
166/// use ironflow_engine::handler::{WorkflowHandler, HandlerFuture, WorkflowInfo};
167/// use ironflow_engine::context::WorkflowContext;
168/// use ironflow_store::memory::InMemoryStore;
169/// use ironflow_store::models::TriggerKind;
170/// use ironflow_core::providers::claude::ClaudeCodeProvider;
171/// use serde_json::json;
172///
173/// struct CiWorkflow;
174/// impl WorkflowHandler for CiWorkflow {
175///     fn name(&self) -> &str { "ci" }
176///     fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
177///         Box::pin(async move {
178///             ctx.shell("test", ShellConfig::new("cargo test")).await?;
179///             Ok(())
180///         })
181///     }
182/// }
183///
184/// # async fn example() -> Result<(), ironflow_engine::error::EngineError> {
185/// let store = Arc::new(InMemoryStore::new());
186/// let provider = Arc::new(ClaudeCodeProvider::new());
187/// let mut engine = Engine::new(store, provider);
188/// engine.register(CiWorkflow)?;
189///
190/// let result = engine.run_handler("ci", TriggerKind::Manual, json!({})).await?;
191/// tracing::info!(run_id = %result.run.id, status = ?result.run.status, steps = result.steps.len(), "run completed");
192/// # Ok(())
193/// # }
194/// ```
195pub struct Engine {
196    store: Arc<dyn Store>,
197    provider: Arc<dyn AgentProvider>,
198    handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
199    event_publisher: EventPublisher,
200    log_sender: Option<LogSender>,
201    budget: BudgetConfig,
202    artifact_sink: Option<Arc<dyn ArtifactSink>>,
203    guard_config: Option<WorkflowGuardConfig>,
204    event_bus: Option<WorkflowEventBus>,
205    decision_provider: Option<Arc<dyn DecisionProvider>>,
206    step_interceptor: Option<Arc<dyn StepInterceptor>>,
207    execution_mode: ExecutionMode,
208}
209
210/// Validate a workflow category path.
211///
212/// A category is a `/`-separated list of non-empty segments. This function
213/// rejects empty paths, leading or trailing `/`, consecutive `/`, and
214/// segments containing only whitespace.
215///
216/// # Errors
217///
218/// Returns [`EngineError::InvalidWorkflow`] when the category is malformed.
219fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
220    let reject = |reason: &str| {
221        Err(EngineError::InvalidWorkflow(format!(
222            "handler '{handler_name}' has invalid category '{category}': {reason}"
223        )))
224    };
225
226    if category.is_empty() {
227        return reject("empty category");
228    }
229    if category.starts_with('/') {
230        return reject("leading '/'");
231    }
232    if category.ends_with('/') {
233        return reject("trailing '/'");
234    }
235    for segment in category.split('/') {
236        if segment.is_empty() {
237            return reject("empty segment (double '/')");
238        }
239        if segment.trim().is_empty() {
240            return reject("whitespace-only segment");
241        }
242    }
243    Ok(())
244}
245
246/// Read a run id from the label `key` of a sub-workflow child run.
247///
248/// `None` when the run was not started by a `Workflow` step, when the label
249/// is missing or invalid, or when it points back at the run itself.
250fn chain_label(run: &Run, key: &str) -> Option<Uuid> {
251    if !matches!(run.trigger, TriggerKind::Workflow) {
252        return None;
253    }
254    let id = Uuid::parse_str(run.labels.get(key)?).ok()?;
255    (id != run.id).then_some(id)
256}
257
258/// The root run of the chain a sub-workflow child run belongs to.
259///
260/// Resuming a child resumes this root instead: the root replays its steps
261/// and re-enters the same child run through its open `Workflow` step.
262pub(crate) fn chain_root(run: &Run) -> Option<Uuid> {
263    chain_label(run, LABEL_ROOT_RUN_ID)
264}
265
266/// The run whose `Workflow` step started this sub-workflow child run.
267fn chain_parent(run: &Run) -> Option<Uuid> {
268    chain_label(run, PARENT_RUN_ID_LABEL)
269}
270
271impl Engine {
272    /// Create a new engine with the given store and agent provider.
273    ///
274    /// # Examples
275    ///
276    /// ```no_run
277    /// use std::sync::Arc;
278    /// use ironflow_engine::engine::Engine;
279    /// use ironflow_store::memory::InMemoryStore;
280    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
281    ///
282    /// let engine = Engine::new(
283    ///     Arc::new(InMemoryStore::new()),
284    ///     Arc::new(ClaudeCodeProvider::new()),
285    /// );
286    /// ```
287    pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
288        Self {
289            store,
290            provider,
291            handlers: HashMap::new(),
292            event_publisher: EventPublisher::new(),
293            log_sender: None,
294            budget: BudgetConfig::new(),
295            artifact_sink: None,
296            guard_config: None,
297            event_bus: None,
298            decision_provider: None,
299            step_interceptor: None,
300            execution_mode: ExecutionMode::default(),
301        }
302    }
303
304    /// Wire a [`DecisionProvider`] backend for `ctx.decision(...)` steps.
305    ///
306    /// Without this, a workflow that reaches a decision step fails with
307    /// [`EngineError::NoDecisionProvider`].
308    ///
309    /// # Examples
310    ///
311    /// ```no_run
312    /// use std::sync::Arc;
313    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
314    /// use ironflow_core::providers::record_replay_decision::RecordReplayDecisionProvider;
315    /// use ironflow_engine::engine::Engine;
316    /// use ironflow_store::memory::InMemoryStore;
317    ///
318    /// let engine = Engine::new(
319    ///     Arc::new(InMemoryStore::new()),
320    ///     Arc::new(ClaudeCodeProvider::new()),
321    /// )
322    /// .with_decision_provider(Arc::new(RecordReplayDecisionProvider::replay("tests/fixtures")));
323    /// ```
324    pub fn with_decision_provider(mut self, provider: Arc<dyn DecisionProvider>) -> Self {
325        self.decision_provider = Some(provider);
326        self
327    }
328
329    /// Wire a [`StepInterceptor`] that resolves steps without executing them.
330    ///
331    /// Used by [`crate::testing::TestEngine`] to mock shell, HTTP and approval
332    /// steps. `None` (the default) executes every step for real.
333    ///
334    /// # Examples
335    ///
336    /// ```no_run
337    /// use std::sync::Arc;
338    ///
339    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
340    /// use ironflow_engine::engine::Engine;
341    /// use ironflow_engine::executor::StepInterceptor;
342    /// use ironflow_store::memory::InMemoryStore;
343    ///
344    /// # fn example(interceptor: Arc<dyn StepInterceptor>) {
345    /// let engine = Engine::new(
346    ///     Arc::new(InMemoryStore::new()),
347    ///     Arc::new(ClaudeCodeProvider::new()),
348    /// )
349    /// .with_step_interceptor(interceptor);
350    /// # }
351    /// ```
352    pub fn with_step_interceptor(mut self, interceptor: Arc<dyn StepInterceptor>) -> Self {
353        self.step_interceptor = Some(interceptor);
354        self
355    }
356
357    /// The step interceptor wired into this engine, if any.
358    pub fn step_interceptor(&self) -> Option<&Arc<dyn StepInterceptor>> {
359        self.step_interceptor.as_ref()
360    }
361
362    /// Apply cost guardrails to this engine.
363    ///
364    /// Without this, both the per-run cap default and the monthly quota are
365    /// disabled and the engine behaves exactly as before.
366    ///
367    /// # Examples
368    ///
369    /// ```no_run
370    /// use std::sync::Arc;
371    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
372    /// use ironflow_engine::budget::BudgetConfig;
373    /// use ironflow_engine::engine::Engine;
374    /// use ironflow_store::memory::InMemoryStore;
375    ///
376    /// let engine = Engine::new(
377    ///     Arc::new(InMemoryStore::new()),
378    ///     Arc::new(ClaudeCodeProvider::new()),
379    /// )
380    /// .with_budget_config(BudgetConfig::from_env());
381    /// ```
382    pub fn with_budget_config(mut self, budget: BudgetConfig) -> Self {
383        self.budget = budget;
384        self
385    }
386
387    /// Returns the cost guardrails applied by this engine.
388    pub fn budget_config(&self) -> &BudgetConfig {
389        &self.budget
390    }
391
392    /// Apply workflow guard configuration to this engine.
393    ///
394    /// When set, every workflow run created by this engine is protected
395    /// by the guard. Handlers can override this via
396    /// [`WorkflowHandler::guard_config`].
397    ///
398    /// # Examples
399    ///
400    /// ```no_run
401    /// use std::sync::Arc;
402    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
403    /// use ironflow_engine::engine::Engine;
404    /// use ironflow_engine::guard::WorkflowGuardConfig;
405    /// use ironflow_store::memory::InMemoryStore;
406    ///
407    /// let engine = Engine::new(
408    ///     Arc::new(InMemoryStore::new()),
409    ///     Arc::new(ClaudeCodeProvider::new()),
410    /// )
411    /// .with_guard_config(WorkflowGuardConfig::default());
412    /// ```
413    pub fn with_guard_config(mut self, config: WorkflowGuardConfig) -> Self {
414        self.guard_config = Some(config);
415        self
416    }
417
418    /// Returns the workflow guard configuration, if any.
419    pub fn guard_config(&self) -> Option<&WorkflowGuardConfig> {
420        self.guard_config.as_ref()
421    }
422
423    /// Choose where a run resumes once its approval, human input or
424    /// escalation is resolved.
425    ///
426    /// Defaults to [`ExecutionMode::Local`]. Set [`ExecutionMode::Workers`]
427    /// on an API process that delegates execution to workers.
428    ///
429    /// # Examples
430    ///
431    /// ```no_run
432    /// use std::sync::Arc;
433    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
434    /// use ironflow_engine::engine::{Engine, ExecutionMode};
435    /// use ironflow_store::memory::InMemoryStore;
436    ///
437    /// let engine = Engine::new(
438    ///     Arc::new(InMemoryStore::new()),
439    ///     Arc::new(ClaudeCodeProvider::new()),
440    /// )
441    /// .with_execution_mode(ExecutionMode::Workers);
442    /// ```
443    pub fn with_execution_mode(mut self, mode: ExecutionMode) -> Self {
444        self.execution_mode = mode;
445        self
446    }
447
448    /// Returns where a run resumes once its gate is resolved.
449    pub fn execution_mode(&self) -> ExecutionMode {
450        self.execution_mode
451    }
452
453    /// Attach a log sender for real-time step output streaming.
454    ///
455    /// When set, all workflow contexts created by this engine will forward
456    /// step output (shell stdout/stderr, agent system messages) to the
457    /// given sender.
458    pub fn set_log_sender(&mut self, sender: LogSender) {
459        self.log_sender = Some(sender);
460    }
461
462    /// Attach the backend that stores and serves artifact bytes.
463    ///
464    /// Every context this engine builds inherits it. Without one, steps that
465    /// declare artifacts fail explicitly and every other step is unaffected.
466    ///
467    /// # Examples
468    ///
469    /// ```no_run
470    /// use std::sync::Arc;
471    ///
472    /// use ironflow_engine::artifact::ArtifactSink;
473    /// use ironflow_engine::engine::Engine;
474    ///
475    /// # fn example(engine: &mut Engine, sink: Arc<dyn ArtifactSink>) {
476    /// engine.set_artifact_sink(sink);
477    /// # }
478    /// ```
479    pub fn set_artifact_sink(&mut self, sink: Arc<dyn ArtifactSink>) {
480        self.artifact_sink = Some(sink);
481    }
482
483    /// The artifact backend attached to this engine, if any.
484    pub fn artifact_sink(&self) -> Option<&Arc<dyn ArtifactSink>> {
485        self.artifact_sink.as_ref()
486    }
487
488    /// Attach a [`WorkflowEventBus`] for per-run real-time monitoring.
489    ///
490    /// When set, all workflow contexts created by this engine will publish
491    /// step-level events (`StepStarted`, `StepCompleted`, `StepFailed`)
492    /// to the bus.
493    ///
494    /// # Examples
495    ///
496    /// ```no_run
497    /// use ironflow_engine::engine::Engine;
498    /// use ironflow_engine::notify::WorkflowEventBus;
499    ///
500    /// # fn example(engine: &mut Engine) {
501    /// engine.set_event_bus(WorkflowEventBus::new());
502    /// # }
503    /// ```
504    pub fn set_event_bus(&mut self, bus: WorkflowEventBus) {
505        self.event_bus = Some(bus);
506    }
507
508    /// The workflow event bus attached to this engine, if any.
509    pub fn event_bus(&self) -> Option<&WorkflowEventBus> {
510        self.event_bus.as_ref()
511    }
512
513    /// Returns a reference to the backing store.
514    pub fn store(&self) -> &Arc<dyn Store> {
515        &self.store
516    }
517
518    /// Returns a reference to the agent provider.
519    pub fn provider(&self) -> &Arc<dyn AgentProvider> {
520        &self.provider
521    }
522
523    /// Build a [`WorkflowContext`] with access to the handler registry.
524    ///
525    /// The context is seeded with the run's attempt number and the cost and
526    /// duration already accumulated by previous attempts, so a retried run
527    /// reports the total it really consumed rather than only its last attempt.
528    ///
529    /// The run's persisted `max_cost_usd` becomes the context cost cap; `None`
530    /// disables the per-run budget check for that context.
531    fn build_context(&self, run: &Run) -> WorkflowContext {
532        let handlers = self.handlers.clone();
533        let resolver: crate::context::HandlerResolver =
534            Arc::new(move |name: &str| handlers.get(name).cloned());
535        let mut ctx = WorkflowContext::with_handler_resolver(
536            run.id,
537            run.workflow_name.clone(),
538            self.store.clone(),
539            self.provider.clone(),
540            resolver,
541        );
542        ctx.carry_over_run_totals(run.retry_count + 1, run.cost_usd, run.duration_ms);
543        ctx.set_max_cost_usd(run.max_cost_usd);
544        ctx.set_run_created_at(run.created_at);
545        if let Some(ref sender) = self.log_sender {
546            ctx.set_log_sender(sender.clone());
547        }
548        if let Some(ref sink) = self.artifact_sink {
549            ctx.set_artifact_sink(sink.clone());
550        }
551        if let Some(ref bus) = self.event_bus {
552            ctx.set_event_bus(bus.clone());
553        }
554        if let Some(ref provider) = self.decision_provider {
555            ctx.set_decision_provider(provider.clone());
556        }
557        if let Some(ref interceptor) = self.step_interceptor {
558            ctx.set_step_interceptor(interceptor.clone());
559        }
560        ctx
561    }
562
563    /// Build a context with the guard attached.
564    ///
565    /// Uses the handler's `guard_config()` if present, otherwise falls back
566    /// to the engine's global configuration. Creates a fresh shared state
567    /// for each top-level run.
568    fn build_context_with_guard(
569        &self,
570        run: &Run,
571        handler: &dyn WorkflowHandler,
572    ) -> WorkflowContext {
573        let mut ctx = self.build_context(run);
574        let guard_config = handler.guard_config().or_else(|| self.guard_config.clone());
575        if let Some(config) = guard_config {
576            ctx.set_guard(config, new_shared_guard_state());
577        }
578        ctx
579    }
580
581    /// Reject the creation of a new run when the monthly quota is exhausted.
582    ///
583    /// The window is the current calendar month in UTC. Runs already in flight
584    /// are never interrupted -- only creation is refused.
585    ///
586    /// # Errors
587    ///
588    /// Returns [`EngineError::MonthlyBudgetExceeded`] when the accumulated cost
589    /// of the month has reached the configured quota. Returns
590    /// [`EngineError::Store`] when the aggregate query fails.
591    async fn check_monthly_quota(&self, workflow_name: &str) -> Result<(), EngineError> {
592        let Some(limit) = self.budget.monthly_cost_limit_usd else {
593            return Ok(());
594        };
595
596        let stats = self
597            .store
598            .get_stats(RunFilter {
599                created_after: Some(month_start(Utc::now())),
600                ..RunFilter::default()
601            })
602            .await?;
603
604        if stats.total_cost_usd < limit {
605            return Ok(());
606        }
607
608        warn!(
609            workflow = %workflow_name,
610            limit_usd = %limit,
611            spent_usd = %stats.total_cost_usd,
612            "monthly cost quota exhausted, refusing new run"
613        );
614
615        #[cfg(feature = "prometheus")]
616        counter!(
617            RUN_BUDGET_EXCEEDED_TOTAL,
618            "workflow" => workflow_name.to_string(),
619            "scope" => "monthly",
620        )
621        .increment(1);
622
623        Err(EngineError::MonthlyBudgetExceeded {
624            limit_usd: limit,
625            spent_usd: stats.total_cost_usd,
626        })
627    }
628
629    // -----------------------------------------------------------------------
630    // Handler registration
631    // -----------------------------------------------------------------------
632
633    /// Register a [`WorkflowHandler`] for dynamic workflow execution.
634    ///
635    /// The handler is looked up by [`WorkflowHandler::name`] when executing
636    /// or enqueuing.
637    ///
638    /// # Errors
639    ///
640    /// Returns [`EngineError::InvalidWorkflow`] if a handler with the same
641    /// name is already registered.
642    ///
643    /// # Examples
644    ///
645    /// ```no_run
646    /// use std::sync::Arc;
647    /// use ironflow_engine::engine::Engine;
648    /// use ironflow_engine::handler::{WorkflowHandler, HandlerFuture};
649    /// use ironflow_engine::context::WorkflowContext;
650    /// use ironflow_engine::config::ShellConfig;
651    /// use ironflow_store::memory::InMemoryStore;
652    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
653    ///
654    /// struct MyWorkflow;
655    /// impl WorkflowHandler for MyWorkflow {
656    ///     fn name(&self) -> &str { "my-workflow" }
657    ///     fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
658    ///         Box::pin(async move {
659    ///             ctx.shell("step1", ShellConfig::new("echo done")).await?;
660    ///             Ok(())
661    ///         })
662    ///     }
663    /// }
664    ///
665    /// let mut engine = Engine::new(
666    ///     Arc::new(InMemoryStore::new()),
667    ///     Arc::new(ClaudeCodeProvider::new()),
668    /// );
669    /// engine.register(MyWorkflow)?;
670    /// # Ok::<(), ironflow_engine::error::EngineError>(())
671    /// ```
672    pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
673        let name = handler.name().to_string();
674        if self.handlers.contains_key(&name) {
675            return Err(EngineError::InvalidWorkflow(format!(
676                "handler '{}' already registered",
677                name
678            )));
679        }
680        if let Some(category) = handler.category() {
681            validate_category(&name, category)?;
682        }
683        self.handlers.insert(name, Arc::new(handler));
684        Ok(())
685    }
686
687    /// Register a pre-boxed workflow handler.
688    ///
689    /// # Errors
690    ///
691    /// Returns [`EngineError::InvalidWorkflow`] if a handler with the same
692    /// name is already registered or if its category is invalid.
693    pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
694        let name = handler.name().to_string();
695        if self.handlers.contains_key(&name) {
696            return Err(EngineError::InvalidWorkflow(format!(
697                "handler '{}' already registered",
698                name
699            )));
700        }
701        if let Some(category) = handler.category() {
702            validate_category(&name, category)?;
703        }
704        self.handlers.insert(name, Arc::from(handler));
705        Ok(())
706    }
707
708    /// Get a registered handler by name.
709    pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
710        self.handlers.get(name)
711    }
712
713    /// List registered handler names.
714    pub fn handler_names(&self) -> Vec<&str> {
715        self.handlers.keys().map(|s| s.as_str()).collect()
716    }
717
718    /// Get detailed info about a registered workflow handler.
719    pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
720        self.handlers.get(name).map(|h| h.describe())
721    }
722
723    /// List handlers that have a cron schedule configured.
724    ///
725    /// Returns pairs of `(workflow_name, cron_expression)` for all handlers
726    /// where [`WorkflowHandler::schedule`] returns `Some`.
727    ///
728    /// Use this to wire scheduled handlers into a scheduler
729    /// (e.g. the schedule ticker in `ironflow-api`).
730    ///
731    /// # Examples
732    ///
733    /// ```no_run
734    /// # use std::sync::Arc;
735    /// # use ironflow_engine::engine::Engine;
736    /// # use ironflow_store::memory::InMemoryStore;
737    /// # use ironflow_core::providers::claude::ClaudeCodeProvider;
738    /// let engine = Engine::new(
739    ///     Arc::new(InMemoryStore::new()),
740    ///     Arc::new(ClaudeCodeProvider::new()),
741    /// );
742    /// for (name, schedule) in engine.scheduled_handlers() {
743    ///     tracing::info!("{name} runs on schedule: {schedule}");
744    /// }
745    /// ```
746    pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
747        self.handlers
748            .iter()
749            .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
750            .collect()
751    }
752
753    /// Register an event subscriber for domain events.
754    ///
755    /// The subscriber is called only for events whose type is in
756    /// `event_types`. Pass [`Event::ALL`] to receive every event.
757    ///
758    /// # Examples
759    ///
760    /// ```no_run
761    /// use ironflow_engine::engine::Engine;
762    /// use ironflow_engine::notify::{Event, WebhookSubscriber};
763    /// use ironflow_store::memory::InMemoryStore;
764    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
765    /// use std::sync::Arc;
766    ///
767    /// let mut engine = Engine::new(
768    ///     Arc::new(InMemoryStore::new()),
769    ///     Arc::new(ClaudeCodeProvider::new()),
770    /// );
771    ///
772    /// engine.subscribe(
773    ///     WebhookSubscriber::new("https://hooks.example.com/events"),
774    ///     &[Event::RUN_STATUS_CHANGED, Event::STEP_FAILED],
775    /// );
776    /// ```
777    pub fn subscribe(
778        &mut self,
779        subscriber: impl EventSubscriber + 'static,
780        event_types: &[&'static str],
781    ) {
782        self.event_publisher.subscribe(subscriber, event_types);
783    }
784
785    /// Returns a reference to the event publisher.
786    ///
787    /// Useful for publishing events from outside the engine (e.g. auth
788    /// routes in the API layer).
789    pub fn event_publisher(&self) -> &EventPublisher {
790        &self.event_publisher
791    }
792
793    // -----------------------------------------------------------------------
794    // Dynamic workflow execution (WorkflowHandler)
795    // -----------------------------------------------------------------------
796
797    /// Execute a registered handler inline.
798    ///
799    /// Creates a run, builds a [`WorkflowContext`], calls the handler's
800    /// [`execute`](WorkflowHandler::execute), and finalizes the run.
801    ///
802    /// # Errors
803    ///
804    /// Returns [`EngineError::InvalidWorkflow`] if no handler is registered
805    /// with that name. Returns [`EngineError`] if execution fails.
806    ///
807    /// # Examples
808    ///
809    /// ```no_run
810    /// use std::sync::Arc;
811    /// use ironflow_engine::engine::Engine;
812    /// use ironflow_store::memory::InMemoryStore;
813    /// use ironflow_store::models::TriggerKind;
814    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
815    /// use serde_json::json;
816    ///
817    /// # async fn example(engine: &Engine) -> Result<(), ironflow_engine::error::EngineError> {
818    /// let run = engine.run_handler("deploy", TriggerKind::Manual, json!({})).await?;
819    /// # Ok(())
820    /// # }
821    /// ```
822    #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
823    pub async fn run_handler(
824        &self,
825        handler_name: &str,
826        trigger: TriggerKind,
827        payload: Value,
828    ) -> Result<WorkflowResult, EngineError> {
829        let handler = self
830            .handlers
831            .get(handler_name)
832            .ok_or_else(|| {
833                EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
834            })?
835            .clone();
836
837        self.check_monthly_quota(handler_name).await?;
838
839        let handler_version = handler.version().map(str::to_string);
840        let max_cost_usd = self
841            .budget
842            .resolve_run_cap(None, handler.default_max_cost_usd());
843        let run = self
844            .store
845            .create_run(NewRun {
846                created_by: None,
847                workflow_name: handler_name.to_string(),
848                trigger,
849                payload,
850                max_retries: 0,
851                handler_version,
852                labels: handler.default_labels(),
853                scheduled_at: None,
854                idempotency_key: None,
855                concurrency_key: None,
856                max_cost_usd,
857            })
858            .await?
859            .into_run();
860
861        let run_id = run.id;
862        info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
863
864        self.store
865            .update_run_status(run_id, RunStatus::Running)
866            .await?;
867
868        #[cfg(feature = "prometheus")]
869        gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
870
871        let run_start = Instant::now();
872        let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
873
874        let result = handler.execute(&mut ctx).await;
875        self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
876            .await
877    }
878
879    /// Build the execution plan for a registered handler without running it.
880    ///
881    /// Executes the handler with every step method in recording mode: no
882    /// command is spawned, no HTTP request is sent, no agent is called,
883    /// nothing is persisted. Conditions declared with
884    /// [`WorkflowContext::when`](crate::context::WorkflowContext::when) are
885    /// evaluated against `payload`; those declared with
886    /// [`WorkflowContext::when_dynamic`](crate::context::WorkflowContext::when_dynamic)
887    /// are reported as unevaluable.
888    ///
889    /// Step outputs are synthetic and success-shaped, so the plan follows the
890    /// nominal branch. A handler that unwraps a decision answer, or that
891    /// deserializes `ctx.input::<T>()` against a payload it does not match,
892    /// aborts the plan: the partial plan is returned with
893    /// [`ExecutionPlan::incomplete_reason`] set.
894    ///
895    /// # Errors
896    ///
897    /// Returns [`EngineError::InvalidWorkflow`] when no handler is registered
898    /// under `handler_name` or when `options.max_depth` is zero. Returns
899    /// [`EngineError::Store`] when the duration-history query fails. A handler
900    /// that errors mid-plan does **not** fail this call.
901    ///
902    /// # Examples
903    ///
904    /// ```no_run
905    /// use ironflow_engine::engine::Engine;
906    /// use ironflow_engine::error::EngineError;
907    /// use ironflow_engine::plan::PlanOptions;
908    /// use serde_json::json;
909    ///
910    /// # async fn example(engine: &Engine) -> Result<(), EngineError> {
911    /// let plan = engine
912    ///     .plan_handler("deploy", json!({"env": "prod"}), PlanOptions::default())
913    ///     .await?;
914    /// for step in &plan.steps {
915    ///     println!("{} ({:?})", step.name, step.kind);
916    /// }
917    /// # Ok(())
918    /// # }
919    /// ```
920    #[tracing::instrument(name = "engine.plan_handler", skip_all, fields(workflow = %handler_name))]
921    pub async fn plan_handler(
922        &self,
923        handler_name: &str,
924        payload: Value,
925        options: PlanOptions,
926    ) -> Result<ExecutionPlan, EngineError> {
927        if options.max_depth == 0 {
928            return Err(EngineError::InvalidWorkflow(
929                "max_depth must be at least 1".to_string(),
930            ));
931        }
932
933        let handler = self
934            .handlers
935            .get(handler_name)
936            .ok_or_else(|| {
937                EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
938            })?
939            .clone();
940
941        let estimates = if options.estimate_durations {
942            estimate_durations(&self.store, handler_name, options.sample_runs).await?
943        } else {
944            HashMap::new()
945        };
946
947        let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
948            handler_name.to_string(),
949            payload,
950            options.max_depth,
951            estimates,
952        )));
953
954        // Deliberately bare: no guard, no event bus, no log sender, no artifact
955        // sink and no budget. Planning produces no side effect to report.
956        let handlers = self.handlers.clone();
957        let resolver: crate::context::HandlerResolver =
958            Arc::new(move |name: &str| handlers.get(name).cloned());
959        let mut ctx = WorkflowContext::with_handler_resolver(
960            Uuid::now_v7(),
961            handler_name.to_string(),
962            self.store.clone(),
963            self.provider.clone(),
964            resolver,
965        );
966        ctx.set_plan(shared.clone());
967
968        if let Err(err) = handler.execute(&mut ctx).await {
969            lock_plan(&shared).fail(err.to_string());
970        }
971        drop(ctx);
972
973        let plan = match Arc::try_unwrap(shared) {
974            Ok(mutex) => mutex
975                .into_inner()
976                .unwrap_or_else(|poisoned| poisoned.into_inner())
977                .into_plan(),
978            Err(shared) => lock_plan(&shared).snapshot(),
979        };
980
981        info!(
982            workflow = %handler_name,
983            steps = plan.steps.len(),
984            truncated = plan.truncated,
985            "execution plan built"
986        );
987
988        Ok(plan)
989    }
990
991    /// Enqueue a handler-based workflow for worker execution.
992    ///
993    /// The workflow name is stored in the run. The worker looks up the
994    /// handler by name when executing.
995    ///
996    /// # Errors
997    ///
998    /// Returns [`EngineError::InvalidWorkflow`] if no handler is registered.
999    /// Returns [`EngineError::MonthlyBudgetExceeded`] if the monthly cost quota
1000    /// is exhausted.
1001    #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
1002    pub async fn enqueue_handler(
1003        &self,
1004        handler_name: &str,
1005        trigger: TriggerKind,
1006        payload: Value,
1007        max_retries: u32,
1008    ) -> Result<Run, EngineError> {
1009        self.enqueue_handler_with_options(
1010            handler_name,
1011            trigger,
1012            payload,
1013            EnqueueOptions {
1014                max_retries,
1015                ..Default::default()
1016            },
1017        )
1018        .await
1019        .map(RunCreation::into_run)
1020    }
1021
1022    /// Enqueue a handler-based workflow with labels, deferred scheduling, an
1023    /// optional cost cap, an optional author, and an optional idempotency key.
1024    ///
1025    /// See [`EnqueueOptions`] for the individual settings.
1026    ///
1027    /// When [`EnqueueOptions::idempotency_key`] is set and already bound to a run
1028    /// created within
1029    /// [`IDEMPOTENCY_WINDOW`](ironflow_store::entities::IDEMPOTENCY_WINDOW), nothing
1030    /// is enqueued and the original run is returned as [`RunCreation::Existing`].
1031    ///
1032    /// # Errors
1033    ///
1034    /// Returns [`EngineError::InvalidWorkflow`] if no handler is registered.
1035    /// Returns [`EngineError::MonthlyBudgetExceeded`] if the monthly cost quota
1036    /// is exhausted. Returns [`EngineError::ConcurrencyConflict`] if
1037    /// [`EnqueueOptions::concurrency_key`] is held by another non-terminal run.
1038    /// Returns [`EngineError::Store`] if the run cannot be persisted.
1039    ///
1040    /// # Examples
1041    ///
1042    /// ```no_run
1043    /// use ironflow_engine::engine::{Engine, EnqueueOptions};
1044    /// use ironflow_store::models::TriggerKind;
1045    /// use serde_json::json;
1046    ///
1047    /// # async fn example(engine: &Engine) -> Result<(), ironflow_engine::error::EngineError> {
1048    /// let creation = engine
1049    ///     .enqueue_handler_with_options(
1050    ///         "deploy",
1051    ///         TriggerKind::Api,
1052    ///         json!({"env": "prod"}),
1053    ///         EnqueueOptions {
1054    ///             max_retries: 3,
1055    ///             idempotency_key: Some("github:abc-123".to_string()),
1056    ///             ..Default::default()
1057    ///         },
1058    ///     )
1059    ///     .await?;
1060    ///
1061    /// if creation.is_created() {
1062    ///     println!("enqueued {}", creation.run().id);
1063    /// }
1064    /// # Ok(())
1065    /// # }
1066    /// ```
1067    #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
1068    pub async fn enqueue_handler_with_options(
1069        &self,
1070        handler_name: &str,
1071        trigger: TriggerKind,
1072        payload: Value,
1073        options: EnqueueOptions,
1074    ) -> Result<RunCreation, EngineError> {
1075        let EnqueueOptions {
1076            max_retries,
1077            labels,
1078            scheduled_at,
1079            max_cost_usd,
1080            created_by,
1081            idempotency_key,
1082            concurrency_key,
1083        } = options;
1084
1085        let handler = self.handlers.get(handler_name).ok_or_else(|| {
1086            EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1087        })?;
1088
1089        self.check_monthly_quota(handler_name).await?;
1090
1091        let handler_version = handler.version().map(str::to_string);
1092        let mut merged_labels = handler.default_labels();
1093        merged_labels.extend(labels);
1094        let resolved_cap = self
1095            .budget
1096            .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
1097
1098        let creation = self
1099            .store
1100            .create_run(NewRun {
1101                workflow_name: handler_name.to_string(),
1102                trigger,
1103                payload,
1104                max_retries,
1105                handler_version,
1106                labels: merged_labels,
1107                scheduled_at,
1108                created_by,
1109                idempotency_key,
1110                concurrency_key,
1111                max_cost_usd: resolved_cap,
1112            })
1113            .await?;
1114
1115        match &creation {
1116            RunCreation::Created(run) => info!(
1117                run_id = %run.id,
1118                workflow = %handler_name,
1119                max_cost_usd = ?resolved_cap,
1120                "handler run enqueued"
1121            ),
1122            RunCreation::Existing(run) => info!(
1123                run_id = %run.id,
1124                workflow = %handler_name,
1125                "idempotent replay, nothing enqueued"
1126            ),
1127        }
1128
1129        Ok(creation)
1130    }
1131
1132    /// Execute a handler-based run (used by the worker after pick_next_pending).
1133    ///
1134    /// Looks up the handler by the run's `workflow_name` and executes it
1135    /// with a fresh [`WorkflowContext`], after
1136    /// [`AgentProvider::release_run`] has stopped whatever a previous
1137    /// execution of the run left running.
1138    ///
1139    /// # Errors
1140    ///
1141    /// Returns [`EngineError::InvalidWorkflow`] if no handler matches. A
1142    /// failed release fails the execution with [`EngineError::Operation`],
1143    /// replayed while the run has retries left. Returns
1144    /// [`EngineError::HandlerVersionMismatch`] when the handler's current
1145    /// version is incompatible with the run's `handler_version` -- checked
1146    /// before any step is replayed.
1147    ///
1148    /// A suspended child run of a sub-workflow is never executed on its own:
1149    /// its root run is resumed instead, and re-enters the child (see
1150    /// [`resume_run`](Self::resume_run)).
1151    #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
1152    pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1153        let run = self
1154            .store
1155            .get_run(run_id)
1156            .await?
1157            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1158
1159        if let Some(root_run_id) = chain_root(&run) {
1160            return self.resume_chain(run_id, root_run_id).await;
1161        }
1162
1163        let handler = self
1164            .handlers
1165            .get(&run.workflow_name)
1166            .ok_or_else(|| {
1167                EngineError::InvalidWorkflow(format!(
1168                    "no handler registered: {}",
1169                    run.workflow_name
1170                ))
1171            })?
1172            .clone();
1173
1174        #[cfg(feature = "prometheus")]
1175        gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
1176
1177        let run_start = Instant::now();
1178        let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1179
1180        // Replay the steps already persisted for this run: on a retry, an approval
1181        // a human already granted must not be asked again, and a run requeued to
1182        // `Pending` after its approval, human input or escalation resolved
1183        // (`ExecutionMode::Workers`, `retry_count` unchanged) must not re-run
1184        // completed steps. A brand-new run has no steps, so this is a no-op.
1185        //
1186        // The handler version is checked first: replaying an incompatible
1187        // handler's steps risks serving one step's cached output to another
1188        // (`EngineError::ReplayDivergence`), so no step is replayed at all
1189        // when the handler changed incompatibly since the run was created.
1190        let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1191            ctx.load_replay_steps().await?;
1192            self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1193                .await
1194        } else {
1195            Err(EngineError::HandlerVersionMismatch {
1196                run_id,
1197                workflow_name: run.workflow_name.clone(),
1198                run_version: run
1199                    .handler_version
1200                    .clone()
1201                    .unwrap_or_else(|| "unknown".to_string()),
1202                current_version: handler
1203                    .version()
1204                    .map(str::to_string)
1205                    .unwrap_or_else(|| "unknown".to_string()),
1206            })
1207        };
1208
1209        self.finalize_run(
1210            run_id,
1211            &run.workflow_name,
1212            result,
1213            &ctx,
1214            run_start,
1215            run.labels,
1216        )
1217        .await
1218    }
1219
1220    /// Execute a run by its ID (used by the worker after pick_next_pending).
1221    ///
1222    /// Delegates to [`execute_handler_run`](Self::execute_handler_run).
1223    ///
1224    /// # Errors
1225    ///
1226    /// Returns [`EngineError`] if the run is not found or execution fails.
1227    #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
1228    pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1229        self.execute_handler_run(run_id).await
1230    }
1231
1232    /// Resume a run after human approval.
1233    ///
1234    /// Re-executes the handler with step replay: completed steps return
1235    /// cached output, approved approval steps are skipped, and execution
1236    /// continues from the first unexecuted step.
1237    ///
1238    /// Supports multiple approval gates -- each resume replays all prior
1239    /// steps and stops at the next approval (or completes the run).
1240    ///
1241    /// When `run_id` is a child run of a sub-workflow, the root run of its
1242    /// chain is moved back to `Running` and resumed instead: it replays, and
1243    /// its open `Workflow` step re-enters the same child run. The returned
1244    /// result is the root run's.
1245    ///
1246    /// Like [`execute_handler_run`](Self::execute_handler_run), the handler
1247    /// only starts once [`AgentProvider::release_run`] has stopped whatever
1248    /// a previous execution of the run left running.
1249    ///
1250    /// # Errors
1251    ///
1252    /// Returns [`EngineError::InvalidWorkflow`] if no handler matches, or when
1253    /// the root run of a child cannot be resumed (it is running or finished:
1254    /// the child is then failed).
1255    /// Returns [`EngineError`] if execution fails or hits another approval.
1256    /// A failed release fails the execution with [`EngineError::Operation`].
1257    /// Returns [`EngineError::HandlerVersionMismatch`] when the handler's
1258    /// current version is incompatible with the run's `handler_version` --
1259    /// checked before any step is replayed.
1260    #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
1261    pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1262        let run = self
1263            .store
1264            .get_run(run_id)
1265            .await?
1266            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1267
1268        if let Some(root_run_id) = chain_root(&run) {
1269            return self.resume_chain(run_id, root_run_id).await;
1270        }
1271
1272        self.resume_loaded_run(run).await
1273    }
1274
1275    /// Resume the root run of a suspended child run.
1276    ///
1277    /// The root waits without `scheduled_at` while its child is suspended, so
1278    /// nothing but this path ever wakes it. It is moved to `Running` and
1279    /// resumed; its replay re-enters the child. A root that is not suspended
1280    /// (already running, or finished) cannot take the child back: the child
1281    /// is failed.
1282    async fn resume_chain(
1283        &self,
1284        child_run_id: Uuid,
1285        root_run_id: Uuid,
1286    ) -> Result<WorkflowResult, EngineError> {
1287        let root = self
1288            .store
1289            .get_run(root_run_id)
1290            .await?
1291            .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1292
1293        match root.status.state {
1294            RunStatus::AwaitingApproval | RunStatus::Pending => {
1295                self.store
1296                    .update_run_status(root_run_id, RunStatus::Running)
1297                    .await?;
1298            }
1299            RunStatus::Sleeping => {
1300                self.store
1301                    .update_run_status(root_run_id, RunStatus::Pending)
1302                    .await?;
1303                self.store
1304                    .update_run_status(root_run_id, RunStatus::Running)
1305                    .await?;
1306            }
1307            other => {
1308                let reason = format!(
1309                    "cannot resume child run {child_run_id}: root run {root_run_id} is {other}"
1310                );
1311                if let Err(err) = self
1312                    .fail_or_schedule_retry(child_run_id, &reason, false, None, None)
1313                    .await
1314                {
1315                    error!(
1316                        run_id = %child_run_id,
1317                        error = %err,
1318                        "failed to fail a child run whose root cannot resume"
1319                    );
1320                }
1321                return Err(EngineError::InvalidWorkflow(reason));
1322            }
1323        }
1324
1325        info!(
1326            run_id = %child_run_id,
1327            root_run_id = %root_run_id,
1328            "child run resumed through its root run"
1329        );
1330
1331        let root = self
1332            .store
1333            .get_run(root_run_id)
1334            .await?
1335            .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1336        self.resume_loaded_run(root).await
1337    }
1338
1339    /// Resume `run`, already loaded and already `Running`.
1340    async fn resume_loaded_run(&self, run: Run) -> Result<WorkflowResult, EngineError> {
1341        let run_id = run.id;
1342        let handler = self
1343            .handlers
1344            .get(&run.workflow_name)
1345            .ok_or_else(|| {
1346                EngineError::InvalidWorkflow(format!(
1347                    "no handler registered: {}",
1348                    run.workflow_name
1349                ))
1350            })?
1351            .clone();
1352
1353        info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
1354
1355        let run_start = Instant::now();
1356        let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1357
1358        let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1359            ctx.load_replay_steps().await?;
1360            self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1361                .await
1362        } else {
1363            Err(EngineError::HandlerVersionMismatch {
1364                run_id,
1365                workflow_name: run.workflow_name.clone(),
1366                run_version: run
1367                    .handler_version
1368                    .clone()
1369                    .unwrap_or_else(|| "unknown".to_string()),
1370                current_version: handler
1371                    .version()
1372                    .map(str::to_string)
1373                    .unwrap_or_else(|| "unknown".to_string()),
1374            })
1375        };
1376
1377        self.finalize_run(
1378            run_id,
1379            &run.workflow_name,
1380            result,
1381            &ctx,
1382            run_start,
1383            run.labels,
1384        )
1385        .await
1386    }
1387
1388    /// Store a signal and resolve every step waiting for its `(name, key)`.
1389    ///
1390    /// Each waiting step validates the payload against the JSON schema it
1391    /// stored when it opened. A matching step is completed with the payload
1392    /// and its run, if `Sleeping`, goes back to `Pending`: under
1393    /// [`ExecutionMode::Local`] it resumes in a background task, under
1394    /// [`ExecutionMode::Workers`] a worker picks it up. A step whose schema
1395    /// the payload does not match keeps waiting and is listed in
1396    /// [`SignalDelivery::rejected`].
1397    ///
1398    /// The signal is stored even when nobody waits for it: a run opening its
1399    /// wait step later still finds it. A signal whose `idempotency_id` was
1400    /// already used is not stored nor delivered again, and comes back with
1401    /// [`SignalDelivery::duplicate`] set.
1402    ///
1403    /// Publishes [`Event::SignalReceived`] for every stored signal.
1404    ///
1405    /// # Errors
1406    ///
1407    /// Returns [`EngineError::InvalidSignal`] when `name` or `key` is empty,
1408    /// and [`EngineError::Store`] when the signal cannot be stored or its
1409    /// waiters cannot be listed.
1410    ///
1411    /// # Examples
1412    ///
1413    /// ```no_run
1414    /// use std::sync::Arc;
1415    /// use ironflow_engine::engine::Engine;
1416    /// use ironflow_engine::error::EngineError;
1417    /// use ironflow_store::entities::NewSignal;
1418    /// use serde_json::json;
1419    ///
1420    /// # async fn example(engine: Arc<Engine>) -> Result<(), EngineError> {
1421    /// let delivery = engine
1422    ///     .deliver_signal(NewSignal {
1423    ///         name: "ci.pipeline_finished".to_string(),
1424    ///         key: "4f2a9c1".to_string(),
1425    ///         payload: json!({"status": "success"}),
1426    ///         idempotency_id: Some("delivery-42".to_string()),
1427    ///     })
1428    ///     .await?;
1429    /// println!("{} runs resumed", delivery.resumed.len());
1430    /// # Ok(())
1431    /// # }
1432    /// ```
1433    pub async fn deliver_signal(
1434        self: &Arc<Self>,
1435        signal: NewSignal,
1436    ) -> Result<SignalDelivery, EngineError> {
1437        if signal.name.trim().is_empty() {
1438            return Err(EngineError::InvalidSignal(
1439                "signal name must not be empty".to_string(),
1440            ));
1441        }
1442        if signal.key.trim().is_empty() {
1443            return Err(EngineError::InvalidSignal(
1444                "signal key must not be empty".to_string(),
1445            ));
1446        }
1447
1448        let stored = match self.store.insert_signal(signal).await? {
1449            SignalInsert::Created(stored) => stored,
1450            SignalInsert::Duplicate(existing) => {
1451                info!(
1452                    signal_id = %existing.id,
1453                    signal = %existing.name,
1454                    key = %existing.key,
1455                    "duplicate signal ignored"
1456                );
1457                return Ok(SignalDelivery {
1458                    signal_id: existing.id,
1459                    duplicate: true,
1460                    resumed: Vec::new(),
1461                    rejected: Vec::new(),
1462                });
1463            }
1464        };
1465
1466        let waiters = self
1467            .store
1468            .list_signal_waiters(&stored.name, &stored.key)
1469            .await?;
1470        let mut resumed = Vec::new();
1471        let mut rejected = Vec::new();
1472
1473        for step in waiters {
1474            if let Err(error) = validate_step_payload(step.input.as_ref(), &stored.payload) {
1475                rejected.push(SignalRejected {
1476                    run_id: step.run_id,
1477                    step_id: step.id,
1478                    error,
1479                });
1480                continue;
1481            }
1482
1483            match self
1484                .store
1485                .resolve_signal_step(step.id, received_output(&stored))
1486                .await
1487            {
1488                Ok(SignalStepResolution::Resolved {
1489                    run_id,
1490                    run_resumed,
1491                }) => {
1492                    resumed.push(SignalResumed {
1493                        run_id,
1494                        step_id: step.id,
1495                    });
1496                    if run_resumed && self.execution_mode == ExecutionMode::Local {
1497                        self.spawn_local_resume(run_id);
1498                    }
1499                }
1500                // A concurrent delivery or the timeout resolved it first.
1501                Ok(SignalStepResolution::NotWaiting { .. }) => {}
1502                Err(err) => {
1503                    error!(
1504                        run_id = %step.run_id,
1505                        step_id = %step.id,
1506                        error = %err,
1507                        "failed to resolve a waiting signal step"
1508                    );
1509                    rejected.push(SignalRejected {
1510                        run_id: step.run_id,
1511                        step_id: step.id,
1512                        error: err.to_string(),
1513                    });
1514                }
1515            }
1516        }
1517
1518        info!(
1519            signal_id = %stored.id,
1520            signal = %stored.name,
1521            key = %stored.key,
1522            resumed = resumed.len(),
1523            rejected = rejected.len(),
1524            "signal received"
1525        );
1526        self.event_publisher
1527            .publish(Event::SignalReceived(SignalReceivedEvent {
1528                signal_id: stored.id,
1529                name: stored.name.clone(),
1530                key: stored.key.clone(),
1531                resumed_runs: resumed.iter().map(|r| r.run_id).collect(),
1532                at: stored.received_at,
1533            }));
1534
1535        Ok(SignalDelivery {
1536            signal_id: stored.id,
1537            duplicate: false,
1538            resumed,
1539            rejected,
1540        })
1541    }
1542
1543    /// Send a typed signal: shorthand for [`deliver_signal`](Self::deliver_signal)
1544    /// with `S::NAME` as the name and `signal` as the payload.
1545    ///
1546    /// `key` identifies the occurrence (a commit SHA, an order ID).
1547    /// `idempotency_id`, when set, makes a redelivery of the same event (a
1548    /// webhook retried by its sender) a no-op.
1549    ///
1550    /// # Errors
1551    ///
1552    /// Returns [`EngineError::Serialization`] when `signal` cannot be
1553    /// serialized, and every error of [`deliver_signal`](Self::deliver_signal).
1554    ///
1555    /// # Examples
1556    ///
1557    /// ```no_run
1558    /// use std::sync::Arc;
1559    /// use ironflow_engine::engine::Engine;
1560    /// use ironflow_engine::error::EngineError;
1561    /// use ironflow_engine::signal::Signal;
1562    /// use schemars::JsonSchema;
1563    /// use serde::{Deserialize, Serialize};
1564    ///
1565    /// #[derive(Serialize, Deserialize, JsonSchema)]
1566    /// struct PipelineFinished {
1567    ///     status: String,
1568    /// }
1569    ///
1570    /// impl Signal for PipelineFinished {
1571    ///     const NAME: &'static str = "ci.pipeline_finished";
1572    /// }
1573    ///
1574    /// # async fn example(engine: Arc<Engine>) -> Result<(), EngineError> {
1575    /// let finished = PipelineFinished { status: "success".to_string() };
1576    /// engine.send_signal(&finished, "4f2a9c1", Some("delivery-42")).await?;
1577    /// # Ok(())
1578    /// # }
1579    /// ```
1580    pub async fn send_signal<S: Signal>(
1581        self: &Arc<Self>,
1582        signal: &S,
1583        key: &str,
1584        idempotency_id: Option<&str>,
1585    ) -> Result<SignalDelivery, EngineError> {
1586        let payload = to_value(signal)?;
1587        self.deliver_signal(NewSignal {
1588            name: S::NAME.to_string(),
1589            key: key.to_string(),
1590            payload,
1591            idempotency_id: idempotency_id.map(str::to_string),
1592        })
1593        .await
1594    }
1595
1596    /// Resume a run requeued to `Pending` in a background task.
1597    ///
1598    /// Used under [`ExecutionMode::Local`], where no worker would pick the
1599    /// run up. The state change already happened, so a failed resume is
1600    /// logged, not rolled back.
1601    pub(crate) fn spawn_local_resume(self: &Arc<Self>, run_id: Uuid) {
1602        let engine = Arc::clone(self);
1603        spawn(async move {
1604            if let Err(err) = engine
1605                .store
1606                .update_run_status(run_id, RunStatus::Running)
1607                .await
1608            {
1609                error!(run_id = %run_id, error = %err, "failed to restart a woken run");
1610                return;
1611            }
1612            if let Err(err) = engine.resume_run(run_id).await {
1613                error!(run_id = %run_id, error = %err, "failed to resume a woken run");
1614            }
1615        });
1616    }
1617
1618    /// Record a run failure, replaying the run later when retries remain.
1619    ///
1620    /// This is the single place where a failed run's fate is decided. When
1621    /// `retryable` is true and the run has not exhausted `max_retries`, the run
1622    /// moves to [`RunStatus::Retrying`] with `scheduled_at` set to
1623    /// `now + backoff` -- [`pick_next_pending`](ironflow_store::store::RunStore::pick_next_pending)
1624    /// picks it up again once that time has passed. Otherwise the run moves to
1625    /// [`RunStatus::Failed`].
1626    ///
1627    /// Either way, steps left non-terminal by the failed attempt are closed via
1628    /// [`fail_orphaned_steps`](Self::fail_orphaned_steps) so they are never
1629    /// confused with the next attempt's steps.
1630    ///
1631    /// Callers pass `retryable` explicitly rather than an error value, because
1632    /// the worker classifies failures it observes from the outside (a timeout, a
1633    /// panicked task) that never produce an [`EngineError`]. Use
1634    /// [`is_run_retryable`] to classify an
1635    /// [`EngineError`].
1636    ///
1637    /// `cost_usd` and `duration_ms` are `None` when the caller does not know the
1638    /// totals (a timeout or a panic observed from outside the handler); the
1639    /// values already stored on the run are then left untouched.
1640    ///
1641    /// Returns the status the run was moved to.
1642    ///
1643    /// # Errors
1644    ///
1645    /// Returns [`EngineError::Store`] if the run does not exist or the update
1646    /// cannot be persisted.
1647    ///
1648    /// # Examples
1649    ///
1650    /// ```no_run
1651    /// use ironflow_engine::engine::Engine;
1652    /// use ironflow_engine::error::EngineError;
1653    /// use uuid::Uuid;
1654    ///
1655    /// # async fn example(engine: &Engine, run_id: Uuid) -> Result<(), EngineError> {
1656    /// let status = engine
1657    ///     .fail_or_schedule_retry(run_id, "run timed out after 300s", true, None, None)
1658    ///     .await?;
1659    /// # Ok(())
1660    /// # }
1661    /// ```
1662    pub async fn fail_or_schedule_retry(
1663        &self,
1664        run_id: Uuid,
1665        error: &str,
1666        retryable: bool,
1667        cost_usd: Option<Decimal>,
1668        duration_ms: Option<u64>,
1669    ) -> Result<RunStatus, EngineError> {
1670        let run = self
1671            .store
1672            .get_run(run_id)
1673            .await?
1674            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1675
1676        let has_attempts_left = run.retry_count < run.max_retries;
1677        let update = if retryable && has_attempts_left {
1678            let backoff = backoff_for_retry(run.retry_count);
1679            let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1680
1681            info!(
1682                run_id = %run_id,
1683                workflow = %run.workflow_name,
1684                attempt = run.retry_count + 1,
1685                max_retries = run.max_retries,
1686                backoff_secs = backoff.as_secs(),
1687                scheduled_at = %scheduled_at,
1688                "run failed, scheduling retry"
1689            );
1690
1691            RunUpdate {
1692                status: Some(RunStatus::Retrying),
1693                error: Some(error.to_string()),
1694                increment_retry: true,
1695                cost_usd,
1696                duration_ms,
1697                scheduled_at: Some(scheduled_at),
1698                ..RunUpdate::default()
1699            }
1700        } else {
1701            RunUpdate {
1702                status: Some(RunStatus::Failed),
1703                error: Some(error.to_string()),
1704                cost_usd,
1705                duration_ms,
1706                completed_at: Some(Utc::now()),
1707                ..RunUpdate::default()
1708            }
1709        };
1710
1711        let status = update.status.unwrap_or(RunStatus::Failed);
1712        self.store.update_run(run_id, update).await?;
1713        self.fail_orphaned_steps(run_id, error).await?;
1714
1715        Ok(status)
1716    }
1717
1718    /// Fail all non-terminal steps for a run.
1719    ///
1720    /// Called after a run is marked as failed (timeout, error, panic) to clean up
1721    /// orphaned steps that are still in `Running`, `Pending`, or `AwaitingApproval`.
1722    ///
1723    /// - `Running` / `AwaitingApproval` steps are marked `Failed`.
1724    /// - `Pending` steps are marked `Skipped` (FSM does not allow Pending -> Failed).
1725    ///
1726    /// Errors from individual step updates are logged but do not abort the cleanup.
1727    ///
1728    /// # Errors
1729    ///
1730    /// Returns [`EngineError`] if listing steps fails.
1731    pub async fn fail_orphaned_steps(
1732        &self,
1733        run_id: Uuid,
1734        error_message: &str,
1735    ) -> Result<(), EngineError> {
1736        let steps = self.store.list_steps(run_id).await?;
1737        let now = Utc::now();
1738
1739        for step in steps {
1740            if step.status.state.is_terminal() {
1741                continue;
1742            }
1743
1744            let (target_status, error) = match step.status.state {
1745                StepStatus::Running | StepStatus::AwaitingApproval => {
1746                    let err = if step.error.is_some() {
1747                        None
1748                    } else {
1749                        Some(error_message.to_string())
1750                    };
1751                    (StepStatus::Failed, err)
1752                }
1753                StepStatus::Pending => (StepStatus::Skipped, None),
1754                _ => continue,
1755            };
1756
1757            if let Err(e) = self
1758                .store
1759                .update_step(
1760                    step.id,
1761                    StepUpdate {
1762                        status: Some(target_status),
1763                        error,
1764                        completed_at: Some(now),
1765                        ..StepUpdate::default()
1766                    },
1767                )
1768                .await
1769            {
1770                warn!(
1771                    run_id = %run_id,
1772                    step_id = %step.id,
1773                    step_name = %step.name,
1774                    error = %e,
1775                    "failed to cleanup orphaned step"
1776                );
1777            } else {
1778                info!(
1779                    run_id = %run_id,
1780                    step_id = %step.id,
1781                    step_name = %step.name,
1782                    from = %step.status.state,
1783                    to = %target_status,
1784                    "cleaned up orphaned step"
1785                );
1786            }
1787        }
1788
1789        Ok(())
1790    }
1791
1792    /// Execute `handler` once [`AgentProvider::release_run`] has stopped
1793    /// whatever a previous execution of the run left running (an agent pod
1794    /// writing to a shared worktree). A failed release fails the execution
1795    /// before its first step, with [`EngineError::Operation`].
1796    async fn release_then_execute(
1797        &self,
1798        run_id: Uuid,
1799        handler: &dyn WorkflowHandler,
1800        ctx: &mut WorkflowContext,
1801    ) -> Result<(), EngineError> {
1802        match self.provider.release_run(&run_id.to_string()).await {
1803            Ok(()) => handler.execute(ctx).await,
1804            Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
1805        }
1806    }
1807
1808    /// Finalize a run with the given result and context.
1809    ///
1810    /// On success: updates run to Completed with cost, duration, and completed_at.
1811    /// On failure: updates run to Failed with error, cost, duration, and completed_at.
1812    /// Always: fetches and returns the final Run.
1813    async fn finalize_run(
1814        &self,
1815        run_id: Uuid,
1816        workflow_name: &str,
1817        result: Result<(), EngineError>,
1818        ctx: &WorkflowContext,
1819        run_start: Instant,
1820        run_labels: HashMap<String, String>,
1821    ) -> Result<WorkflowResult, EngineError> {
1822        // Covers the whole run: previous attempts plus this one, so a retried
1823        // run reports the time it really consumed.
1824        let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
1825        let completed_at = Utc::now();
1826
1827        let final_status;
1828        let final_run;
1829
1830        match result {
1831            Ok(()) => {
1832                final_status = if ctx.has_allowed_failure() {
1833                    RunStatus::Warning
1834                } else {
1835                    RunStatus::Completed
1836                };
1837                final_run = self
1838                    .store
1839                    .update_run_returning(
1840                        run_id,
1841                        RunUpdate {
1842                            status: Some(final_status),
1843                            cost_usd: Some(ctx.total_cost_usd()),
1844                            duration_ms: Some(total_duration),
1845                            completed_at: Some(completed_at),
1846                            output: ctx.output().cloned(),
1847                            ..RunUpdate::default()
1848                        },
1849                    )
1850                    .await?;
1851
1852                info!(
1853                    run_id = %run_id,
1854                    status = %final_status,
1855                    cost_usd = %ctx.total_cost_usd(),
1856                    duration_ms = total_duration,
1857                    "run completed"
1858                );
1859            }
1860            Err(EngineError::ApprovalRequired {
1861                run_id: approval_run_id,
1862                step_id,
1863                ref message,
1864            }) => {
1865                final_status = RunStatus::AwaitingApproval;
1866                final_run = self
1867                    .store
1868                    .update_run_returning(
1869                        run_id,
1870                        RunUpdate {
1871                            status: Some(RunStatus::AwaitingApproval),
1872                            cost_usd: Some(ctx.total_cost_usd()),
1873                            duration_ms: Some(total_duration),
1874                            ..RunUpdate::default()
1875                        },
1876                    )
1877                    .await?;
1878
1879                info!(
1880                    run_id = %approval_run_id,
1881                    step_id = %step_id,
1882                    message = %message,
1883                    "run awaiting approval"
1884                );
1885
1886                self.publish_approval_requested(approval_run_id, step_id, message)
1887                    .await?;
1888            }
1889            Err(EngineError::ChildSuspended {
1890                run_id: child_run_id,
1891                ref cause,
1892            }) => {
1893                final_status = cause.suspension_status();
1894                // No `scheduled_at`: the suspended descendant owns the wake-up
1895                // and resumes this run through `resume_chain`. A wake-up armed
1896                // here too would resume the chain twice.
1897                final_run = self
1898                    .store
1899                    .update_run_returning(
1900                        run_id,
1901                        RunUpdate {
1902                            status: Some(final_status),
1903                            cost_usd: Some(ctx.total_cost_usd()),
1904                            duration_ms: Some(total_duration),
1905                            ..RunUpdate::default()
1906                        },
1907                    )
1908                    .await?;
1909
1910                let leaf = cause.suspension_leaf();
1911                info!(
1912                    run_id = %run_id,
1913                    child_run_id = %child_run_id,
1914                    status = %final_status,
1915                    cause = %leaf,
1916                    "run suspended with its child run"
1917                );
1918
1919                match leaf {
1920                    EngineError::ApprovalRequired {
1921                        run_id: approval_run_id,
1922                        step_id,
1923                        message,
1924                    } => {
1925                        self.publish_approval_requested(*approval_run_id, *step_id, message)
1926                            .await?;
1927                    }
1928                    EngineError::SignalWaiting {
1929                        run_id: wait_run_id,
1930                        step_id,
1931                        step_name,
1932                        name,
1933                        key,
1934                        deadline_at,
1935                    } => {
1936                        self.event_publisher
1937                            .publish(Event::SignalAwaited(SignalAwaitedEvent {
1938                                run_id: *wait_run_id,
1939                                step_id: *step_id,
1940                                step_name: step_name.clone(),
1941                                name: name.clone(),
1942                                key: key.clone(),
1943                                deadline_at: *deadline_at,
1944                                at: Utc::now(),
1945                            }));
1946                    }
1947                    // A human input or a delay publishes no suspension event,
1948                    // like on a top-level run.
1949                    _ => {}
1950                }
1951            }
1952            Err(EngineError::HumanInputRequired {
1953                run_id: input_run_id,
1954                step_id,
1955                ref message,
1956            }) => {
1957                final_status = RunStatus::AwaitingApproval;
1958                final_run = self
1959                    .store
1960                    .update_run_returning(
1961                        run_id,
1962                        RunUpdate {
1963                            status: Some(RunStatus::AwaitingApproval),
1964                            cost_usd: Some(ctx.total_cost_usd()),
1965                            duration_ms: Some(total_duration),
1966                            ..RunUpdate::default()
1967                        },
1968                    )
1969                    .await?;
1970
1971                // No `ApprovalRequested` event: a human input is not an approval.
1972                info!(
1973                    run_id = %input_run_id,
1974                    step_id = %step_id,
1975                    message = %message,
1976                    "run awaiting human input"
1977                );
1978            }
1979            Err(EngineError::DelaySleeping {
1980                run_id: delay_run_id,
1981                step_id,
1982                wake_at,
1983            }) => {
1984                final_status = RunStatus::Sleeping;
1985                final_run = self
1986                    .store
1987                    .update_run_returning(
1988                        run_id,
1989                        RunUpdate {
1990                            status: Some(RunStatus::Sleeping),
1991                            cost_usd: Some(ctx.total_cost_usd()),
1992                            duration_ms: Some(total_duration),
1993                            scheduled_at: Some(wake_at),
1994                            ..RunUpdate::default()
1995                        },
1996                    )
1997                    .await?;
1998
1999                info!(
2000                    run_id = %delay_run_id,
2001                    step_id = %step_id,
2002                    wake_at = %wake_at,
2003                    "run sleeping until delay elapses"
2004                );
2005            }
2006            Err(EngineError::SignalWaiting {
2007                run_id: wait_run_id,
2008                step_id,
2009                ref step_name,
2010                ref name,
2011                ref key,
2012                deadline_at,
2013            }) => {
2014                final_status = RunStatus::Sleeping;
2015                // Atomic with the step lock: a signal delivered since the step
2016                // opened leaves the run due right away instead of until the
2017                // deadline.
2018                let waiting = self
2019                    .store
2020                    .suspend_run_on_signal(run_id, step_id, deadline_at)
2021                    .await?;
2022                final_run = self
2023                    .store
2024                    .update_run_returning(
2025                        run_id,
2026                        RunUpdate {
2027                            cost_usd: Some(ctx.total_cost_usd()),
2028                            duration_ms: Some(total_duration),
2029                            ..RunUpdate::default()
2030                        },
2031                    )
2032                    .await?;
2033
2034                if waiting {
2035                    self.event_publisher
2036                        .publish(Event::SignalAwaited(SignalAwaitedEvent {
2037                            run_id: wait_run_id,
2038                            step_id,
2039                            step_name: step_name.clone(),
2040                            name: name.clone(),
2041                            key: key.clone(),
2042                            deadline_at,
2043                            at: Utc::now(),
2044                        }));
2045                }
2046
2047                info!(
2048                    run_id = %wait_run_id,
2049                    step_id = %step_id,
2050                    signal = %name,
2051                    key = %key,
2052                    deadline_at = %deadline_at,
2053                    waiting,
2054                    "run sleeping until a signal arrives"
2055                );
2056            }
2057            Err(err) => {
2058                // A guardrail stop (budget or workflow guard) is deliberate,
2059                // not a breakage: the run is cancelled, never failed and
2060                // never replayed.
2061                let guardrail_stop = matches!(
2062                    err,
2063                    EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2064                );
2065
2066                final_status = if guardrail_stop {
2067                    if let Err(store_err) = self
2068                        .store
2069                        .update_run(
2070                            run_id,
2071                            RunUpdate {
2072                                status: Some(RunStatus::Cancelled),
2073                                error: Some(err.to_string()),
2074                                cost_usd: Some(ctx.total_cost_usd()),
2075                                duration_ms: Some(total_duration),
2076                                completed_at: Some(completed_at),
2077                                output: ctx.output().cloned(),
2078                                ..RunUpdate::default()
2079                            },
2080                        )
2081                        .await
2082                    {
2083                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2084                    }
2085                    if let Err(cleanup_err) = self
2086                        .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2087                        .await
2088                    {
2089                        error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2090                    }
2091                    RunStatus::Cancelled
2092                } else {
2093                    // Written before the failure so a failed run keeps the
2094                    // verdict the handler set before returning its error.
2095                    if let Some(output) = ctx.output()
2096                        && let Err(store_err) = self
2097                            .store
2098                            .update_run(
2099                                run_id,
2100                                RunUpdate {
2101                                    output: Some(output.clone()),
2102                                    ..RunUpdate::default()
2103                                },
2104                            )
2105                            .await
2106                    {
2107                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run output");
2108                    }
2109                    self.fail_or_schedule_retry(
2110                        run_id,
2111                        &err.to_string(),
2112                        is_run_retryable(&err),
2113                        Some(ctx.total_cost_usd()),
2114                        Some(total_duration),
2115                    )
2116                    .await
2117                    .unwrap_or_else(|store_err| {
2118                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2119                        RunStatus::Failed
2120                    })
2121                };
2122
2123                if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2124                    self.on_run_budget_exceeded(workflow_name, run_id, &err);
2125                }
2126
2127                error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2128
2129                self.publish_run_status_changed(
2130                    workflow_name,
2131                    run_id,
2132                    final_status,
2133                    Some(err.to_string()),
2134                    ctx,
2135                    total_duration,
2136                    run_labels,
2137                );
2138
2139                #[cfg(feature = "prometheus")]
2140                self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2141
2142                return Err(err);
2143            }
2144        }
2145
2146        self.publish_run_status_changed(
2147            workflow_name,
2148            run_id,
2149            final_status,
2150            None,
2151            ctx,
2152            total_duration,
2153            run_labels,
2154        );
2155
2156        #[cfg(feature = "prometheus")]
2157        self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2158
2159        Ok(WorkflowResult {
2160            run: final_run,
2161            steps: ctx.step_results().to_vec(),
2162        })
2163    }
2164
2165    /// Publish [`Event::ApprovalRequested`] for an approval gate that opened
2166    /// on `run_id`, with the requirement recorded when the gate opened.
2167    async fn publish_approval_requested(
2168        &self,
2169        run_id: Uuid,
2170        step_id: Uuid,
2171        message: &str,
2172    ) -> Result<(), EngineError> {
2173        let requirement = self
2174            .store
2175            .get_step(step_id)
2176            .await?
2177            .and_then(|s| s.approval_requirement);
2178        self.event_publisher
2179            .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2180                run_id,
2181                step_id,
2182                message: message.to_string(),
2183                requirement,
2184                at: Utc::now(),
2185            }));
2186        Ok(())
2187    }
2188
2189    /// Fail every ancestor of a child run, closest first.
2190    ///
2191    /// Used when a gate inside a sub-workflow is rejected: the child run
2192    /// fails, and the runs suspended with it (its parent, up to the root)
2193    /// must not stay suspended on a child that will never resume. Each
2194    /// ancestor goes through [`fail_or_schedule_retry`](Self::fail_or_schedule_retry)
2195    /// without a retry, which also fails its open `Workflow` step. A run that
2196    /// is not a child of a sub-workflow has no ancestor: nothing happens.
2197    ///
2198    /// # Errors
2199    ///
2200    /// Returns [`EngineError::Store`] if a run of the chain does not exist or
2201    /// its failure cannot be persisted.
2202    ///
2203    /// # Examples
2204    ///
2205    /// ```no_run
2206    /// use ironflow_engine::engine::Engine;
2207    /// use ironflow_engine::error::EngineError;
2208    /// use uuid::Uuid;
2209    ///
2210    /// # async fn example(engine: &Engine, child_run_id: Uuid) -> Result<(), EngineError> {
2211    /// engine
2212    ///     .fail_ancestors(child_run_id, "approval rejected in a sub-workflow")
2213    ///     .await?;
2214    /// # Ok(())
2215    /// # }
2216    /// ```
2217    pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2218        let mut current = self
2219            .store
2220            .get_run(run_id)
2221            .await?
2222            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2223        // Labels are data: a chain that loops back on itself stops there.
2224        let mut visited = HashSet::from([run_id]);
2225
2226        while let Some(parent_id) = chain_parent(&current) {
2227            if !visited.insert(parent_id) {
2228                break;
2229            }
2230            let status = self
2231                .fail_or_schedule_retry(parent_id, reason, false, None, None)
2232                .await?;
2233            info!(
2234                run_id = %run_id,
2235                ancestor_run_id = %parent_id,
2236                status = %status,
2237                "ancestor run failed with its child"
2238            );
2239            current = self
2240                .store
2241                .get_run(parent_id)
2242                .await?
2243                .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2244        }
2245
2246        Ok(())
2247    }
2248
2249    /// Emit Prometheus metrics for a completed run.
2250    #[cfg(feature = "prometheus")]
2251    fn emit_run_metrics(
2252        &self,
2253        workflow_name: &str,
2254        status: RunStatus,
2255        duration_ms: u64,
2256        ctx: &WorkflowContext,
2257    ) {
2258        let status_str = status.to_string();
2259        let wf = workflow_name.to_string();
2260
2261        counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2262        histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2263            .record(duration_ms as f64 / 1000.0);
2264        histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2265            ctx.total_cost_usd()
2266                .to_string()
2267                .parse::<f64>()
2268                .unwrap_or(0.0),
2269        );
2270        gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2271    }
2272
2273    /// Record the metric and publish the audit event for a run that hit its
2274    /// cost cap.
2275    ///
2276    /// A non-[`RunBudgetExceeded`](EngineError::RunBudgetExceeded) error is
2277    /// ignored, so callers can pass the error unconditionally.
2278    fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2279        let EngineError::RunBudgetExceeded {
2280            limit_usd,
2281            spent_usd,
2282            step_budget_usd,
2283            ..
2284        } = err
2285        else {
2286            return;
2287        };
2288
2289        #[cfg(feature = "prometheus")]
2290        counter!(
2291            RUN_BUDGET_EXCEEDED_TOTAL,
2292            "workflow" => workflow_name.to_string(),
2293            "scope" => "run",
2294        )
2295        .increment(1);
2296
2297        self.event_publisher
2298            .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2299                run_id,
2300                workflow_name: workflow_name.to_string(),
2301                limit_usd: *limit_usd,
2302                spent_usd: *spent_usd,
2303                step_budget_usd: *step_budget_usd,
2304                at: Utc::now(),
2305            }));
2306    }
2307
2308    /// Publish a run status changed event to all registered subscribers.
2309    ///
2310    /// `from` is always `Running` because `finalize_run` is only called
2311    /// from a running state.
2312    #[allow(clippy::too_many_arguments)]
2313    fn publish_run_status_changed(
2314        &self,
2315        workflow_name: &str,
2316        run_id: Uuid,
2317        to: RunStatus,
2318        error: Option<String>,
2319        ctx: &WorkflowContext,
2320        duration_ms: u64,
2321        labels: HashMap<String, String>,
2322    ) {
2323        let now = Utc::now();
2324        let cost_usd = ctx.total_cost_usd();
2325        let wf = workflow_name.to_string();
2326
2327        self.event_publisher
2328            .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2329                run_id,
2330                workflow_name: wf.clone(),
2331                from: RunStatus::Running,
2332                to,
2333                error: error.clone(),
2334                cost_usd,
2335                duration_ms,
2336                labels: labels.clone(),
2337                at: now,
2338            }));
2339
2340        if to == RunStatus::Failed {
2341            self.event_publisher
2342                .publish(Event::RunFailed(RunFailedEvent {
2343                    run_id,
2344                    workflow_name: wf,
2345                    error,
2346                    cost_usd,
2347                    duration_ms,
2348                    labels,
2349                    at: now,
2350                }));
2351        }
2352    }
2353}
2354
2355impl fmt::Debug for Engine {
2356    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2357        f.debug_struct("Engine")
2358            .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2359            .finish_non_exhaustive()
2360    }
2361}
2362
2363#[cfg(test)]
2364mod tests {
2365    use super::*;
2366    use crate::config::ShellConfig;
2367    use crate::handler::{HandlerFuture, WorkflowHandler};
2368    use ironflow_core::providers::claude::ClaudeCodeProvider;
2369    use ironflow_core::providers::record_replay::RecordReplayProvider;
2370    use ironflow_store::memory::InMemoryStore;
2371    use ironflow_store::models::StepStatus;
2372    use serde_json::json;
2373
2374    // Test handler that echoes a message via shell
2375    struct EchoWorkflow;
2376
2377    impl WorkflowHandler for EchoWorkflow {
2378        fn name(&self) -> &str {
2379            "echo-workflow"
2380        }
2381
2382        fn describe(&self) -> WorkflowInfo {
2383            WorkflowInfo {
2384                description: "A simple workflow that echoes hello".to_string(),
2385                source_code: None,
2386                sub_workflows: Vec::new(),
2387                category: None,
2388                version: self.version().map(str::to_string),
2389                compatible_versions: Vec::new(),
2390                input_schema: None,
2391                default_labels: HashMap::new(),
2392                schedule: self.schedule().cloned(),
2393                default_max_cost_usd: self.default_max_cost_usd(),
2394            }
2395        }
2396
2397        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2398            Box::pin(async move {
2399                ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2400                Ok(())
2401            })
2402        }
2403    }
2404
2405    // Test handler that fails
2406    struct FailingWorkflow;
2407
2408    impl WorkflowHandler for FailingWorkflow {
2409        fn name(&self) -> &str {
2410            "failing-workflow"
2411        }
2412
2413        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2414            Box::pin(async move {
2415                ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2416                Ok(())
2417            })
2418        }
2419    }
2420
2421    fn create_test_engine() -> Engine {
2422        let store = Arc::new(InMemoryStore::new());
2423        let inner = ClaudeCodeProvider::new();
2424        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2425            inner,
2426            "/tmp/ironflow-fixtures",
2427        ));
2428        Engine::new(store, provider)
2429    }
2430
2431    #[test]
2432    fn engine_new_creates_instance() {
2433        let engine = create_test_engine();
2434        assert_eq!(engine.handler_names().len(), 0);
2435    }
2436
2437    #[test]
2438    fn execution_mode_defaults_to_local() {
2439        let engine = create_test_engine();
2440        assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2441    }
2442
2443    #[test]
2444    fn with_execution_mode_overrides_the_default() {
2445        let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2446        assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2447    }
2448
2449    #[test]
2450    fn engine_register_handler() {
2451        let mut engine = create_test_engine();
2452        let result = engine.register(EchoWorkflow);
2453        assert!(result.is_ok());
2454        assert_eq!(engine.handler_names().len(), 1);
2455        assert!(engine.handler_names().contains(&"echo-workflow"));
2456    }
2457
2458    #[test]
2459    fn engine_register_duplicate_returns_error() {
2460        let mut engine = create_test_engine();
2461        engine.register(EchoWorkflow).unwrap();
2462        let result = engine.register(EchoWorkflow);
2463        assert!(result.is_err());
2464    }
2465
2466    #[test]
2467    fn engine_get_handler_found() {
2468        let mut engine = create_test_engine();
2469        engine.register(EchoWorkflow).unwrap();
2470        let handler = engine.get_handler("echo-workflow");
2471        assert!(handler.is_some());
2472    }
2473
2474    #[test]
2475    fn engine_get_handler_not_found() {
2476        let engine = create_test_engine();
2477        let handler = engine.get_handler("nonexistent");
2478        assert!(handler.is_none());
2479    }
2480
2481    #[test]
2482    fn engine_handler_names_lists_all() {
2483        let mut engine = create_test_engine();
2484        engine.register(EchoWorkflow).unwrap();
2485        engine.register(FailingWorkflow).unwrap();
2486        let names = engine.handler_names();
2487        assert_eq!(names.len(), 2);
2488        assert!(names.contains(&"echo-workflow"));
2489        assert!(names.contains(&"failing-workflow"));
2490    }
2491
2492    #[test]
2493    fn engine_handler_info_returns_description() {
2494        let mut engine = create_test_engine();
2495        engine.register(EchoWorkflow).unwrap();
2496        let info = engine.handler_info("echo-workflow");
2497        assert!(info.is_some());
2498        let info = info.unwrap();
2499        assert_eq!(info.description, "A simple workflow that echoes hello");
2500    }
2501
2502    struct CategorizedWorkflow;
2503
2504    impl WorkflowHandler for CategorizedWorkflow {
2505        fn name(&self) -> &str {
2506            "categorized"
2507        }
2508        fn category(&self) -> Option<&str> {
2509            Some("data/etl")
2510        }
2511        fn execute<'a>(
2512            &'a self,
2513            _ctx: &'a mut WorkflowContext,
2514        ) -> crate::handler::HandlerFuture<'a> {
2515            Box::pin(async move { Ok(()) })
2516        }
2517    }
2518
2519    #[test]
2520    fn engine_default_describe_propagates_category() {
2521        let mut engine = create_test_engine();
2522        engine.register(CategorizedWorkflow).unwrap();
2523        let info = engine.handler_info("categorized").unwrap();
2524        assert_eq!(info.category.as_deref(), Some("data/etl"));
2525    }
2526
2527    #[test]
2528    fn engine_default_describe_without_category() {
2529        let mut engine = create_test_engine();
2530        engine.register(EchoWorkflow).unwrap();
2531        let info = engine.handler_info("echo-workflow").unwrap();
2532        assert!(info.category.is_none());
2533    }
2534
2535    // -----------------------------------------------------------------------
2536    // Schedule tests
2537    // -----------------------------------------------------------------------
2538
2539    struct ScheduledWorkflow {
2540        schedule: CronSchedule,
2541    }
2542
2543    impl ScheduledWorkflow {
2544        fn new() -> Self {
2545            Self {
2546                schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2547            }
2548        }
2549    }
2550
2551    impl WorkflowHandler for ScheduledWorkflow {
2552        fn name(&self) -> &str {
2553            "scheduled"
2554        }
2555        fn schedule(&self) -> Option<&CronSchedule> {
2556            Some(&self.schedule)
2557        }
2558        fn execute<'a>(
2559            &'a self,
2560            _ctx: &'a mut WorkflowContext,
2561        ) -> crate::handler::HandlerFuture<'a> {
2562            Box::pin(async move { Ok(()) })
2563        }
2564    }
2565
2566    #[test]
2567    fn engine_default_describe_propagates_schedule() {
2568        let mut engine = create_test_engine();
2569        engine.register(ScheduledWorkflow::new()).unwrap();
2570        let info = engine.handler_info("scheduled").unwrap();
2571        assert_eq!(
2572            info.schedule.as_ref().map(|s| s.as_str()),
2573            Some("0 0 * * * *")
2574        );
2575    }
2576
2577    #[test]
2578    fn engine_default_describe_without_schedule() {
2579        let mut engine = create_test_engine();
2580        engine.register(EchoWorkflow).unwrap();
2581        let info = engine.handler_info("echo-workflow").unwrap();
2582        assert!(info.schedule.is_none());
2583    }
2584
2585    #[test]
2586    fn scheduled_handlers_returns_only_scheduled() {
2587        let mut engine = create_test_engine();
2588        engine.register(EchoWorkflow).unwrap();
2589        engine.register(ScheduledWorkflow::new()).unwrap();
2590        engine.register(FailingWorkflow).unwrap();
2591
2592        let scheduled = engine.scheduled_handlers();
2593        assert_eq!(scheduled.len(), 1);
2594        assert_eq!(scheduled[0].0, "scheduled");
2595        assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
2596    }
2597
2598    #[test]
2599    fn scheduled_handlers_empty_when_none_scheduled() {
2600        let mut engine = create_test_engine();
2601        engine.register(EchoWorkflow).unwrap();
2602        engine.register(FailingWorkflow).unwrap();
2603
2604        let scheduled = engine.scheduled_handlers();
2605        assert!(scheduled.is_empty());
2606    }
2607
2608    struct BadCategoryWorkflow(&'static str);
2609
2610    impl WorkflowHandler for BadCategoryWorkflow {
2611        fn name(&self) -> &str {
2612            "bad-category"
2613        }
2614        fn category(&self) -> Option<&str> {
2615            Some(self.0)
2616        }
2617        fn execute<'a>(
2618            &'a self,
2619            _ctx: &'a mut WorkflowContext,
2620        ) -> crate::handler::HandlerFuture<'a> {
2621            Box::pin(async move { Ok(()) })
2622        }
2623    }
2624
2625    #[test]
2626    fn engine_register_rejects_empty_category() {
2627        let mut engine = create_test_engine();
2628        let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2629        match err {
2630            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2631            other => panic!("expected InvalidWorkflow, got {other:?}"),
2632        }
2633    }
2634
2635    #[test]
2636    fn engine_register_rejects_leading_slash_category() {
2637        let mut engine = create_test_engine();
2638        let err = engine
2639            .register(BadCategoryWorkflow("/data/etl"))
2640            .unwrap_err();
2641        match err {
2642            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2643            other => panic!("expected InvalidWorkflow, got {other:?}"),
2644        }
2645    }
2646
2647    #[test]
2648    fn engine_register_rejects_trailing_slash_category() {
2649        let mut engine = create_test_engine();
2650        let err = engine
2651            .register(BadCategoryWorkflow("data/etl/"))
2652            .unwrap_err();
2653        match err {
2654            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2655            other => panic!("expected InvalidWorkflow, got {other:?}"),
2656        }
2657    }
2658
2659    #[test]
2660    fn engine_register_rejects_double_slash_category() {
2661        let mut engine = create_test_engine();
2662        let err = engine
2663            .register(BadCategoryWorkflow("data//etl"))
2664            .unwrap_err();
2665        match err {
2666            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2667            other => panic!("expected InvalidWorkflow, got {other:?}"),
2668        }
2669    }
2670
2671    #[test]
2672    fn engine_register_rejects_whitespace_only_segment_category() {
2673        let mut engine = create_test_engine();
2674        let err = engine
2675            .register(BadCategoryWorkflow("data/ /etl"))
2676            .unwrap_err();
2677        match err {
2678            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2679            other => panic!("expected InvalidWorkflow, got {other:?}"),
2680        }
2681    }
2682
2683    #[test]
2684    fn engine_register_accepts_valid_nested_category() {
2685        let mut engine = create_test_engine();
2686        assert!(engine.register(CategorizedWorkflow).is_ok());
2687    }
2688
2689    #[tokio::test]
2690    async fn engine_unknown_workflow_returns_error() {
2691        let engine = create_test_engine();
2692        let result = engine
2693            .run_handler("unknown", TriggerKind::Manual, json!({}))
2694            .await;
2695        assert!(result.is_err());
2696        match result {
2697            Err(EngineError::InvalidWorkflow(msg)) => {
2698                assert!(msg.contains("no handler registered"));
2699            }
2700            _ => panic!("expected InvalidWorkflow error"),
2701        }
2702    }
2703
2704    #[tokio::test]
2705    async fn engine_enqueue_handler_creates_pending_run() {
2706        let mut engine = create_test_engine();
2707        engine.register(EchoWorkflow).unwrap();
2708
2709        let run = engine
2710            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2711            .await
2712            .unwrap();
2713        assert_eq!(run.status.state, RunStatus::Pending);
2714        assert_eq!(run.workflow_name, "echo-workflow");
2715    }
2716
2717    #[tokio::test]
2718    async fn enqueue_handler_leaves_the_run_unattributed() {
2719        let mut engine = create_test_engine();
2720        engine.register(EchoWorkflow).unwrap();
2721
2722        let run = engine
2723            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2724            .await
2725            .unwrap();
2726
2727        assert!(run.created_by.is_none());
2728    }
2729
2730    #[tokio::test]
2731    async fn enqueue_handler_with_options_records_the_author() {
2732        let mut engine = create_test_engine();
2733        engine.register(EchoWorkflow).unwrap();
2734        let actor = RunActor::User {
2735            user_id: Uuid::now_v7(),
2736        };
2737
2738        let run = engine
2739            .enqueue_handler_with_options(
2740                "echo-workflow",
2741                TriggerKind::Api,
2742                json!({}),
2743                EnqueueOptions {
2744                    created_by: Some(actor.clone()),
2745                    ..Default::default()
2746                },
2747            )
2748            .await
2749            .unwrap()
2750            .into_run();
2751
2752        assert_eq!(run.created_by, Some(actor));
2753    }
2754
2755    #[tokio::test]
2756    async fn enqueue_handler_with_options_accepts_no_author() {
2757        let mut engine = create_test_engine();
2758        engine.register(EchoWorkflow).unwrap();
2759
2760        let run = engine
2761            .enqueue_handler_with_options(
2762                "echo-workflow",
2763                TriggerKind::Cron {
2764                    schedule: "0 * * * * *".to_string(),
2765                },
2766                json!({}),
2767                EnqueueOptions::default(),
2768            )
2769            .await
2770            .unwrap()
2771            .into_run();
2772
2773        assert!(run.created_by.is_none());
2774    }
2775
2776    #[tokio::test]
2777    async fn run_handler_leaves_the_run_unattributed() {
2778        let mut engine = create_test_engine();
2779        engine.register(EchoWorkflow).unwrap();
2780
2781        let run = engine
2782            .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
2783            .await
2784            .unwrap()
2785            .run;
2786
2787        assert!(run.created_by.is_none());
2788    }
2789
2790    #[tokio::test]
2791    async fn engine_register_boxed() {
2792        let mut engine = create_test_engine();
2793        let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
2794        let result = engine.register_boxed(handler);
2795        assert!(result.is_ok());
2796        assert_eq!(engine.handler_names().len(), 1);
2797    }
2798
2799    #[tokio::test]
2800    async fn engine_store_and_provider_accessors() {
2801        let store = Arc::new(InMemoryStore::new());
2802        let inner = ClaudeCodeProvider::new();
2803        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2804            inner,
2805            "/tmp/ironflow-fixtures",
2806        ));
2807        let engine = Engine::new(store.clone(), provider.clone());
2808
2809        // Verify accessors return references
2810        let _ = engine.store();
2811        let _ = engine.provider();
2812    }
2813
2814    // -----------------------------------------------------------------------
2815    // Operation trait tests
2816    // -----------------------------------------------------------------------
2817
2818    use crate::operation::{Operation, OperationContext};
2819    use async_trait::async_trait;
2820    use ironflow_core::error::OperationError;
2821    use ironflow_store::models::StepKind;
2822
2823    struct FakeGitlabOp {
2824        project_id: u64,
2825        title: String,
2826    }
2827
2828    #[async_trait]
2829    impl Operation for FakeGitlabOp {
2830        fn kind(&self) -> &str {
2831            "gitlab"
2832        }
2833
2834        async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2835            Ok(json!({
2836                "issue_id": 42,
2837                "project_id": self.project_id,
2838                "title": self.title,
2839            }))
2840        }
2841
2842        fn input(&self) -> Option<Value> {
2843            Some(json!({
2844                "project_id": self.project_id,
2845                "title": self.title,
2846            }))
2847        }
2848    }
2849
2850    struct FailingOp;
2851
2852    #[async_trait]
2853    impl Operation for FailingOp {
2854        fn kind(&self) -> &str {
2855            "broken-service"
2856        }
2857
2858        async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2859            Err(OperationError::Http {
2860                status: None,
2861                message: "service unavailable".to_string(),
2862            })
2863        }
2864    }
2865
2866    struct OperationWorkflow;
2867
2868    impl WorkflowHandler for OperationWorkflow {
2869        fn name(&self) -> &str {
2870            "operation-workflow"
2871        }
2872
2873        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2874            Box::pin(async move {
2875                let op = FakeGitlabOp {
2876                    project_id: 123,
2877                    title: "Bug report".to_string(),
2878                };
2879                ctx.operation("create-issue", &op).await?;
2880                Ok(())
2881            })
2882        }
2883    }
2884
2885    struct FailingOperationWorkflow;
2886
2887    impl WorkflowHandler for FailingOperationWorkflow {
2888        fn name(&self) -> &str {
2889            "failing-operation-workflow"
2890        }
2891
2892        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2893            Box::pin(async move {
2894                ctx.operation("broken-call", &FailingOp).await?;
2895                Ok(())
2896            })
2897        }
2898    }
2899
2900    struct MixedWorkflow;
2901
2902    impl WorkflowHandler for MixedWorkflow {
2903        fn name(&self) -> &str {
2904            "mixed-workflow"
2905        }
2906
2907        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2908            Box::pin(async move {
2909                ctx.shell("build", ShellConfig::new("echo built")).await?;
2910                let op = FakeGitlabOp {
2911                    project_id: 456,
2912                    title: "Deploy done".to_string(),
2913                };
2914                let result = ctx.operation("notify-gitlab", &op).await?;
2915                assert_eq!(result.output["issue_id"], 42);
2916                Ok(())
2917            })
2918        }
2919    }
2920
2921    #[tokio::test]
2922    async fn operation_step_happy_path() {
2923        let mut engine = create_test_engine();
2924        engine.register(OperationWorkflow).unwrap();
2925
2926        let run = engine
2927            .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
2928            .await
2929            .unwrap()
2930            .run;
2931
2932        assert_eq!(run.status.state, RunStatus::Completed);
2933
2934        let steps = engine.store().list_steps(run.id).await.unwrap();
2935
2936        assert_eq!(steps.len(), 1);
2937        assert_eq!(steps[0].name, "create-issue");
2938        assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
2939        assert_eq!(
2940            steps[0].status.state,
2941            ironflow_store::models::StepStatus::Completed
2942        );
2943
2944        let output = steps[0].output.as_ref().unwrap();
2945        assert_eq!(output["issue_id"], 42);
2946        assert_eq!(output["project_id"], 123);
2947
2948        let input = steps[0].input.as_ref().unwrap();
2949        assert_eq!(input["project_id"], 123);
2950        assert_eq!(input["title"], "Bug report");
2951    }
2952
2953    #[tokio::test]
2954    async fn operation_step_failure_marks_run_failed() {
2955        let mut engine = create_test_engine();
2956        engine.register(FailingOperationWorkflow).unwrap();
2957
2958        let result = engine
2959            .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
2960            .await;
2961
2962        assert!(result.is_err());
2963    }
2964
2965    #[tokio::test]
2966    async fn operation_mixed_with_shell_steps() {
2967        let mut engine = create_test_engine();
2968        engine.register(MixedWorkflow).unwrap();
2969
2970        let run = engine
2971            .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
2972            .await
2973            .unwrap()
2974            .run;
2975
2976        assert_eq!(run.status.state, RunStatus::Completed);
2977
2978        let steps = engine.store().list_steps(run.id).await.unwrap();
2979
2980        assert_eq!(steps.len(), 2);
2981        assert_eq!(steps[0].kind, StepKind::Shell);
2982        assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
2983        assert_eq!(steps[0].position, 0);
2984        assert_eq!(steps[1].position, 1);
2985    }
2986
2987    // -----------------------------------------------------------------------
2988    // Approval + resume tests
2989    // -----------------------------------------------------------------------
2990
2991    use crate::config::ApprovalConfig;
2992
2993    struct SingleApprovalWorkflow;
2994
2995    impl WorkflowHandler for SingleApprovalWorkflow {
2996        fn name(&self) -> &str {
2997            "single-approval"
2998        }
2999
3000        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3001            Box::pin(async move {
3002                ctx.shell("build", ShellConfig::new("echo built")).await?;
3003                ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
3004                ctx.shell("deploy", ShellConfig::new("echo deployed"))
3005                    .await?;
3006                Ok(())
3007            })
3008        }
3009    }
3010
3011    struct DoubleApprovalWorkflow;
3012
3013    impl WorkflowHandler for DoubleApprovalWorkflow {
3014        fn name(&self) -> &str {
3015            "double-approval"
3016        }
3017
3018        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3019            Box::pin(async move {
3020                ctx.shell("build", ShellConfig::new("echo built")).await?;
3021                ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
3022                    .await?;
3023                ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3024                    .await?;
3025                ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3026                    .await?;
3027                ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3028                    .await?;
3029                Ok(())
3030            })
3031        }
3032    }
3033
3034    #[tokio::test]
3035    async fn approval_pauses_run() {
3036        let mut engine = create_test_engine();
3037        engine.register(SingleApprovalWorkflow).unwrap();
3038
3039        let run = engine
3040            .run_handler("single-approval", TriggerKind::Manual, json!({}))
3041            .await
3042            .unwrap()
3043            .run;
3044
3045        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3046
3047        let steps = engine.store().list_steps(run.id).await.unwrap();
3048        assert_eq!(steps.len(), 2); // build + approval gate
3049        assert_eq!(steps[0].kind, StepKind::Shell);
3050        assert_eq!(steps[0].status.state, StepStatus::Completed);
3051        assert_eq!(steps[1].kind, StepKind::Approval);
3052        assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3053    }
3054
3055    #[tokio::test]
3056    async fn approval_resume_completes_run() {
3057        let mut engine = create_test_engine();
3058        engine.register(SingleApprovalWorkflow).unwrap();
3059
3060        // First execution: pauses at approval
3061        let run = engine
3062            .run_handler("single-approval", TriggerKind::Manual, json!({}))
3063            .await
3064            .unwrap()
3065            .run;
3066        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3067
3068        // Simulate approval: transition to Running
3069        engine
3070            .store()
3071            .update_run_status(run.id, RunStatus::Running)
3072            .await
3073            .unwrap();
3074
3075        // Resume: replays build, skips approval, executes deploy
3076        let resumed = engine.resume_run(run.id).await.unwrap().run;
3077        assert_eq!(resumed.status.state, RunStatus::Completed);
3078
3079        let steps = engine.store().list_steps(run.id).await.unwrap();
3080        assert_eq!(steps.len(), 3); // build + approval + deploy
3081        assert_eq!(steps[0].name, "build");
3082        assert_eq!(steps[0].status.state, StepStatus::Completed);
3083        assert_eq!(steps[1].name, "gate");
3084        assert_eq!(steps[1].kind, StepKind::Approval);
3085        assert_eq!(steps[1].status.state, StepStatus::Completed);
3086        assert_eq!(steps[2].name, "deploy");
3087        assert_eq!(steps[2].status.state, StepStatus::Completed);
3088    }
3089
3090    #[tokio::test]
3091    async fn double_approval_two_resumes() {
3092        let mut engine = create_test_engine();
3093        engine.register(DoubleApprovalWorkflow).unwrap();
3094
3095        // First execution: pauses at staging-gate
3096        let run = engine
3097            .run_handler("double-approval", TriggerKind::Manual, json!({}))
3098            .await
3099            .unwrap()
3100            .run;
3101        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3102
3103        let steps = engine.store().list_steps(run.id).await.unwrap();
3104        assert_eq!(steps.len(), 2); // build + staging-gate
3105
3106        // First approval
3107        engine
3108            .store()
3109            .update_run_status(run.id, RunStatus::Running)
3110            .await
3111            .unwrap();
3112
3113        let resumed = engine.resume_run(run.id).await.unwrap().run;
3114        assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3115
3116        let steps = engine.store().list_steps(run.id).await.unwrap();
3117        assert_eq!(steps.len(), 4); // build + staging-gate + deploy-staging + prod-gate
3118
3119        // Second approval
3120        engine
3121            .store()
3122            .update_run_status(run.id, RunStatus::Running)
3123            .await
3124            .unwrap();
3125
3126        let final_run = engine.resume_run(run.id).await.unwrap().run;
3127        assert_eq!(final_run.status.state, RunStatus::Completed);
3128
3129        let steps = engine.store().list_steps(run.id).await.unwrap();
3130        assert_eq!(steps.len(), 5);
3131        assert_eq!(steps[0].name, "build");
3132        assert_eq!(steps[1].name, "staging-gate");
3133        assert_eq!(steps[2].name, "deploy-staging");
3134        assert_eq!(steps[3].name, "prod-gate");
3135        assert_eq!(steps[4].name, "deploy-prod");
3136
3137        for step in &steps {
3138            assert_eq!(step.status.state, StepStatus::Completed);
3139        }
3140    }
3141
3142    // -----------------------------------------------------------------------
3143    // fail_orphaned_steps tests
3144    // -----------------------------------------------------------------------
3145
3146    use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3147
3148    async fn create_step_with_status(
3149        store: &Arc<dyn Store>,
3150        run_id: Uuid,
3151        name: &str,
3152        position: u32,
3153        status: StepStatus,
3154    ) -> ironflow_store::models::Step {
3155        let step = store
3156            .create_step(NewStep {
3157                run_id,
3158                trace_id: step_trace_id(run_id, name, position),
3159                name: name.to_string(),
3160                kind: StepKind::Shell,
3161                position,
3162                input: None,
3163                is_error_handler: false,
3164            })
3165            .await
3166            .unwrap();
3167
3168        match status {
3169            StepStatus::Pending => {}
3170            StepStatus::Running => {
3171                store
3172                    .update_step(
3173                        step.id,
3174                        StepUpdate {
3175                            status: Some(StepStatus::Running),
3176                            ..StepUpdate::default()
3177                        },
3178                    )
3179                    .await
3180                    .unwrap();
3181            }
3182            StepStatus::Completed => {
3183                store
3184                    .update_step(
3185                        step.id,
3186                        StepUpdate {
3187                            status: Some(StepStatus::Running),
3188                            ..StepUpdate::default()
3189                        },
3190                    )
3191                    .await
3192                    .unwrap();
3193                store
3194                    .update_step(
3195                        step.id,
3196                        StepUpdate {
3197                            status: Some(StepStatus::Completed),
3198                            ..StepUpdate::default()
3199                        },
3200                    )
3201                    .await
3202                    .unwrap();
3203            }
3204            StepStatus::AwaitingApproval => {
3205                store
3206                    .update_step(
3207                        step.id,
3208                        StepUpdate {
3209                            status: Some(StepStatus::Running),
3210                            ..StepUpdate::default()
3211                        },
3212                    )
3213                    .await
3214                    .unwrap();
3215                store
3216                    .update_step(
3217                        step.id,
3218                        StepUpdate {
3219                            status: Some(StepStatus::AwaitingApproval),
3220                            ..StepUpdate::default()
3221                        },
3222                    )
3223                    .await
3224                    .unwrap();
3225            }
3226            _ => panic!("unsupported status for test helper: {status}"),
3227        }
3228
3229        store.get_step(step.id).await.unwrap().unwrap()
3230    }
3231
3232    #[tokio::test]
3233    async fn fail_orphaned_steps_marks_running_as_failed() {
3234        let engine = create_test_engine();
3235        let run = engine
3236            .store()
3237            .create_run(NewRun {
3238                created_by: None,
3239                workflow_name: "test".to_string(),
3240                trigger: TriggerKind::Manual,
3241                payload: json!({}),
3242                max_retries: 0,
3243                handler_version: None,
3244                labels: HashMap::new(),
3245                scheduled_at: None,
3246                idempotency_key: None,
3247                concurrency_key: None,
3248                max_cost_usd: None,
3249            })
3250            .await
3251            .unwrap()
3252            .into_run();
3253
3254        let step = create_step_with_status(
3255            engine.store(),
3256            run.id,
3257            "running-step",
3258            0,
3259            StepStatus::Running,
3260        )
3261        .await;
3262
3263        engine
3264            .fail_orphaned_steps(run.id, "parent run timed out")
3265            .await
3266            .unwrap();
3267
3268        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3269        assert_eq!(updated.status.state, StepStatus::Failed);
3270        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3271        assert!(updated.completed_at.is_some());
3272    }
3273
3274    #[tokio::test]
3275    async fn fail_orphaned_steps_marks_pending_as_skipped() {
3276        let engine = create_test_engine();
3277        let run = engine
3278            .store()
3279            .create_run(NewRun {
3280                created_by: None,
3281                workflow_name: "test".to_string(),
3282                trigger: TriggerKind::Manual,
3283                payload: json!({}),
3284                max_retries: 0,
3285                handler_version: None,
3286                labels: HashMap::new(),
3287                scheduled_at: None,
3288                idempotency_key: None,
3289                concurrency_key: None,
3290                max_cost_usd: None,
3291            })
3292            .await
3293            .unwrap()
3294            .into_run();
3295
3296        let step = create_step_with_status(
3297            engine.store(),
3298            run.id,
3299            "pending-step",
3300            0,
3301            StepStatus::Pending,
3302        )
3303        .await;
3304
3305        engine
3306            .fail_orphaned_steps(run.id, "parent run timed out")
3307            .await
3308            .unwrap();
3309
3310        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3311        assert_eq!(updated.status.state, StepStatus::Skipped);
3312        assert!(updated.error.is_none());
3313        assert!(updated.completed_at.is_some());
3314    }
3315
3316    #[tokio::test]
3317    async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
3318        let engine = create_test_engine();
3319        let run = engine
3320            .store()
3321            .create_run(NewRun {
3322                created_by: None,
3323                workflow_name: "test".to_string(),
3324                trigger: TriggerKind::Manual,
3325                payload: json!({}),
3326                max_retries: 0,
3327                handler_version: None,
3328                labels: HashMap::new(),
3329                scheduled_at: None,
3330                idempotency_key: None,
3331                concurrency_key: None,
3332                max_cost_usd: None,
3333            })
3334            .await
3335            .unwrap()
3336            .into_run();
3337
3338        let step = create_step_with_status(
3339            engine.store(),
3340            run.id,
3341            "approval-step",
3342            0,
3343            StepStatus::AwaitingApproval,
3344        )
3345        .await;
3346
3347        engine
3348            .fail_orphaned_steps(run.id, "parent run timed out")
3349            .await
3350            .unwrap();
3351
3352        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3353        assert_eq!(updated.status.state, StepStatus::Failed);
3354        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3355        assert!(updated.completed_at.is_some());
3356    }
3357
3358    #[tokio::test]
3359    async fn fail_orphaned_steps_skips_terminal_steps() {
3360        let engine = create_test_engine();
3361        let run = engine
3362            .store()
3363            .create_run(NewRun {
3364                created_by: None,
3365                workflow_name: "test".to_string(),
3366                trigger: TriggerKind::Manual,
3367                payload: json!({}),
3368                max_retries: 0,
3369                handler_version: None,
3370                labels: HashMap::new(),
3371                scheduled_at: None,
3372                idempotency_key: None,
3373                concurrency_key: None,
3374                max_cost_usd: None,
3375            })
3376            .await
3377            .unwrap()
3378            .into_run();
3379
3380        let completed_step =
3381            create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
3382        let running_step =
3383            create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
3384                .await;
3385
3386        engine
3387            .fail_orphaned_steps(run.id, "parent run timed out")
3388            .await
3389            .unwrap();
3390
3391        let completed = engine
3392            .store()
3393            .get_step(completed_step.id)
3394            .await
3395            .unwrap()
3396            .unwrap();
3397        assert_eq!(completed.status.state, StepStatus::Completed);
3398
3399        let failed = engine
3400            .store()
3401            .get_step(running_step.id)
3402            .await
3403            .unwrap()
3404            .unwrap();
3405        assert_eq!(failed.status.state, StepStatus::Failed);
3406    }
3407
3408    #[tokio::test]
3409    async fn fail_orphaned_steps_mixed_states() {
3410        let engine = create_test_engine();
3411        let run = engine
3412            .store()
3413            .create_run(NewRun {
3414                created_by: None,
3415                workflow_name: "test".to_string(),
3416                trigger: TriggerKind::Manual,
3417                payload: json!({}),
3418                max_retries: 0,
3419                handler_version: None,
3420                labels: HashMap::new(),
3421                scheduled_at: None,
3422                idempotency_key: None,
3423                concurrency_key: None,
3424                max_cost_usd: None,
3425            })
3426            .await
3427            .unwrap()
3428            .into_run();
3429
3430        let s_completed =
3431            create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
3432                .await;
3433        let s_running =
3434            create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
3435        let s_pending =
3436            create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
3437
3438        engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
3439
3440        let r_completed = engine
3441            .store()
3442            .get_step(s_completed.id)
3443            .await
3444            .unwrap()
3445            .unwrap();
3446        assert_eq!(r_completed.status.state, StepStatus::Completed);
3447
3448        let r_running = engine
3449            .store()
3450            .get_step(s_running.id)
3451            .await
3452            .unwrap()
3453            .unwrap();
3454        assert_eq!(r_running.status.state, StepStatus::Failed);
3455        assert_eq!(r_running.error.as_deref(), Some("timeout"));
3456
3457        let r_pending = engine
3458            .store()
3459            .get_step(s_pending.id)
3460            .await
3461            .unwrap()
3462            .unwrap();
3463        assert_eq!(r_pending.status.state, StepStatus::Skipped);
3464        assert!(r_pending.error.is_none());
3465    }
3466
3467    #[tokio::test]
3468    async fn fail_orphaned_steps_no_steps_is_noop() {
3469        let engine = create_test_engine();
3470        let run = engine
3471            .store()
3472            .create_run(NewRun {
3473                created_by: None,
3474                workflow_name: "test".to_string(),
3475                trigger: TriggerKind::Manual,
3476                payload: json!({}),
3477                max_retries: 0,
3478                handler_version: None,
3479                labels: HashMap::new(),
3480                scheduled_at: None,
3481                idempotency_key: None,
3482                concurrency_key: None,
3483                max_cost_usd: None,
3484            })
3485            .await
3486            .unwrap()
3487            .into_run();
3488
3489        let result = engine.fail_orphaned_steps(run.id, "timeout").await;
3490        assert!(result.is_ok());
3491    }
3492
3493    #[tokio::test]
3494    async fn fail_orphaned_steps_preserves_existing_error() {
3495        let engine = create_test_engine();
3496        let run = engine
3497            .store()
3498            .create_run(NewRun {
3499                created_by: None,
3500                workflow_name: "test".to_string(),
3501                trigger: TriggerKind::Manual,
3502                payload: json!({}),
3503                max_retries: 0,
3504                handler_version: None,
3505                labels: HashMap::new(),
3506                scheduled_at: None,
3507                idempotency_key: None,
3508                concurrency_key: None,
3509                max_cost_usd: None,
3510            })
3511            .await
3512            .unwrap()
3513            .into_run();
3514
3515        let step_with_error = create_step_with_status(
3516            engine.store(),
3517            run.id,
3518            "already-errored",
3519            0,
3520            StepStatus::Running,
3521        )
3522        .await;
3523
3524        engine
3525            .store()
3526            .update_step(
3527                step_with_error.id,
3528                StepUpdate {
3529                    error: Some("real error from provider".to_string()),
3530                    ..StepUpdate::default()
3531                },
3532            )
3533            .await
3534            .unwrap();
3535
3536        let step_no_error = create_step_with_status(
3537            engine.store(),
3538            run.id,
3539            "no-error-yet",
3540            1,
3541            StepStatus::Running,
3542        )
3543        .await;
3544
3545        engine
3546            .fail_orphaned_steps(run.id, "parent run failed")
3547            .await
3548            .unwrap();
3549
3550        let updated_with = engine
3551            .store()
3552            .get_step(step_with_error.id)
3553            .await
3554            .unwrap()
3555            .unwrap();
3556        assert_eq!(updated_with.status.state, StepStatus::Failed);
3557        assert_eq!(
3558            updated_with.error.as_deref(),
3559            Some("real error from provider"),
3560        );
3561
3562        let updated_without = engine
3563            .store()
3564            .get_step(step_no_error.id)
3565            .await
3566            .unwrap()
3567            .unwrap();
3568        assert_eq!(updated_without.status.state, StepStatus::Failed);
3569        assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
3570    }
3571}