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