Skip to main content

ironflow_engine/
engine.rs

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