Skip to main content

ironflow_engine/
engine.rs

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