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                            ..RunUpdate::default()
1847                        },
1848                    )
1849                    .await?;
1850
1851                info!(
1852                    run_id = %run_id,
1853                    status = %final_status,
1854                    cost_usd = %ctx.total_cost_usd(),
1855                    duration_ms = total_duration,
1856                    "run completed"
1857                );
1858            }
1859            Err(EngineError::ApprovalRequired {
1860                run_id: approval_run_id,
1861                step_id,
1862                ref message,
1863            }) => {
1864                final_status = RunStatus::AwaitingApproval;
1865                final_run = self
1866                    .store
1867                    .update_run_returning(
1868                        run_id,
1869                        RunUpdate {
1870                            status: Some(RunStatus::AwaitingApproval),
1871                            cost_usd: Some(ctx.total_cost_usd()),
1872                            duration_ms: Some(total_duration),
1873                            ..RunUpdate::default()
1874                        },
1875                    )
1876                    .await?;
1877
1878                info!(
1879                    run_id = %approval_run_id,
1880                    step_id = %step_id,
1881                    message = %message,
1882                    "run awaiting approval"
1883                );
1884
1885                self.publish_approval_requested(approval_run_id, step_id, message)
1886                    .await?;
1887            }
1888            Err(EngineError::ChildSuspended {
1889                run_id: child_run_id,
1890                ref cause,
1891            }) => {
1892                final_status = cause.suspension_status();
1893                // No `scheduled_at`: the suspended descendant owns the wake-up
1894                // and resumes this run through `resume_chain`. A wake-up armed
1895                // here too would resume the chain twice.
1896                final_run = self
1897                    .store
1898                    .update_run_returning(
1899                        run_id,
1900                        RunUpdate {
1901                            status: Some(final_status),
1902                            cost_usd: Some(ctx.total_cost_usd()),
1903                            duration_ms: Some(total_duration),
1904                            ..RunUpdate::default()
1905                        },
1906                    )
1907                    .await?;
1908
1909                let leaf = cause.suspension_leaf();
1910                info!(
1911                    run_id = %run_id,
1912                    child_run_id = %child_run_id,
1913                    status = %final_status,
1914                    cause = %leaf,
1915                    "run suspended with its child run"
1916                );
1917
1918                match leaf {
1919                    EngineError::ApprovalRequired {
1920                        run_id: approval_run_id,
1921                        step_id,
1922                        message,
1923                    } => {
1924                        self.publish_approval_requested(*approval_run_id, *step_id, message)
1925                            .await?;
1926                    }
1927                    EngineError::SignalWaiting {
1928                        run_id: wait_run_id,
1929                        step_id,
1930                        step_name,
1931                        name,
1932                        key,
1933                        deadline_at,
1934                    } => {
1935                        self.event_publisher
1936                            .publish(Event::SignalAwaited(SignalAwaitedEvent {
1937                                run_id: *wait_run_id,
1938                                step_id: *step_id,
1939                                step_name: step_name.clone(),
1940                                name: name.clone(),
1941                                key: key.clone(),
1942                                deadline_at: *deadline_at,
1943                                at: Utc::now(),
1944                            }));
1945                    }
1946                    // A human input or a delay publishes no suspension event,
1947                    // like on a top-level run.
1948                    _ => {}
1949                }
1950            }
1951            Err(EngineError::HumanInputRequired {
1952                run_id: input_run_id,
1953                step_id,
1954                ref message,
1955            }) => {
1956                final_status = RunStatus::AwaitingApproval;
1957                final_run = self
1958                    .store
1959                    .update_run_returning(
1960                        run_id,
1961                        RunUpdate {
1962                            status: Some(RunStatus::AwaitingApproval),
1963                            cost_usd: Some(ctx.total_cost_usd()),
1964                            duration_ms: Some(total_duration),
1965                            ..RunUpdate::default()
1966                        },
1967                    )
1968                    .await?;
1969
1970                // No `ApprovalRequested` event: a human input is not an approval.
1971                info!(
1972                    run_id = %input_run_id,
1973                    step_id = %step_id,
1974                    message = %message,
1975                    "run awaiting human input"
1976                );
1977            }
1978            Err(EngineError::DelaySleeping {
1979                run_id: delay_run_id,
1980                step_id,
1981                wake_at,
1982            }) => {
1983                final_status = RunStatus::Sleeping;
1984                final_run = self
1985                    .store
1986                    .update_run_returning(
1987                        run_id,
1988                        RunUpdate {
1989                            status: Some(RunStatus::Sleeping),
1990                            cost_usd: Some(ctx.total_cost_usd()),
1991                            duration_ms: Some(total_duration),
1992                            scheduled_at: Some(wake_at),
1993                            ..RunUpdate::default()
1994                        },
1995                    )
1996                    .await?;
1997
1998                info!(
1999                    run_id = %delay_run_id,
2000                    step_id = %step_id,
2001                    wake_at = %wake_at,
2002                    "run sleeping until delay elapses"
2003                );
2004            }
2005            Err(EngineError::SignalWaiting {
2006                run_id: wait_run_id,
2007                step_id,
2008                ref step_name,
2009                ref name,
2010                ref key,
2011                deadline_at,
2012            }) => {
2013                final_status = RunStatus::Sleeping;
2014                // Atomic with the step lock: a signal delivered since the step
2015                // opened leaves the run due right away instead of until the
2016                // deadline.
2017                let waiting = self
2018                    .store
2019                    .suspend_run_on_signal(run_id, step_id, deadline_at)
2020                    .await?;
2021                final_run = self
2022                    .store
2023                    .update_run_returning(
2024                        run_id,
2025                        RunUpdate {
2026                            cost_usd: Some(ctx.total_cost_usd()),
2027                            duration_ms: Some(total_duration),
2028                            ..RunUpdate::default()
2029                        },
2030                    )
2031                    .await?;
2032
2033                if waiting {
2034                    self.event_publisher
2035                        .publish(Event::SignalAwaited(SignalAwaitedEvent {
2036                            run_id: wait_run_id,
2037                            step_id,
2038                            step_name: step_name.clone(),
2039                            name: name.clone(),
2040                            key: key.clone(),
2041                            deadline_at,
2042                            at: Utc::now(),
2043                        }));
2044                }
2045
2046                info!(
2047                    run_id = %wait_run_id,
2048                    step_id = %step_id,
2049                    signal = %name,
2050                    key = %key,
2051                    deadline_at = %deadline_at,
2052                    waiting,
2053                    "run sleeping until a signal arrives"
2054                );
2055            }
2056            Err(err) => {
2057                // A guardrail stop (budget or workflow guard) is deliberate,
2058                // not a breakage: the run is cancelled, never failed and
2059                // never replayed.
2060                let guardrail_stop = matches!(
2061                    err,
2062                    EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2063                );
2064
2065                final_status = if guardrail_stop {
2066                    if let Err(store_err) = self
2067                        .store
2068                        .update_run(
2069                            run_id,
2070                            RunUpdate {
2071                                status: Some(RunStatus::Cancelled),
2072                                error: Some(err.to_string()),
2073                                cost_usd: Some(ctx.total_cost_usd()),
2074                                duration_ms: Some(total_duration),
2075                                completed_at: Some(completed_at),
2076                                ..RunUpdate::default()
2077                            },
2078                        )
2079                        .await
2080                    {
2081                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2082                    }
2083                    if let Err(cleanup_err) = self
2084                        .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2085                        .await
2086                    {
2087                        error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2088                    }
2089                    RunStatus::Cancelled
2090                } else {
2091                    self.fail_or_schedule_retry(
2092                        run_id,
2093                        &err.to_string(),
2094                        is_run_retryable(&err),
2095                        Some(ctx.total_cost_usd()),
2096                        Some(total_duration),
2097                    )
2098                    .await
2099                    .unwrap_or_else(|store_err| {
2100                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2101                        RunStatus::Failed
2102                    })
2103                };
2104
2105                if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2106                    self.on_run_budget_exceeded(workflow_name, run_id, &err);
2107                }
2108
2109                error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2110
2111                self.publish_run_status_changed(
2112                    workflow_name,
2113                    run_id,
2114                    final_status,
2115                    Some(err.to_string()),
2116                    ctx,
2117                    total_duration,
2118                    run_labels,
2119                );
2120
2121                #[cfg(feature = "prometheus")]
2122                self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2123
2124                return Err(err);
2125            }
2126        }
2127
2128        self.publish_run_status_changed(
2129            workflow_name,
2130            run_id,
2131            final_status,
2132            None,
2133            ctx,
2134            total_duration,
2135            run_labels,
2136        );
2137
2138        #[cfg(feature = "prometheus")]
2139        self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2140
2141        Ok(WorkflowResult {
2142            run: final_run,
2143            steps: ctx.step_results().to_vec(),
2144        })
2145    }
2146
2147    /// Publish [`Event::ApprovalRequested`] for an approval gate that opened
2148    /// on `run_id`, with the requirement recorded when the gate opened.
2149    async fn publish_approval_requested(
2150        &self,
2151        run_id: Uuid,
2152        step_id: Uuid,
2153        message: &str,
2154    ) -> Result<(), EngineError> {
2155        let requirement = self
2156            .store
2157            .get_step(step_id)
2158            .await?
2159            .and_then(|s| s.approval_requirement);
2160        self.event_publisher
2161            .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2162                run_id,
2163                step_id,
2164                message: message.to_string(),
2165                requirement,
2166                at: Utc::now(),
2167            }));
2168        Ok(())
2169    }
2170
2171    /// Fail every ancestor of a child run, closest first.
2172    ///
2173    /// Used when a gate inside a sub-workflow is rejected: the child run
2174    /// fails, and the runs suspended with it (its parent, up to the root)
2175    /// must not stay suspended on a child that will never resume. Each
2176    /// ancestor goes through [`fail_or_schedule_retry`](Self::fail_or_schedule_retry)
2177    /// without a retry, which also fails its open `Workflow` step. A run that
2178    /// is not a child of a sub-workflow has no ancestor: nothing happens.
2179    ///
2180    /// # Errors
2181    ///
2182    /// Returns [`EngineError::Store`] if a run of the chain does not exist or
2183    /// its failure cannot be persisted.
2184    ///
2185    /// # Examples
2186    ///
2187    /// ```no_run
2188    /// use ironflow_engine::engine::Engine;
2189    /// use ironflow_engine::error::EngineError;
2190    /// use uuid::Uuid;
2191    ///
2192    /// # async fn example(engine: &Engine, child_run_id: Uuid) -> Result<(), EngineError> {
2193    /// engine
2194    ///     .fail_ancestors(child_run_id, "approval rejected in a sub-workflow")
2195    ///     .await?;
2196    /// # Ok(())
2197    /// # }
2198    /// ```
2199    pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2200        let mut current = self
2201            .store
2202            .get_run(run_id)
2203            .await?
2204            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2205        // Labels are data: a chain that loops back on itself stops there.
2206        let mut visited = HashSet::from([run_id]);
2207
2208        while let Some(parent_id) = chain_parent(&current) {
2209            if !visited.insert(parent_id) {
2210                break;
2211            }
2212            let status = self
2213                .fail_or_schedule_retry(parent_id, reason, false, None, None)
2214                .await?;
2215            info!(
2216                run_id = %run_id,
2217                ancestor_run_id = %parent_id,
2218                status = %status,
2219                "ancestor run failed with its child"
2220            );
2221            current = self
2222                .store
2223                .get_run(parent_id)
2224                .await?
2225                .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2226        }
2227
2228        Ok(())
2229    }
2230
2231    /// Emit Prometheus metrics for a completed run.
2232    #[cfg(feature = "prometheus")]
2233    fn emit_run_metrics(
2234        &self,
2235        workflow_name: &str,
2236        status: RunStatus,
2237        duration_ms: u64,
2238        ctx: &WorkflowContext,
2239    ) {
2240        let status_str = status.to_string();
2241        let wf = workflow_name.to_string();
2242
2243        counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2244        histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2245            .record(duration_ms as f64 / 1000.0);
2246        histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2247            ctx.total_cost_usd()
2248                .to_string()
2249                .parse::<f64>()
2250                .unwrap_or(0.0),
2251        );
2252        gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2253    }
2254
2255    /// Record the metric and publish the audit event for a run that hit its
2256    /// cost cap.
2257    ///
2258    /// A non-[`RunBudgetExceeded`](EngineError::RunBudgetExceeded) error is
2259    /// ignored, so callers can pass the error unconditionally.
2260    fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2261        let EngineError::RunBudgetExceeded {
2262            limit_usd,
2263            spent_usd,
2264            step_budget_usd,
2265            ..
2266        } = err
2267        else {
2268            return;
2269        };
2270
2271        #[cfg(feature = "prometheus")]
2272        counter!(
2273            RUN_BUDGET_EXCEEDED_TOTAL,
2274            "workflow" => workflow_name.to_string(),
2275            "scope" => "run",
2276        )
2277        .increment(1);
2278
2279        self.event_publisher
2280            .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2281                run_id,
2282                workflow_name: workflow_name.to_string(),
2283                limit_usd: *limit_usd,
2284                spent_usd: *spent_usd,
2285                step_budget_usd: *step_budget_usd,
2286                at: Utc::now(),
2287            }));
2288    }
2289
2290    /// Publish a run status changed event to all registered subscribers.
2291    ///
2292    /// `from` is always `Running` because `finalize_run` is only called
2293    /// from a running state.
2294    #[allow(clippy::too_many_arguments)]
2295    fn publish_run_status_changed(
2296        &self,
2297        workflow_name: &str,
2298        run_id: Uuid,
2299        to: RunStatus,
2300        error: Option<String>,
2301        ctx: &WorkflowContext,
2302        duration_ms: u64,
2303        labels: HashMap<String, String>,
2304    ) {
2305        let now = Utc::now();
2306        let cost_usd = ctx.total_cost_usd();
2307        let wf = workflow_name.to_string();
2308
2309        self.event_publisher
2310            .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2311                run_id,
2312                workflow_name: wf.clone(),
2313                from: RunStatus::Running,
2314                to,
2315                error: error.clone(),
2316                cost_usd,
2317                duration_ms,
2318                labels: labels.clone(),
2319                at: now,
2320            }));
2321
2322        if to == RunStatus::Failed {
2323            self.event_publisher
2324                .publish(Event::RunFailed(RunFailedEvent {
2325                    run_id,
2326                    workflow_name: wf,
2327                    error,
2328                    cost_usd,
2329                    duration_ms,
2330                    labels,
2331                    at: now,
2332                }));
2333        }
2334    }
2335}
2336
2337impl fmt::Debug for Engine {
2338    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2339        f.debug_struct("Engine")
2340            .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2341            .finish_non_exhaustive()
2342    }
2343}
2344
2345#[cfg(test)]
2346mod tests {
2347    use super::*;
2348    use crate::config::ShellConfig;
2349    use crate::handler::{HandlerFuture, WorkflowHandler};
2350    use ironflow_core::providers::claude::ClaudeCodeProvider;
2351    use ironflow_core::providers::record_replay::RecordReplayProvider;
2352    use ironflow_store::memory::InMemoryStore;
2353    use ironflow_store::models::StepStatus;
2354    use serde_json::json;
2355
2356    // Test handler that echoes a message via shell
2357    struct EchoWorkflow;
2358
2359    impl WorkflowHandler for EchoWorkflow {
2360        fn name(&self) -> &str {
2361            "echo-workflow"
2362        }
2363
2364        fn describe(&self) -> WorkflowInfo {
2365            WorkflowInfo {
2366                description: "A simple workflow that echoes hello".to_string(),
2367                source_code: None,
2368                sub_workflows: Vec::new(),
2369                category: None,
2370                version: self.version().map(str::to_string),
2371                compatible_versions: Vec::new(),
2372                input_schema: None,
2373                default_labels: HashMap::new(),
2374                schedule: self.schedule().cloned(),
2375                default_max_cost_usd: self.default_max_cost_usd(),
2376            }
2377        }
2378
2379        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2380            Box::pin(async move {
2381                ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2382                Ok(())
2383            })
2384        }
2385    }
2386
2387    // Test handler that fails
2388    struct FailingWorkflow;
2389
2390    impl WorkflowHandler for FailingWorkflow {
2391        fn name(&self) -> &str {
2392            "failing-workflow"
2393        }
2394
2395        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2396            Box::pin(async move {
2397                ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2398                Ok(())
2399            })
2400        }
2401    }
2402
2403    fn create_test_engine() -> Engine {
2404        let store = Arc::new(InMemoryStore::new());
2405        let inner = ClaudeCodeProvider::new();
2406        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2407            inner,
2408            "/tmp/ironflow-fixtures",
2409        ));
2410        Engine::new(store, provider)
2411    }
2412
2413    #[test]
2414    fn engine_new_creates_instance() {
2415        let engine = create_test_engine();
2416        assert_eq!(engine.handler_names().len(), 0);
2417    }
2418
2419    #[test]
2420    fn execution_mode_defaults_to_local() {
2421        let engine = create_test_engine();
2422        assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2423    }
2424
2425    #[test]
2426    fn with_execution_mode_overrides_the_default() {
2427        let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2428        assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2429    }
2430
2431    #[test]
2432    fn engine_register_handler() {
2433        let mut engine = create_test_engine();
2434        let result = engine.register(EchoWorkflow);
2435        assert!(result.is_ok());
2436        assert_eq!(engine.handler_names().len(), 1);
2437        assert!(engine.handler_names().contains(&"echo-workflow"));
2438    }
2439
2440    #[test]
2441    fn engine_register_duplicate_returns_error() {
2442        let mut engine = create_test_engine();
2443        engine.register(EchoWorkflow).unwrap();
2444        let result = engine.register(EchoWorkflow);
2445        assert!(result.is_err());
2446    }
2447
2448    #[test]
2449    fn engine_get_handler_found() {
2450        let mut engine = create_test_engine();
2451        engine.register(EchoWorkflow).unwrap();
2452        let handler = engine.get_handler("echo-workflow");
2453        assert!(handler.is_some());
2454    }
2455
2456    #[test]
2457    fn engine_get_handler_not_found() {
2458        let engine = create_test_engine();
2459        let handler = engine.get_handler("nonexistent");
2460        assert!(handler.is_none());
2461    }
2462
2463    #[test]
2464    fn engine_handler_names_lists_all() {
2465        let mut engine = create_test_engine();
2466        engine.register(EchoWorkflow).unwrap();
2467        engine.register(FailingWorkflow).unwrap();
2468        let names = engine.handler_names();
2469        assert_eq!(names.len(), 2);
2470        assert!(names.contains(&"echo-workflow"));
2471        assert!(names.contains(&"failing-workflow"));
2472    }
2473
2474    #[test]
2475    fn engine_handler_info_returns_description() {
2476        let mut engine = create_test_engine();
2477        engine.register(EchoWorkflow).unwrap();
2478        let info = engine.handler_info("echo-workflow");
2479        assert!(info.is_some());
2480        let info = info.unwrap();
2481        assert_eq!(info.description, "A simple workflow that echoes hello");
2482    }
2483
2484    struct CategorizedWorkflow;
2485
2486    impl WorkflowHandler for CategorizedWorkflow {
2487        fn name(&self) -> &str {
2488            "categorized"
2489        }
2490        fn category(&self) -> Option<&str> {
2491            Some("data/etl")
2492        }
2493        fn execute<'a>(
2494            &'a self,
2495            _ctx: &'a mut WorkflowContext,
2496        ) -> crate::handler::HandlerFuture<'a> {
2497            Box::pin(async move { Ok(()) })
2498        }
2499    }
2500
2501    #[test]
2502    fn engine_default_describe_propagates_category() {
2503        let mut engine = create_test_engine();
2504        engine.register(CategorizedWorkflow).unwrap();
2505        let info = engine.handler_info("categorized").unwrap();
2506        assert_eq!(info.category.as_deref(), Some("data/etl"));
2507    }
2508
2509    #[test]
2510    fn engine_default_describe_without_category() {
2511        let mut engine = create_test_engine();
2512        engine.register(EchoWorkflow).unwrap();
2513        let info = engine.handler_info("echo-workflow").unwrap();
2514        assert!(info.category.is_none());
2515    }
2516
2517    // -----------------------------------------------------------------------
2518    // Schedule tests
2519    // -----------------------------------------------------------------------
2520
2521    struct ScheduledWorkflow {
2522        schedule: CronSchedule,
2523    }
2524
2525    impl ScheduledWorkflow {
2526        fn new() -> Self {
2527            Self {
2528                schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2529            }
2530        }
2531    }
2532
2533    impl WorkflowHandler for ScheduledWorkflow {
2534        fn name(&self) -> &str {
2535            "scheduled"
2536        }
2537        fn schedule(&self) -> Option<&CronSchedule> {
2538            Some(&self.schedule)
2539        }
2540        fn execute<'a>(
2541            &'a self,
2542            _ctx: &'a mut WorkflowContext,
2543        ) -> crate::handler::HandlerFuture<'a> {
2544            Box::pin(async move { Ok(()) })
2545        }
2546    }
2547
2548    #[test]
2549    fn engine_default_describe_propagates_schedule() {
2550        let mut engine = create_test_engine();
2551        engine.register(ScheduledWorkflow::new()).unwrap();
2552        let info = engine.handler_info("scheduled").unwrap();
2553        assert_eq!(
2554            info.schedule.as_ref().map(|s| s.as_str()),
2555            Some("0 0 * * * *")
2556        );
2557    }
2558
2559    #[test]
2560    fn engine_default_describe_without_schedule() {
2561        let mut engine = create_test_engine();
2562        engine.register(EchoWorkflow).unwrap();
2563        let info = engine.handler_info("echo-workflow").unwrap();
2564        assert!(info.schedule.is_none());
2565    }
2566
2567    #[test]
2568    fn scheduled_handlers_returns_only_scheduled() {
2569        let mut engine = create_test_engine();
2570        engine.register(EchoWorkflow).unwrap();
2571        engine.register(ScheduledWorkflow::new()).unwrap();
2572        engine.register(FailingWorkflow).unwrap();
2573
2574        let scheduled = engine.scheduled_handlers();
2575        assert_eq!(scheduled.len(), 1);
2576        assert_eq!(scheduled[0].0, "scheduled");
2577        assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
2578    }
2579
2580    #[test]
2581    fn scheduled_handlers_empty_when_none_scheduled() {
2582        let mut engine = create_test_engine();
2583        engine.register(EchoWorkflow).unwrap();
2584        engine.register(FailingWorkflow).unwrap();
2585
2586        let scheduled = engine.scheduled_handlers();
2587        assert!(scheduled.is_empty());
2588    }
2589
2590    struct BadCategoryWorkflow(&'static str);
2591
2592    impl WorkflowHandler for BadCategoryWorkflow {
2593        fn name(&self) -> &str {
2594            "bad-category"
2595        }
2596        fn category(&self) -> Option<&str> {
2597            Some(self.0)
2598        }
2599        fn execute<'a>(
2600            &'a self,
2601            _ctx: &'a mut WorkflowContext,
2602        ) -> crate::handler::HandlerFuture<'a> {
2603            Box::pin(async move { Ok(()) })
2604        }
2605    }
2606
2607    #[test]
2608    fn engine_register_rejects_empty_category() {
2609        let mut engine = create_test_engine();
2610        let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2611        match err {
2612            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2613            other => panic!("expected InvalidWorkflow, got {other:?}"),
2614        }
2615    }
2616
2617    #[test]
2618    fn engine_register_rejects_leading_slash_category() {
2619        let mut engine = create_test_engine();
2620        let err = engine
2621            .register(BadCategoryWorkflow("/data/etl"))
2622            .unwrap_err();
2623        match err {
2624            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2625            other => panic!("expected InvalidWorkflow, got {other:?}"),
2626        }
2627    }
2628
2629    #[test]
2630    fn engine_register_rejects_trailing_slash_category() {
2631        let mut engine = create_test_engine();
2632        let err = engine
2633            .register(BadCategoryWorkflow("data/etl/"))
2634            .unwrap_err();
2635        match err {
2636            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2637            other => panic!("expected InvalidWorkflow, got {other:?}"),
2638        }
2639    }
2640
2641    #[test]
2642    fn engine_register_rejects_double_slash_category() {
2643        let mut engine = create_test_engine();
2644        let err = engine
2645            .register(BadCategoryWorkflow("data//etl"))
2646            .unwrap_err();
2647        match err {
2648            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2649            other => panic!("expected InvalidWorkflow, got {other:?}"),
2650        }
2651    }
2652
2653    #[test]
2654    fn engine_register_rejects_whitespace_only_segment_category() {
2655        let mut engine = create_test_engine();
2656        let err = engine
2657            .register(BadCategoryWorkflow("data/ /etl"))
2658            .unwrap_err();
2659        match err {
2660            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2661            other => panic!("expected InvalidWorkflow, got {other:?}"),
2662        }
2663    }
2664
2665    #[test]
2666    fn engine_register_accepts_valid_nested_category() {
2667        let mut engine = create_test_engine();
2668        assert!(engine.register(CategorizedWorkflow).is_ok());
2669    }
2670
2671    #[tokio::test]
2672    async fn engine_unknown_workflow_returns_error() {
2673        let engine = create_test_engine();
2674        let result = engine
2675            .run_handler("unknown", TriggerKind::Manual, json!({}))
2676            .await;
2677        assert!(result.is_err());
2678        match result {
2679            Err(EngineError::InvalidWorkflow(msg)) => {
2680                assert!(msg.contains("no handler registered"));
2681            }
2682            _ => panic!("expected InvalidWorkflow error"),
2683        }
2684    }
2685
2686    #[tokio::test]
2687    async fn engine_enqueue_handler_creates_pending_run() {
2688        let mut engine = create_test_engine();
2689        engine.register(EchoWorkflow).unwrap();
2690
2691        let run = engine
2692            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2693            .await
2694            .unwrap();
2695        assert_eq!(run.status.state, RunStatus::Pending);
2696        assert_eq!(run.workflow_name, "echo-workflow");
2697    }
2698
2699    #[tokio::test]
2700    async fn enqueue_handler_leaves_the_run_unattributed() {
2701        let mut engine = create_test_engine();
2702        engine.register(EchoWorkflow).unwrap();
2703
2704        let run = engine
2705            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2706            .await
2707            .unwrap();
2708
2709        assert!(run.created_by.is_none());
2710    }
2711
2712    #[tokio::test]
2713    async fn enqueue_handler_with_options_records_the_author() {
2714        let mut engine = create_test_engine();
2715        engine.register(EchoWorkflow).unwrap();
2716        let actor = RunActor::User {
2717            user_id: Uuid::now_v7(),
2718        };
2719
2720        let run = engine
2721            .enqueue_handler_with_options(
2722                "echo-workflow",
2723                TriggerKind::Api,
2724                json!({}),
2725                EnqueueOptions {
2726                    created_by: Some(actor.clone()),
2727                    ..Default::default()
2728                },
2729            )
2730            .await
2731            .unwrap()
2732            .into_run();
2733
2734        assert_eq!(run.created_by, Some(actor));
2735    }
2736
2737    #[tokio::test]
2738    async fn enqueue_handler_with_options_accepts_no_author() {
2739        let mut engine = create_test_engine();
2740        engine.register(EchoWorkflow).unwrap();
2741
2742        let run = engine
2743            .enqueue_handler_with_options(
2744                "echo-workflow",
2745                TriggerKind::Cron {
2746                    schedule: "0 * * * * *".to_string(),
2747                },
2748                json!({}),
2749                EnqueueOptions::default(),
2750            )
2751            .await
2752            .unwrap()
2753            .into_run();
2754
2755        assert!(run.created_by.is_none());
2756    }
2757
2758    #[tokio::test]
2759    async fn run_handler_leaves_the_run_unattributed() {
2760        let mut engine = create_test_engine();
2761        engine.register(EchoWorkflow).unwrap();
2762
2763        let run = engine
2764            .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
2765            .await
2766            .unwrap()
2767            .run;
2768
2769        assert!(run.created_by.is_none());
2770    }
2771
2772    #[tokio::test]
2773    async fn engine_register_boxed() {
2774        let mut engine = create_test_engine();
2775        let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
2776        let result = engine.register_boxed(handler);
2777        assert!(result.is_ok());
2778        assert_eq!(engine.handler_names().len(), 1);
2779    }
2780
2781    #[tokio::test]
2782    async fn engine_store_and_provider_accessors() {
2783        let store = Arc::new(InMemoryStore::new());
2784        let inner = ClaudeCodeProvider::new();
2785        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2786            inner,
2787            "/tmp/ironflow-fixtures",
2788        ));
2789        let engine = Engine::new(store.clone(), provider.clone());
2790
2791        // Verify accessors return references
2792        let _ = engine.store();
2793        let _ = engine.provider();
2794    }
2795
2796    // -----------------------------------------------------------------------
2797    // Operation trait tests
2798    // -----------------------------------------------------------------------
2799
2800    use crate::operation::{Operation, OperationContext};
2801    use async_trait::async_trait;
2802    use ironflow_core::error::OperationError;
2803    use ironflow_store::models::StepKind;
2804
2805    struct FakeGitlabOp {
2806        project_id: u64,
2807        title: String,
2808    }
2809
2810    #[async_trait]
2811    impl Operation for FakeGitlabOp {
2812        fn kind(&self) -> &str {
2813            "gitlab"
2814        }
2815
2816        async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2817            Ok(json!({
2818                "issue_id": 42,
2819                "project_id": self.project_id,
2820                "title": self.title,
2821            }))
2822        }
2823
2824        fn input(&self) -> Option<Value> {
2825            Some(json!({
2826                "project_id": self.project_id,
2827                "title": self.title,
2828            }))
2829        }
2830    }
2831
2832    struct FailingOp;
2833
2834    #[async_trait]
2835    impl Operation for FailingOp {
2836        fn kind(&self) -> &str {
2837            "broken-service"
2838        }
2839
2840        async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2841            Err(OperationError::Http {
2842                status: None,
2843                message: "service unavailable".to_string(),
2844            })
2845        }
2846    }
2847
2848    struct OperationWorkflow;
2849
2850    impl WorkflowHandler for OperationWorkflow {
2851        fn name(&self) -> &str {
2852            "operation-workflow"
2853        }
2854
2855        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2856            Box::pin(async move {
2857                let op = FakeGitlabOp {
2858                    project_id: 123,
2859                    title: "Bug report".to_string(),
2860                };
2861                ctx.operation("create-issue", &op).await?;
2862                Ok(())
2863            })
2864        }
2865    }
2866
2867    struct FailingOperationWorkflow;
2868
2869    impl WorkflowHandler for FailingOperationWorkflow {
2870        fn name(&self) -> &str {
2871            "failing-operation-workflow"
2872        }
2873
2874        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2875            Box::pin(async move {
2876                ctx.operation("broken-call", &FailingOp).await?;
2877                Ok(())
2878            })
2879        }
2880    }
2881
2882    struct MixedWorkflow;
2883
2884    impl WorkflowHandler for MixedWorkflow {
2885        fn name(&self) -> &str {
2886            "mixed-workflow"
2887        }
2888
2889        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2890            Box::pin(async move {
2891                ctx.shell("build", ShellConfig::new("echo built")).await?;
2892                let op = FakeGitlabOp {
2893                    project_id: 456,
2894                    title: "Deploy done".to_string(),
2895                };
2896                let result = ctx.operation("notify-gitlab", &op).await?;
2897                assert_eq!(result.output["issue_id"], 42);
2898                Ok(())
2899            })
2900        }
2901    }
2902
2903    #[tokio::test]
2904    async fn operation_step_happy_path() {
2905        let mut engine = create_test_engine();
2906        engine.register(OperationWorkflow).unwrap();
2907
2908        let run = engine
2909            .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
2910            .await
2911            .unwrap()
2912            .run;
2913
2914        assert_eq!(run.status.state, RunStatus::Completed);
2915
2916        let steps = engine.store().list_steps(run.id).await.unwrap();
2917
2918        assert_eq!(steps.len(), 1);
2919        assert_eq!(steps[0].name, "create-issue");
2920        assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
2921        assert_eq!(
2922            steps[0].status.state,
2923            ironflow_store::models::StepStatus::Completed
2924        );
2925
2926        let output = steps[0].output.as_ref().unwrap();
2927        assert_eq!(output["issue_id"], 42);
2928        assert_eq!(output["project_id"], 123);
2929
2930        let input = steps[0].input.as_ref().unwrap();
2931        assert_eq!(input["project_id"], 123);
2932        assert_eq!(input["title"], "Bug report");
2933    }
2934
2935    #[tokio::test]
2936    async fn operation_step_failure_marks_run_failed() {
2937        let mut engine = create_test_engine();
2938        engine.register(FailingOperationWorkflow).unwrap();
2939
2940        let result = engine
2941            .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
2942            .await;
2943
2944        assert!(result.is_err());
2945    }
2946
2947    #[tokio::test]
2948    async fn operation_mixed_with_shell_steps() {
2949        let mut engine = create_test_engine();
2950        engine.register(MixedWorkflow).unwrap();
2951
2952        let run = engine
2953            .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
2954            .await
2955            .unwrap()
2956            .run;
2957
2958        assert_eq!(run.status.state, RunStatus::Completed);
2959
2960        let steps = engine.store().list_steps(run.id).await.unwrap();
2961
2962        assert_eq!(steps.len(), 2);
2963        assert_eq!(steps[0].kind, StepKind::Shell);
2964        assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
2965        assert_eq!(steps[0].position, 0);
2966        assert_eq!(steps[1].position, 1);
2967    }
2968
2969    // -----------------------------------------------------------------------
2970    // Approval + resume tests
2971    // -----------------------------------------------------------------------
2972
2973    use crate::config::ApprovalConfig;
2974
2975    struct SingleApprovalWorkflow;
2976
2977    impl WorkflowHandler for SingleApprovalWorkflow {
2978        fn name(&self) -> &str {
2979            "single-approval"
2980        }
2981
2982        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2983            Box::pin(async move {
2984                ctx.shell("build", ShellConfig::new("echo built")).await?;
2985                ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
2986                ctx.shell("deploy", ShellConfig::new("echo deployed"))
2987                    .await?;
2988                Ok(())
2989            })
2990        }
2991    }
2992
2993    struct DoubleApprovalWorkflow;
2994
2995    impl WorkflowHandler for DoubleApprovalWorkflow {
2996        fn name(&self) -> &str {
2997            "double-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("staging-gate", ApprovalConfig::new("Deploy staging?"))
3004                    .await?;
3005                ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3006                    .await?;
3007                ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3008                    .await?;
3009                ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3010                    .await?;
3011                Ok(())
3012            })
3013        }
3014    }
3015
3016    #[tokio::test]
3017    async fn approval_pauses_run() {
3018        let mut engine = create_test_engine();
3019        engine.register(SingleApprovalWorkflow).unwrap();
3020
3021        let run = engine
3022            .run_handler("single-approval", TriggerKind::Manual, json!({}))
3023            .await
3024            .unwrap()
3025            .run;
3026
3027        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3028
3029        let steps = engine.store().list_steps(run.id).await.unwrap();
3030        assert_eq!(steps.len(), 2); // build + approval gate
3031        assert_eq!(steps[0].kind, StepKind::Shell);
3032        assert_eq!(steps[0].status.state, StepStatus::Completed);
3033        assert_eq!(steps[1].kind, StepKind::Approval);
3034        assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3035    }
3036
3037    #[tokio::test]
3038    async fn approval_resume_completes_run() {
3039        let mut engine = create_test_engine();
3040        engine.register(SingleApprovalWorkflow).unwrap();
3041
3042        // First execution: pauses at approval
3043        let run = engine
3044            .run_handler("single-approval", TriggerKind::Manual, json!({}))
3045            .await
3046            .unwrap()
3047            .run;
3048        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3049
3050        // Simulate approval: transition to Running
3051        engine
3052            .store()
3053            .update_run_status(run.id, RunStatus::Running)
3054            .await
3055            .unwrap();
3056
3057        // Resume: replays build, skips approval, executes deploy
3058        let resumed = engine.resume_run(run.id).await.unwrap().run;
3059        assert_eq!(resumed.status.state, RunStatus::Completed);
3060
3061        let steps = engine.store().list_steps(run.id).await.unwrap();
3062        assert_eq!(steps.len(), 3); // build + approval + deploy
3063        assert_eq!(steps[0].name, "build");
3064        assert_eq!(steps[0].status.state, StepStatus::Completed);
3065        assert_eq!(steps[1].name, "gate");
3066        assert_eq!(steps[1].kind, StepKind::Approval);
3067        assert_eq!(steps[1].status.state, StepStatus::Completed);
3068        assert_eq!(steps[2].name, "deploy");
3069        assert_eq!(steps[2].status.state, StepStatus::Completed);
3070    }
3071
3072    #[tokio::test]
3073    async fn double_approval_two_resumes() {
3074        let mut engine = create_test_engine();
3075        engine.register(DoubleApprovalWorkflow).unwrap();
3076
3077        // First execution: pauses at staging-gate
3078        let run = engine
3079            .run_handler("double-approval", TriggerKind::Manual, json!({}))
3080            .await
3081            .unwrap()
3082            .run;
3083        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3084
3085        let steps = engine.store().list_steps(run.id).await.unwrap();
3086        assert_eq!(steps.len(), 2); // build + staging-gate
3087
3088        // First approval
3089        engine
3090            .store()
3091            .update_run_status(run.id, RunStatus::Running)
3092            .await
3093            .unwrap();
3094
3095        let resumed = engine.resume_run(run.id).await.unwrap().run;
3096        assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3097
3098        let steps = engine.store().list_steps(run.id).await.unwrap();
3099        assert_eq!(steps.len(), 4); // build + staging-gate + deploy-staging + prod-gate
3100
3101        // Second approval
3102        engine
3103            .store()
3104            .update_run_status(run.id, RunStatus::Running)
3105            .await
3106            .unwrap();
3107
3108        let final_run = engine.resume_run(run.id).await.unwrap().run;
3109        assert_eq!(final_run.status.state, RunStatus::Completed);
3110
3111        let steps = engine.store().list_steps(run.id).await.unwrap();
3112        assert_eq!(steps.len(), 5);
3113        assert_eq!(steps[0].name, "build");
3114        assert_eq!(steps[1].name, "staging-gate");
3115        assert_eq!(steps[2].name, "deploy-staging");
3116        assert_eq!(steps[3].name, "prod-gate");
3117        assert_eq!(steps[4].name, "deploy-prod");
3118
3119        for step in &steps {
3120            assert_eq!(step.status.state, StepStatus::Completed);
3121        }
3122    }
3123
3124    // -----------------------------------------------------------------------
3125    // fail_orphaned_steps tests
3126    // -----------------------------------------------------------------------
3127
3128    use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3129
3130    async fn create_step_with_status(
3131        store: &Arc<dyn Store>,
3132        run_id: Uuid,
3133        name: &str,
3134        position: u32,
3135        status: StepStatus,
3136    ) -> ironflow_store::models::Step {
3137        let step = store
3138            .create_step(NewStep {
3139                run_id,
3140                trace_id: step_trace_id(run_id, name, position),
3141                name: name.to_string(),
3142                kind: StepKind::Shell,
3143                position,
3144                input: None,
3145                is_error_handler: false,
3146            })
3147            .await
3148            .unwrap();
3149
3150        match status {
3151            StepStatus::Pending => {}
3152            StepStatus::Running => {
3153                store
3154                    .update_step(
3155                        step.id,
3156                        StepUpdate {
3157                            status: Some(StepStatus::Running),
3158                            ..StepUpdate::default()
3159                        },
3160                    )
3161                    .await
3162                    .unwrap();
3163            }
3164            StepStatus::Completed => {
3165                store
3166                    .update_step(
3167                        step.id,
3168                        StepUpdate {
3169                            status: Some(StepStatus::Running),
3170                            ..StepUpdate::default()
3171                        },
3172                    )
3173                    .await
3174                    .unwrap();
3175                store
3176                    .update_step(
3177                        step.id,
3178                        StepUpdate {
3179                            status: Some(StepStatus::Completed),
3180                            ..StepUpdate::default()
3181                        },
3182                    )
3183                    .await
3184                    .unwrap();
3185            }
3186            StepStatus::AwaitingApproval => {
3187                store
3188                    .update_step(
3189                        step.id,
3190                        StepUpdate {
3191                            status: Some(StepStatus::Running),
3192                            ..StepUpdate::default()
3193                        },
3194                    )
3195                    .await
3196                    .unwrap();
3197                store
3198                    .update_step(
3199                        step.id,
3200                        StepUpdate {
3201                            status: Some(StepStatus::AwaitingApproval),
3202                            ..StepUpdate::default()
3203                        },
3204                    )
3205                    .await
3206                    .unwrap();
3207            }
3208            _ => panic!("unsupported status for test helper: {status}"),
3209        }
3210
3211        store.get_step(step.id).await.unwrap().unwrap()
3212    }
3213
3214    #[tokio::test]
3215    async fn fail_orphaned_steps_marks_running_as_failed() {
3216        let engine = create_test_engine();
3217        let run = engine
3218            .store()
3219            .create_run(NewRun {
3220                created_by: None,
3221                workflow_name: "test".to_string(),
3222                trigger: TriggerKind::Manual,
3223                payload: json!({}),
3224                max_retries: 0,
3225                handler_version: None,
3226                labels: HashMap::new(),
3227                scheduled_at: None,
3228                idempotency_key: None,
3229                concurrency_key: None,
3230                max_cost_usd: None,
3231            })
3232            .await
3233            .unwrap()
3234            .into_run();
3235
3236        let step = create_step_with_status(
3237            engine.store(),
3238            run.id,
3239            "running-step",
3240            0,
3241            StepStatus::Running,
3242        )
3243        .await;
3244
3245        engine
3246            .fail_orphaned_steps(run.id, "parent run timed out")
3247            .await
3248            .unwrap();
3249
3250        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3251        assert_eq!(updated.status.state, StepStatus::Failed);
3252        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3253        assert!(updated.completed_at.is_some());
3254    }
3255
3256    #[tokio::test]
3257    async fn fail_orphaned_steps_marks_pending_as_skipped() {
3258        let engine = create_test_engine();
3259        let run = engine
3260            .store()
3261            .create_run(NewRun {
3262                created_by: None,
3263                workflow_name: "test".to_string(),
3264                trigger: TriggerKind::Manual,
3265                payload: json!({}),
3266                max_retries: 0,
3267                handler_version: None,
3268                labels: HashMap::new(),
3269                scheduled_at: None,
3270                idempotency_key: None,
3271                concurrency_key: None,
3272                max_cost_usd: None,
3273            })
3274            .await
3275            .unwrap()
3276            .into_run();
3277
3278        let step = create_step_with_status(
3279            engine.store(),
3280            run.id,
3281            "pending-step",
3282            0,
3283            StepStatus::Pending,
3284        )
3285        .await;
3286
3287        engine
3288            .fail_orphaned_steps(run.id, "parent run timed out")
3289            .await
3290            .unwrap();
3291
3292        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3293        assert_eq!(updated.status.state, StepStatus::Skipped);
3294        assert!(updated.error.is_none());
3295        assert!(updated.completed_at.is_some());
3296    }
3297
3298    #[tokio::test]
3299    async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
3300        let engine = create_test_engine();
3301        let run = engine
3302            .store()
3303            .create_run(NewRun {
3304                created_by: None,
3305                workflow_name: "test".to_string(),
3306                trigger: TriggerKind::Manual,
3307                payload: json!({}),
3308                max_retries: 0,
3309                handler_version: None,
3310                labels: HashMap::new(),
3311                scheduled_at: None,
3312                idempotency_key: None,
3313                concurrency_key: None,
3314                max_cost_usd: None,
3315            })
3316            .await
3317            .unwrap()
3318            .into_run();
3319
3320        let step = create_step_with_status(
3321            engine.store(),
3322            run.id,
3323            "approval-step",
3324            0,
3325            StepStatus::AwaitingApproval,
3326        )
3327        .await;
3328
3329        engine
3330            .fail_orphaned_steps(run.id, "parent run timed out")
3331            .await
3332            .unwrap();
3333
3334        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3335        assert_eq!(updated.status.state, StepStatus::Failed);
3336        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3337        assert!(updated.completed_at.is_some());
3338    }
3339
3340    #[tokio::test]
3341    async fn fail_orphaned_steps_skips_terminal_steps() {
3342        let engine = create_test_engine();
3343        let run = engine
3344            .store()
3345            .create_run(NewRun {
3346                created_by: None,
3347                workflow_name: "test".to_string(),
3348                trigger: TriggerKind::Manual,
3349                payload: json!({}),
3350                max_retries: 0,
3351                handler_version: None,
3352                labels: HashMap::new(),
3353                scheduled_at: None,
3354                idempotency_key: None,
3355                concurrency_key: None,
3356                max_cost_usd: None,
3357            })
3358            .await
3359            .unwrap()
3360            .into_run();
3361
3362        let completed_step =
3363            create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
3364        let running_step =
3365            create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
3366                .await;
3367
3368        engine
3369            .fail_orphaned_steps(run.id, "parent run timed out")
3370            .await
3371            .unwrap();
3372
3373        let completed = engine
3374            .store()
3375            .get_step(completed_step.id)
3376            .await
3377            .unwrap()
3378            .unwrap();
3379        assert_eq!(completed.status.state, StepStatus::Completed);
3380
3381        let failed = engine
3382            .store()
3383            .get_step(running_step.id)
3384            .await
3385            .unwrap()
3386            .unwrap();
3387        assert_eq!(failed.status.state, StepStatus::Failed);
3388    }
3389
3390    #[tokio::test]
3391    async fn fail_orphaned_steps_mixed_states() {
3392        let engine = create_test_engine();
3393        let run = engine
3394            .store()
3395            .create_run(NewRun {
3396                created_by: None,
3397                workflow_name: "test".to_string(),
3398                trigger: TriggerKind::Manual,
3399                payload: json!({}),
3400                max_retries: 0,
3401                handler_version: None,
3402                labels: HashMap::new(),
3403                scheduled_at: None,
3404                idempotency_key: None,
3405                concurrency_key: None,
3406                max_cost_usd: None,
3407            })
3408            .await
3409            .unwrap()
3410            .into_run();
3411
3412        let s_completed =
3413            create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
3414                .await;
3415        let s_running =
3416            create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
3417        let s_pending =
3418            create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
3419
3420        engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
3421
3422        let r_completed = engine
3423            .store()
3424            .get_step(s_completed.id)
3425            .await
3426            .unwrap()
3427            .unwrap();
3428        assert_eq!(r_completed.status.state, StepStatus::Completed);
3429
3430        let r_running = engine
3431            .store()
3432            .get_step(s_running.id)
3433            .await
3434            .unwrap()
3435            .unwrap();
3436        assert_eq!(r_running.status.state, StepStatus::Failed);
3437        assert_eq!(r_running.error.as_deref(), Some("timeout"));
3438
3439        let r_pending = engine
3440            .store()
3441            .get_step(s_pending.id)
3442            .await
3443            .unwrap()
3444            .unwrap();
3445        assert_eq!(r_pending.status.state, StepStatus::Skipped);
3446        assert!(r_pending.error.is_none());
3447    }
3448
3449    #[tokio::test]
3450    async fn fail_orphaned_steps_no_steps_is_noop() {
3451        let engine = create_test_engine();
3452        let run = engine
3453            .store()
3454            .create_run(NewRun {
3455                created_by: None,
3456                workflow_name: "test".to_string(),
3457                trigger: TriggerKind::Manual,
3458                payload: json!({}),
3459                max_retries: 0,
3460                handler_version: None,
3461                labels: HashMap::new(),
3462                scheduled_at: None,
3463                idempotency_key: None,
3464                concurrency_key: None,
3465                max_cost_usd: None,
3466            })
3467            .await
3468            .unwrap()
3469            .into_run();
3470
3471        let result = engine.fail_orphaned_steps(run.id, "timeout").await;
3472        assert!(result.is_ok());
3473    }
3474
3475    #[tokio::test]
3476    async fn fail_orphaned_steps_preserves_existing_error() {
3477        let engine = create_test_engine();
3478        let run = engine
3479            .store()
3480            .create_run(NewRun {
3481                created_by: None,
3482                workflow_name: "test".to_string(),
3483                trigger: TriggerKind::Manual,
3484                payload: json!({}),
3485                max_retries: 0,
3486                handler_version: None,
3487                labels: HashMap::new(),
3488                scheduled_at: None,
3489                idempotency_key: None,
3490                concurrency_key: None,
3491                max_cost_usd: None,
3492            })
3493            .await
3494            .unwrap()
3495            .into_run();
3496
3497        let step_with_error = create_step_with_status(
3498            engine.store(),
3499            run.id,
3500            "already-errored",
3501            0,
3502            StepStatus::Running,
3503        )
3504        .await;
3505
3506        engine
3507            .store()
3508            .update_step(
3509                step_with_error.id,
3510                StepUpdate {
3511                    error: Some("real error from provider".to_string()),
3512                    ..StepUpdate::default()
3513                },
3514            )
3515            .await
3516            .unwrap();
3517
3518        let step_no_error = create_step_with_status(
3519            engine.store(),
3520            run.id,
3521            "no-error-yet",
3522            1,
3523            StepStatus::Running,
3524        )
3525        .await;
3526
3527        engine
3528            .fail_orphaned_steps(run.id, "parent run failed")
3529            .await
3530            .unwrap();
3531
3532        let updated_with = engine
3533            .store()
3534            .get_step(step_with_error.id)
3535            .await
3536            .unwrap()
3537            .unwrap();
3538        assert_eq!(updated_with.status.state, StepStatus::Failed);
3539        assert_eq!(
3540            updated_with.error.as_deref(),
3541            Some("real error from provider"),
3542        );
3543
3544        let updated_without = engine
3545            .store()
3546            .get_step(step_no_error.id)
3547            .await
3548            .unwrap()
3549            .unwrap();
3550        assert_eq!(updated_without.status.state, StepStatus::Failed);
3551        assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
3552    }
3553}