Skip to main content

ironflow_engine/
engine.rs

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