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.
1651    ///
1652    /// Callers pass `retryable` explicitly rather than an error value, because
1653    /// the worker classifies failures it observes from the outside (a timeout, a
1654    /// panicked task) that never produce an [`EngineError`]. Use
1655    /// [`is_run_retryable`] to classify an
1656    /// [`EngineError`].
1657    ///
1658    /// `cost_usd` and `duration_ms` are `None` when the caller does not know the
1659    /// totals (a timeout or a panic observed from outside the handler); the
1660    /// values already stored on the run are then left untouched.
1661    ///
1662    /// Returns the status the run was moved to.
1663    ///
1664    /// # Errors
1665    ///
1666    /// Returns [`EngineError::Store`] if the run does not exist or the update
1667    /// cannot be persisted.
1668    ///
1669    /// # Examples
1670    ///
1671    /// ```no_run
1672    /// use ironflow_engine::engine::Engine;
1673    /// use ironflow_engine::error::EngineError;
1674    /// use uuid::Uuid;
1675    ///
1676    /// # async fn example(engine: &Engine, run_id: Uuid) -> Result<(), EngineError> {
1677    /// let status = engine
1678    ///     .fail_or_schedule_retry(run_id, "run timed out after 300s", true, None, None)
1679    ///     .await?;
1680    /// # Ok(())
1681    /// # }
1682    /// ```
1683    pub async fn fail_or_schedule_retry(
1684        &self,
1685        run_id: Uuid,
1686        error: &str,
1687        retryable: bool,
1688        cost_usd: Option<Decimal>,
1689        duration_ms: Option<u64>,
1690    ) -> Result<RunStatus, EngineError> {
1691        let run = self
1692            .store
1693            .get_run(run_id)
1694            .await?
1695            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1696
1697        let has_attempts_left = run.retry_count < run.max_retries;
1698        let update = if retryable && has_attempts_left {
1699            let backoff = backoff_for_retry(run.retry_count);
1700            let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
1701
1702            info!(
1703                run_id = %run_id,
1704                workflow = %run.workflow_name,
1705                attempt = run.retry_count + 1,
1706                max_retries = run.max_retries,
1707                backoff_secs = backoff.as_secs(),
1708                scheduled_at = %scheduled_at,
1709                "run failed, scheduling retry"
1710            );
1711
1712            RunUpdate {
1713                status: Some(RunStatus::Retrying),
1714                error: Some(error.to_string()),
1715                increment_retry: true,
1716                cost_usd,
1717                duration_ms,
1718                scheduled_at: Some(scheduled_at),
1719                ..RunUpdate::default()
1720            }
1721        } else {
1722            RunUpdate {
1723                status: Some(RunStatus::Failed),
1724                error: Some(error.to_string()),
1725                cost_usd,
1726                duration_ms,
1727                completed_at: Some(Utc::now()),
1728                ..RunUpdate::default()
1729            }
1730        };
1731
1732        let status = update.status.unwrap_or(RunStatus::Failed);
1733        self.store.update_run(run_id, update).await?;
1734        self.fail_orphaned_steps(run_id, error).await?;
1735
1736        Ok(status)
1737    }
1738
1739    /// Mark the steps of a run requeued after a lost worker lease as interrupted.
1740    ///
1741    /// Every `Running` step is marked `Failed` with
1742    /// [`STEP_INTERRUPTED_ERROR`](ironflow_store::store::STEP_INTERRUPTED_ERROR).
1743    /// The run keeps its attempt number, so when it is picked up again its
1744    /// finished steps are replayed and each interrupted step is executed again
1745    /// at the same position; an interrupted `Workflow` step re-enters the same
1746    /// child run. `Pending` and `AwaitingApproval` steps are left as they are,
1747    /// unlike [`fail_orphaned_steps`](Self::fail_orphaned_steps), which ends
1748    /// the run's steps for good.
1749    ///
1750    /// Errors from individual step updates are logged but do not abort the cleanup.
1751    ///
1752    /// # Errors
1753    ///
1754    /// Returns [`EngineError`] if listing steps fails.
1755    ///
1756    /// # Examples
1757    ///
1758    /// ```no_run
1759    /// use ironflow_engine::engine::Engine;
1760    /// use ironflow_engine::error::EngineError;
1761    /// use uuid::Uuid;
1762    ///
1763    /// # async fn example(engine: &Engine, run_id: Uuid) -> Result<(), EngineError> {
1764    /// engine.interrupt_running_steps(run_id).await?;
1765    /// # Ok(())
1766    /// # }
1767    /// ```
1768    pub async fn interrupt_running_steps(&self, run_id: Uuid) -> Result<(), EngineError> {
1769        interrupt_running_steps(self.store.as_ref(), run_id).await
1770    }
1771
1772    /// Fail all non-terminal steps for a run.
1773    ///
1774    /// Called after a run is marked as failed (timeout, error, panic) to clean up
1775    /// orphaned steps that are still in `Running`, `Pending`, or `AwaitingApproval`.
1776    ///
1777    /// - `Running` / `AwaitingApproval` steps are marked `Failed`.
1778    /// - `Pending` steps are marked `Skipped` (FSM does not allow Pending -> Failed).
1779    ///
1780    /// Errors from individual step updates are logged but do not abort the cleanup.
1781    ///
1782    /// # Errors
1783    ///
1784    /// Returns [`EngineError`] if listing steps fails.
1785    pub async fn fail_orphaned_steps(
1786        &self,
1787        run_id: Uuid,
1788        error_message: &str,
1789    ) -> Result<(), EngineError> {
1790        let steps = self.store.list_steps(run_id).await?;
1791        let now = Utc::now();
1792
1793        for step in steps {
1794            if step.status.state.is_terminal() {
1795                continue;
1796            }
1797
1798            let (target_status, error) = match step.status.state {
1799                StepStatus::Running | StepStatus::AwaitingApproval => {
1800                    let err = if step.error.is_some() {
1801                        None
1802                    } else {
1803                        Some(error_message.to_string())
1804                    };
1805                    (StepStatus::Failed, err)
1806                }
1807                StepStatus::Pending => (StepStatus::Skipped, None),
1808                _ => continue,
1809            };
1810
1811            if let Err(e) = self
1812                .store
1813                .update_step(
1814                    step.id,
1815                    StepUpdate {
1816                        status: Some(target_status),
1817                        error,
1818                        completed_at: Some(now),
1819                        ..StepUpdate::default()
1820                    },
1821                )
1822                .await
1823            {
1824                warn!(
1825                    run_id = %run_id,
1826                    step_id = %step.id,
1827                    step_name = %step.name,
1828                    error = %e,
1829                    "failed to cleanup orphaned step"
1830                );
1831            } else {
1832                info!(
1833                    run_id = %run_id,
1834                    step_id = %step.id,
1835                    step_name = %step.name,
1836                    from = %step.status.state,
1837                    to = %target_status,
1838                    "cleaned up orphaned step"
1839                );
1840            }
1841        }
1842
1843        Ok(())
1844    }
1845
1846    /// Execute `handler` once [`AgentProvider::release_run`] has stopped
1847    /// whatever a previous execution of the run left running (an agent pod
1848    /// writing to a shared worktree). A failed release fails the execution
1849    /// before its first step, with [`EngineError::Operation`].
1850    async fn release_then_execute(
1851        &self,
1852        run_id: Uuid,
1853        handler: &dyn WorkflowHandler,
1854        ctx: &mut WorkflowContext,
1855    ) -> Result<(), EngineError> {
1856        match self.provider.release_run(&run_id.to_string()).await {
1857            Ok(()) => handler.execute(ctx).await,
1858            Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
1859        }
1860    }
1861
1862    /// Finalize a run with the given result and context.
1863    ///
1864    /// On success: updates run to Completed with cost, duration, and completed_at.
1865    /// On failure: updates run to Failed with error, cost, duration, and completed_at.
1866    /// Always: fetches and returns the final Run.
1867    async fn finalize_run(
1868        &self,
1869        run_id: Uuid,
1870        workflow_name: &str,
1871        result: Result<(), EngineError>,
1872        ctx: &WorkflowContext,
1873        run_start: Instant,
1874        run_labels: HashMap<String, String>,
1875    ) -> Result<WorkflowResult, EngineError> {
1876        // Covers the whole run: previous attempts plus this one, so a retried
1877        // run reports the time it really consumed.
1878        let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
1879        let completed_at = Utc::now();
1880
1881        let final_status;
1882        let final_run;
1883
1884        match result {
1885            Ok(()) => {
1886                final_status = if ctx.has_allowed_failure() {
1887                    RunStatus::Warning
1888                } else {
1889                    RunStatus::Completed
1890                };
1891                final_run = self
1892                    .store
1893                    .update_run_returning(
1894                        run_id,
1895                        RunUpdate {
1896                            status: Some(final_status),
1897                            cost_usd: Some(ctx.total_cost_usd()),
1898                            duration_ms: Some(total_duration),
1899                            completed_at: Some(completed_at),
1900                            output: ctx.output().cloned(),
1901                            ..RunUpdate::default()
1902                        },
1903                    )
1904                    .await?;
1905
1906                info!(
1907                    run_id = %run_id,
1908                    status = %final_status,
1909                    cost_usd = %ctx.total_cost_usd(),
1910                    duration_ms = total_duration,
1911                    "run completed"
1912                );
1913            }
1914            Err(EngineError::ApprovalRequired {
1915                run_id: approval_run_id,
1916                step_id,
1917                ref message,
1918            }) => {
1919                final_status = RunStatus::AwaitingApproval;
1920                final_run = self
1921                    .store
1922                    .update_run_returning(
1923                        run_id,
1924                        RunUpdate {
1925                            status: Some(RunStatus::AwaitingApproval),
1926                            cost_usd: Some(ctx.total_cost_usd()),
1927                            duration_ms: Some(total_duration),
1928                            ..RunUpdate::default()
1929                        },
1930                    )
1931                    .await?;
1932
1933                info!(
1934                    run_id = %approval_run_id,
1935                    step_id = %step_id,
1936                    message = %message,
1937                    "run awaiting approval"
1938                );
1939
1940                self.publish_approval_requested(approval_run_id, step_id, message)
1941                    .await?;
1942            }
1943            Err(EngineError::ChildSuspended {
1944                run_id: child_run_id,
1945                ref cause,
1946            }) => {
1947                final_status = cause.suspension_status();
1948                // No `scheduled_at`: the suspended descendant owns the wake-up
1949                // and resumes this run through `resume_chain`. A wake-up armed
1950                // here too would resume the chain twice.
1951                final_run = self
1952                    .store
1953                    .update_run_returning(
1954                        run_id,
1955                        RunUpdate {
1956                            status: Some(final_status),
1957                            cost_usd: Some(ctx.total_cost_usd()),
1958                            duration_ms: Some(total_duration),
1959                            ..RunUpdate::default()
1960                        },
1961                    )
1962                    .await?;
1963
1964                let leaf = cause.suspension_leaf();
1965                info!(
1966                    run_id = %run_id,
1967                    child_run_id = %child_run_id,
1968                    status = %final_status,
1969                    cause = %leaf,
1970                    "run suspended with its child run"
1971                );
1972
1973                match leaf {
1974                    EngineError::ApprovalRequired {
1975                        run_id: approval_run_id,
1976                        step_id,
1977                        message,
1978                    } => {
1979                        self.publish_approval_requested(*approval_run_id, *step_id, message)
1980                            .await?;
1981                    }
1982                    EngineError::SignalWaiting {
1983                        run_id: wait_run_id,
1984                        step_id,
1985                        step_name,
1986                        name,
1987                        key,
1988                        deadline_at,
1989                    } => {
1990                        self.event_publisher
1991                            .publish(Event::SignalAwaited(SignalAwaitedEvent {
1992                                run_id: *wait_run_id,
1993                                step_id: *step_id,
1994                                step_name: step_name.clone(),
1995                                name: name.clone(),
1996                                key: key.clone(),
1997                                deadline_at: *deadline_at,
1998                                at: Utc::now(),
1999                            }));
2000                    }
2001                    // A human input or a delay publishes no suspension event,
2002                    // like on a top-level run.
2003                    _ => {}
2004                }
2005            }
2006            Err(EngineError::HumanInputRequired {
2007                run_id: input_run_id,
2008                step_id,
2009                ref message,
2010            }) => {
2011                final_status = RunStatus::AwaitingApproval;
2012                final_run = self
2013                    .store
2014                    .update_run_returning(
2015                        run_id,
2016                        RunUpdate {
2017                            status: Some(RunStatus::AwaitingApproval),
2018                            cost_usd: Some(ctx.total_cost_usd()),
2019                            duration_ms: Some(total_duration),
2020                            ..RunUpdate::default()
2021                        },
2022                    )
2023                    .await?;
2024
2025                // No `ApprovalRequested` event: a human input is not an approval.
2026                info!(
2027                    run_id = %input_run_id,
2028                    step_id = %step_id,
2029                    message = %message,
2030                    "run awaiting human input"
2031                );
2032            }
2033            Err(EngineError::DelaySleeping {
2034                run_id: delay_run_id,
2035                step_id,
2036                wake_at,
2037            }) => {
2038                final_status = RunStatus::Sleeping;
2039                final_run = self
2040                    .store
2041                    .update_run_returning(
2042                        run_id,
2043                        RunUpdate {
2044                            status: Some(RunStatus::Sleeping),
2045                            cost_usd: Some(ctx.total_cost_usd()),
2046                            duration_ms: Some(total_duration),
2047                            scheduled_at: Some(wake_at),
2048                            ..RunUpdate::default()
2049                        },
2050                    )
2051                    .await?;
2052
2053                info!(
2054                    run_id = %delay_run_id,
2055                    step_id = %step_id,
2056                    wake_at = %wake_at,
2057                    "run sleeping until delay elapses"
2058                );
2059            }
2060            Err(EngineError::SignalWaiting {
2061                run_id: wait_run_id,
2062                step_id,
2063                ref step_name,
2064                ref name,
2065                ref key,
2066                deadline_at,
2067            }) => {
2068                final_status = RunStatus::Sleeping;
2069                // Atomic with the step lock: a signal delivered since the step
2070                // opened leaves the run due right away instead of until the
2071                // deadline.
2072                let waiting = self
2073                    .store
2074                    .suspend_run_on_signal(run_id, step_id, deadline_at)
2075                    .await?;
2076                final_run = self
2077                    .store
2078                    .update_run_returning(
2079                        run_id,
2080                        RunUpdate {
2081                            cost_usd: Some(ctx.total_cost_usd()),
2082                            duration_ms: Some(total_duration),
2083                            ..RunUpdate::default()
2084                        },
2085                    )
2086                    .await?;
2087
2088                if waiting {
2089                    self.event_publisher
2090                        .publish(Event::SignalAwaited(SignalAwaitedEvent {
2091                            run_id: wait_run_id,
2092                            step_id,
2093                            step_name: step_name.clone(),
2094                            name: name.clone(),
2095                            key: key.clone(),
2096                            deadline_at,
2097                            at: Utc::now(),
2098                        }));
2099                }
2100
2101                info!(
2102                    run_id = %wait_run_id,
2103                    step_id = %step_id,
2104                    signal = %name,
2105                    key = %key,
2106                    deadline_at = %deadline_at,
2107                    waiting,
2108                    "run sleeping until a signal arrives"
2109                );
2110            }
2111            Err(err) => {
2112                // A guardrail stop (budget or workflow guard) is deliberate,
2113                // not a breakage: the run is cancelled, never failed and
2114                // never replayed.
2115                let guardrail_stop = matches!(
2116                    err,
2117                    EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2118                );
2119
2120                final_status = if guardrail_stop {
2121                    if let Err(store_err) = self
2122                        .store
2123                        .update_run(
2124                            run_id,
2125                            RunUpdate {
2126                                status: Some(RunStatus::Cancelled),
2127                                error: Some(err.to_string()),
2128                                cost_usd: Some(ctx.total_cost_usd()),
2129                                duration_ms: Some(total_duration),
2130                                completed_at: Some(completed_at),
2131                                output: ctx.output().cloned(),
2132                                ..RunUpdate::default()
2133                            },
2134                        )
2135                        .await
2136                    {
2137                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2138                    }
2139                    if let Err(cleanup_err) = self
2140                        .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2141                        .await
2142                    {
2143                        error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2144                    }
2145                    RunStatus::Cancelled
2146                } else {
2147                    // Written before the failure so a failed run keeps the
2148                    // verdict the handler set before returning its error.
2149                    if let Some(output) = ctx.output()
2150                        && let Err(store_err) = self
2151                            .store
2152                            .update_run(
2153                                run_id,
2154                                RunUpdate {
2155                                    output: Some(output.clone()),
2156                                    ..RunUpdate::default()
2157                                },
2158                            )
2159                            .await
2160                    {
2161                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run output");
2162                    }
2163                    self.fail_or_schedule_retry(
2164                        run_id,
2165                        &err.to_string(),
2166                        is_run_retryable(&err),
2167                        Some(ctx.total_cost_usd()),
2168                        Some(total_duration),
2169                    )
2170                    .await
2171                    .unwrap_or_else(|store_err| {
2172                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2173                        RunStatus::Failed
2174                    })
2175                };
2176
2177                if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2178                    self.on_run_budget_exceeded(workflow_name, run_id, &err);
2179                }
2180
2181                error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2182
2183                self.publish_run_status_changed(
2184                    workflow_name,
2185                    run_id,
2186                    final_status,
2187                    Some(err.to_string()),
2188                    ctx,
2189                    total_duration,
2190                    run_labels,
2191                );
2192
2193                #[cfg(feature = "prometheus")]
2194                self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2195
2196                return Err(err);
2197            }
2198        }
2199
2200        self.publish_run_status_changed(
2201            workflow_name,
2202            run_id,
2203            final_status,
2204            None,
2205            ctx,
2206            total_duration,
2207            run_labels,
2208        );
2209
2210        #[cfg(feature = "prometheus")]
2211        self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2212
2213        Ok(WorkflowResult {
2214            run: final_run,
2215            steps: ctx.step_results().to_vec(),
2216        })
2217    }
2218
2219    /// Publish [`Event::ApprovalRequested`] for an approval gate that opened
2220    /// on `run_id`, with the requirement recorded when the gate opened.
2221    async fn publish_approval_requested(
2222        &self,
2223        run_id: Uuid,
2224        step_id: Uuid,
2225        message: &str,
2226    ) -> Result<(), EngineError> {
2227        let requirement = self
2228            .store
2229            .get_step(step_id)
2230            .await?
2231            .and_then(|s| s.approval_requirement);
2232        self.event_publisher
2233            .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2234                run_id,
2235                step_id,
2236                message: message.to_string(),
2237                requirement,
2238                at: Utc::now(),
2239            }));
2240        Ok(())
2241    }
2242
2243    /// Fail every ancestor of a child run, closest first.
2244    ///
2245    /// Used when a gate inside a sub-workflow is rejected: the child run
2246    /// fails, and the runs suspended with it (its parent, up to the root)
2247    /// must not stay suspended on a child that will never resume. Each
2248    /// ancestor goes through [`fail_or_schedule_retry`](Self::fail_or_schedule_retry)
2249    /// without a retry, which also fails its open `Workflow` step. A run that
2250    /// is not a child of a sub-workflow has no ancestor: nothing happens.
2251    ///
2252    /// # Errors
2253    ///
2254    /// Returns [`EngineError::Store`] if a run of the chain does not exist or
2255    /// its failure cannot be persisted.
2256    ///
2257    /// # Examples
2258    ///
2259    /// ```no_run
2260    /// use ironflow_engine::engine::Engine;
2261    /// use ironflow_engine::error::EngineError;
2262    /// use uuid::Uuid;
2263    ///
2264    /// # async fn example(engine: &Engine, child_run_id: Uuid) -> Result<(), EngineError> {
2265    /// engine
2266    ///     .fail_ancestors(child_run_id, "approval rejected in a sub-workflow")
2267    ///     .await?;
2268    /// # Ok(())
2269    /// # }
2270    /// ```
2271    pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2272        let mut current = self
2273            .store
2274            .get_run(run_id)
2275            .await?
2276            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2277        // Labels are data: a chain that loops back on itself stops there.
2278        let mut visited = HashSet::from([run_id]);
2279
2280        while let Some(parent_id) = chain_parent(&current) {
2281            if !visited.insert(parent_id) {
2282                break;
2283            }
2284            let status = self
2285                .fail_or_schedule_retry(parent_id, reason, false, None, None)
2286                .await?;
2287            info!(
2288                run_id = %run_id,
2289                ancestor_run_id = %parent_id,
2290                status = %status,
2291                "ancestor run failed with its child"
2292            );
2293            current = self
2294                .store
2295                .get_run(parent_id)
2296                .await?
2297                .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2298        }
2299
2300        Ok(())
2301    }
2302
2303    /// Emit Prometheus metrics for a completed run.
2304    #[cfg(feature = "prometheus")]
2305    fn emit_run_metrics(
2306        &self,
2307        workflow_name: &str,
2308        status: RunStatus,
2309        duration_ms: u64,
2310        ctx: &WorkflowContext,
2311    ) {
2312        let status_str = status.to_string();
2313        let wf = workflow_name.to_string();
2314
2315        counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2316        histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2317            .record(duration_ms as f64 / 1000.0);
2318        histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2319            ctx.total_cost_usd()
2320                .to_string()
2321                .parse::<f64>()
2322                .unwrap_or(0.0),
2323        );
2324        gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2325    }
2326
2327    /// Record the metric and publish the audit event for a run that hit its
2328    /// cost cap.
2329    ///
2330    /// A non-[`RunBudgetExceeded`](EngineError::RunBudgetExceeded) error is
2331    /// ignored, so callers can pass the error unconditionally.
2332    fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2333        let EngineError::RunBudgetExceeded {
2334            limit_usd,
2335            spent_usd,
2336            step_budget_usd,
2337            ..
2338        } = err
2339        else {
2340            return;
2341        };
2342
2343        #[cfg(feature = "prometheus")]
2344        counter!(
2345            RUN_BUDGET_EXCEEDED_TOTAL,
2346            "workflow" => workflow_name.to_string(),
2347            "scope" => "run",
2348        )
2349        .increment(1);
2350
2351        self.event_publisher
2352            .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2353                run_id,
2354                workflow_name: workflow_name.to_string(),
2355                limit_usd: *limit_usd,
2356                spent_usd: *spent_usd,
2357                step_budget_usd: *step_budget_usd,
2358                at: Utc::now(),
2359            }));
2360    }
2361
2362    /// Publish a run status changed event to all registered subscribers.
2363    ///
2364    /// `from` is always `Running` because `finalize_run` is only called
2365    /// from a running state.
2366    #[allow(clippy::too_many_arguments)]
2367    fn publish_run_status_changed(
2368        &self,
2369        workflow_name: &str,
2370        run_id: Uuid,
2371        to: RunStatus,
2372        error: Option<String>,
2373        ctx: &WorkflowContext,
2374        duration_ms: u64,
2375        labels: HashMap<String, String>,
2376    ) {
2377        let now = Utc::now();
2378        let cost_usd = ctx.total_cost_usd();
2379        let wf = workflow_name.to_string();
2380
2381        self.event_publisher
2382            .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2383                run_id,
2384                workflow_name: wf.clone(),
2385                from: RunStatus::Running,
2386                to,
2387                error: error.clone(),
2388                cost_usd,
2389                duration_ms,
2390                labels: labels.clone(),
2391                at: now,
2392            }));
2393
2394        if to == RunStatus::Failed {
2395            self.event_publisher
2396                .publish(Event::RunFailed(RunFailedEvent {
2397                    run_id,
2398                    workflow_name: wf,
2399                    error,
2400                    cost_usd,
2401                    duration_ms,
2402                    labels,
2403                    at: now,
2404                }));
2405        }
2406    }
2407}
2408
2409impl fmt::Debug for Engine {
2410    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2411        f.debug_struct("Engine")
2412            .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2413            .finish_non_exhaustive()
2414    }
2415}
2416
2417#[cfg(test)]
2418mod tests {
2419    use super::*;
2420    use crate::config::ShellConfig;
2421    use crate::handler::{HandlerFuture, WorkflowHandler};
2422    use ironflow_core::providers::claude::ClaudeCodeProvider;
2423    use ironflow_core::providers::record_replay::RecordReplayProvider;
2424    use ironflow_store::memory::InMemoryStore;
2425    use ironflow_store::models::StepStatus;
2426    use serde_json::json;
2427
2428    // Test handler that echoes a message via shell
2429    struct EchoWorkflow;
2430
2431    impl WorkflowHandler for EchoWorkflow {
2432        fn name(&self) -> &str {
2433            "echo-workflow"
2434        }
2435
2436        fn describe(&self) -> WorkflowInfo {
2437            WorkflowInfo {
2438                description: "A simple workflow that echoes hello".to_string(),
2439                source_code: None,
2440                sub_workflows: Vec::new(),
2441                category: None,
2442                version: self.version().map(str::to_string),
2443                compatible_versions: Vec::new(),
2444                input_schema: None,
2445                default_labels: HashMap::new(),
2446                schedule: self.schedule().cloned(),
2447                default_max_cost_usd: self.default_max_cost_usd(),
2448            }
2449        }
2450
2451        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2452            Box::pin(async move {
2453                ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2454                Ok(())
2455            })
2456        }
2457    }
2458
2459    // Test handler that fails
2460    struct FailingWorkflow;
2461
2462    impl WorkflowHandler for FailingWorkflow {
2463        fn name(&self) -> &str {
2464            "failing-workflow"
2465        }
2466
2467        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2468            Box::pin(async move {
2469                ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2470                Ok(())
2471            })
2472        }
2473    }
2474
2475    fn create_test_engine() -> Engine {
2476        let store = Arc::new(InMemoryStore::new());
2477        let inner = ClaudeCodeProvider::new();
2478        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2479            inner,
2480            "/tmp/ironflow-fixtures",
2481        ));
2482        Engine::new(store, provider)
2483    }
2484
2485    #[test]
2486    fn engine_new_creates_instance() {
2487        let engine = create_test_engine();
2488        assert_eq!(engine.handler_names().len(), 0);
2489    }
2490
2491    #[test]
2492    fn execution_mode_defaults_to_local() {
2493        let engine = create_test_engine();
2494        assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2495    }
2496
2497    #[test]
2498    fn with_execution_mode_overrides_the_default() {
2499        let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2500        assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2501    }
2502
2503    #[test]
2504    fn engine_register_handler() {
2505        let mut engine = create_test_engine();
2506        let result = engine.register(EchoWorkflow);
2507        assert!(result.is_ok());
2508        assert_eq!(engine.handler_names().len(), 1);
2509        assert!(engine.handler_names().contains(&"echo-workflow"));
2510    }
2511
2512    #[test]
2513    fn engine_register_duplicate_returns_error() {
2514        let mut engine = create_test_engine();
2515        engine.register(EchoWorkflow).unwrap();
2516        let result = engine.register(EchoWorkflow);
2517        assert!(result.is_err());
2518    }
2519
2520    #[test]
2521    fn engine_get_handler_found() {
2522        let mut engine = create_test_engine();
2523        engine.register(EchoWorkflow).unwrap();
2524        let handler = engine.get_handler("echo-workflow");
2525        assert!(handler.is_some());
2526    }
2527
2528    #[test]
2529    fn engine_get_handler_not_found() {
2530        let engine = create_test_engine();
2531        let handler = engine.get_handler("nonexistent");
2532        assert!(handler.is_none());
2533    }
2534
2535    #[test]
2536    fn engine_handler_names_lists_all() {
2537        let mut engine = create_test_engine();
2538        engine.register(EchoWorkflow).unwrap();
2539        engine.register(FailingWorkflow).unwrap();
2540        let names = engine.handler_names();
2541        assert_eq!(names.len(), 2);
2542        assert!(names.contains(&"echo-workflow"));
2543        assert!(names.contains(&"failing-workflow"));
2544    }
2545
2546    #[test]
2547    fn engine_handler_info_returns_description() {
2548        let mut engine = create_test_engine();
2549        engine.register(EchoWorkflow).unwrap();
2550        let info = engine.handler_info("echo-workflow");
2551        assert!(info.is_some());
2552        let info = info.unwrap();
2553        assert_eq!(info.description, "A simple workflow that echoes hello");
2554    }
2555
2556    struct CategorizedWorkflow;
2557
2558    impl WorkflowHandler for CategorizedWorkflow {
2559        fn name(&self) -> &str {
2560            "categorized"
2561        }
2562        fn category(&self) -> Option<&str> {
2563            Some("data/etl")
2564        }
2565        fn execute<'a>(
2566            &'a self,
2567            _ctx: &'a mut WorkflowContext,
2568        ) -> crate::handler::HandlerFuture<'a> {
2569            Box::pin(async move { Ok(()) })
2570        }
2571    }
2572
2573    #[test]
2574    fn engine_default_describe_propagates_category() {
2575        let mut engine = create_test_engine();
2576        engine.register(CategorizedWorkflow).unwrap();
2577        let info = engine.handler_info("categorized").unwrap();
2578        assert_eq!(info.category.as_deref(), Some("data/etl"));
2579    }
2580
2581    #[test]
2582    fn engine_default_describe_without_category() {
2583        let mut engine = create_test_engine();
2584        engine.register(EchoWorkflow).unwrap();
2585        let info = engine.handler_info("echo-workflow").unwrap();
2586        assert!(info.category.is_none());
2587    }
2588
2589    // -----------------------------------------------------------------------
2590    // Schedule tests
2591    // -----------------------------------------------------------------------
2592
2593    struct ScheduledWorkflow {
2594        schedule: CronSchedule,
2595    }
2596
2597    impl ScheduledWorkflow {
2598        fn new() -> Self {
2599            Self {
2600                schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2601            }
2602        }
2603    }
2604
2605    impl WorkflowHandler for ScheduledWorkflow {
2606        fn name(&self) -> &str {
2607            "scheduled"
2608        }
2609        fn schedule(&self) -> Option<&CronSchedule> {
2610            Some(&self.schedule)
2611        }
2612        fn execute<'a>(
2613            &'a self,
2614            _ctx: &'a mut WorkflowContext,
2615        ) -> crate::handler::HandlerFuture<'a> {
2616            Box::pin(async move { Ok(()) })
2617        }
2618    }
2619
2620    #[test]
2621    fn engine_default_describe_propagates_schedule() {
2622        let mut engine = create_test_engine();
2623        engine.register(ScheduledWorkflow::new()).unwrap();
2624        let info = engine.handler_info("scheduled").unwrap();
2625        assert_eq!(
2626            info.schedule.as_ref().map(|s| s.as_str()),
2627            Some("0 0 * * * *")
2628        );
2629    }
2630
2631    #[test]
2632    fn engine_default_describe_without_schedule() {
2633        let mut engine = create_test_engine();
2634        engine.register(EchoWorkflow).unwrap();
2635        let info = engine.handler_info("echo-workflow").unwrap();
2636        assert!(info.schedule.is_none());
2637    }
2638
2639    #[test]
2640    fn scheduled_handlers_returns_only_scheduled() {
2641        let mut engine = create_test_engine();
2642        engine.register(EchoWorkflow).unwrap();
2643        engine.register(ScheduledWorkflow::new()).unwrap();
2644        engine.register(FailingWorkflow).unwrap();
2645
2646        let scheduled = engine.scheduled_handlers();
2647        assert_eq!(scheduled.len(), 1);
2648        assert_eq!(scheduled[0].0, "scheduled");
2649        assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
2650    }
2651
2652    #[test]
2653    fn scheduled_handlers_empty_when_none_scheduled() {
2654        let mut engine = create_test_engine();
2655        engine.register(EchoWorkflow).unwrap();
2656        engine.register(FailingWorkflow).unwrap();
2657
2658        let scheduled = engine.scheduled_handlers();
2659        assert!(scheduled.is_empty());
2660    }
2661
2662    struct BadCategoryWorkflow(&'static str);
2663
2664    impl WorkflowHandler for BadCategoryWorkflow {
2665        fn name(&self) -> &str {
2666            "bad-category"
2667        }
2668        fn category(&self) -> Option<&str> {
2669            Some(self.0)
2670        }
2671        fn execute<'a>(
2672            &'a self,
2673            _ctx: &'a mut WorkflowContext,
2674        ) -> crate::handler::HandlerFuture<'a> {
2675            Box::pin(async move { Ok(()) })
2676        }
2677    }
2678
2679    #[test]
2680    fn engine_register_rejects_empty_category() {
2681        let mut engine = create_test_engine();
2682        let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
2683        match err {
2684            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
2685            other => panic!("expected InvalidWorkflow, got {other:?}"),
2686        }
2687    }
2688
2689    #[test]
2690    fn engine_register_rejects_leading_slash_category() {
2691        let mut engine = create_test_engine();
2692        let err = engine
2693            .register(BadCategoryWorkflow("/data/etl"))
2694            .unwrap_err();
2695        match err {
2696            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
2697            other => panic!("expected InvalidWorkflow, got {other:?}"),
2698        }
2699    }
2700
2701    #[test]
2702    fn engine_register_rejects_trailing_slash_category() {
2703        let mut engine = create_test_engine();
2704        let err = engine
2705            .register(BadCategoryWorkflow("data/etl/"))
2706            .unwrap_err();
2707        match err {
2708            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
2709            other => panic!("expected InvalidWorkflow, got {other:?}"),
2710        }
2711    }
2712
2713    #[test]
2714    fn engine_register_rejects_double_slash_category() {
2715        let mut engine = create_test_engine();
2716        let err = engine
2717            .register(BadCategoryWorkflow("data//etl"))
2718            .unwrap_err();
2719        match err {
2720            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
2721            other => panic!("expected InvalidWorkflow, got {other:?}"),
2722        }
2723    }
2724
2725    #[test]
2726    fn engine_register_rejects_whitespace_only_segment_category() {
2727        let mut engine = create_test_engine();
2728        let err = engine
2729            .register(BadCategoryWorkflow("data/ /etl"))
2730            .unwrap_err();
2731        match err {
2732            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
2733            other => panic!("expected InvalidWorkflow, got {other:?}"),
2734        }
2735    }
2736
2737    #[test]
2738    fn engine_register_accepts_valid_nested_category() {
2739        let mut engine = create_test_engine();
2740        assert!(engine.register(CategorizedWorkflow).is_ok());
2741    }
2742
2743    #[tokio::test]
2744    async fn engine_unknown_workflow_returns_error() {
2745        let engine = create_test_engine();
2746        let result = engine
2747            .run_handler("unknown", TriggerKind::Manual, json!({}))
2748            .await;
2749        assert!(result.is_err());
2750        match result {
2751            Err(EngineError::InvalidWorkflow(msg)) => {
2752                assert!(msg.contains("no handler registered"));
2753            }
2754            _ => panic!("expected InvalidWorkflow error"),
2755        }
2756    }
2757
2758    #[tokio::test]
2759    async fn engine_enqueue_handler_creates_pending_run() {
2760        let mut engine = create_test_engine();
2761        engine.register(EchoWorkflow).unwrap();
2762
2763        let run = engine
2764            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2765            .await
2766            .unwrap();
2767        assert_eq!(run.status.state, RunStatus::Pending);
2768        assert_eq!(run.workflow_name, "echo-workflow");
2769    }
2770
2771    #[tokio::test]
2772    async fn enqueue_handler_leaves_the_run_unattributed() {
2773        let mut engine = create_test_engine();
2774        engine.register(EchoWorkflow).unwrap();
2775
2776        let run = engine
2777            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
2778            .await
2779            .unwrap();
2780
2781        assert!(run.created_by.is_none());
2782    }
2783
2784    #[tokio::test]
2785    async fn enqueue_handler_with_options_records_the_author() {
2786        let mut engine = create_test_engine();
2787        engine.register(EchoWorkflow).unwrap();
2788        let actor = RunActor::User {
2789            user_id: Uuid::now_v7(),
2790        };
2791
2792        let run = engine
2793            .enqueue_handler_with_options(
2794                "echo-workflow",
2795                TriggerKind::Api,
2796                json!({}),
2797                EnqueueOptions {
2798                    created_by: Some(actor.clone()),
2799                    ..Default::default()
2800                },
2801            )
2802            .await
2803            .unwrap()
2804            .into_run();
2805
2806        assert_eq!(run.created_by, Some(actor));
2807    }
2808
2809    #[tokio::test]
2810    async fn enqueue_handler_with_options_accepts_no_author() {
2811        let mut engine = create_test_engine();
2812        engine.register(EchoWorkflow).unwrap();
2813
2814        let run = engine
2815            .enqueue_handler_with_options(
2816                "echo-workflow",
2817                TriggerKind::Cron {
2818                    schedule: "0 * * * * *".to_string(),
2819                },
2820                json!({}),
2821                EnqueueOptions::default(),
2822            )
2823            .await
2824            .unwrap()
2825            .into_run();
2826
2827        assert!(run.created_by.is_none());
2828    }
2829
2830    #[tokio::test]
2831    async fn enqueue_handler_with_options_stores_concurrency_limits() {
2832        let mut engine = create_test_engine();
2833        engine.register(EchoWorkflow).unwrap();
2834        let limits = vec![
2835            ConcurrencyLimit::new("repo:acme", 2),
2836            ConcurrencyLimit::new("tenant:42", 5),
2837        ];
2838
2839        let run = engine
2840            .enqueue_handler_with_options(
2841                "echo-workflow",
2842                TriggerKind::Api,
2843                json!({}),
2844                EnqueueOptions {
2845                    concurrency_limits: limits.clone(),
2846                    ..Default::default()
2847                },
2848            )
2849            .await
2850            .unwrap()
2851            .into_run();
2852
2853        assert_eq!(run.concurrency_limits, limits);
2854    }
2855
2856    #[tokio::test]
2857    async fn enqueue_rejects_invalid_concurrency_limits() {
2858        let mut engine = create_test_engine();
2859        engine.register(EchoWorkflow).unwrap();
2860
2861        let invalid = [
2862            vec![ConcurrencyLimit::new("repo:acme", 0)],
2863            vec![ConcurrencyLimit::new("", 1)],
2864            vec![
2865                ConcurrencyLimit::new("repo:acme", 1),
2866                ConcurrencyLimit::new("repo:acme", 2),
2867            ],
2868        ];
2869        for concurrency_limits in invalid {
2870            let err = engine
2871                .enqueue_handler_with_options(
2872                    "echo-workflow",
2873                    TriggerKind::Api,
2874                    json!({}),
2875                    EnqueueOptions {
2876                        concurrency_limits,
2877                        ..Default::default()
2878                    },
2879                )
2880                .await
2881                .unwrap_err();
2882            assert!(
2883                matches!(err, EngineError::InvalidConcurrencyLimit(_)),
2884                "{err:?}"
2885            );
2886        }
2887
2888        // Validated before the handler lookup.
2889        let err = engine
2890            .enqueue_handler_with_options(
2891                "not-registered",
2892                TriggerKind::Api,
2893                json!({}),
2894                EnqueueOptions {
2895                    concurrency_limits: vec![ConcurrencyLimit::new("repo:acme", 0)],
2896                    ..Default::default()
2897                },
2898            )
2899            .await
2900            .unwrap_err();
2901        assert!(
2902            matches!(err, EngineError::InvalidConcurrencyLimit(_)),
2903            "{err:?}"
2904        );
2905
2906        let page = engine
2907            .store()
2908            .list_runs(RunFilter::default(), 1, 10)
2909            .await
2910            .unwrap();
2911        assert_eq!(page.total, 0, "no run may be created");
2912    }
2913
2914    #[tokio::test]
2915    async fn run_handler_leaves_the_run_unattributed() {
2916        let mut engine = create_test_engine();
2917        engine.register(EchoWorkflow).unwrap();
2918
2919        let run = engine
2920            .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
2921            .await
2922            .unwrap()
2923            .run;
2924
2925        assert!(run.created_by.is_none());
2926    }
2927
2928    #[tokio::test]
2929    async fn engine_register_boxed() {
2930        let mut engine = create_test_engine();
2931        let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
2932        let result = engine.register_boxed(handler);
2933        assert!(result.is_ok());
2934        assert_eq!(engine.handler_names().len(), 1);
2935    }
2936
2937    #[tokio::test]
2938    async fn engine_store_and_provider_accessors() {
2939        let store = Arc::new(InMemoryStore::new());
2940        let inner = ClaudeCodeProvider::new();
2941        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2942            inner,
2943            "/tmp/ironflow-fixtures",
2944        ));
2945        let engine = Engine::new(store.clone(), provider.clone());
2946
2947        // Verify accessors return references
2948        let _ = engine.store();
2949        let _ = engine.provider();
2950    }
2951
2952    // -----------------------------------------------------------------------
2953    // Operation trait tests
2954    // -----------------------------------------------------------------------
2955
2956    use crate::operation::{Operation, OperationContext};
2957    use async_trait::async_trait;
2958    use ironflow_core::error::OperationError;
2959    use ironflow_store::models::StepKind;
2960
2961    struct FakeGitlabOp {
2962        project_id: u64,
2963        title: String,
2964    }
2965
2966    #[async_trait]
2967    impl Operation for FakeGitlabOp {
2968        fn kind(&self) -> &str {
2969            "gitlab"
2970        }
2971
2972        async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2973            Ok(json!({
2974                "issue_id": 42,
2975                "project_id": self.project_id,
2976                "title": self.title,
2977            }))
2978        }
2979
2980        fn input(&self) -> Option<Value> {
2981            Some(json!({
2982                "project_id": self.project_id,
2983                "title": self.title,
2984            }))
2985        }
2986    }
2987
2988    struct FailingOp;
2989
2990    #[async_trait]
2991    impl Operation for FailingOp {
2992        fn kind(&self) -> &str {
2993            "broken-service"
2994        }
2995
2996        async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2997            Err(OperationError::Http {
2998                status: None,
2999                message: "service unavailable".to_string(),
3000            })
3001        }
3002    }
3003
3004    struct OperationWorkflow;
3005
3006    impl WorkflowHandler for OperationWorkflow {
3007        fn name(&self) -> &str {
3008            "operation-workflow"
3009        }
3010
3011        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3012            Box::pin(async move {
3013                let op = FakeGitlabOp {
3014                    project_id: 123,
3015                    title: "Bug report".to_string(),
3016                };
3017                ctx.operation("create-issue", &op).await?;
3018                Ok(())
3019            })
3020        }
3021    }
3022
3023    struct FailingOperationWorkflow;
3024
3025    impl WorkflowHandler for FailingOperationWorkflow {
3026        fn name(&self) -> &str {
3027            "failing-operation-workflow"
3028        }
3029
3030        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3031            Box::pin(async move {
3032                ctx.operation("broken-call", &FailingOp).await?;
3033                Ok(())
3034            })
3035        }
3036    }
3037
3038    struct MixedWorkflow;
3039
3040    impl WorkflowHandler for MixedWorkflow {
3041        fn name(&self) -> &str {
3042            "mixed-workflow"
3043        }
3044
3045        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3046            Box::pin(async move {
3047                ctx.shell("build", ShellConfig::new("echo built")).await?;
3048                let op = FakeGitlabOp {
3049                    project_id: 456,
3050                    title: "Deploy done".to_string(),
3051                };
3052                let result = ctx.operation("notify-gitlab", &op).await?;
3053                assert_eq!(result.output["issue_id"], 42);
3054                Ok(())
3055            })
3056        }
3057    }
3058
3059    #[tokio::test]
3060    async fn operation_step_happy_path() {
3061        let mut engine = create_test_engine();
3062        engine.register(OperationWorkflow).unwrap();
3063
3064        let run = engine
3065            .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
3066            .await
3067            .unwrap()
3068            .run;
3069
3070        assert_eq!(run.status.state, RunStatus::Completed);
3071
3072        let steps = engine.store().list_steps(run.id).await.unwrap();
3073
3074        assert_eq!(steps.len(), 1);
3075        assert_eq!(steps[0].name, "create-issue");
3076        assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
3077        assert_eq!(
3078            steps[0].status.state,
3079            ironflow_store::models::StepStatus::Completed
3080        );
3081
3082        let output = steps[0].output.as_ref().unwrap();
3083        assert_eq!(output["issue_id"], 42);
3084        assert_eq!(output["project_id"], 123);
3085
3086        let input = steps[0].input.as_ref().unwrap();
3087        assert_eq!(input["project_id"], 123);
3088        assert_eq!(input["title"], "Bug report");
3089    }
3090
3091    #[tokio::test]
3092    async fn operation_step_failure_marks_run_failed() {
3093        let mut engine = create_test_engine();
3094        engine.register(FailingOperationWorkflow).unwrap();
3095
3096        let result = engine
3097            .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
3098            .await;
3099
3100        assert!(result.is_err());
3101    }
3102
3103    #[tokio::test]
3104    async fn operation_mixed_with_shell_steps() {
3105        let mut engine = create_test_engine();
3106        engine.register(MixedWorkflow).unwrap();
3107
3108        let run = engine
3109            .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
3110            .await
3111            .unwrap()
3112            .run;
3113
3114        assert_eq!(run.status.state, RunStatus::Completed);
3115
3116        let steps = engine.store().list_steps(run.id).await.unwrap();
3117
3118        assert_eq!(steps.len(), 2);
3119        assert_eq!(steps[0].kind, StepKind::Shell);
3120        assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
3121        assert_eq!(steps[0].position, 0);
3122        assert_eq!(steps[1].position, 1);
3123    }
3124
3125    // -----------------------------------------------------------------------
3126    // Approval + resume tests
3127    // -----------------------------------------------------------------------
3128
3129    use crate::config::ApprovalConfig;
3130
3131    struct SingleApprovalWorkflow;
3132
3133    impl WorkflowHandler for SingleApprovalWorkflow {
3134        fn name(&self) -> &str {
3135            "single-approval"
3136        }
3137
3138        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3139            Box::pin(async move {
3140                ctx.shell("build", ShellConfig::new("echo built")).await?;
3141                ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
3142                ctx.shell("deploy", ShellConfig::new("echo deployed"))
3143                    .await?;
3144                Ok(())
3145            })
3146        }
3147    }
3148
3149    struct DoubleApprovalWorkflow;
3150
3151    impl WorkflowHandler for DoubleApprovalWorkflow {
3152        fn name(&self) -> &str {
3153            "double-approval"
3154        }
3155
3156        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3157            Box::pin(async move {
3158                ctx.shell("build", ShellConfig::new("echo built")).await?;
3159                ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
3160                    .await?;
3161                ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3162                    .await?;
3163                ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3164                    .await?;
3165                ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3166                    .await?;
3167                Ok(())
3168            })
3169        }
3170    }
3171
3172    #[tokio::test]
3173    async fn approval_pauses_run() {
3174        let mut engine = create_test_engine();
3175        engine.register(SingleApprovalWorkflow).unwrap();
3176
3177        let run = engine
3178            .run_handler("single-approval", TriggerKind::Manual, json!({}))
3179            .await
3180            .unwrap()
3181            .run;
3182
3183        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3184
3185        let steps = engine.store().list_steps(run.id).await.unwrap();
3186        assert_eq!(steps.len(), 2); // build + approval gate
3187        assert_eq!(steps[0].kind, StepKind::Shell);
3188        assert_eq!(steps[0].status.state, StepStatus::Completed);
3189        assert_eq!(steps[1].kind, StepKind::Approval);
3190        assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3191    }
3192
3193    #[tokio::test]
3194    async fn approval_resume_completes_run() {
3195        let mut engine = create_test_engine();
3196        engine.register(SingleApprovalWorkflow).unwrap();
3197
3198        // First execution: pauses at approval
3199        let run = engine
3200            .run_handler("single-approval", TriggerKind::Manual, json!({}))
3201            .await
3202            .unwrap()
3203            .run;
3204        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3205
3206        // Simulate approval: transition to Running
3207        engine
3208            .store()
3209            .update_run_status(run.id, RunStatus::Running)
3210            .await
3211            .unwrap();
3212
3213        // Resume: replays build, skips approval, executes deploy
3214        let resumed = engine.resume_run(run.id).await.unwrap().run;
3215        assert_eq!(resumed.status.state, RunStatus::Completed);
3216
3217        let steps = engine.store().list_steps(run.id).await.unwrap();
3218        assert_eq!(steps.len(), 3); // build + approval + deploy
3219        assert_eq!(steps[0].name, "build");
3220        assert_eq!(steps[0].status.state, StepStatus::Completed);
3221        assert_eq!(steps[1].name, "gate");
3222        assert_eq!(steps[1].kind, StepKind::Approval);
3223        assert_eq!(steps[1].status.state, StepStatus::Completed);
3224        assert_eq!(steps[2].name, "deploy");
3225        assert_eq!(steps[2].status.state, StepStatus::Completed);
3226    }
3227
3228    #[tokio::test]
3229    async fn double_approval_two_resumes() {
3230        let mut engine = create_test_engine();
3231        engine.register(DoubleApprovalWorkflow).unwrap();
3232
3233        // First execution: pauses at staging-gate
3234        let run = engine
3235            .run_handler("double-approval", TriggerKind::Manual, json!({}))
3236            .await
3237            .unwrap()
3238            .run;
3239        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3240
3241        let steps = engine.store().list_steps(run.id).await.unwrap();
3242        assert_eq!(steps.len(), 2); // build + staging-gate
3243
3244        // First approval
3245        engine
3246            .store()
3247            .update_run_status(run.id, RunStatus::Running)
3248            .await
3249            .unwrap();
3250
3251        let resumed = engine.resume_run(run.id).await.unwrap().run;
3252        assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3253
3254        let steps = engine.store().list_steps(run.id).await.unwrap();
3255        assert_eq!(steps.len(), 4); // build + staging-gate + deploy-staging + prod-gate
3256
3257        // Second approval
3258        engine
3259            .store()
3260            .update_run_status(run.id, RunStatus::Running)
3261            .await
3262            .unwrap();
3263
3264        let final_run = engine.resume_run(run.id).await.unwrap().run;
3265        assert_eq!(final_run.status.state, RunStatus::Completed);
3266
3267        let steps = engine.store().list_steps(run.id).await.unwrap();
3268        assert_eq!(steps.len(), 5);
3269        assert_eq!(steps[0].name, "build");
3270        assert_eq!(steps[1].name, "staging-gate");
3271        assert_eq!(steps[2].name, "deploy-staging");
3272        assert_eq!(steps[3].name, "prod-gate");
3273        assert_eq!(steps[4].name, "deploy-prod");
3274
3275        for step in &steps {
3276            assert_eq!(step.status.state, StepStatus::Completed);
3277        }
3278    }
3279
3280    // -----------------------------------------------------------------------
3281    // fail_orphaned_steps tests
3282    // -----------------------------------------------------------------------
3283
3284    use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3285
3286    async fn create_step_with_status(
3287        store: &Arc<dyn Store>,
3288        run_id: Uuid,
3289        name: &str,
3290        position: u32,
3291        status: StepStatus,
3292    ) -> ironflow_store::models::Step {
3293        let step = store
3294            .create_step(NewStep {
3295                run_id,
3296                trace_id: step_trace_id(run_id, name, position),
3297                name: name.to_string(),
3298                kind: StepKind::Shell,
3299                position,
3300                input: None,
3301                is_error_handler: false,
3302            })
3303            .await
3304            .unwrap();
3305
3306        match status {
3307            StepStatus::Pending => {}
3308            StepStatus::Running => {
3309                store
3310                    .update_step(
3311                        step.id,
3312                        StepUpdate {
3313                            status: Some(StepStatus::Running),
3314                            ..StepUpdate::default()
3315                        },
3316                    )
3317                    .await
3318                    .unwrap();
3319            }
3320            StepStatus::Completed => {
3321                store
3322                    .update_step(
3323                        step.id,
3324                        StepUpdate {
3325                            status: Some(StepStatus::Running),
3326                            ..StepUpdate::default()
3327                        },
3328                    )
3329                    .await
3330                    .unwrap();
3331                store
3332                    .update_step(
3333                        step.id,
3334                        StepUpdate {
3335                            status: Some(StepStatus::Completed),
3336                            ..StepUpdate::default()
3337                        },
3338                    )
3339                    .await
3340                    .unwrap();
3341            }
3342            StepStatus::AwaitingApproval => {
3343                store
3344                    .update_step(
3345                        step.id,
3346                        StepUpdate {
3347                            status: Some(StepStatus::Running),
3348                            ..StepUpdate::default()
3349                        },
3350                    )
3351                    .await
3352                    .unwrap();
3353                store
3354                    .update_step(
3355                        step.id,
3356                        StepUpdate {
3357                            status: Some(StepStatus::AwaitingApproval),
3358                            ..StepUpdate::default()
3359                        },
3360                    )
3361                    .await
3362                    .unwrap();
3363            }
3364            _ => panic!("unsupported status for test helper: {status}"),
3365        }
3366
3367        store.get_step(step.id).await.unwrap().unwrap()
3368    }
3369
3370    #[tokio::test]
3371    async fn fail_orphaned_steps_marks_running_as_failed() {
3372        let engine = create_test_engine();
3373        let run = engine
3374            .store()
3375            .create_run(NewRun {
3376                created_by: None,
3377                workflow_name: "test".to_string(),
3378                trigger: TriggerKind::Manual,
3379                payload: json!({}),
3380                max_retries: 0,
3381                handler_version: None,
3382                labels: HashMap::new(),
3383                scheduled_at: None,
3384                idempotency_key: None,
3385                concurrency_key: None,
3386                concurrency_limits: Vec::new(),
3387                max_cost_usd: None,
3388            })
3389            .await
3390            .unwrap()
3391            .into_run();
3392
3393        let step = create_step_with_status(
3394            engine.store(),
3395            run.id,
3396            "running-step",
3397            0,
3398            StepStatus::Running,
3399        )
3400        .await;
3401
3402        engine
3403            .fail_orphaned_steps(run.id, "parent run timed out")
3404            .await
3405            .unwrap();
3406
3407        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3408        assert_eq!(updated.status.state, StepStatus::Failed);
3409        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3410        assert!(updated.completed_at.is_some());
3411    }
3412
3413    #[tokio::test]
3414    async fn fail_orphaned_steps_marks_pending_as_skipped() {
3415        let engine = create_test_engine();
3416        let run = engine
3417            .store()
3418            .create_run(NewRun {
3419                created_by: None,
3420                workflow_name: "test".to_string(),
3421                trigger: TriggerKind::Manual,
3422                payload: json!({}),
3423                max_retries: 0,
3424                handler_version: None,
3425                labels: HashMap::new(),
3426                scheduled_at: None,
3427                idempotency_key: None,
3428                concurrency_key: None,
3429                concurrency_limits: Vec::new(),
3430                max_cost_usd: None,
3431            })
3432            .await
3433            .unwrap()
3434            .into_run();
3435
3436        let step = create_step_with_status(
3437            engine.store(),
3438            run.id,
3439            "pending-step",
3440            0,
3441            StepStatus::Pending,
3442        )
3443        .await;
3444
3445        engine
3446            .fail_orphaned_steps(run.id, "parent run timed out")
3447            .await
3448            .unwrap();
3449
3450        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3451        assert_eq!(updated.status.state, StepStatus::Skipped);
3452        assert!(updated.error.is_none());
3453        assert!(updated.completed_at.is_some());
3454    }
3455
3456    #[tokio::test]
3457    async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
3458        let engine = create_test_engine();
3459        let run = engine
3460            .store()
3461            .create_run(NewRun {
3462                created_by: None,
3463                workflow_name: "test".to_string(),
3464                trigger: TriggerKind::Manual,
3465                payload: json!({}),
3466                max_retries: 0,
3467                handler_version: None,
3468                labels: HashMap::new(),
3469                scheduled_at: None,
3470                idempotency_key: None,
3471                concurrency_key: None,
3472                concurrency_limits: Vec::new(),
3473                max_cost_usd: None,
3474            })
3475            .await
3476            .unwrap()
3477            .into_run();
3478
3479        let step = create_step_with_status(
3480            engine.store(),
3481            run.id,
3482            "approval-step",
3483            0,
3484            StepStatus::AwaitingApproval,
3485        )
3486        .await;
3487
3488        engine
3489            .fail_orphaned_steps(run.id, "parent run timed out")
3490            .await
3491            .unwrap();
3492
3493        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
3494        assert_eq!(updated.status.state, StepStatus::Failed);
3495        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
3496        assert!(updated.completed_at.is_some());
3497    }
3498
3499    #[tokio::test]
3500    async fn fail_orphaned_steps_skips_terminal_steps() {
3501        let engine = create_test_engine();
3502        let run = engine
3503            .store()
3504            .create_run(NewRun {
3505                created_by: None,
3506                workflow_name: "test".to_string(),
3507                trigger: TriggerKind::Manual,
3508                payload: json!({}),
3509                max_retries: 0,
3510                handler_version: None,
3511                labels: HashMap::new(),
3512                scheduled_at: None,
3513                idempotency_key: None,
3514                concurrency_key: None,
3515                concurrency_limits: Vec::new(),
3516                max_cost_usd: None,
3517            })
3518            .await
3519            .unwrap()
3520            .into_run();
3521
3522        let completed_step =
3523            create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
3524        let running_step =
3525            create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
3526                .await;
3527
3528        engine
3529            .fail_orphaned_steps(run.id, "parent run timed out")
3530            .await
3531            .unwrap();
3532
3533        let completed = engine
3534            .store()
3535            .get_step(completed_step.id)
3536            .await
3537            .unwrap()
3538            .unwrap();
3539        assert_eq!(completed.status.state, StepStatus::Completed);
3540
3541        let failed = engine
3542            .store()
3543            .get_step(running_step.id)
3544            .await
3545            .unwrap()
3546            .unwrap();
3547        assert_eq!(failed.status.state, StepStatus::Failed);
3548    }
3549
3550    #[tokio::test]
3551    async fn fail_orphaned_steps_mixed_states() {
3552        let engine = create_test_engine();
3553        let run = engine
3554            .store()
3555            .create_run(NewRun {
3556                created_by: None,
3557                workflow_name: "test".to_string(),
3558                trigger: TriggerKind::Manual,
3559                payload: json!({}),
3560                max_retries: 0,
3561                handler_version: None,
3562                labels: HashMap::new(),
3563                scheduled_at: None,
3564                idempotency_key: None,
3565                concurrency_key: None,
3566                concurrency_limits: Vec::new(),
3567                max_cost_usd: None,
3568            })
3569            .await
3570            .unwrap()
3571            .into_run();
3572
3573        let s_completed =
3574            create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
3575                .await;
3576        let s_running =
3577            create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
3578        let s_pending =
3579            create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
3580
3581        engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
3582
3583        let r_completed = engine
3584            .store()
3585            .get_step(s_completed.id)
3586            .await
3587            .unwrap()
3588            .unwrap();
3589        assert_eq!(r_completed.status.state, StepStatus::Completed);
3590
3591        let r_running = engine
3592            .store()
3593            .get_step(s_running.id)
3594            .await
3595            .unwrap()
3596            .unwrap();
3597        assert_eq!(r_running.status.state, StepStatus::Failed);
3598        assert_eq!(r_running.error.as_deref(), Some("timeout"));
3599
3600        let r_pending = engine
3601            .store()
3602            .get_step(s_pending.id)
3603            .await
3604            .unwrap()
3605            .unwrap();
3606        assert_eq!(r_pending.status.state, StepStatus::Skipped);
3607        assert!(r_pending.error.is_none());
3608    }
3609
3610    #[tokio::test]
3611    async fn fail_orphaned_steps_no_steps_is_noop() {
3612        let engine = create_test_engine();
3613        let run = engine
3614            .store()
3615            .create_run(NewRun {
3616                created_by: None,
3617                workflow_name: "test".to_string(),
3618                trigger: TriggerKind::Manual,
3619                payload: json!({}),
3620                max_retries: 0,
3621                handler_version: None,
3622                labels: HashMap::new(),
3623                scheduled_at: None,
3624                idempotency_key: None,
3625                concurrency_key: None,
3626                concurrency_limits: Vec::new(),
3627                max_cost_usd: None,
3628            })
3629            .await
3630            .unwrap()
3631            .into_run();
3632
3633        let result = engine.fail_orphaned_steps(run.id, "timeout").await;
3634        assert!(result.is_ok());
3635    }
3636
3637    #[tokio::test]
3638    async fn fail_orphaned_steps_preserves_existing_error() {
3639        let engine = create_test_engine();
3640        let run = engine
3641            .store()
3642            .create_run(NewRun {
3643                created_by: None,
3644                workflow_name: "test".to_string(),
3645                trigger: TriggerKind::Manual,
3646                payload: json!({}),
3647                max_retries: 0,
3648                handler_version: None,
3649                labels: HashMap::new(),
3650                scheduled_at: None,
3651                idempotency_key: None,
3652                concurrency_key: None,
3653                concurrency_limits: Vec::new(),
3654                max_cost_usd: None,
3655            })
3656            .await
3657            .unwrap()
3658            .into_run();
3659
3660        let step_with_error = create_step_with_status(
3661            engine.store(),
3662            run.id,
3663            "already-errored",
3664            0,
3665            StepStatus::Running,
3666        )
3667        .await;
3668
3669        engine
3670            .store()
3671            .update_step(
3672                step_with_error.id,
3673                StepUpdate {
3674                    error: Some("real error from provider".to_string()),
3675                    ..StepUpdate::default()
3676                },
3677            )
3678            .await
3679            .unwrap();
3680
3681        let step_no_error = create_step_with_status(
3682            engine.store(),
3683            run.id,
3684            "no-error-yet",
3685            1,
3686            StepStatus::Running,
3687        )
3688        .await;
3689
3690        engine
3691            .fail_orphaned_steps(run.id, "parent run failed")
3692            .await
3693            .unwrap();
3694
3695        let updated_with = engine
3696            .store()
3697            .get_step(step_with_error.id)
3698            .await
3699            .unwrap()
3700            .unwrap();
3701        assert_eq!(updated_with.status.state, StepStatus::Failed);
3702        assert_eq!(
3703            updated_with.error.as_deref(),
3704            Some("real error from provider"),
3705        );
3706
3707        let updated_without = engine
3708            .store()
3709            .get_step(step_no_error.id)
3710            .await
3711            .unwrap()
3712            .unwrap();
3713        assert_eq!(updated_without.status.state, StepStatus::Failed);
3714        assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
3715    }
3716}