Skip to main content

ironflow_engine/
engine.rs

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