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, Run, RunActor, RunCreation, RunFilter,
30    RunStatus, RunUpdate, SignalInsert, SignalStepResolution, StepStatus, StepUpdate, TriggerKind,
31    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::SignalWaiting {
2137                run_id: wait_run_id,
2138                step_id,
2139                ref step_name,
2140                ref name,
2141                ref key,
2142                deadline_at,
2143            }) => {
2144                final_status = RunStatus::Sleeping;
2145                // Atomic with the step lock: a signal delivered since the step
2146                // opened leaves the run due right away instead of until the
2147                // deadline.
2148                let waiting = self
2149                    .store
2150                    .suspend_run_on_signal(run_id, step_id, deadline_at)
2151                    .await?;
2152                final_run = self
2153                    .store
2154                    .update_run_returning(
2155                        run_id,
2156                        RunUpdate {
2157                            cost_usd: Some(ctx.total_cost_usd()),
2158                            duration_ms: Some(total_duration),
2159                            ..RunUpdate::default()
2160                        },
2161                    )
2162                    .await?;
2163
2164                if waiting {
2165                    self.event_publisher
2166                        .publish(Event::SignalAwaited(SignalAwaitedEvent {
2167                            run_id: wait_run_id,
2168                            step_id,
2169                            step_name: step_name.clone(),
2170                            name: name.clone(),
2171                            key: key.clone(),
2172                            deadline_at,
2173                            at: Utc::now(),
2174                        }));
2175                }
2176
2177                info!(
2178                    run_id = %wait_run_id,
2179                    step_id = %step_id,
2180                    signal = %name,
2181                    key = %key,
2182                    deadline_at = %deadline_at,
2183                    waiting,
2184                    "run sleeping until a signal arrives"
2185                );
2186            }
2187            Err(err) => {
2188                // A guardrail stop (budget or workflow guard) is deliberate,
2189                // not a breakage: the run is cancelled, never failed and
2190                // never replayed.
2191                let guardrail_stop = matches!(
2192                    err,
2193                    EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2194                );
2195
2196                final_status = if guardrail_stop {
2197                    if let Err(store_err) = self
2198                        .store
2199                        .update_run(
2200                            run_id,
2201                            RunUpdate {
2202                                status: Some(RunStatus::Cancelled),
2203                                error: Some(err.to_string()),
2204                                cost_usd: Some(ctx.total_cost_usd()),
2205                                duration_ms: Some(total_duration),
2206                                completed_at: Some(completed_at),
2207                                output: ctx.output().cloned(),
2208                                ..RunUpdate::default()
2209                            },
2210                        )
2211                        .await
2212                    {
2213                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2214                    }
2215                    if let Err(cleanup_err) = self
2216                        .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2217                        .await
2218                    {
2219                        error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2220                    }
2221                    RunStatus::Cancelled
2222                } else {
2223                    // Written before the failure so a failed run keeps the
2224                    // verdict the handler set before returning its error.
2225                    if let Some(output) = ctx.output()
2226                        && let Err(store_err) = self
2227                            .store
2228                            .update_run(
2229                                run_id,
2230                                RunUpdate {
2231                                    output: Some(output.clone()),
2232                                    ..RunUpdate::default()
2233                                },
2234                            )
2235                            .await
2236                    {
2237                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run output");
2238                    }
2239                    self.fail_or_schedule_retry(
2240                        run_id,
2241                        &err.to_string(),
2242                        is_run_retryable(&err),
2243                        Some(ctx.total_cost_usd()),
2244                        Some(total_duration),
2245                    )
2246                    .await
2247                    .unwrap_or_else(|store_err| {
2248                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2249                        RunStatus::Failed
2250                    })
2251                };
2252
2253                if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2254                    self.on_run_budget_exceeded(workflow_name, run_id, &err);
2255                }
2256
2257                error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2258
2259                self.publish_run_status_changed(
2260                    workflow_name,
2261                    run_id,
2262                    final_status,
2263                    Some(err.to_string()),
2264                    ctx,
2265                    total_duration,
2266                    run_labels,
2267                );
2268
2269                #[cfg(feature = "prometheus")]
2270                self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2271
2272                return Err(err);
2273            }
2274        }
2275
2276        self.publish_run_status_changed(
2277            workflow_name,
2278            run_id,
2279            final_status,
2280            None,
2281            ctx,
2282            total_duration,
2283            run_labels,
2284        );
2285
2286        #[cfg(feature = "prometheus")]
2287        self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2288
2289        Ok(WorkflowResult {
2290            run: final_run,
2291            steps: ctx.step_results().to_vec(),
2292        })
2293    }
2294
2295    /// Publish [`Event::ApprovalRequested`] for an approval gate that opened
2296    /// on `run_id`, with the requirement recorded when the gate opened.
2297    async fn publish_approval_requested(
2298        &self,
2299        run_id: Uuid,
2300        step_id: Uuid,
2301        message: &str,
2302    ) -> Result<(), EngineError> {
2303        let requirement = self
2304            .store
2305            .get_step(step_id)
2306            .await?
2307            .and_then(|s| s.approval_requirement);
2308        self.event_publisher
2309            .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2310                run_id,
2311                step_id,
2312                message: message.to_string(),
2313                requirement,
2314                at: Utc::now(),
2315            }));
2316        Ok(())
2317    }
2318
2319    /// Fail every ancestor of a child run, closest first.
2320    ///
2321    /// Used when a gate inside a sub-workflow is rejected: the child run
2322    /// fails, and the runs suspended with it (its parent, up to the root)
2323    /// must not stay suspended on a child that will never resume. Each
2324    /// ancestor goes through [`fail_or_schedule_retry`](Self::fail_or_schedule_retry)
2325    /// without a retry, which also fails its open `Workflow` step. A run that
2326    /// is not a child of a sub-workflow has no ancestor: nothing happens.
2327    ///
2328    /// # Errors
2329    ///
2330    /// Returns [`EngineError::Store`] if a run of the chain does not exist or
2331    /// its failure cannot be persisted.
2332    ///
2333    /// # Examples
2334    ///
2335    /// ```no_run
2336    /// use ironflow_engine::engine::Engine;
2337    /// use ironflow_engine::error::EngineError;
2338    /// use uuid::Uuid;
2339    ///
2340    /// # async fn example(engine: &Engine, child_run_id: Uuid) -> Result<(), EngineError> {
2341    /// engine
2342    ///     .fail_ancestors(child_run_id, "approval rejected in a sub-workflow")
2343    ///     .await?;
2344    /// # Ok(())
2345    /// # }
2346    /// ```
2347    pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2348        let mut current = self
2349            .store
2350            .get_run(run_id)
2351            .await?
2352            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2353        // Labels are data: a chain that loops back on itself stops there.
2354        let mut visited = HashSet::from([run_id]);
2355
2356        while let Some(parent_id) = chain_parent(&current) {
2357            if !visited.insert(parent_id) {
2358                break;
2359            }
2360            let status = self
2361                .fail_or_schedule_retry(parent_id, reason, false, None, None)
2362                .await?;
2363            info!(
2364                run_id = %run_id,
2365                ancestor_run_id = %parent_id,
2366                status = %status,
2367                "ancestor run failed with its child"
2368            );
2369            current = self
2370                .store
2371                .get_run(parent_id)
2372                .await?
2373                .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2374        }
2375
2376        Ok(())
2377    }
2378
2379    /// Emit Prometheus metrics for a completed run.
2380    #[cfg(feature = "prometheus")]
2381    fn emit_run_metrics(
2382        &self,
2383        workflow_name: &str,
2384        status: RunStatus,
2385        duration_ms: u64,
2386        ctx: &WorkflowContext,
2387    ) {
2388        let status_str = status.to_string();
2389        let wf = workflow_name.to_string();
2390
2391        counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2392        histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2393            .record(duration_ms as f64 / 1000.0);
2394        histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2395            ctx.total_cost_usd()
2396                .to_string()
2397                .parse::<f64>()
2398                .unwrap_or(0.0),
2399        );
2400        gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2401    }
2402
2403    /// Record the metric and publish the audit event for a run that hit its
2404    /// cost cap.
2405    ///
2406    /// A non-[`RunBudgetExceeded`](EngineError::RunBudgetExceeded) error is
2407    /// ignored, so callers can pass the error unconditionally.
2408    fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2409        let EngineError::RunBudgetExceeded {
2410            limit_usd,
2411            spent_usd,
2412            step_budget_usd,
2413            ..
2414        } = err
2415        else {
2416            return;
2417        };
2418
2419        #[cfg(feature = "prometheus")]
2420        counter!(
2421            RUN_BUDGET_EXCEEDED_TOTAL,
2422            "workflow" => workflow_name.to_string(),
2423            "scope" => "run",
2424        )
2425        .increment(1);
2426
2427        self.event_publisher
2428            .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2429                run_id,
2430                workflow_name: workflow_name.to_string(),
2431                limit_usd: *limit_usd,
2432                spent_usd: *spent_usd,
2433                step_budget_usd: *step_budget_usd,
2434                at: Utc::now(),
2435            }));
2436    }
2437
2438    /// Publish a run status changed event to all registered subscribers.
2439    ///
2440    /// `from` is always `Running` because `finalize_run` is only called
2441    /// from a running state.
2442    #[allow(clippy::too_many_arguments)]
2443    fn publish_run_status_changed(
2444        &self,
2445        workflow_name: &str,
2446        run_id: Uuid,
2447        to: RunStatus,
2448        error: Option<String>,
2449        ctx: &WorkflowContext,
2450        duration_ms: u64,
2451        labels: HashMap<String, String>,
2452    ) {
2453        let now = Utc::now();
2454        let cost_usd = ctx.total_cost_usd();
2455        let wf = workflow_name.to_string();
2456
2457        self.event_publisher
2458            .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2459                run_id,
2460                workflow_name: wf.clone(),
2461                from: RunStatus::Running,
2462                to,
2463                error: error.clone(),
2464                cost_usd,
2465                duration_ms,
2466                labels: labels.clone(),
2467                at: now,
2468            }));
2469
2470        if to == RunStatus::Failed {
2471            self.event_publisher
2472                .publish(Event::RunFailed(RunFailedEvent {
2473                    run_id,
2474                    workflow_name: wf,
2475                    error,
2476                    cost_usd,
2477                    duration_ms,
2478                    labels,
2479                    at: now,
2480                }));
2481        }
2482    }
2483}
2484
2485impl fmt::Debug for Engine {
2486    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2487        f.debug_struct("Engine")
2488            .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2489            .finish_non_exhaustive()
2490    }
2491}
2492
2493#[cfg(test)]
2494mod tests {
2495    use super::*;
2496    use crate::config::ShellConfig;
2497    use crate::handler::{HandlerFuture, WorkflowHandler};
2498    use ironflow_core::providers::claude::ClaudeCodeProvider;
2499    use ironflow_core::providers::record_replay::RecordReplayProvider;
2500    use ironflow_store::memory::InMemoryStore;
2501    use ironflow_store::models::StepStatus;
2502    use serde_json::json;
2503
2504    // Test handler that echoes a message via shell
2505    struct EchoWorkflow;
2506
2507    impl WorkflowHandler for EchoWorkflow {
2508        fn name(&self) -> &str {
2509            "echo-workflow"
2510        }
2511
2512        fn describe(&self) -> WorkflowInfo {
2513            WorkflowInfo {
2514                description: "A simple workflow that echoes hello".to_string(),
2515                source_code: None,
2516                sub_workflows: Vec::new(),
2517                category: None,
2518                version: self.version().map(str::to_string),
2519                compatible_versions: Vec::new(),
2520                input_schema: None,
2521                default_labels: HashMap::new(),
2522                schedule: self.schedule().cloned(),
2523                default_max_cost_usd: self.default_max_cost_usd(),
2524            }
2525        }
2526
2527        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2528            Box::pin(async move {
2529                ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2530                Ok(())
2531            })
2532        }
2533    }
2534
2535    // Test handler that fails
2536    struct FailingWorkflow;
2537
2538    impl WorkflowHandler for FailingWorkflow {
2539        fn name(&self) -> &str {
2540            "failing-workflow"
2541        }
2542
2543        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2544            Box::pin(async move {
2545                ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2546                Ok(())
2547            })
2548        }
2549    }
2550
2551    fn create_test_engine() -> Engine {
2552        let store = Arc::new(InMemoryStore::new());
2553        let inner = ClaudeCodeProvider::new();
2554        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2555            inner,
2556            "/tmp/ironflow-fixtures",
2557        ));
2558        Engine::new(store, provider)
2559    }
2560
2561    #[test]
2562    fn engine_new_creates_instance() {
2563        let engine = create_test_engine();
2564        assert_eq!(engine.handler_names().len(), 0);
2565    }
2566
2567    #[test]
2568    fn execution_mode_defaults_to_local() {
2569        let engine = create_test_engine();
2570        assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2571    }
2572
2573    #[test]
2574    fn with_execution_mode_overrides_the_default() {
2575        let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2576        assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2577    }
2578
2579    #[test]
2580    fn engine_register_handler() {
2581        let mut engine = create_test_engine();
2582        let result = engine.register(EchoWorkflow);
2583        assert!(result.is_ok());
2584        assert_eq!(engine.handler_names().len(), 1);
2585        assert!(engine.handler_names().contains(&"echo-workflow"));
2586    }
2587
2588    #[test]
2589    fn engine_register_duplicate_returns_error() {
2590        let mut engine = create_test_engine();
2591        engine.register(EchoWorkflow).unwrap();
2592        let result = engine.register(EchoWorkflow);
2593        assert!(result.is_err());
2594    }
2595
2596    #[test]
2597    fn engine_get_handler_found() {
2598        let mut engine = create_test_engine();
2599        engine.register(EchoWorkflow).unwrap();
2600        let handler = engine.get_handler("echo-workflow");
2601        assert!(handler.is_some());
2602    }
2603
2604    #[test]
2605    fn engine_get_handler_not_found() {
2606        let engine = create_test_engine();
2607        let handler = engine.get_handler("nonexistent");
2608        assert!(handler.is_none());
2609    }
2610
2611    #[test]
2612    fn engine_handler_names_lists_all() {
2613        let mut engine = create_test_engine();
2614        engine.register(EchoWorkflow).unwrap();
2615        engine.register(FailingWorkflow).unwrap();
2616        let names = engine.handler_names();
2617        assert_eq!(names.len(), 2);
2618        assert!(names.contains(&"echo-workflow"));
2619        assert!(names.contains(&"failing-workflow"));
2620    }
2621
2622    #[test]
2623    fn engine_handler_info_returns_description() {
2624        let mut engine = create_test_engine();
2625        engine.register(EchoWorkflow).unwrap();
2626        let info = engine.handler_info("echo-workflow");
2627        assert!(info.is_some());
2628        let info = info.unwrap();
2629        assert_eq!(info.description, "A simple workflow that echoes hello");
2630    }
2631
2632    struct CategorizedWorkflow;
2633
2634    impl WorkflowHandler for CategorizedWorkflow {
2635        fn name(&self) -> &str {
2636            "categorized"
2637        }
2638        fn category(&self) -> Option<&str> {
2639            Some("data/etl")
2640        }
2641        fn execute<'a>(
2642            &'a self,
2643            _ctx: &'a mut WorkflowContext,
2644        ) -> crate::handler::HandlerFuture<'a> {
2645            Box::pin(async move { Ok(()) })
2646        }
2647    }
2648
2649    #[test]
2650    fn engine_default_describe_propagates_category() {
2651        let mut engine = create_test_engine();
2652        engine.register(CategorizedWorkflow).unwrap();
2653        let info = engine.handler_info("categorized").unwrap();
2654        assert_eq!(info.category.as_deref(), Some("data/etl"));
2655    }
2656
2657    #[test]
2658    fn engine_default_describe_without_category() {
2659        let mut engine = create_test_engine();
2660        engine.register(EchoWorkflow).unwrap();
2661        let info = engine.handler_info("echo-workflow").unwrap();
2662        assert!(info.category.is_none());
2663    }
2664
2665    // -----------------------------------------------------------------------
2666    // Schedule tests
2667    // -----------------------------------------------------------------------
2668
2669    struct ScheduledWorkflow {
2670        schedule: CronSchedule,
2671    }
2672
2673    impl ScheduledWorkflow {
2674        fn new() -> Self {
2675            Self {
2676                schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2677            }
2678        }
2679    }
2680
2681    impl WorkflowHandler for ScheduledWorkflow {
2682        fn name(&self) -> &str {
2683            "scheduled"
2684        }
2685        fn schedule(&self) -> Option<&CronSchedule> {
2686            Some(&self.schedule)
2687        }
2688        fn execute<'a>(
2689            &'a self,
2690            _ctx: &'a mut WorkflowContext,
2691        ) -> crate::handler::HandlerFuture<'a> {
2692            Box::pin(async move { Ok(()) })
2693        }
2694    }
2695
2696    #[test]
2697    fn engine_default_describe_propagates_schedule() {
2698        let mut engine = create_test_engine();
2699        engine.register(ScheduledWorkflow::new()).unwrap();
2700        let info = engine.handler_info("scheduled").unwrap();
2701        assert_eq!(
2702            info.schedule.as_ref().map(|s| s.as_str()),
2703            Some("0 0 * * * *")
2704        );
2705    }
2706
2707    #[test]
2708    fn engine_default_describe_without_schedule() {
2709        let mut engine = create_test_engine();
2710        engine.register(EchoWorkflow).unwrap();
2711        let info = engine.handler_info("echo-workflow").unwrap();
2712        assert!(info.schedule.is_none());
2713    }
2714
2715    #[test]
2716    fn scheduled_handlers_returns_only_scheduled() {
2717        let mut engine = create_test_engine();
2718        engine.register(EchoWorkflow).unwrap();
2719        engine.register(ScheduledWorkflow::new()).unwrap();
2720        engine.register(FailingWorkflow).unwrap();
2721
2722        let scheduled = engine.scheduled_handlers();
2723        assert_eq!(scheduled.len(), 1);
2724        assert_eq!(scheduled[0].0, "scheduled");
2725        assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
2726    }
2727
2728    #[test]
2729    fn scheduled_handlers_empty_when_none_scheduled() {
2730        let mut engine = create_test_engine();
2731        engine.register(EchoWorkflow).unwrap();
2732        engine.register(FailingWorkflow).unwrap();
2733
2734        let scheduled = engine.scheduled_handlers();
2735        assert!(scheduled.is_empty());
2736    }
2737
2738    struct BadCategoryWorkflow(&'static str);
2739
2740    impl WorkflowHandler for BadCategoryWorkflow {
2741        fn name(&self) -> &str {
2742            "bad-category"
2743        }
2744        fn category(&self) -> Option<&str> {
2745            Some(self.0)
2746        }
2747        fn execute<'a>(
2748            &'a self,
2749            _ctx: &'a mut WorkflowContext,
2750        ) -> crate::handler::HandlerFuture<'a> {
2751            Box::pin(async move { Ok(()) })
2752        }
2753    }
2754
2755    #[test]
2756    fn engine_register_rejects_empty_category() {
2757        let mut engine = create_test_engine();
2758        let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2759        match err {
2760            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2761            other => panic!("expected InvalidWorkflow, got {other:?}"),
2762        }
2763    }
2764
2765    #[test]
2766    fn engine_register_rejects_leading_slash_category() {
2767        let mut engine = create_test_engine();
2768        let err = engine
2769            .register(BadCategoryWorkflow("/data/etl"))
2770            .unwrap_err();
2771        match err {
2772            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2773            other => panic!("expected InvalidWorkflow, got {other:?}"),
2774        }
2775    }
2776
2777    #[test]
2778    fn engine_register_rejects_trailing_slash_category() {
2779        let mut engine = create_test_engine();
2780        let err = engine
2781            .register(BadCategoryWorkflow("data/etl/"))
2782            .unwrap_err();
2783        match err {
2784            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2785            other => panic!("expected InvalidWorkflow, got {other:?}"),
2786        }
2787    }
2788
2789    #[test]
2790    fn engine_register_rejects_double_slash_category() {
2791        let mut engine = create_test_engine();
2792        let err = engine
2793            .register(BadCategoryWorkflow("data//etl"))
2794            .unwrap_err();
2795        match err {
2796            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2797            other => panic!("expected InvalidWorkflow, got {other:?}"),
2798        }
2799    }
2800
2801    #[test]
2802    fn engine_register_rejects_whitespace_only_segment_category() {
2803        let mut engine = create_test_engine();
2804        let err = engine
2805            .register(BadCategoryWorkflow("data/ /etl"))
2806            .unwrap_err();
2807        match err {
2808            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2809            other => panic!("expected InvalidWorkflow, got {other:?}"),
2810        }
2811    }
2812
2813    #[test]
2814    fn engine_register_accepts_valid_nested_category() {
2815        let mut engine = create_test_engine();
2816        assert!(engine.register(CategorizedWorkflow).is_ok());
2817    }
2818
2819    #[tokio::test]
2820    async fn engine_unknown_workflow_returns_error() {
2821        let engine = create_test_engine();
2822        let result = engine
2823            .run_handler("unknown", TriggerKind::Manual, json!({}))
2824            .await;
2825        assert!(result.is_err());
2826        match result {
2827            Err(EngineError::InvalidWorkflow(msg)) => {
2828                assert!(msg.contains("no handler registered"));
2829            }
2830            _ => panic!("expected InvalidWorkflow error"),
2831        }
2832    }
2833
2834    #[tokio::test]
2835    async fn engine_enqueue_handler_creates_pending_run() {
2836        let mut engine = create_test_engine();
2837        engine.register(EchoWorkflow).unwrap();
2838
2839        let run = engine
2840            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2841            .await
2842            .unwrap();
2843        assert_eq!(run.status.state, RunStatus::Pending);
2844        assert_eq!(run.workflow_name, "echo-workflow");
2845    }
2846
2847    #[tokio::test]
2848    async fn enqueue_handler_leaves_the_run_unattributed() {
2849        let mut engine = create_test_engine();
2850        engine.register(EchoWorkflow).unwrap();
2851
2852        let run = engine
2853            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2854            .await
2855            .unwrap();
2856
2857        assert!(run.created_by.is_none());
2858    }
2859
2860    #[tokio::test]
2861    async fn enqueue_handler_with_options_records_the_author() {
2862        let mut engine = create_test_engine();
2863        engine.register(EchoWorkflow).unwrap();
2864        let actor = RunActor::User {
2865            user_id: Uuid::now_v7(),
2866        };
2867
2868        let run = engine
2869            .enqueue_handler_with_options(
2870                "echo-workflow",
2871                TriggerKind::Api,
2872                json!({}),
2873                EnqueueOptions {
2874                    created_by: Some(actor.clone()),
2875                    ..Default::default()
2876                },
2877            )
2878            .await
2879            .unwrap()
2880            .into_run();
2881
2882        assert_eq!(run.created_by, Some(actor));
2883    }
2884
2885    #[tokio::test]
2886    async fn enqueue_handler_with_options_accepts_no_author() {
2887        let mut engine = create_test_engine();
2888        engine.register(EchoWorkflow).unwrap();
2889
2890        let run = engine
2891            .enqueue_handler_with_options(
2892                "echo-workflow",
2893                TriggerKind::Cron {
2894                    schedule: "0 * * * * *".to_string(),
2895                },
2896                json!({}),
2897                EnqueueOptions::default(),
2898            )
2899            .await
2900            .unwrap()
2901            .into_run();
2902
2903        assert!(run.created_by.is_none());
2904    }
2905
2906    #[tokio::test]
2907    async fn enqueue_handler_with_options_stores_concurrency_limits() {
2908        let mut engine = create_test_engine();
2909        engine.register(EchoWorkflow).unwrap();
2910        let limits = vec![
2911            ConcurrencyLimit::new("repo:acme", 2),
2912            ConcurrencyLimit::new("tenant:42", 5),
2913        ];
2914
2915        let run = engine
2916            .enqueue_handler_with_options(
2917                "echo-workflow",
2918                TriggerKind::Api,
2919                json!({}),
2920                EnqueueOptions {
2921                    concurrency_limits: limits.clone(),
2922                    ..Default::default()
2923                },
2924            )
2925            .await
2926            .unwrap()
2927            .into_run();
2928
2929        assert_eq!(run.concurrency_limits, limits);
2930    }
2931
2932    #[tokio::test]
2933    async fn enqueue_rejects_invalid_concurrency_limits() {
2934        let mut engine = create_test_engine();
2935        engine.register(EchoWorkflow).unwrap();
2936
2937        let invalid = [
2938            vec![ConcurrencyLimit::new("repo:acme", 0)],
2939            vec![ConcurrencyLimit::new("", 1)],
2940            vec![
2941                ConcurrencyLimit::new("repo:acme", 1),
2942                ConcurrencyLimit::new("repo:acme", 2),
2943            ],
2944        ];
2945        for concurrency_limits in invalid {
2946            let err = engine
2947                .enqueue_handler_with_options(
2948                    "echo-workflow",
2949                    TriggerKind::Api,
2950                    json!({}),
2951                    EnqueueOptions {
2952                        concurrency_limits,
2953                        ..Default::default()
2954                    },
2955                )
2956                .await
2957                .unwrap_err();
2958            assert!(
2959                matches!(err, EngineError::InvalidConcurrencyLimit(_)),
2960                "{err:?}"
2961            );
2962        }
2963
2964        // Validated before the handler lookup.
2965        let err = engine
2966            .enqueue_handler_with_options(
2967                "not-registered",
2968                TriggerKind::Api,
2969                json!({}),
2970                EnqueueOptions {
2971                    concurrency_limits: vec![ConcurrencyLimit::new("repo:acme", 0)],
2972                    ..Default::default()
2973                },
2974            )
2975            .await
2976            .unwrap_err();
2977        assert!(
2978            matches!(err, EngineError::InvalidConcurrencyLimit(_)),
2979            "{err:?}"
2980        );
2981
2982        let page = engine
2983            .store()
2984            .list_runs(RunFilter::default(), 1, 10)
2985            .await
2986            .unwrap();
2987        assert_eq!(page.total, 0, "no run may be created");
2988    }
2989
2990    #[tokio::test]
2991    async fn run_handler_leaves_the_run_unattributed() {
2992        let mut engine = create_test_engine();
2993        engine.register(EchoWorkflow).unwrap();
2994
2995        let run = engine
2996            .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
2997            .await
2998            .unwrap()
2999            .run;
3000
3001        assert!(run.created_by.is_none());
3002    }
3003
3004    #[tokio::test]
3005    async fn engine_register_boxed() {
3006        let mut engine = create_test_engine();
3007        let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
3008        let result = engine.register_boxed(handler);
3009        assert!(result.is_ok());
3010        assert_eq!(engine.handler_names().len(), 1);
3011    }
3012
3013    #[tokio::test]
3014    async fn engine_store_and_provider_accessors() {
3015        let store = Arc::new(InMemoryStore::new());
3016        let inner = ClaudeCodeProvider::new();
3017        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
3018            inner,
3019            "/tmp/ironflow-fixtures",
3020        ));
3021        let engine = Engine::new(store.clone(), provider.clone());
3022
3023        // Verify accessors return references
3024        let _ = engine.store();
3025        let _ = engine.provider();
3026    }
3027
3028    // -----------------------------------------------------------------------
3029    // Operation trait tests
3030    // -----------------------------------------------------------------------
3031
3032    use crate::operation::{Operation, OperationContext};
3033    use async_trait::async_trait;
3034    use ironflow_core::error::OperationError;
3035    use ironflow_store::models::StepKind;
3036
3037    struct FakeGitlabOp {
3038        project_id: u64,
3039        title: String,
3040    }
3041
3042    #[async_trait]
3043    impl Operation for FakeGitlabOp {
3044        fn kind(&self) -> &str {
3045            "gitlab"
3046        }
3047
3048        async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3049            Ok(json!({
3050                "issue_id": 42,
3051                "project_id": self.project_id,
3052                "title": self.title,
3053            }))
3054        }
3055
3056        fn input(&self) -> Option<Value> {
3057            Some(json!({
3058                "project_id": self.project_id,
3059                "title": self.title,
3060            }))
3061        }
3062    }
3063
3064    struct FailingOp;
3065
3066    #[async_trait]
3067    impl Operation for FailingOp {
3068        fn kind(&self) -> &str {
3069            "broken-service"
3070        }
3071
3072        async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3073            Err(OperationError::Http {
3074                status: None,
3075                message: "service unavailable".to_string(),
3076            })
3077        }
3078    }
3079
3080    struct OperationWorkflow;
3081
3082    impl WorkflowHandler for OperationWorkflow {
3083        fn name(&self) -> &str {
3084            "operation-workflow"
3085        }
3086
3087        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3088            Box::pin(async move {
3089                let op = FakeGitlabOp {
3090                    project_id: 123,
3091                    title: "Bug report".to_string(),
3092                };
3093                ctx.operation("create-issue", &op).await?;
3094                Ok(())
3095            })
3096        }
3097    }
3098
3099    struct FailingOperationWorkflow;
3100
3101    impl WorkflowHandler for FailingOperationWorkflow {
3102        fn name(&self) -> &str {
3103            "failing-operation-workflow"
3104        }
3105
3106        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3107            Box::pin(async move {
3108                ctx.operation("broken-call", &FailingOp).await?;
3109                Ok(())
3110            })
3111        }
3112    }
3113
3114    struct MixedWorkflow;
3115
3116    impl WorkflowHandler for MixedWorkflow {
3117        fn name(&self) -> &str {
3118            "mixed-workflow"
3119        }
3120
3121        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3122            Box::pin(async move {
3123                ctx.shell("build", ShellConfig::new("echo built")).await?;
3124                let op = FakeGitlabOp {
3125                    project_id: 456,
3126                    title: "Deploy done".to_string(),
3127                };
3128                let result = ctx.operation("notify-gitlab", &op).await?;
3129                assert_eq!(result.output["issue_id"], 42);
3130                Ok(())
3131            })
3132        }
3133    }
3134
3135    #[tokio::test]
3136    async fn operation_step_happy_path() {
3137        let mut engine = create_test_engine();
3138        engine.register(OperationWorkflow).unwrap();
3139
3140        let run = engine
3141            .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
3142            .await
3143            .unwrap()
3144            .run;
3145
3146        assert_eq!(run.status.state, RunStatus::Completed);
3147
3148        let steps = engine.store().list_steps(run.id).await.unwrap();
3149
3150        assert_eq!(steps.len(), 1);
3151        assert_eq!(steps[0].name, "create-issue");
3152        assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
3153        assert_eq!(
3154            steps[0].status.state,
3155            ironflow_store::models::StepStatus::Completed
3156        );
3157
3158        let output = steps[0].output.as_ref().unwrap();
3159        assert_eq!(output["issue_id"], 42);
3160        assert_eq!(output["project_id"], 123);
3161
3162        let input = steps[0].input.as_ref().unwrap();
3163        assert_eq!(input["project_id"], 123);
3164        assert_eq!(input["title"], "Bug report");
3165    }
3166
3167    #[tokio::test]
3168    async fn operation_step_failure_marks_run_failed() {
3169        let mut engine = create_test_engine();
3170        engine.register(FailingOperationWorkflow).unwrap();
3171
3172        let result = engine
3173            .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
3174            .await;
3175
3176        assert!(result.is_err());
3177    }
3178
3179    #[tokio::test]
3180    async fn operation_mixed_with_shell_steps() {
3181        let mut engine = create_test_engine();
3182        engine.register(MixedWorkflow).unwrap();
3183
3184        let run = engine
3185            .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
3186            .await
3187            .unwrap()
3188            .run;
3189
3190        assert_eq!(run.status.state, RunStatus::Completed);
3191
3192        let steps = engine.store().list_steps(run.id).await.unwrap();
3193
3194        assert_eq!(steps.len(), 2);
3195        assert_eq!(steps[0].kind, StepKind::Shell);
3196        assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
3197        assert_eq!(steps[0].position, 0);
3198        assert_eq!(steps[1].position, 1);
3199    }
3200
3201    // -----------------------------------------------------------------------
3202    // Approval + resume tests
3203    // -----------------------------------------------------------------------
3204
3205    use crate::config::ApprovalConfig;
3206
3207    struct SingleApprovalWorkflow;
3208
3209    impl WorkflowHandler for SingleApprovalWorkflow {
3210        fn name(&self) -> &str {
3211            "single-approval"
3212        }
3213
3214        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3215            Box::pin(async move {
3216                ctx.shell("build", ShellConfig::new("echo built")).await?;
3217                ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
3218                ctx.shell("deploy", ShellConfig::new("echo deployed"))
3219                    .await?;
3220                Ok(())
3221            })
3222        }
3223    }
3224
3225    struct DoubleApprovalWorkflow;
3226
3227    impl WorkflowHandler for DoubleApprovalWorkflow {
3228        fn name(&self) -> &str {
3229            "double-approval"
3230        }
3231
3232        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3233            Box::pin(async move {
3234                ctx.shell("build", ShellConfig::new("echo built")).await?;
3235                ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
3236                    .await?;
3237                ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3238                    .await?;
3239                ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3240                    .await?;
3241                ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3242                    .await?;
3243                Ok(())
3244            })
3245        }
3246    }
3247
3248    #[tokio::test]
3249    async fn approval_pauses_run() {
3250        let mut engine = create_test_engine();
3251        engine.register(SingleApprovalWorkflow).unwrap();
3252
3253        let run = engine
3254            .run_handler("single-approval", TriggerKind::Manual, json!({}))
3255            .await
3256            .unwrap()
3257            .run;
3258
3259        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3260
3261        let steps = engine.store().list_steps(run.id).await.unwrap();
3262        assert_eq!(steps.len(), 2); // build + approval gate
3263        assert_eq!(steps[0].kind, StepKind::Shell);
3264        assert_eq!(steps[0].status.state, StepStatus::Completed);
3265        assert_eq!(steps[1].kind, StepKind::Approval);
3266        assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3267    }
3268
3269    #[tokio::test]
3270    async fn approval_resume_completes_run() {
3271        let mut engine = create_test_engine();
3272        engine.register(SingleApprovalWorkflow).unwrap();
3273
3274        // First execution: pauses at approval
3275        let run = engine
3276            .run_handler("single-approval", TriggerKind::Manual, json!({}))
3277            .await
3278            .unwrap()
3279            .run;
3280        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3281
3282        // Simulate approval: transition to Running
3283        engine
3284            .store()
3285            .update_run_status(run.id, RunStatus::Running)
3286            .await
3287            .unwrap();
3288
3289        // Resume: replays build, skips approval, executes deploy
3290        let resumed = engine.resume_run(run.id).await.unwrap().run;
3291        assert_eq!(resumed.status.state, RunStatus::Completed);
3292
3293        let steps = engine.store().list_steps(run.id).await.unwrap();
3294        assert_eq!(steps.len(), 3); // build + approval + deploy
3295        assert_eq!(steps[0].name, "build");
3296        assert_eq!(steps[0].status.state, StepStatus::Completed);
3297        assert_eq!(steps[1].name, "gate");
3298        assert_eq!(steps[1].kind, StepKind::Approval);
3299        assert_eq!(steps[1].status.state, StepStatus::Completed);
3300        assert_eq!(steps[2].name, "deploy");
3301        assert_eq!(steps[2].status.state, StepStatus::Completed);
3302    }
3303
3304    #[tokio::test]
3305    async fn double_approval_two_resumes() {
3306        let mut engine = create_test_engine();
3307        engine.register(DoubleApprovalWorkflow).unwrap();
3308
3309        // First execution: pauses at staging-gate
3310        let run = engine
3311            .run_handler("double-approval", TriggerKind::Manual, json!({}))
3312            .await
3313            .unwrap()
3314            .run;
3315        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3316
3317        let steps = engine.store().list_steps(run.id).await.unwrap();
3318        assert_eq!(steps.len(), 2); // build + staging-gate
3319
3320        // First approval
3321        engine
3322            .store()
3323            .update_run_status(run.id, RunStatus::Running)
3324            .await
3325            .unwrap();
3326
3327        let resumed = engine.resume_run(run.id).await.unwrap().run;
3328        assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3329
3330        let steps = engine.store().list_steps(run.id).await.unwrap();
3331        assert_eq!(steps.len(), 4); // build + staging-gate + deploy-staging + prod-gate
3332
3333        // Second approval
3334        engine
3335            .store()
3336            .update_run_status(run.id, RunStatus::Running)
3337            .await
3338            .unwrap();
3339
3340        let final_run = engine.resume_run(run.id).await.unwrap().run;
3341        assert_eq!(final_run.status.state, RunStatus::Completed);
3342
3343        let steps = engine.store().list_steps(run.id).await.unwrap();
3344        assert_eq!(steps.len(), 5);
3345        assert_eq!(steps[0].name, "build");
3346        assert_eq!(steps[1].name, "staging-gate");
3347        assert_eq!(steps[2].name, "deploy-staging");
3348        assert_eq!(steps[3].name, "prod-gate");
3349        assert_eq!(steps[4].name, "deploy-prod");
3350
3351        for step in &steps {
3352            assert_eq!(step.status.state, StepStatus::Completed);
3353        }
3354    }
3355
3356    // -----------------------------------------------------------------------
3357    // fail_orphaned_steps tests
3358    // -----------------------------------------------------------------------
3359
3360    use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3361
3362    async fn create_step_with_status(
3363        store: &Arc<dyn Store>,
3364        run_id: Uuid,
3365        name: &str,
3366        position: u32,
3367        status: StepStatus,
3368    ) -> ironflow_store::models::Step {
3369        let step = store
3370            .create_step(NewStep {
3371                run_id,
3372                trace_id: step_trace_id(run_id, name, position),
3373                name: name.to_string(),
3374                kind: StepKind::Shell,
3375                position,
3376                input: None,
3377                is_error_handler: false,
3378            })
3379            .await
3380            .unwrap();
3381
3382        match status {
3383            StepStatus::Pending => {}
3384            StepStatus::Running => {
3385                store
3386                    .update_step(
3387                        step.id,
3388                        StepUpdate {
3389                            status: Some(StepStatus::Running),
3390                            ..StepUpdate::default()
3391                        },
3392                    )
3393                    .await
3394                    .unwrap();
3395            }
3396            StepStatus::Completed => {
3397                store
3398                    .update_step(
3399                        step.id,
3400                        StepUpdate {
3401                            status: Some(StepStatus::Running),
3402                            ..StepUpdate::default()
3403                        },
3404                    )
3405                    .await
3406                    .unwrap();
3407                store
3408                    .update_step(
3409                        step.id,
3410                        StepUpdate {
3411                            status: Some(StepStatus::Completed),
3412                            ..StepUpdate::default()
3413                        },
3414                    )
3415                    .await
3416                    .unwrap();
3417            }
3418            StepStatus::AwaitingApproval => {
3419                store
3420                    .update_step(
3421                        step.id,
3422                        StepUpdate {
3423                            status: Some(StepStatus::Running),
3424                            ..StepUpdate::default()
3425                        },
3426                    )
3427                    .await
3428                    .unwrap();
3429                store
3430                    .update_step(
3431                        step.id,
3432                        StepUpdate {
3433                            status: Some(StepStatus::AwaitingApproval),
3434                            ..StepUpdate::default()
3435                        },
3436                    )
3437                    .await
3438                    .unwrap();
3439            }
3440            _ => panic!("unsupported status for test helper: {status}"),
3441        }
3442
3443        store.get_step(step.id).await.unwrap().unwrap()
3444    }
3445
3446    #[tokio::test]
3447    async fn fail_orphaned_steps_marks_running_as_failed() {
3448        let engine = create_test_engine();
3449        let run = engine
3450            .store()
3451            .create_run(NewRun {
3452                created_by: None,
3453                workflow_name: "test".to_string(),
3454                trigger: TriggerKind::Manual,
3455                payload: json!({}),
3456                max_retries: 0,
3457                handler_version: None,
3458                labels: HashMap::new(),
3459                scheduled_at: None,
3460                idempotency_key: None,
3461                concurrency_key: None,
3462                concurrency_limits: Vec::new(),
3463                max_cost_usd: None,
3464            })
3465            .await
3466            .unwrap()
3467            .into_run();
3468
3469        let step = create_step_with_status(
3470            engine.store(),
3471            run.id,
3472            "running-step",
3473            0,
3474            StepStatus::Running,
3475        )
3476        .await;
3477
3478        engine
3479            .fail_orphaned_steps(run.id, "parent run timed out")
3480            .await
3481            .unwrap();
3482
3483        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3484        assert_eq!(updated.status.state, StepStatus::Failed);
3485        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3486        assert!(updated.completed_at.is_some());
3487    }
3488
3489    #[tokio::test]
3490    async fn fail_orphaned_steps_marks_pending_as_skipped() {
3491        let engine = create_test_engine();
3492        let run = engine
3493            .store()
3494            .create_run(NewRun {
3495                created_by: None,
3496                workflow_name: "test".to_string(),
3497                trigger: TriggerKind::Manual,
3498                payload: json!({}),
3499                max_retries: 0,
3500                handler_version: None,
3501                labels: HashMap::new(),
3502                scheduled_at: None,
3503                idempotency_key: None,
3504                concurrency_key: None,
3505                concurrency_limits: Vec::new(),
3506                max_cost_usd: None,
3507            })
3508            .await
3509            .unwrap()
3510            .into_run();
3511
3512        let step = create_step_with_status(
3513            engine.store(),
3514            run.id,
3515            "pending-step",
3516            0,
3517            StepStatus::Pending,
3518        )
3519        .await;
3520
3521        engine
3522            .fail_orphaned_steps(run.id, "parent run timed out")
3523            .await
3524            .unwrap();
3525
3526        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3527        assert_eq!(updated.status.state, StepStatus::Skipped);
3528        assert!(updated.error.is_none());
3529        assert!(updated.completed_at.is_some());
3530    }
3531
3532    #[tokio::test]
3533    async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
3534        let engine = create_test_engine();
3535        let run = engine
3536            .store()
3537            .create_run(NewRun {
3538                created_by: None,
3539                workflow_name: "test".to_string(),
3540                trigger: TriggerKind::Manual,
3541                payload: json!({}),
3542                max_retries: 0,
3543                handler_version: None,
3544                labels: HashMap::new(),
3545                scheduled_at: None,
3546                idempotency_key: None,
3547                concurrency_key: None,
3548                concurrency_limits: Vec::new(),
3549                max_cost_usd: None,
3550            })
3551            .await
3552            .unwrap()
3553            .into_run();
3554
3555        let step = create_step_with_status(
3556            engine.store(),
3557            run.id,
3558            "approval-step",
3559            0,
3560            StepStatus::AwaitingApproval,
3561        )
3562        .await;
3563
3564        engine
3565            .fail_orphaned_steps(run.id, "parent run timed out")
3566            .await
3567            .unwrap();
3568
3569        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3570        assert_eq!(updated.status.state, StepStatus::Failed);
3571        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3572        assert!(updated.completed_at.is_some());
3573    }
3574
3575    #[tokio::test]
3576    async fn fail_orphaned_steps_skips_terminal_steps() {
3577        let engine = create_test_engine();
3578        let run = engine
3579            .store()
3580            .create_run(NewRun {
3581                created_by: None,
3582                workflow_name: "test".to_string(),
3583                trigger: TriggerKind::Manual,
3584                payload: json!({}),
3585                max_retries: 0,
3586                handler_version: None,
3587                labels: HashMap::new(),
3588                scheduled_at: None,
3589                idempotency_key: None,
3590                concurrency_key: None,
3591                concurrency_limits: Vec::new(),
3592                max_cost_usd: None,
3593            })
3594            .await
3595            .unwrap()
3596            .into_run();
3597
3598        let completed_step =
3599            create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
3600        let running_step =
3601            create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
3602                .await;
3603
3604        engine
3605            .fail_orphaned_steps(run.id, "parent run timed out")
3606            .await
3607            .unwrap();
3608
3609        let completed = engine
3610            .store()
3611            .get_step(completed_step.id)
3612            .await
3613            .unwrap()
3614            .unwrap();
3615        assert_eq!(completed.status.state, StepStatus::Completed);
3616
3617        let failed = engine
3618            .store()
3619            .get_step(running_step.id)
3620            .await
3621            .unwrap()
3622            .unwrap();
3623        assert_eq!(failed.status.state, StepStatus::Failed);
3624    }
3625
3626    #[tokio::test]
3627    async fn fail_orphaned_steps_mixed_states() {
3628        let engine = create_test_engine();
3629        let run = engine
3630            .store()
3631            .create_run(NewRun {
3632                created_by: None,
3633                workflow_name: "test".to_string(),
3634                trigger: TriggerKind::Manual,
3635                payload: json!({}),
3636                max_retries: 0,
3637                handler_version: None,
3638                labels: HashMap::new(),
3639                scheduled_at: None,
3640                idempotency_key: None,
3641                concurrency_key: None,
3642                concurrency_limits: Vec::new(),
3643                max_cost_usd: None,
3644            })
3645            .await
3646            .unwrap()
3647            .into_run();
3648
3649        let s_completed =
3650            create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
3651                .await;
3652        let s_running =
3653            create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
3654        let s_pending =
3655            create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
3656
3657        engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
3658
3659        let r_completed = engine
3660            .store()
3661            .get_step(s_completed.id)
3662            .await
3663            .unwrap()
3664            .unwrap();
3665        assert_eq!(r_completed.status.state, StepStatus::Completed);
3666
3667        let r_running = engine
3668            .store()
3669            .get_step(s_running.id)
3670            .await
3671            .unwrap()
3672            .unwrap();
3673        assert_eq!(r_running.status.state, StepStatus::Failed);
3674        assert_eq!(r_running.error.as_deref(), Some("timeout"));
3675
3676        let r_pending = engine
3677            .store()
3678            .get_step(s_pending.id)
3679            .await
3680            .unwrap()
3681            .unwrap();
3682        assert_eq!(r_pending.status.state, StepStatus::Skipped);
3683        assert!(r_pending.error.is_none());
3684    }
3685
3686    #[tokio::test]
3687    async fn fail_orphaned_steps_no_steps_is_noop() {
3688        let engine = create_test_engine();
3689        let run = engine
3690            .store()
3691            .create_run(NewRun {
3692                created_by: None,
3693                workflow_name: "test".to_string(),
3694                trigger: TriggerKind::Manual,
3695                payload: json!({}),
3696                max_retries: 0,
3697                handler_version: None,
3698                labels: HashMap::new(),
3699                scheduled_at: None,
3700                idempotency_key: None,
3701                concurrency_key: None,
3702                concurrency_limits: Vec::new(),
3703                max_cost_usd: None,
3704            })
3705            .await
3706            .unwrap()
3707            .into_run();
3708
3709        let result = engine.fail_orphaned_steps(run.id, "timeout").await;
3710        assert!(result.is_ok());
3711    }
3712
3713    #[tokio::test]
3714    async fn fail_orphaned_steps_preserves_existing_error() {
3715        let engine = create_test_engine();
3716        let run = engine
3717            .store()
3718            .create_run(NewRun {
3719                created_by: None,
3720                workflow_name: "test".to_string(),
3721                trigger: TriggerKind::Manual,
3722                payload: json!({}),
3723                max_retries: 0,
3724                handler_version: None,
3725                labels: HashMap::new(),
3726                scheduled_at: None,
3727                idempotency_key: None,
3728                concurrency_key: None,
3729                concurrency_limits: Vec::new(),
3730                max_cost_usd: None,
3731            })
3732            .await
3733            .unwrap()
3734            .into_run();
3735
3736        let step_with_error = create_step_with_status(
3737            engine.store(),
3738            run.id,
3739            "already-errored",
3740            0,
3741            StepStatus::Running,
3742        )
3743        .await;
3744
3745        engine
3746            .store()
3747            .update_step(
3748                step_with_error.id,
3749                StepUpdate {
3750                    error: Some("real error from provider".to_string()),
3751                    ..StepUpdate::default()
3752                },
3753            )
3754            .await
3755            .unwrap();
3756
3757        let step_no_error = create_step_with_status(
3758            engine.store(),
3759            run.id,
3760            "no-error-yet",
3761            1,
3762            StepStatus::Running,
3763        )
3764        .await;
3765
3766        engine
3767            .fail_orphaned_steps(run.id, "parent run failed")
3768            .await
3769            .unwrap();
3770
3771        let updated_with = engine
3772            .store()
3773            .get_step(step_with_error.id)
3774            .await
3775            .unwrap()
3776            .unwrap();
3777        assert_eq!(updated_with.status.state, StepStatus::Failed);
3778        assert_eq!(
3779            updated_with.error.as_deref(),
3780            Some("real error from provider"),
3781        );
3782
3783        let updated_without = engine
3784            .store()
3785            .get_step(step_no_error.id)
3786            .await
3787            .unwrap()
3788            .unwrap();
3789        assert_eq!(updated_without.status.state, StepStatus::Failed);
3790        assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
3791    }
3792}