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