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