Skip to main content

ironflow_engine/
engine.rs

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