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