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    /// Returns the status the run was moved to.
1825    ///
1826    /// # Errors
1827    ///
1828    /// Returns [`EngineError::Store`] if the run does not exist or the update
1829    /// cannot be persisted.
1830    ///
1831    /// # Examples
1832    ///
1833    /// ```no_run
1834    /// use ironflow_engine::engine::Engine;
1835    /// use ironflow_engine::error::EngineError;
1836    /// use uuid::Uuid;
1837    ///
1838    /// # async fn example(engine: &Engine, run_id: Uuid) -> Result<(), EngineError> {
1839    /// let status = engine
1840    ///     .fail_or_schedule_retry(run_id, "run timed out after 300s", true, None, None)
1841    ///     .await?;
1842    /// # Ok(())
1843    /// # }
1844    /// ```
1845    pub async fn fail_or_schedule_retry(
1846        &self,
1847        run_id: Uuid,
1848        error: &str,
1849        retryable: bool,
1850        cost_usd: Option<Decimal>,
1851        duration_ms: Option<u64>,
1852    ) -> Result<RunStatus, EngineError> {
1853        let run = self
1854            .store
1855            .get_run(run_id)
1856            .await?
1857            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1858
1859        let has_attempts_left = run.retry_count < run.max_retries;
1860        let update = if retryable && has_attempts_left {
1861            let backoff = backoff_for_retry(run.retry_count);
1862            let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1863
1864            info!(
1865                run_id = %run_id,
1866                workflow = %run.workflow_name,
1867                attempt = run.retry_count + 1,
1868                max_retries = run.max_retries,
1869                backoff_secs = backoff.as_secs(),
1870                scheduled_at = %scheduled_at,
1871                "run failed, scheduling retry"
1872            );
1873
1874            RunUpdate {
1875                status: Some(RunStatus::Retrying),
1876                error: Some(error.to_string()),
1877                increment_retry: true,
1878                cost_usd,
1879                duration_ms,
1880                scheduled_at: Some(scheduled_at),
1881                ..RunUpdate::default()
1882            }
1883        } else {
1884            RunUpdate {
1885                status: Some(RunStatus::Failed),
1886                error: Some(error.to_string()),
1887                cost_usd,
1888                duration_ms,
1889                completed_at: Some(Utc::now()),
1890                ..RunUpdate::default()
1891            }
1892        };
1893
1894        let status = update.status.unwrap_or(RunStatus::Failed);
1895        self.store.update_run(run_id, update).await?;
1896        self.fail_orphaned_steps(run_id, error).await?;
1897        // The attempt is over: a retry starts new children, and nothing drives
1898        // those this attempt left running.
1899        self.cancel_descendants_of_stopped_run(run_id, error).await;
1900
1901        Ok(status)
1902    }
1903
1904    /// Mark the steps of a run requeued after a lost worker lease as interrupted.
1905    ///
1906    /// Every `Running` step is marked `Failed` with
1907    /// [`STEP_INTERRUPTED_ERROR`](ironflow_store::store::STEP_INTERRUPTED_ERROR).
1908    /// The run keeps its attempt number, so when it is picked up again its
1909    /// finished steps are replayed and each interrupted step is executed again
1910    /// at the same position; an interrupted `Workflow` step re-enters the same
1911    /// child run. `Pending` and `AwaitingApproval` steps are left as they are,
1912    /// unlike [`fail_orphaned_steps`](Self::fail_orphaned_steps), which ends
1913    /// the run's steps for good.
1914    ///
1915    /// Errors from individual step updates are logged but do not abort the cleanup.
1916    ///
1917    /// # Errors
1918    ///
1919    /// Returns [`EngineError`] if listing steps fails.
1920    ///
1921    /// # Examples
1922    ///
1923    /// ```no_run
1924    /// use ironflow_engine::engine::Engine;
1925    /// use ironflow_engine::error::EngineError;
1926    /// use uuid::Uuid;
1927    ///
1928    /// # async fn example(engine: &Engine, run_id: Uuid) -> Result<(), EngineError> {
1929    /// engine.interrupt_running_steps(run_id).await?;
1930    /// # Ok(())
1931    /// # }
1932    /// ```
1933    pub async fn interrupt_running_steps(&self, run_id: Uuid) -> Result<(), EngineError> {
1934        interrupt_running_steps(self.store.as_ref(), run_id).await
1935    }
1936
1937    /// Fail all non-terminal steps for a run.
1938    ///
1939    /// Called after a run is marked as failed (timeout, error, panic) to clean up
1940    /// orphaned steps that are still in `Running`, `Pending`, or `AwaitingApproval`.
1941    ///
1942    /// - `Running` / `AwaitingApproval` steps are marked `Failed`.
1943    /// - `Pending` steps are marked `Skipped` (FSM does not allow Pending -> Failed).
1944    ///
1945    /// Errors from individual step updates are logged but do not abort the cleanup.
1946    ///
1947    /// # Errors
1948    ///
1949    /// Returns [`EngineError`] if listing steps fails.
1950    pub async fn fail_orphaned_steps(
1951        &self,
1952        run_id: Uuid,
1953        error_message: &str,
1954    ) -> Result<(), EngineError> {
1955        let steps = self.store.list_steps(run_id).await?;
1956        let now = Utc::now();
1957
1958        for step in steps {
1959            if step.status.state.is_terminal() {
1960                continue;
1961            }
1962
1963            let (target_status, error) = match step.status.state {
1964                StepStatus::Running | StepStatus::AwaitingApproval => {
1965                    let err = if step.error.is_some() {
1966                        None
1967                    } else {
1968                        Some(error_message.to_string())
1969                    };
1970                    (StepStatus::Failed, err)
1971                }
1972                StepStatus::Pending => (StepStatus::Skipped, None),
1973                _ => continue,
1974            };
1975
1976            if let Err(e) = self
1977                .store
1978                .update_step(
1979                    step.id,
1980                    StepUpdate {
1981                        status: Some(target_status),
1982                        error,
1983                        completed_at: Some(now),
1984                        ..StepUpdate::default()
1985                    },
1986                )
1987                .await
1988            {
1989                warn!(
1990                    run_id = %run_id,
1991                    step_id = %step.id,
1992                    step_name = %step.name,
1993                    error = %e,
1994                    "failed to cleanup orphaned step"
1995                );
1996            } else {
1997                info!(
1998                    run_id = %run_id,
1999                    step_id = %step.id,
2000                    step_name = %step.name,
2001                    from = %step.status.state,
2002                    to = %target_status,
2003                    "cleaned up orphaned step"
2004                );
2005            }
2006        }
2007
2008        Ok(())
2009    }
2010
2011    /// Execute `handler` once [`AgentProvider::release_run`] has stopped
2012    /// whatever a previous execution of the run left running (an agent pod
2013    /// writing to a shared worktree). A failed release fails the execution
2014    /// before its first step, with [`EngineError::Operation`].
2015    async fn release_then_execute(
2016        &self,
2017        run_id: Uuid,
2018        handler: &dyn WorkflowHandler,
2019        ctx: &mut WorkflowContext,
2020    ) -> Result<(), EngineError> {
2021        match self.provider.release_run(&run_id.to_string()).await {
2022            Ok(()) => handler.execute(ctx).await,
2023            Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
2024        }
2025    }
2026
2027    /// Finalize a run with the given result and context.
2028    ///
2029    /// On success: updates run to Completed with cost, duration, and completed_at.
2030    /// On failure: updates run to Failed with error, cost, duration, and completed_at.
2031    /// Always: fetches and returns the final Run.
2032    async fn finalize_run(
2033        &self,
2034        run_id: Uuid,
2035        workflow_name: &str,
2036        result: Result<(), EngineError>,
2037        ctx: &WorkflowContext,
2038        run_start: Instant,
2039        run_labels: HashMap<String, String>,
2040    ) -> Result<WorkflowResult, EngineError> {
2041        // Covers the whole run: previous attempts plus this one, so a retried
2042        // run reports the time it really consumed.
2043        let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
2044        let completed_at = Utc::now();
2045
2046        let final_status;
2047        let final_run;
2048
2049        match result {
2050            Ok(()) => {
2051                final_status = if ctx.has_allowed_failure() {
2052                    RunStatus::Warning
2053                } else {
2054                    RunStatus::Completed
2055                };
2056                final_run = self
2057                    .store
2058                    .update_run_returning(
2059                        run_id,
2060                        RunUpdate {
2061                            status: Some(final_status),
2062                            cost_usd: Some(ctx.total_cost_usd()),
2063                            duration_ms: Some(total_duration),
2064                            completed_at: Some(completed_at),
2065                            output: ctx.output().cloned(),
2066                            ..RunUpdate::default()
2067                        },
2068                    )
2069                    .await?;
2070
2071                info!(
2072                    run_id = %run_id,
2073                    status = %final_status,
2074                    cost_usd = %ctx.total_cost_usd(),
2075                    duration_ms = total_duration,
2076                    "run completed"
2077                );
2078            }
2079            Err(EngineError::ApprovalRequired {
2080                run_id: approval_run_id,
2081                step_id,
2082                ref message,
2083            }) => {
2084                final_status = RunStatus::AwaitingApproval;
2085                final_run = self
2086                    .store
2087                    .update_run_returning(
2088                        run_id,
2089                        RunUpdate {
2090                            status: Some(RunStatus::AwaitingApproval),
2091                            cost_usd: Some(ctx.total_cost_usd()),
2092                            duration_ms: Some(total_duration),
2093                            ..RunUpdate::default()
2094                        },
2095                    )
2096                    .await?;
2097
2098                info!(
2099                    run_id = %approval_run_id,
2100                    step_id = %step_id,
2101                    message = %message,
2102                    "run awaiting approval"
2103                );
2104
2105                self.publish_approval_requested(approval_run_id, step_id, message)
2106                    .await?;
2107            }
2108            Err(EngineError::ChildSuspended {
2109                run_id: child_run_id,
2110                ref cause,
2111            }) => {
2112                final_status = cause.suspension_status();
2113                // No `scheduled_at`: the suspended descendant owns the wake-up
2114                // and resumes this run through `resume_chain`. A wake-up armed
2115                // here too would resume the chain twice.
2116                final_run = self
2117                    .store
2118                    .update_run_returning(
2119                        run_id,
2120                        RunUpdate {
2121                            status: Some(final_status),
2122                            cost_usd: Some(ctx.total_cost_usd()),
2123                            duration_ms: Some(total_duration),
2124                            ..RunUpdate::default()
2125                        },
2126                    )
2127                    .await?;
2128
2129                let leaf = cause.suspension_leaf();
2130                info!(
2131                    run_id = %run_id,
2132                    child_run_id = %child_run_id,
2133                    status = %final_status,
2134                    cause = %leaf,
2135                    "run suspended with its child run"
2136                );
2137
2138                match leaf {
2139                    EngineError::ApprovalRequired {
2140                        run_id: approval_run_id,
2141                        step_id,
2142                        message,
2143                    } => {
2144                        self.publish_approval_requested(*approval_run_id, *step_id, message)
2145                            .await?;
2146                    }
2147                    EngineError::SignalWaiting {
2148                        run_id: wait_run_id,
2149                        step_id,
2150                        step_name,
2151                        name,
2152                        key,
2153                        deadline_at,
2154                    } => {
2155                        self.event_publisher
2156                            .publish(Event::SignalAwaited(SignalAwaitedEvent {
2157                                run_id: *wait_run_id,
2158                                step_id: *step_id,
2159                                step_name: step_name.clone(),
2160                                name: name.clone(),
2161                                key: key.clone(),
2162                                deadline_at: *deadline_at,
2163                                at: Utc::now(),
2164                            }));
2165                    }
2166                    // A human input or a delay publishes no suspension event,
2167                    // like on a top-level run.
2168                    _ => {}
2169                }
2170            }
2171            Err(EngineError::HumanInputRequired {
2172                run_id: input_run_id,
2173                step_id,
2174                ref message,
2175            }) => {
2176                final_status = RunStatus::AwaitingApproval;
2177                final_run = self
2178                    .store
2179                    .update_run_returning(
2180                        run_id,
2181                        RunUpdate {
2182                            status: Some(RunStatus::AwaitingApproval),
2183                            cost_usd: Some(ctx.total_cost_usd()),
2184                            duration_ms: Some(total_duration),
2185                            ..RunUpdate::default()
2186                        },
2187                    )
2188                    .await?;
2189
2190                // No `ApprovalRequested` event: a human input is not an approval.
2191                info!(
2192                    run_id = %input_run_id,
2193                    step_id = %step_id,
2194                    message = %message,
2195                    "run awaiting human input"
2196                );
2197            }
2198            Err(EngineError::DelaySleeping {
2199                run_id: delay_run_id,
2200                step_id,
2201                wake_at,
2202            }) => {
2203                final_status = RunStatus::Sleeping;
2204                final_run = self
2205                    .store
2206                    .update_run_returning(
2207                        run_id,
2208                        RunUpdate {
2209                            status: Some(RunStatus::Sleeping),
2210                            cost_usd: Some(ctx.total_cost_usd()),
2211                            duration_ms: Some(total_duration),
2212                            scheduled_at: Some(wake_at),
2213                            ..RunUpdate::default()
2214                        },
2215                    )
2216                    .await?;
2217
2218                info!(
2219                    run_id = %delay_run_id,
2220                    step_id = %step_id,
2221                    wake_at = %wake_at,
2222                    "run sleeping until delay elapses"
2223                );
2224            }
2225            Err(EngineError::CapacitySleeping {
2226                run_id: capacity_run_id,
2227                step_id,
2228                ref kind,
2229                wake_at,
2230            }) => {
2231                final_status = RunStatus::Sleeping;
2232                final_run = self
2233                    .store
2234                    .update_run_returning(
2235                        run_id,
2236                        RunUpdate {
2237                            status: Some(RunStatus::Sleeping),
2238                            cost_usd: Some(ctx.total_cost_usd()),
2239                            duration_ms: Some(total_duration),
2240                            scheduled_at: Some(wake_at),
2241                            capacity_wait_kind: Some(ProviderKind::new(kind.as_str())),
2242                            ..RunUpdate::default()
2243                        },
2244                    )
2245                    .await?;
2246
2247                info!(
2248                    run_id = %capacity_run_id,
2249                    step_id = %step_id,
2250                    kind = %kind,
2251                    wake_at = %wake_at,
2252                    "run sleeping until provider capacity returns"
2253                );
2254            }
2255            Err(EngineError::SignalWaiting {
2256                run_id: wait_run_id,
2257                step_id,
2258                ref step_name,
2259                ref name,
2260                ref key,
2261                deadline_at,
2262            }) => {
2263                final_status = RunStatus::Sleeping;
2264                // Atomic with the step lock: a signal delivered since the step
2265                // opened leaves the run due right away instead of until the
2266                // deadline.
2267                let waiting = self
2268                    .store
2269                    .suspend_run_on_signal(run_id, step_id, deadline_at)
2270                    .await?;
2271                final_run = self
2272                    .store
2273                    .update_run_returning(
2274                        run_id,
2275                        RunUpdate {
2276                            cost_usd: Some(ctx.total_cost_usd()),
2277                            duration_ms: Some(total_duration),
2278                            ..RunUpdate::default()
2279                        },
2280                    )
2281                    .await?;
2282
2283                if waiting {
2284                    self.event_publisher
2285                        .publish(Event::SignalAwaited(SignalAwaitedEvent {
2286                            run_id: wait_run_id,
2287                            step_id,
2288                            step_name: step_name.clone(),
2289                            name: name.clone(),
2290                            key: key.clone(),
2291                            deadline_at,
2292                            at: Utc::now(),
2293                        }));
2294                }
2295
2296                info!(
2297                    run_id = %wait_run_id,
2298                    step_id = %step_id,
2299                    signal = %name,
2300                    key = %key,
2301                    deadline_at = %deadline_at,
2302                    waiting,
2303                    "run sleeping until a signal arrives"
2304                );
2305            }
2306            Err(err) => {
2307                // A guardrail stop (budget or workflow guard) is deliberate,
2308                // not a breakage: the run is cancelled, never failed and
2309                // never replayed.
2310                let guardrail_stop = matches!(
2311                    err,
2312                    EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2313                );
2314
2315                final_status = if guardrail_stop {
2316                    if let Err(store_err) = self
2317                        .store
2318                        .update_run(
2319                            run_id,
2320                            RunUpdate {
2321                                status: Some(RunStatus::Cancelled),
2322                                error: Some(err.to_string()),
2323                                cost_usd: Some(ctx.total_cost_usd()),
2324                                duration_ms: Some(total_duration),
2325                                completed_at: Some(completed_at),
2326                                output: ctx.output().cloned(),
2327                                ..RunUpdate::default()
2328                            },
2329                        )
2330                        .await
2331                    {
2332                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2333                    }
2334                    if let Err(cleanup_err) = self
2335                        .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2336                        .await
2337                    {
2338                        error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2339                    }
2340                    RunStatus::Cancelled
2341                } else {
2342                    // Written before the failure so a failed run keeps the
2343                    // verdict the handler set before returning its error.
2344                    if let Some(output) = ctx.output()
2345                        && let Err(store_err) = self
2346                            .store
2347                            .update_run(
2348                                run_id,
2349                                RunUpdate {
2350                                    output: Some(output.clone()),
2351                                    ..RunUpdate::default()
2352                                },
2353                            )
2354                            .await
2355                    {
2356                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run output");
2357                    }
2358                    self.fail_or_schedule_retry(
2359                        run_id,
2360                        &err.to_string(),
2361                        is_run_retryable(&err),
2362                        Some(ctx.total_cost_usd()),
2363                        Some(total_duration),
2364                    )
2365                    .await
2366                    .unwrap_or_else(|store_err| {
2367                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2368                        RunStatus::Failed
2369                    })
2370                };
2371
2372                if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2373                    self.on_run_budget_exceeded(workflow_name, run_id, &err);
2374                }
2375
2376                error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2377
2378                self.publish_run_status_changed(
2379                    workflow_name,
2380                    run_id,
2381                    final_status,
2382                    Some(err.to_string()),
2383                    ctx,
2384                    total_duration,
2385                    run_labels,
2386                );
2387
2388                #[cfg(feature = "prometheus")]
2389                self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2390
2391                return Err(err);
2392            }
2393        }
2394
2395        self.publish_run_status_changed(
2396            workflow_name,
2397            run_id,
2398            final_status,
2399            None,
2400            ctx,
2401            total_duration,
2402            run_labels,
2403        );
2404
2405        #[cfg(feature = "prometheus")]
2406        self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2407
2408        Ok(WorkflowResult {
2409            run: final_run,
2410            steps: ctx.step_results().to_vec(),
2411        })
2412    }
2413
2414    /// Publish [`Event::ApprovalRequested`] for an approval gate that opened
2415    /// on `run_id`, with the requirement recorded when the gate opened.
2416    async fn publish_approval_requested(
2417        &self,
2418        run_id: Uuid,
2419        step_id: Uuid,
2420        message: &str,
2421    ) -> Result<(), EngineError> {
2422        let requirement = self
2423            .store
2424            .get_step(step_id)
2425            .await?
2426            .and_then(|s| s.approval_requirement);
2427        self.event_publisher
2428            .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2429                run_id,
2430                step_id,
2431                message: message.to_string(),
2432                requirement,
2433                at: Utc::now(),
2434            }));
2435        Ok(())
2436    }
2437
2438    /// Fail every ancestor of a child run, closest first.
2439    ///
2440    /// Used when a gate inside a sub-workflow is rejected: the child run
2441    /// fails, and the runs suspended with it (its parent, up to the root)
2442    /// must not stay suspended on a child that will never resume. Each
2443    /// ancestor goes through [`fail_or_schedule_retry`](Self::fail_or_schedule_retry)
2444    /// without a retry, which also fails its open `Workflow` step. A run that
2445    /// is not a child of a sub-workflow has no ancestor: nothing happens.
2446    ///
2447    /// # Errors
2448    ///
2449    /// Returns [`EngineError::Store`] if a run of the chain does not exist or
2450    /// its failure cannot be persisted.
2451    ///
2452    /// # Examples
2453    ///
2454    /// ```no_run
2455    /// use ironflow_engine::engine::Engine;
2456    /// use ironflow_engine::error::EngineError;
2457    /// use uuid::Uuid;
2458    ///
2459    /// # async fn example(engine: &Engine, child_run_id: Uuid) -> Result<(), EngineError> {
2460    /// engine
2461    ///     .fail_ancestors(child_run_id, "approval rejected in a sub-workflow")
2462    ///     .await?;
2463    /// # Ok(())
2464    /// # }
2465    /// ```
2466    pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2467        let mut current = self
2468            .store
2469            .get_run(run_id)
2470            .await?
2471            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2472        // Labels are data: a chain that loops back on itself stops there.
2473        let mut visited = HashSet::from([run_id]);
2474
2475        while let Some(parent_id) = chain_parent(&current) {
2476            if !visited.insert(parent_id) {
2477                break;
2478            }
2479            let status = self
2480                .fail_or_schedule_retry(parent_id, reason, false, None, None)
2481                .await?;
2482            info!(
2483                run_id = %run_id,
2484                ancestor_run_id = %parent_id,
2485                status = %status,
2486                "ancestor run failed with its child"
2487            );
2488            current = self
2489                .store
2490                .get_run(parent_id)
2491                .await?
2492                .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2493        }
2494
2495        Ok(())
2496    }
2497
2498    /// Emit Prometheus metrics for a completed run.
2499    #[cfg(feature = "prometheus")]
2500    fn emit_run_metrics(
2501        &self,
2502        workflow_name: &str,
2503        status: RunStatus,
2504        duration_ms: u64,
2505        ctx: &WorkflowContext,
2506    ) {
2507        let status_str = status.to_string();
2508        let wf = workflow_name.to_string();
2509
2510        counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2511        histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2512            .record(duration_ms as f64 / 1000.0);
2513        histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2514            ctx.total_cost_usd()
2515                .to_string()
2516                .parse::<f64>()
2517                .unwrap_or(0.0),
2518        );
2519        gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2520    }
2521
2522    /// Record the metric and publish the audit event for a run that hit its
2523    /// cost cap.
2524    ///
2525    /// A non-[`RunBudgetExceeded`](EngineError::RunBudgetExceeded) error is
2526    /// ignored, so callers can pass the error unconditionally.
2527    fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2528        let EngineError::RunBudgetExceeded {
2529            limit_usd,
2530            spent_usd,
2531            step_budget_usd,
2532            ..
2533        } = err
2534        else {
2535            return;
2536        };
2537
2538        #[cfg(feature = "prometheus")]
2539        counter!(
2540            RUN_BUDGET_EXCEEDED_TOTAL,
2541            "workflow" => workflow_name.to_string(),
2542            "scope" => "run",
2543        )
2544        .increment(1);
2545
2546        self.event_publisher
2547            .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2548                run_id,
2549                workflow_name: workflow_name.to_string(),
2550                limit_usd: *limit_usd,
2551                spent_usd: *spent_usd,
2552                step_budget_usd: *step_budget_usd,
2553                at: Utc::now(),
2554            }));
2555    }
2556
2557    /// Publish a run status changed event to all registered subscribers.
2558    ///
2559    /// `from` is always `Running` because `finalize_run` is only called
2560    /// from a running state.
2561    #[allow(clippy::too_many_arguments)]
2562    fn publish_run_status_changed(
2563        &self,
2564        workflow_name: &str,
2565        run_id: Uuid,
2566        to: RunStatus,
2567        error: Option<String>,
2568        ctx: &WorkflowContext,
2569        duration_ms: u64,
2570        labels: HashMap<String, String>,
2571    ) {
2572        let now = Utc::now();
2573        let cost_usd = ctx.total_cost_usd();
2574        let wf = workflow_name.to_string();
2575
2576        self.event_publisher
2577            .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2578                run_id,
2579                workflow_name: wf.clone(),
2580                from: RunStatus::Running,
2581                to,
2582                error: error.clone(),
2583                cost_usd,
2584                duration_ms,
2585                labels: labels.clone(),
2586                at: now,
2587            }));
2588
2589        if to == RunStatus::Failed {
2590            self.event_publisher
2591                .publish(Event::RunFailed(RunFailedEvent {
2592                    run_id,
2593                    workflow_name: wf,
2594                    error,
2595                    cost_usd,
2596                    duration_ms,
2597                    labels,
2598                    at: now,
2599                }));
2600        }
2601    }
2602}
2603
2604impl fmt::Debug for Engine {
2605    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2606        f.debug_struct("Engine")
2607            .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2608            .finish_non_exhaustive()
2609    }
2610}
2611
2612#[cfg(test)]
2613mod tests {
2614    use super::*;
2615    use crate::config::ShellConfig;
2616    use crate::handler::{HandlerFuture, WorkflowHandler};
2617    use ironflow_core::providers::claude::ClaudeCodeProvider;
2618    use ironflow_core::providers::record_replay::RecordReplayProvider;
2619    use ironflow_store::memory::InMemoryStore;
2620    use ironflow_store::models::{MAX_PRIORITY, MIN_PRIORITY, StepStatus};
2621    use serde_json::json;
2622
2623    // Test handler that echoes a message via shell
2624    struct EchoWorkflow;
2625
2626    impl WorkflowHandler for EchoWorkflow {
2627        fn name(&self) -> &str {
2628            "echo-workflow"
2629        }
2630
2631        fn describe(&self) -> WorkflowInfo {
2632            WorkflowInfo {
2633                description: "A simple workflow that echoes hello".to_string(),
2634                source_code: None,
2635                sub_workflows: Vec::new(),
2636                category: None,
2637                version: self.version().map(str::to_string),
2638                compatible_versions: Vec::new(),
2639                input_schema: None,
2640                default_labels: HashMap::new(),
2641                schedule: self.schedule().cloned(),
2642                default_max_cost_usd: self.default_max_cost_usd(),
2643                priority: self.priority(),
2644            }
2645        }
2646
2647        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2648            Box::pin(async move {
2649                ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2650                Ok(())
2651            })
2652        }
2653    }
2654
2655    // Test handler that fails
2656    struct FailingWorkflow;
2657
2658    impl WorkflowHandler for FailingWorkflow {
2659        fn name(&self) -> &str {
2660            "failing-workflow"
2661        }
2662
2663        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2664            Box::pin(async move {
2665                ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2666                Ok(())
2667            })
2668        }
2669    }
2670
2671    fn create_test_engine() -> Engine {
2672        let store = Arc::new(InMemoryStore::new());
2673        let inner = ClaudeCodeProvider::new();
2674        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2675            inner,
2676            "/tmp/ironflow-fixtures",
2677        ));
2678        Engine::new(store, provider)
2679    }
2680
2681    #[test]
2682    fn engine_new_creates_instance() {
2683        let engine = create_test_engine();
2684        assert_eq!(engine.handler_names().len(), 0);
2685    }
2686
2687    #[test]
2688    fn execution_mode_defaults_to_local() {
2689        let engine = create_test_engine();
2690        assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2691    }
2692
2693    #[test]
2694    fn with_execution_mode_overrides_the_default() {
2695        let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2696        assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2697    }
2698
2699    #[test]
2700    fn engine_register_handler() {
2701        let mut engine = create_test_engine();
2702        let result = engine.register(EchoWorkflow);
2703        assert!(result.is_ok());
2704        assert_eq!(engine.handler_names().len(), 1);
2705        assert!(engine.handler_names().contains(&"echo-workflow"));
2706    }
2707
2708    #[test]
2709    fn engine_register_duplicate_returns_error() {
2710        let mut engine = create_test_engine();
2711        engine.register(EchoWorkflow).unwrap();
2712        let result = engine.register(EchoWorkflow);
2713        assert!(result.is_err());
2714    }
2715
2716    #[test]
2717    fn engine_get_handler_found() {
2718        let mut engine = create_test_engine();
2719        engine.register(EchoWorkflow).unwrap();
2720        let handler = engine.get_handler("echo-workflow");
2721        assert!(handler.is_some());
2722    }
2723
2724    #[test]
2725    fn engine_get_handler_not_found() {
2726        let engine = create_test_engine();
2727        let handler = engine.get_handler("nonexistent");
2728        assert!(handler.is_none());
2729    }
2730
2731    #[test]
2732    fn engine_handler_names_lists_all() {
2733        let mut engine = create_test_engine();
2734        engine.register(EchoWorkflow).unwrap();
2735        engine.register(FailingWorkflow).unwrap();
2736        let names = engine.handler_names();
2737        assert_eq!(names.len(), 2);
2738        assert!(names.contains(&"echo-workflow"));
2739        assert!(names.contains(&"failing-workflow"));
2740    }
2741
2742    #[test]
2743    fn engine_handler_info_returns_description() {
2744        let mut engine = create_test_engine();
2745        engine.register(EchoWorkflow).unwrap();
2746        let info = engine.handler_info("echo-workflow");
2747        assert!(info.is_some());
2748        let info = info.unwrap();
2749        assert_eq!(info.description, "A simple workflow that echoes hello");
2750    }
2751
2752    struct CategorizedWorkflow;
2753
2754    impl WorkflowHandler for CategorizedWorkflow {
2755        fn name(&self) -> &str {
2756            "categorized"
2757        }
2758        fn category(&self) -> Option<&str> {
2759            Some("data/etl")
2760        }
2761        fn execute<'a>(
2762            &'a self,
2763            _ctx: &'a mut WorkflowContext,
2764        ) -> crate::handler::HandlerFuture<'a> {
2765            Box::pin(async move { Ok(()) })
2766        }
2767    }
2768
2769    #[test]
2770    fn engine_default_describe_propagates_category() {
2771        let mut engine = create_test_engine();
2772        engine.register(CategorizedWorkflow).unwrap();
2773        let info = engine.handler_info("categorized").unwrap();
2774        assert_eq!(info.category.as_deref(), Some("data/etl"));
2775    }
2776
2777    #[test]
2778    fn engine_default_describe_without_category() {
2779        let mut engine = create_test_engine();
2780        engine.register(EchoWorkflow).unwrap();
2781        let info = engine.handler_info("echo-workflow").unwrap();
2782        assert!(info.category.is_none());
2783    }
2784
2785    // -----------------------------------------------------------------------
2786    // Schedule tests
2787    // -----------------------------------------------------------------------
2788
2789    struct ScheduledWorkflow {
2790        schedule: CronSchedule,
2791    }
2792
2793    impl ScheduledWorkflow {
2794        fn new() -> Self {
2795            Self {
2796                schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2797            }
2798        }
2799    }
2800
2801    impl WorkflowHandler for ScheduledWorkflow {
2802        fn name(&self) -> &str {
2803            "scheduled"
2804        }
2805        fn schedule(&self) -> Option<&CronSchedule> {
2806            Some(&self.schedule)
2807        }
2808        fn execute<'a>(
2809            &'a self,
2810            _ctx: &'a mut WorkflowContext,
2811        ) -> crate::handler::HandlerFuture<'a> {
2812            Box::pin(async move { Ok(()) })
2813        }
2814    }
2815
2816    #[test]
2817    fn engine_default_describe_propagates_schedule() {
2818        let mut engine = create_test_engine();
2819        engine.register(ScheduledWorkflow::new()).unwrap();
2820        let info = engine.handler_info("scheduled").unwrap();
2821        assert_eq!(
2822            info.schedule.as_ref().map(|s| s.as_str()),
2823            Some("0 0 * * * *")
2824        );
2825    }
2826
2827    #[test]
2828    fn engine_default_describe_without_schedule() {
2829        let mut engine = create_test_engine();
2830        engine.register(EchoWorkflow).unwrap();
2831        let info = engine.handler_info("echo-workflow").unwrap();
2832        assert!(info.schedule.is_none());
2833    }
2834
2835    #[test]
2836    fn scheduled_handlers_returns_only_scheduled() {
2837        let mut engine = create_test_engine();
2838        engine.register(EchoWorkflow).unwrap();
2839        engine.register(ScheduledWorkflow::new()).unwrap();
2840        engine.register(FailingWorkflow).unwrap();
2841
2842        let scheduled = engine.scheduled_handlers();
2843        assert_eq!(scheduled.len(), 1);
2844        assert_eq!(scheduled[0].0, "scheduled");
2845        assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
2846    }
2847
2848    #[test]
2849    fn scheduled_handlers_empty_when_none_scheduled() {
2850        let mut engine = create_test_engine();
2851        engine.register(EchoWorkflow).unwrap();
2852        engine.register(FailingWorkflow).unwrap();
2853
2854        let scheduled = engine.scheduled_handlers();
2855        assert!(scheduled.is_empty());
2856    }
2857
2858    struct BadCategoryWorkflow(&'static str);
2859
2860    impl WorkflowHandler for BadCategoryWorkflow {
2861        fn name(&self) -> &str {
2862            "bad-category"
2863        }
2864        fn category(&self) -> Option<&str> {
2865            Some(self.0)
2866        }
2867        fn execute<'a>(
2868            &'a self,
2869            _ctx: &'a mut WorkflowContext,
2870        ) -> crate::handler::HandlerFuture<'a> {
2871            Box::pin(async move { Ok(()) })
2872        }
2873    }
2874
2875    #[test]
2876    fn engine_register_rejects_empty_category() {
2877        let mut engine = create_test_engine();
2878        let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2879        match err {
2880            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2881            other => panic!("expected InvalidWorkflow, got {other:?}"),
2882        }
2883    }
2884
2885    #[test]
2886    fn engine_register_rejects_leading_slash_category() {
2887        let mut engine = create_test_engine();
2888        let err = engine
2889            .register(BadCategoryWorkflow("/data/etl"))
2890            .unwrap_err();
2891        match err {
2892            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2893            other => panic!("expected InvalidWorkflow, got {other:?}"),
2894        }
2895    }
2896
2897    #[test]
2898    fn engine_register_rejects_trailing_slash_category() {
2899        let mut engine = create_test_engine();
2900        let err = engine
2901            .register(BadCategoryWorkflow("data/etl/"))
2902            .unwrap_err();
2903        match err {
2904            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2905            other => panic!("expected InvalidWorkflow, got {other:?}"),
2906        }
2907    }
2908
2909    #[test]
2910    fn engine_register_rejects_double_slash_category() {
2911        let mut engine = create_test_engine();
2912        let err = engine
2913            .register(BadCategoryWorkflow("data//etl"))
2914            .unwrap_err();
2915        match err {
2916            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2917            other => panic!("expected InvalidWorkflow, got {other:?}"),
2918        }
2919    }
2920
2921    #[test]
2922    fn engine_register_rejects_whitespace_only_segment_category() {
2923        let mut engine = create_test_engine();
2924        let err = engine
2925            .register(BadCategoryWorkflow("data/ /etl"))
2926            .unwrap_err();
2927        match err {
2928            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2929            other => panic!("expected InvalidWorkflow, got {other:?}"),
2930        }
2931    }
2932
2933    #[test]
2934    fn engine_register_accepts_valid_nested_category() {
2935        let mut engine = create_test_engine();
2936        assert!(engine.register(CategorizedWorkflow).is_ok());
2937    }
2938
2939    #[tokio::test]
2940    async fn engine_unknown_workflow_returns_error() {
2941        let engine = create_test_engine();
2942        let result = engine
2943            .run_handler("unknown", TriggerKind::Manual, json!({}))
2944            .await;
2945        assert!(result.is_err());
2946        match result {
2947            Err(EngineError::InvalidWorkflow(msg)) => {
2948                assert!(msg.contains("no handler registered"));
2949            }
2950            _ => panic!("expected InvalidWorkflow error"),
2951        }
2952    }
2953
2954    #[tokio::test]
2955    async fn engine_enqueue_handler_creates_pending_run() {
2956        let mut engine = create_test_engine();
2957        engine.register(EchoWorkflow).unwrap();
2958
2959        let run = engine
2960            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2961            .await
2962            .unwrap();
2963        assert_eq!(run.status.state, RunStatus::Pending);
2964        assert_eq!(run.workflow_name, "echo-workflow");
2965    }
2966
2967    #[tokio::test]
2968    async fn enqueue_handler_leaves_the_run_unattributed() {
2969        let mut engine = create_test_engine();
2970        engine.register(EchoWorkflow).unwrap();
2971
2972        let run = engine
2973            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2974            .await
2975            .unwrap();
2976
2977        assert!(run.created_by.is_none());
2978    }
2979
2980    #[tokio::test]
2981    async fn enqueue_handler_with_options_records_the_author() {
2982        let mut engine = create_test_engine();
2983        engine.register(EchoWorkflow).unwrap();
2984        let actor = RunActor::User {
2985            user_id: Uuid::now_v7(),
2986        };
2987
2988        let run = engine
2989            .enqueue_handler_with_options(
2990                "echo-workflow",
2991                TriggerKind::Api,
2992                json!({}),
2993                EnqueueOptions {
2994                    created_by: Some(actor.clone()),
2995                    ..Default::default()
2996                },
2997            )
2998            .await
2999            .unwrap()
3000            .into_run();
3001
3002        assert_eq!(run.created_by, Some(actor));
3003    }
3004
3005    #[tokio::test]
3006    async fn enqueue_handler_with_options_accepts_no_author() {
3007        let mut engine = create_test_engine();
3008        engine.register(EchoWorkflow).unwrap();
3009
3010        let run = engine
3011            .enqueue_handler_with_options(
3012                "echo-workflow",
3013                TriggerKind::Cron {
3014                    schedule: "0 * * * * *".to_string(),
3015                },
3016                json!({}),
3017                EnqueueOptions::default(),
3018            )
3019            .await
3020            .unwrap()
3021            .into_run();
3022
3023        assert!(run.created_by.is_none());
3024    }
3025
3026    #[tokio::test]
3027    async fn enqueue_handler_with_options_stores_concurrency_limits() {
3028        let mut engine = create_test_engine();
3029        engine.register(EchoWorkflow).unwrap();
3030        let limits = vec![
3031            ConcurrencyLimit::new("repo:acme", 2),
3032            ConcurrencyLimit::new("tenant:42", 5),
3033        ];
3034
3035        let run = engine
3036            .enqueue_handler_with_options(
3037                "echo-workflow",
3038                TriggerKind::Api,
3039                json!({}),
3040                EnqueueOptions {
3041                    concurrency_limits: limits.clone(),
3042                    ..Default::default()
3043                },
3044            )
3045            .await
3046            .unwrap()
3047            .into_run();
3048
3049        assert_eq!(run.concurrency_limits, limits);
3050    }
3051
3052    #[tokio::test]
3053    async fn enqueue_rejects_invalid_concurrency_limits() {
3054        let mut engine = create_test_engine();
3055        engine.register(EchoWorkflow).unwrap();
3056
3057        let invalid = [
3058            vec![ConcurrencyLimit::new("repo:acme", 0)],
3059            vec![ConcurrencyLimit::new("", 1)],
3060            vec![
3061                ConcurrencyLimit::new("repo:acme", 1),
3062                ConcurrencyLimit::new("repo:acme", 2),
3063            ],
3064        ];
3065        for concurrency_limits in invalid {
3066            let err = engine
3067                .enqueue_handler_with_options(
3068                    "echo-workflow",
3069                    TriggerKind::Api,
3070                    json!({}),
3071                    EnqueueOptions {
3072                        concurrency_limits,
3073                        ..Default::default()
3074                    },
3075                )
3076                .await
3077                .unwrap_err();
3078            assert!(
3079                matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3080                "{err:?}"
3081            );
3082        }
3083
3084        // Validated before the handler lookup.
3085        let err = engine
3086            .enqueue_handler_with_options(
3087                "not-registered",
3088                TriggerKind::Api,
3089                json!({}),
3090                EnqueueOptions {
3091                    concurrency_limits: vec![ConcurrencyLimit::new("repo:acme", 0)],
3092                    ..Default::default()
3093                },
3094            )
3095            .await
3096            .unwrap_err();
3097        assert!(
3098            matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3099            "{err:?}"
3100        );
3101
3102        let page = engine
3103            .store()
3104            .list_runs(RunFilter::default(), 1, 10)
3105            .await
3106            .unwrap();
3107        assert_eq!(page.total, 0, "no run may be created");
3108    }
3109
3110    struct UrgentWorkflow;
3111
3112    impl WorkflowHandler for UrgentWorkflow {
3113        fn name(&self) -> &str {
3114            "urgent-workflow"
3115        }
3116
3117        fn priority(&self) -> i16 {
3118            60
3119        }
3120
3121        fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3122            Box::pin(async { Ok(()) })
3123        }
3124    }
3125
3126    #[tokio::test]
3127    async fn enqueue_priority_defaults_to_the_handler_priority() {
3128        let mut engine = create_test_engine();
3129        engine.register(EchoWorkflow).unwrap();
3130        engine.register(UrgentWorkflow).unwrap();
3131
3132        let echo = engine
3133            .enqueue_handler_with_options(
3134                "echo-workflow",
3135                TriggerKind::Api,
3136                json!({}),
3137                EnqueueOptions::default(),
3138            )
3139            .await
3140            .unwrap()
3141            .into_run();
3142        assert_eq!(echo.priority, 0);
3143
3144        let urgent = engine
3145            .enqueue_handler_with_options(
3146                "urgent-workflow",
3147                TriggerKind::Api,
3148                json!({}),
3149                EnqueueOptions::default(),
3150            )
3151            .await
3152            .unwrap()
3153            .into_run();
3154        assert_eq!(urgent.priority, 60);
3155    }
3156
3157    #[tokio::test]
3158    async fn enqueue_priority_explicit_value_overrides_the_handler() {
3159        let mut engine = create_test_engine();
3160        engine.register(UrgentWorkflow).unwrap();
3161
3162        let run = engine
3163            .enqueue_handler_with_options(
3164                "urgent-workflow",
3165                TriggerKind::Api,
3166                json!({}),
3167                EnqueueOptions {
3168                    priority: Some(-20),
3169                    ..Default::default()
3170                },
3171            )
3172            .await
3173            .unwrap()
3174            .into_run();
3175        assert_eq!(run.priority, -20);
3176
3177        let stored = engine.store().get_run(run.id).await.unwrap().unwrap();
3178        assert_eq!(stored.priority, -20);
3179    }
3180
3181    #[tokio::test]
3182    async fn enqueue_priority_out_of_range_is_rejected() {
3183        let mut engine = create_test_engine();
3184        engine.register(EchoWorkflow).unwrap();
3185
3186        for priority in [MAX_PRIORITY + 1, MIN_PRIORITY - 1] {
3187            let err = engine
3188                .enqueue_handler_with_options(
3189                    "echo-workflow",
3190                    TriggerKind::Api,
3191                    json!({}),
3192                    EnqueueOptions {
3193                        priority: Some(priority),
3194                        ..Default::default()
3195                    },
3196                )
3197                .await
3198                .unwrap_err();
3199            assert!(matches!(err, EngineError::InvalidPriority(_)), "{err:?}");
3200        }
3201
3202        // Validated before the handler lookup.
3203        let err = engine
3204            .enqueue_handler_with_options(
3205                "not-registered",
3206                TriggerKind::Api,
3207                json!({}),
3208                EnqueueOptions {
3209                    priority: Some(MAX_PRIORITY + 1),
3210                    ..Default::default()
3211                },
3212            )
3213            .await
3214            .unwrap_err();
3215        assert!(matches!(err, EngineError::InvalidPriority(_)), "{err:?}");
3216
3217        let page = engine
3218            .store()
3219            .list_runs(RunFilter::default(), 1, 10)
3220            .await
3221            .unwrap();
3222        assert_eq!(page.total, 0, "no run may be created");
3223    }
3224
3225    #[tokio::test]
3226    async fn enqueue_priority_bounds_are_accepted() {
3227        let mut engine = create_test_engine();
3228        engine.register(EchoWorkflow).unwrap();
3229
3230        for priority in [MIN_PRIORITY, MAX_PRIORITY] {
3231            let run = engine
3232                .enqueue_handler_with_options(
3233                    "echo-workflow",
3234                    TriggerKind::Api,
3235                    json!({}),
3236                    EnqueueOptions {
3237                        priority: Some(priority),
3238                        ..Default::default()
3239                    },
3240                )
3241                .await
3242                .unwrap()
3243                .into_run();
3244            assert_eq!(run.priority, priority);
3245        }
3246    }
3247
3248    struct GpuWorkflow;
3249
3250    impl WorkflowHandler for GpuWorkflow {
3251        fn name(&self) -> &str {
3252            "gpu-workflow"
3253        }
3254
3255        fn required_worker_tags(&self) -> Vec<String> {
3256            vec!["gpu".to_string()]
3257        }
3258
3259        fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3260            Box::pin(async move { Ok(()) })
3261        }
3262    }
3263
3264    #[tokio::test]
3265    async fn enqueue_merges_handler_and_request_worker_tags() {
3266        let mut engine = create_test_engine();
3267        engine.register(GpuWorkflow).unwrap();
3268
3269        let run = engine
3270            .enqueue_handler_with_options(
3271                "gpu-workflow",
3272                TriggerKind::Api,
3273                json!({}),
3274                EnqueueOptions {
3275                    worker_tags: vec!["region:eu".to_string(), "gpu".to_string()],
3276                    ..Default::default()
3277                },
3278            )
3279            .await
3280            .unwrap()
3281            .into_run();
3282
3283        assert_eq!(
3284            run.worker_tags,
3285            vec!["gpu".to_string(), "region:eu".to_string()]
3286        );
3287    }
3288
3289    #[tokio::test]
3290    async fn enqueue_without_worker_tags_keeps_handler_tags() {
3291        let mut engine = create_test_engine();
3292        engine.register(GpuWorkflow).unwrap();
3293        engine.register(EchoWorkflow).unwrap();
3294
3295        let gpu = engine
3296            .enqueue_handler("gpu-workflow", TriggerKind::Api, json!({}), 0)
3297            .await
3298            .unwrap();
3299        assert_eq!(gpu.worker_tags, vec!["gpu".to_string()]);
3300
3301        let echo = engine
3302            .enqueue_handler("echo-workflow", TriggerKind::Api, json!({}), 0)
3303            .await
3304            .unwrap();
3305        assert!(echo.worker_tags.is_empty());
3306    }
3307
3308    #[tokio::test]
3309    async fn enqueue_rejects_invalid_worker_tags() {
3310        let mut engine = create_test_engine();
3311        engine.register(EchoWorkflow).unwrap();
3312
3313        for worker_tags in [
3314            vec!["bad,tag".to_string()],
3315            vec![" ".to_string()],
3316            vec!["x".repeat(65)],
3317        ] {
3318            let err = engine
3319                .enqueue_handler_with_options(
3320                    "echo-workflow",
3321                    TriggerKind::Api,
3322                    json!({}),
3323                    EnqueueOptions {
3324                        worker_tags,
3325                        ..Default::default()
3326                    },
3327                )
3328                .await
3329                .unwrap_err();
3330            assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3331        }
3332
3333        // Validated before the handler lookup.
3334        let err = engine
3335            .enqueue_handler_with_options(
3336                "not-registered",
3337                TriggerKind::Api,
3338                json!({}),
3339                EnqueueOptions {
3340                    worker_tags: vec!["bad,tag".to_string()],
3341                    ..Default::default()
3342                },
3343            )
3344            .await
3345            .unwrap_err();
3346        assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3347    }
3348
3349    #[test]
3350    fn worker_tags_are_unset_by_default() {
3351        let engine = create_test_engine();
3352        assert!(engine.worker_tags().is_none());
3353    }
3354
3355    #[test]
3356    fn set_worker_tags_stores_the_tags() {
3357        let mut engine = create_test_engine();
3358        engine.set_worker_tags(vec!["gpu".to_string()]);
3359        assert_eq!(engine.worker_tags(), Some(&["gpu".to_string()][..]));
3360
3361        engine.set_worker_tags(Vec::new());
3362        assert_eq!(engine.worker_tags(), Some(&[][..]));
3363    }
3364
3365    #[tokio::test]
3366    async fn run_handler_records_handler_worker_tags() {
3367        let mut engine = create_test_engine();
3368        engine.register(GpuWorkflow).unwrap();
3369
3370        let result = engine
3371            .run_handler("gpu-workflow", TriggerKind::Manual, json!({}))
3372            .await
3373            .unwrap();
3374        assert_eq!(result.run.worker_tags, vec!["gpu".to_string()]);
3375    }
3376
3377    #[tokio::test]
3378    async fn run_handler_leaves_the_run_unattributed() {
3379        let mut engine = create_test_engine();
3380        engine.register(EchoWorkflow).unwrap();
3381
3382        let run = engine
3383            .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
3384            .await
3385            .unwrap()
3386            .run;
3387
3388        assert!(run.created_by.is_none());
3389    }
3390
3391    #[tokio::test]
3392    async fn run_handler_priority_comes_from_the_handler() {
3393        let mut engine = create_test_engine();
3394        engine.register(UrgentWorkflow).unwrap();
3395
3396        let run = engine
3397            .run_handler("urgent-workflow", TriggerKind::Manual, json!({}))
3398            .await
3399            .unwrap()
3400            .run;
3401
3402        assert_eq!(run.priority, 60);
3403    }
3404
3405    #[tokio::test]
3406    async fn engine_register_boxed() {
3407        let mut engine = create_test_engine();
3408        let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
3409        let result = engine.register_boxed(handler);
3410        assert!(result.is_ok());
3411        assert_eq!(engine.handler_names().len(), 1);
3412    }
3413
3414    #[tokio::test]
3415    async fn engine_store_and_provider_accessors() {
3416        let store = Arc::new(InMemoryStore::new());
3417        let inner = ClaudeCodeProvider::new();
3418        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
3419            inner,
3420            "/tmp/ironflow-fixtures",
3421        ));
3422        let engine = Engine::new(store.clone(), provider.clone());
3423
3424        // Verify accessors return references
3425        let _ = engine.store();
3426        let _ = engine.provider();
3427    }
3428
3429    // -----------------------------------------------------------------------
3430    // Operation trait tests
3431    // -----------------------------------------------------------------------
3432
3433    use crate::operation::{Operation, OperationContext};
3434    use async_trait::async_trait;
3435    use ironflow_core::error::OperationError;
3436    use ironflow_store::models::StepKind;
3437
3438    struct FakeGitlabOp {
3439        project_id: u64,
3440        title: String,
3441    }
3442
3443    #[async_trait]
3444    impl Operation for FakeGitlabOp {
3445        fn kind(&self) -> &str {
3446            "gitlab"
3447        }
3448
3449        async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3450            Ok(json!({
3451                "issue_id": 42,
3452                "project_id": self.project_id,
3453                "title": self.title,
3454            }))
3455        }
3456
3457        fn input(&self) -> Option<Value> {
3458            Some(json!({
3459                "project_id": self.project_id,
3460                "title": self.title,
3461            }))
3462        }
3463    }
3464
3465    struct FailingOp;
3466
3467    #[async_trait]
3468    impl Operation for FailingOp {
3469        fn kind(&self) -> &str {
3470            "broken-service"
3471        }
3472
3473        async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3474            Err(OperationError::Http {
3475                status: None,
3476                message: "service unavailable".to_string(),
3477            })
3478        }
3479    }
3480
3481    struct OperationWorkflow;
3482
3483    impl WorkflowHandler for OperationWorkflow {
3484        fn name(&self) -> &str {
3485            "operation-workflow"
3486        }
3487
3488        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3489            Box::pin(async move {
3490                let op = FakeGitlabOp {
3491                    project_id: 123,
3492                    title: "Bug report".to_string(),
3493                };
3494                ctx.operation("create-issue", &op).await?;
3495                Ok(())
3496            })
3497        }
3498    }
3499
3500    struct FailingOperationWorkflow;
3501
3502    impl WorkflowHandler for FailingOperationWorkflow {
3503        fn name(&self) -> &str {
3504            "failing-operation-workflow"
3505        }
3506
3507        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3508            Box::pin(async move {
3509                ctx.operation("broken-call", &FailingOp).await?;
3510                Ok(())
3511            })
3512        }
3513    }
3514
3515    struct MixedWorkflow;
3516
3517    impl WorkflowHandler for MixedWorkflow {
3518        fn name(&self) -> &str {
3519            "mixed-workflow"
3520        }
3521
3522        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3523            Box::pin(async move {
3524                ctx.shell("build", ShellConfig::new("echo built")).await?;
3525                let op = FakeGitlabOp {
3526                    project_id: 456,
3527                    title: "Deploy done".to_string(),
3528                };
3529                let result = ctx.operation("notify-gitlab", &op).await?;
3530                assert_eq!(result.output["issue_id"], 42);
3531                Ok(())
3532            })
3533        }
3534    }
3535
3536    #[tokio::test]
3537    async fn operation_step_happy_path() {
3538        let mut engine = create_test_engine();
3539        engine.register(OperationWorkflow).unwrap();
3540
3541        let run = engine
3542            .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
3543            .await
3544            .unwrap()
3545            .run;
3546
3547        assert_eq!(run.status.state, RunStatus::Completed);
3548
3549        let steps = engine.store().list_steps(run.id).await.unwrap();
3550
3551        assert_eq!(steps.len(), 1);
3552        assert_eq!(steps[0].name, "create-issue");
3553        assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
3554        assert_eq!(
3555            steps[0].status.state,
3556            ironflow_store::models::StepStatus::Completed
3557        );
3558
3559        let output = steps[0].output.as_ref().unwrap();
3560        assert_eq!(output["issue_id"], 42);
3561        assert_eq!(output["project_id"], 123);
3562
3563        let input = steps[0].input.as_ref().unwrap();
3564        assert_eq!(input["project_id"], 123);
3565        assert_eq!(input["title"], "Bug report");
3566    }
3567
3568    #[tokio::test]
3569    async fn operation_step_failure_marks_run_failed() {
3570        let mut engine = create_test_engine();
3571        engine.register(FailingOperationWorkflow).unwrap();
3572
3573        let result = engine
3574            .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
3575            .await;
3576
3577        assert!(result.is_err());
3578    }
3579
3580    #[tokio::test]
3581    async fn operation_mixed_with_shell_steps() {
3582        let mut engine = create_test_engine();
3583        engine.register(MixedWorkflow).unwrap();
3584
3585        let run = engine
3586            .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
3587            .await
3588            .unwrap()
3589            .run;
3590
3591        assert_eq!(run.status.state, RunStatus::Completed);
3592
3593        let steps = engine.store().list_steps(run.id).await.unwrap();
3594
3595        assert_eq!(steps.len(), 2);
3596        assert_eq!(steps[0].kind, StepKind::Shell);
3597        assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
3598        assert_eq!(steps[0].position, 0);
3599        assert_eq!(steps[1].position, 1);
3600    }
3601
3602    // -----------------------------------------------------------------------
3603    // Approval + resume tests
3604    // -----------------------------------------------------------------------
3605
3606    use crate::config::ApprovalConfig;
3607
3608    struct SingleApprovalWorkflow;
3609
3610    impl WorkflowHandler for SingleApprovalWorkflow {
3611        fn name(&self) -> &str {
3612            "single-approval"
3613        }
3614
3615        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3616            Box::pin(async move {
3617                ctx.shell("build", ShellConfig::new("echo built")).await?;
3618                ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
3619                ctx.shell("deploy", ShellConfig::new("echo deployed"))
3620                    .await?;
3621                Ok(())
3622            })
3623        }
3624    }
3625
3626    struct DoubleApprovalWorkflow;
3627
3628    impl WorkflowHandler for DoubleApprovalWorkflow {
3629        fn name(&self) -> &str {
3630            "double-approval"
3631        }
3632
3633        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3634            Box::pin(async move {
3635                ctx.shell("build", ShellConfig::new("echo built")).await?;
3636                ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
3637                    .await?;
3638                ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3639                    .await?;
3640                ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3641                    .await?;
3642                ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3643                    .await?;
3644                Ok(())
3645            })
3646        }
3647    }
3648
3649    #[tokio::test]
3650    async fn approval_pauses_run() {
3651        let mut engine = create_test_engine();
3652        engine.register(SingleApprovalWorkflow).unwrap();
3653
3654        let run = engine
3655            .run_handler("single-approval", TriggerKind::Manual, json!({}))
3656            .await
3657            .unwrap()
3658            .run;
3659
3660        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3661
3662        let steps = engine.store().list_steps(run.id).await.unwrap();
3663        assert_eq!(steps.len(), 2); // build + approval gate
3664        assert_eq!(steps[0].kind, StepKind::Shell);
3665        assert_eq!(steps[0].status.state, StepStatus::Completed);
3666        assert_eq!(steps[1].kind, StepKind::Approval);
3667        assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3668    }
3669
3670    #[tokio::test]
3671    async fn approval_resume_completes_run() {
3672        let mut engine = create_test_engine();
3673        engine.register(SingleApprovalWorkflow).unwrap();
3674
3675        // First execution: pauses at approval
3676        let run = engine
3677            .run_handler("single-approval", TriggerKind::Manual, json!({}))
3678            .await
3679            .unwrap()
3680            .run;
3681        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3682
3683        // Simulate approval: transition to Running
3684        engine
3685            .store()
3686            .update_run_status(run.id, RunStatus::Running)
3687            .await
3688            .unwrap();
3689
3690        // Resume: replays build, skips approval, executes deploy
3691        let resumed = engine.resume_run(run.id).await.unwrap().run;
3692        assert_eq!(resumed.status.state, RunStatus::Completed);
3693
3694        let steps = engine.store().list_steps(run.id).await.unwrap();
3695        assert_eq!(steps.len(), 3); // build + approval + deploy
3696        assert_eq!(steps[0].name, "build");
3697        assert_eq!(steps[0].status.state, StepStatus::Completed);
3698        assert_eq!(steps[1].name, "gate");
3699        assert_eq!(steps[1].kind, StepKind::Approval);
3700        assert_eq!(steps[1].status.state, StepStatus::Completed);
3701        assert_eq!(steps[2].name, "deploy");
3702        assert_eq!(steps[2].status.state, StepStatus::Completed);
3703    }
3704
3705    #[tokio::test]
3706    async fn double_approval_two_resumes() {
3707        let mut engine = create_test_engine();
3708        engine.register(DoubleApprovalWorkflow).unwrap();
3709
3710        // First execution: pauses at staging-gate
3711        let run = engine
3712            .run_handler("double-approval", TriggerKind::Manual, json!({}))
3713            .await
3714            .unwrap()
3715            .run;
3716        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3717
3718        let steps = engine.store().list_steps(run.id).await.unwrap();
3719        assert_eq!(steps.len(), 2); // build + staging-gate
3720
3721        // First approval
3722        engine
3723            .store()
3724            .update_run_status(run.id, RunStatus::Running)
3725            .await
3726            .unwrap();
3727
3728        let resumed = engine.resume_run(run.id).await.unwrap().run;
3729        assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3730
3731        let steps = engine.store().list_steps(run.id).await.unwrap();
3732        assert_eq!(steps.len(), 4); // build + staging-gate + deploy-staging + prod-gate
3733
3734        // Second approval
3735        engine
3736            .store()
3737            .update_run_status(run.id, RunStatus::Running)
3738            .await
3739            .unwrap();
3740
3741        let final_run = engine.resume_run(run.id).await.unwrap().run;
3742        assert_eq!(final_run.status.state, RunStatus::Completed);
3743
3744        let steps = engine.store().list_steps(run.id).await.unwrap();
3745        assert_eq!(steps.len(), 5);
3746        assert_eq!(steps[0].name, "build");
3747        assert_eq!(steps[1].name, "staging-gate");
3748        assert_eq!(steps[2].name, "deploy-staging");
3749        assert_eq!(steps[3].name, "prod-gate");
3750        assert_eq!(steps[4].name, "deploy-prod");
3751
3752        for step in &steps {
3753            assert_eq!(step.status.state, StepStatus::Completed);
3754        }
3755    }
3756
3757    // -----------------------------------------------------------------------
3758    // fail_orphaned_steps tests
3759    // -----------------------------------------------------------------------
3760
3761    use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3762
3763    async fn create_step_with_status(
3764        store: &Arc<dyn Store>,
3765        run_id: Uuid,
3766        name: &str,
3767        position: u32,
3768        status: StepStatus,
3769    ) -> ironflow_store::models::Step {
3770        let step = store
3771            .create_step(NewStep {
3772                run_id,
3773                trace_id: step_trace_id(run_id, name, position),
3774                name: name.to_string(),
3775                kind: StepKind::Shell,
3776                position,
3777                input: None,
3778                is_error_handler: false,
3779            })
3780            .await
3781            .unwrap();
3782
3783        match status {
3784            StepStatus::Pending => {}
3785            StepStatus::Running => {
3786                store
3787                    .update_step(
3788                        step.id,
3789                        StepUpdate {
3790                            status: Some(StepStatus::Running),
3791                            ..StepUpdate::default()
3792                        },
3793                    )
3794                    .await
3795                    .unwrap();
3796            }
3797            StepStatus::Completed => {
3798                store
3799                    .update_step(
3800                        step.id,
3801                        StepUpdate {
3802                            status: Some(StepStatus::Running),
3803                            ..StepUpdate::default()
3804                        },
3805                    )
3806                    .await
3807                    .unwrap();
3808                store
3809                    .update_step(
3810                        step.id,
3811                        StepUpdate {
3812                            status: Some(StepStatus::Completed),
3813                            ..StepUpdate::default()
3814                        },
3815                    )
3816                    .await
3817                    .unwrap();
3818            }
3819            StepStatus::AwaitingApproval => {
3820                store
3821                    .update_step(
3822                        step.id,
3823                        StepUpdate {
3824                            status: Some(StepStatus::Running),
3825                            ..StepUpdate::default()
3826                        },
3827                    )
3828                    .await
3829                    .unwrap();
3830                store
3831                    .update_step(
3832                        step.id,
3833                        StepUpdate {
3834                            status: Some(StepStatus::AwaitingApproval),
3835                            ..StepUpdate::default()
3836                        },
3837                    )
3838                    .await
3839                    .unwrap();
3840            }
3841            _ => panic!("unsupported status for test helper: {status}"),
3842        }
3843
3844        store.get_step(step.id).await.unwrap().unwrap()
3845    }
3846
3847    #[tokio::test]
3848    async fn fail_orphaned_steps_marks_running_as_failed() {
3849        let engine = create_test_engine();
3850        let run = engine
3851            .store()
3852            .create_run(NewRun {
3853                created_by: None,
3854                workflow_name: "test".to_string(),
3855                trigger: TriggerKind::Manual,
3856                payload: json!({}),
3857                max_retries: 0,
3858                handler_version: None,
3859                labels: HashMap::new(),
3860                scheduled_at: None,
3861                idempotency_key: None,
3862                concurrency_key: None,
3863                priority: 0,
3864                concurrency_limits: Vec::new(),
3865                max_cost_usd: None,
3866                worker_tags: Vec::new(),
3867            })
3868            .await
3869            .unwrap()
3870            .into_run();
3871
3872        let step = create_step_with_status(
3873            engine.store(),
3874            run.id,
3875            "running-step",
3876            0,
3877            StepStatus::Running,
3878        )
3879        .await;
3880
3881        engine
3882            .fail_orphaned_steps(run.id, "parent run timed out")
3883            .await
3884            .unwrap();
3885
3886        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3887        assert_eq!(updated.status.state, StepStatus::Failed);
3888        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3889        assert!(updated.completed_at.is_some());
3890    }
3891
3892    #[tokio::test]
3893    async fn fail_orphaned_steps_marks_pending_as_skipped() {
3894        let engine = create_test_engine();
3895        let run = engine
3896            .store()
3897            .create_run(NewRun {
3898                created_by: None,
3899                workflow_name: "test".to_string(),
3900                trigger: TriggerKind::Manual,
3901                payload: json!({}),
3902                max_retries: 0,
3903                handler_version: None,
3904                labels: HashMap::new(),
3905                scheduled_at: None,
3906                idempotency_key: None,
3907                concurrency_key: None,
3908                priority: 0,
3909                concurrency_limits: Vec::new(),
3910                max_cost_usd: None,
3911                worker_tags: Vec::new(),
3912            })
3913            .await
3914            .unwrap()
3915            .into_run();
3916
3917        let step = create_step_with_status(
3918            engine.store(),
3919            run.id,
3920            "pending-step",
3921            0,
3922            StepStatus::Pending,
3923        )
3924        .await;
3925
3926        engine
3927            .fail_orphaned_steps(run.id, "parent run timed out")
3928            .await
3929            .unwrap();
3930
3931        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3932        assert_eq!(updated.status.state, StepStatus::Skipped);
3933        assert!(updated.error.is_none());
3934        assert!(updated.completed_at.is_some());
3935    }
3936
3937    #[tokio::test]
3938    async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
3939        let engine = create_test_engine();
3940        let run = engine
3941            .store()
3942            .create_run(NewRun {
3943                created_by: None,
3944                workflow_name: "test".to_string(),
3945                trigger: TriggerKind::Manual,
3946                payload: json!({}),
3947                max_retries: 0,
3948                handler_version: None,
3949                labels: HashMap::new(),
3950                scheduled_at: None,
3951                idempotency_key: None,
3952                concurrency_key: None,
3953                priority: 0,
3954                concurrency_limits: Vec::new(),
3955                max_cost_usd: None,
3956                worker_tags: Vec::new(),
3957            })
3958            .await
3959            .unwrap()
3960            .into_run();
3961
3962        let step = create_step_with_status(
3963            engine.store(),
3964            run.id,
3965            "approval-step",
3966            0,
3967            StepStatus::AwaitingApproval,
3968        )
3969        .await;
3970
3971        engine
3972            .fail_orphaned_steps(run.id, "parent run timed out")
3973            .await
3974            .unwrap();
3975
3976        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3977        assert_eq!(updated.status.state, StepStatus::Failed);
3978        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3979        assert!(updated.completed_at.is_some());
3980    }
3981
3982    #[tokio::test]
3983    async fn fail_orphaned_steps_skips_terminal_steps() {
3984        let engine = create_test_engine();
3985        let run = engine
3986            .store()
3987            .create_run(NewRun {
3988                created_by: None,
3989                workflow_name: "test".to_string(),
3990                trigger: TriggerKind::Manual,
3991                payload: json!({}),
3992                max_retries: 0,
3993                handler_version: None,
3994                labels: HashMap::new(),
3995                scheduled_at: None,
3996                idempotency_key: None,
3997                concurrency_key: None,
3998                priority: 0,
3999                concurrency_limits: Vec::new(),
4000                max_cost_usd: None,
4001                worker_tags: Vec::new(),
4002            })
4003            .await
4004            .unwrap()
4005            .into_run();
4006
4007        let completed_step =
4008            create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
4009        let running_step =
4010            create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
4011                .await;
4012
4013        engine
4014            .fail_orphaned_steps(run.id, "parent run timed out")
4015            .await
4016            .unwrap();
4017
4018        let completed = engine
4019            .store()
4020            .get_step(completed_step.id)
4021            .await
4022            .unwrap()
4023            .unwrap();
4024        assert_eq!(completed.status.state, StepStatus::Completed);
4025
4026        let failed = engine
4027            .store()
4028            .get_step(running_step.id)
4029            .await
4030            .unwrap()
4031            .unwrap();
4032        assert_eq!(failed.status.state, StepStatus::Failed);
4033    }
4034
4035    #[tokio::test]
4036    async fn fail_orphaned_steps_mixed_states() {
4037        let engine = create_test_engine();
4038        let run = engine
4039            .store()
4040            .create_run(NewRun {
4041                created_by: None,
4042                workflow_name: "test".to_string(),
4043                trigger: TriggerKind::Manual,
4044                payload: json!({}),
4045                max_retries: 0,
4046                handler_version: None,
4047                labels: HashMap::new(),
4048                scheduled_at: None,
4049                idempotency_key: None,
4050                concurrency_key: None,
4051                priority: 0,
4052                concurrency_limits: Vec::new(),
4053                max_cost_usd: None,
4054                worker_tags: Vec::new(),
4055            })
4056            .await
4057            .unwrap()
4058            .into_run();
4059
4060        let s_completed =
4061            create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
4062                .await;
4063        let s_running =
4064            create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
4065        let s_pending =
4066            create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
4067
4068        engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
4069
4070        let r_completed = engine
4071            .store()
4072            .get_step(s_completed.id)
4073            .await
4074            .unwrap()
4075            .unwrap();
4076        assert_eq!(r_completed.status.state, StepStatus::Completed);
4077
4078        let r_running = engine
4079            .store()
4080            .get_step(s_running.id)
4081            .await
4082            .unwrap()
4083            .unwrap();
4084        assert_eq!(r_running.status.state, StepStatus::Failed);
4085        assert_eq!(r_running.error.as_deref(), Some("timeout"));
4086
4087        let r_pending = engine
4088            .store()
4089            .get_step(s_pending.id)
4090            .await
4091            .unwrap()
4092            .unwrap();
4093        assert_eq!(r_pending.status.state, StepStatus::Skipped);
4094        assert!(r_pending.error.is_none());
4095    }
4096
4097    #[tokio::test]
4098    async fn fail_orphaned_steps_no_steps_is_noop() {
4099        let engine = create_test_engine();
4100        let run = engine
4101            .store()
4102            .create_run(NewRun {
4103                created_by: None,
4104                workflow_name: "test".to_string(),
4105                trigger: TriggerKind::Manual,
4106                payload: json!({}),
4107                max_retries: 0,
4108                handler_version: None,
4109                labels: HashMap::new(),
4110                scheduled_at: None,
4111                idempotency_key: None,
4112                concurrency_key: None,
4113                priority: 0,
4114                concurrency_limits: Vec::new(),
4115                max_cost_usd: None,
4116                worker_tags: Vec::new(),
4117            })
4118            .await
4119            .unwrap()
4120            .into_run();
4121
4122        let result = engine.fail_orphaned_steps(run.id, "timeout").await;
4123        assert!(result.is_ok());
4124    }
4125
4126    #[tokio::test]
4127    async fn fail_orphaned_steps_preserves_existing_error() {
4128        let engine = create_test_engine();
4129        let run = engine
4130            .store()
4131            .create_run(NewRun {
4132                created_by: None,
4133                workflow_name: "test".to_string(),
4134                trigger: TriggerKind::Manual,
4135                payload: json!({}),
4136                max_retries: 0,
4137                handler_version: None,
4138                labels: HashMap::new(),
4139                scheduled_at: None,
4140                idempotency_key: None,
4141                concurrency_key: None,
4142                priority: 0,
4143                concurrency_limits: Vec::new(),
4144                max_cost_usd: None,
4145                worker_tags: Vec::new(),
4146            })
4147            .await
4148            .unwrap()
4149            .into_run();
4150
4151        let step_with_error = create_step_with_status(
4152            engine.store(),
4153            run.id,
4154            "already-errored",
4155            0,
4156            StepStatus::Running,
4157        )
4158        .await;
4159
4160        engine
4161            .store()
4162            .update_step(
4163                step_with_error.id,
4164                StepUpdate {
4165                    error: Some("real error from provider".to_string()),
4166                    ..StepUpdate::default()
4167                },
4168            )
4169            .await
4170            .unwrap();
4171
4172        let step_no_error = create_step_with_status(
4173            engine.store(),
4174            run.id,
4175            "no-error-yet",
4176            1,
4177            StepStatus::Running,
4178        )
4179        .await;
4180
4181        engine
4182            .fail_orphaned_steps(run.id, "parent run failed")
4183            .await
4184            .unwrap();
4185
4186        let updated_with = engine
4187            .store()
4188            .get_step(step_with_error.id)
4189            .await
4190            .unwrap()
4191            .unwrap();
4192        assert_eq!(updated_with.status.state, StepStatus::Failed);
4193        assert_eq!(
4194            updated_with.error.as_deref(),
4195            Some("real error from provider"),
4196        );
4197
4198        let updated_without = engine
4199            .store()
4200            .get_step(step_no_error.id)
4201            .await
4202            .unwrap()
4203            .unwrap();
4204        assert_eq!(updated_without.status.state, StepStatus::Failed);
4205        assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
4206    }
4207}