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    /// When the workflow is paused (see [`Engine::pause_workflow`]) the run is
983    /// created and left `Pending`, with no steps: the returned
984    /// `WorkflowResult.run.status.state` is `Pending`. The hold applies in
985    /// every execution mode; [`Engine::resume_workflow`] runs the held run
986    /// under [`ExecutionMode::Local`].
987    ///
988    /// # Examples
989    ///
990    /// ```no_run
991    /// use std::sync::Arc;
992    /// use ironflow_engine::engine::Engine;
993    /// use ironflow_store::memory::InMemoryStore;
994    /// use ironflow_store::models::TriggerKind;
995    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
996    /// use serde_json::json;
997    ///
998    /// # async fn example(engine: &Engine) -> Result<(), ironflow_engine::error::EngineError> {
999    /// let run = engine.run_handler("deploy", TriggerKind::Manual, json!({})).await?;
1000    /// # Ok(())
1001    /// # }
1002    /// ```
1003    #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
1004    pub async fn run_handler(
1005        &self,
1006        handler_name: &str,
1007        trigger: TriggerKind,
1008        payload: Value,
1009    ) -> Result<WorkflowResult, EngineError> {
1010        let handler = self
1011            .handlers
1012            .get(handler_name)
1013            .ok_or_else(|| {
1014                EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1015            })?
1016            .clone();
1017
1018        self.check_monthly_quota(handler_name).await?;
1019
1020        let handler_version = handler.version().map(str::to_string);
1021        let max_cost_usd = self
1022            .budget
1023            .resolve_run_cap(None, handler.default_max_cost_usd());
1024        let run = self
1025            .store
1026            .create_run(NewRun {
1027                created_by: None,
1028                workflow_name: handler_name.to_string(),
1029                trigger,
1030                payload,
1031                max_retries: 0,
1032                handler_version,
1033                labels: handler.default_labels(),
1034                scheduled_at: None,
1035                idempotency_key: None,
1036                concurrency_key: None,
1037                priority: clamp_priority(handler.priority()),
1038                concurrency_limits: Vec::new(),
1039                max_cost_usd,
1040                worker_tags: normalize_worker_tags(handler.required_worker_tags()),
1041            })
1042            .await?
1043            .into_run();
1044
1045        let run_id = run.id;
1046        info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
1047
1048        if self.is_workflow_paused(handler_name).await? {
1049            info!(run_id = %run_id, workflow = %handler_name, "workflow paused, run left pending");
1050            return Ok(WorkflowResult {
1051                run,
1052                steps: Vec::new(),
1053            });
1054        }
1055
1056        self.store
1057            .update_run_status(run_id, RunStatus::Running)
1058            .await?;
1059
1060        #[cfg(feature = "prometheus")]
1061        gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
1062
1063        let run_start = Instant::now();
1064        let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1065
1066        let result = handler.execute(&mut ctx).await;
1067        self.finalize_run(run_id, handler_name, result, &ctx, run_start, run.labels)
1068            .await
1069    }
1070
1071    /// Build the execution plan for a registered handler without running it.
1072    ///
1073    /// Executes the handler with every step method in recording mode: no
1074    /// command is spawned, no HTTP request is sent, no agent is called,
1075    /// nothing is persisted. Conditions declared with
1076    /// [`WorkflowContext::when`](crate::context::WorkflowContext::when) are
1077    /// evaluated against `payload`; those declared with
1078    /// [`WorkflowContext::when_dynamic`](crate::context::WorkflowContext::when_dynamic)
1079    /// are reported as unevaluable.
1080    ///
1081    /// Step outputs are synthetic and success-shaped, so the plan follows the
1082    /// nominal branch. A handler that unwraps a decision answer, or that
1083    /// deserializes `ctx.input::<T>()` against a payload it does not match,
1084    /// aborts the plan: the partial plan is returned with
1085    /// [`ExecutionPlan::incomplete_reason`] set.
1086    ///
1087    /// # Errors
1088    ///
1089    /// Returns [`EngineError::InvalidWorkflow`] when no handler is registered
1090    /// under `handler_name` or when `options.max_depth` is zero. Returns
1091    /// [`EngineError::Store`] when the duration-history query fails. A handler
1092    /// that errors mid-plan does **not** fail this call.
1093    ///
1094    /// # Examples
1095    ///
1096    /// ```no_run
1097    /// use ironflow_engine::engine::Engine;
1098    /// use ironflow_engine::error::EngineError;
1099    /// use ironflow_engine::plan::PlanOptions;
1100    /// use serde_json::json;
1101    ///
1102    /// # async fn example(engine: &Engine) -> Result<(), EngineError> {
1103    /// let plan = engine
1104    ///     .plan_handler("deploy", json!({"env": "prod"}), PlanOptions::default())
1105    ///     .await?;
1106    /// for step in &plan.steps {
1107    ///     println!("{} ({:?})", step.name, step.kind);
1108    /// }
1109    /// # Ok(())
1110    /// # }
1111    /// ```
1112    #[tracing::instrument(name = "engine.plan_handler", skip_all, fields(workflow = %handler_name))]
1113    pub async fn plan_handler(
1114        &self,
1115        handler_name: &str,
1116        payload: Value,
1117        options: PlanOptions,
1118    ) -> Result<ExecutionPlan, EngineError> {
1119        if options.max_depth == 0 {
1120            return Err(EngineError::InvalidWorkflow(
1121                "max_depth must be at least 1".to_string(),
1122            ));
1123        }
1124
1125        let handler = self
1126            .handlers
1127            .get(handler_name)
1128            .ok_or_else(|| {
1129                EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1130            })?
1131            .clone();
1132
1133        let estimates = if options.estimate_durations {
1134            estimate_durations(&self.store, handler_name, options.sample_runs).await?
1135        } else {
1136            HashMap::new()
1137        };
1138
1139        let shared: SharedPlanRecorder = Arc::new(Mutex::new(PlanRecorder::new(
1140            handler_name.to_string(),
1141            payload,
1142            options.max_depth,
1143            estimates,
1144        )));
1145
1146        // Deliberately bare: no guard, no event bus, no log sender, no artifact
1147        // sink and no budget. Planning produces no side effect to report.
1148        let handlers = self.handlers.clone();
1149        let resolver: crate::context::HandlerResolver =
1150            Arc::new(move |name: &str| handlers.get(name).cloned());
1151        let mut ctx = WorkflowContext::with_handler_resolver(
1152            Uuid::now_v7(),
1153            handler_name.to_string(),
1154            self.store.clone(),
1155            self.provider.clone(),
1156            resolver,
1157        );
1158        ctx.set_plan(shared.clone());
1159
1160        if let Err(err) = handler.execute(&mut ctx).await {
1161            lock_plan(&shared).fail(err.to_string());
1162        }
1163        drop(ctx);
1164
1165        let plan = match Arc::try_unwrap(shared) {
1166            Ok(mutex) => mutex
1167                .into_inner()
1168                .unwrap_or_else(|poisoned| poisoned.into_inner())
1169                .into_plan(),
1170            Err(shared) => lock_plan(&shared).snapshot(),
1171        };
1172
1173        info!(
1174            workflow = %handler_name,
1175            steps = plan.steps.len(),
1176            truncated = plan.truncated,
1177            "execution plan built"
1178        );
1179
1180        Ok(plan)
1181    }
1182
1183    /// Enqueue a handler-based workflow for worker execution.
1184    ///
1185    /// The workflow name is stored in the run. The worker looks up the
1186    /// handler by name when executing.
1187    ///
1188    /// # Errors
1189    ///
1190    /// Returns [`EngineError::InvalidWorkflow`] if no handler is registered.
1191    /// Returns [`EngineError::MonthlyBudgetExceeded`] if the monthly cost quota
1192    /// is exhausted.
1193    #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
1194    pub async fn enqueue_handler(
1195        &self,
1196        handler_name: &str,
1197        trigger: TriggerKind,
1198        payload: Value,
1199        max_retries: u32,
1200    ) -> Result<Run, EngineError> {
1201        self.enqueue_handler_with_options(
1202            handler_name,
1203            trigger,
1204            payload,
1205            EnqueueOptions {
1206                max_retries,
1207                ..Default::default()
1208            },
1209        )
1210        .await
1211        .map(RunCreation::into_run)
1212    }
1213
1214    /// Enqueue a handler-based workflow with labels, deferred scheduling, an
1215    /// optional cost cap, an optional author, and an optional idempotency key.
1216    ///
1217    /// See [`EnqueueOptions`] for the individual settings.
1218    ///
1219    /// When [`EnqueueOptions::idempotency_key`] is set and already bound to a run
1220    /// created within
1221    /// [`IDEMPOTENCY_WINDOW`](ironflow_store::entities::IDEMPOTENCY_WINDOW), nothing
1222    /// is enqueued and the original run is returned as [`RunCreation::Existing`].
1223    ///
1224    /// # Errors
1225    ///
1226    /// Returns [`EngineError::InvalidWorkflow`] if no handler is registered.
1227    /// Returns [`EngineError::MonthlyBudgetExceeded`] if the monthly cost quota
1228    /// is exhausted. Returns [`EngineError::ConcurrencyConflict`] if
1229    /// [`EnqueueOptions::concurrency_key`] is held by another non-terminal run.
1230    /// Returns [`EngineError::InvalidConcurrencyLimit`] if
1231    /// [`EnqueueOptions::concurrency_limits`] is invalid, before any other check.
1232    /// Returns [`EngineError::InvalidPriority`] if [`EnqueueOptions::priority`]
1233    /// is out of range, before any other check.
1234    /// Returns [`EngineError::InvalidWorkerTag`] if [`EnqueueOptions::worker_tags`]
1235    /// holds an invalid tag, checked right after the concurrency limits.
1236    /// Returns [`EngineError::Store`] if the run cannot be persisted.
1237    ///
1238    /// # Examples
1239    ///
1240    /// ```no_run
1241    /// use ironflow_engine::engine::{Engine, EnqueueOptions};
1242    /// use ironflow_store::models::TriggerKind;
1243    /// use serde_json::json;
1244    ///
1245    /// # async fn example(engine: &Engine) -> Result<(), ironflow_engine::error::EngineError> {
1246    /// let creation = engine
1247    ///     .enqueue_handler_with_options(
1248    ///         "deploy",
1249    ///         TriggerKind::Api,
1250    ///         json!({"env": "prod"}),
1251    ///         EnqueueOptions {
1252    ///             max_retries: 3,
1253    ///             idempotency_key: Some("github:abc-123".to_string()),
1254    ///             ..Default::default()
1255    ///         },
1256    ///     )
1257    ///     .await?;
1258    ///
1259    /// if creation.is_created() {
1260    ///     println!("enqueued {}", creation.run().id);
1261    /// }
1262    /// # Ok(())
1263    /// # }
1264    /// ```
1265    #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
1266    pub async fn enqueue_handler_with_options(
1267        &self,
1268        handler_name: &str,
1269        trigger: TriggerKind,
1270        payload: Value,
1271        options: EnqueueOptions,
1272    ) -> Result<RunCreation, EngineError> {
1273        let EnqueueOptions {
1274            max_retries,
1275            labels,
1276            scheduled_at,
1277            max_cost_usd,
1278            created_by,
1279            idempotency_key,
1280            concurrency_key,
1281            concurrency_limits,
1282            priority,
1283            worker_tags,
1284        } = options;
1285
1286        // Checked first: a malformed request is the caller's error, whatever
1287        // the handler or the quota.
1288        validate_concurrency_limits(&concurrency_limits)
1289            .map_err(EngineError::InvalidConcurrencyLimit)?;
1290        if let Some(priority) = priority {
1291            validate_priority(priority).map_err(EngineError::InvalidPriority)?;
1292        }
1293        validate_worker_tags(&worker_tags).map_err(EngineError::InvalidWorkerTag)?;
1294
1295        let handler = self.handlers.get(handler_name).ok_or_else(|| {
1296            EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
1297        })?;
1298
1299        self.check_monthly_quota(handler_name).await?;
1300
1301        let handler_version = handler.version().map(str::to_string);
1302        let mut merged_labels = handler.default_labels();
1303        merged_labels.extend(labels);
1304        let resolved_cap = self
1305            .budget
1306            .resolve_run_cap(max_cost_usd, handler.default_max_cost_usd());
1307        let priority = priority.unwrap_or_else(|| clamp_priority(handler.priority()));
1308        let required_tags = normalize_worker_tags(
1309            handler
1310                .required_worker_tags()
1311                .into_iter()
1312                .chain(worker_tags),
1313        );
1314
1315        let creation = self
1316            .store
1317            .create_run(NewRun {
1318                workflow_name: handler_name.to_string(),
1319                trigger,
1320                payload,
1321                max_retries,
1322                handler_version,
1323                labels: merged_labels,
1324                scheduled_at,
1325                created_by,
1326                idempotency_key,
1327                concurrency_key,
1328                priority,
1329                concurrency_limits,
1330                max_cost_usd: resolved_cap,
1331                worker_tags: required_tags,
1332            })
1333            .await?;
1334
1335        match &creation {
1336            RunCreation::Created(run) => info!(
1337                run_id = %run.id,
1338                workflow = %handler_name,
1339                max_cost_usd = ?resolved_cap,
1340                "handler run enqueued"
1341            ),
1342            RunCreation::Existing(run) => info!(
1343                run_id = %run.id,
1344                workflow = %handler_name,
1345                "idempotent replay, nothing enqueued"
1346            ),
1347        }
1348
1349        Ok(creation)
1350    }
1351
1352    /// Execute a handler-based run (used by the worker after pick_next_pending).
1353    ///
1354    /// Looks up the handler by the run's `workflow_name` and executes it
1355    /// with a fresh [`WorkflowContext`], after
1356    /// [`AgentProvider::release_run`] has stopped whatever a previous
1357    /// execution of the run left running.
1358    ///
1359    /// # Errors
1360    ///
1361    /// Returns [`EngineError::InvalidWorkflow`] if no handler matches. A
1362    /// failed release fails the execution with [`EngineError::Operation`],
1363    /// replayed while the run has retries left. Returns
1364    /// [`EngineError::HandlerVersionMismatch`] when the handler's current
1365    /// version is incompatible with the run's `handler_version` -- checked
1366    /// before any step is replayed.
1367    ///
1368    /// A suspended child run of a sub-workflow is never executed on its own:
1369    /// its root run is resumed instead, and re-enters the child (see
1370    /// [`resume_run`](Self::resume_run)).
1371    #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
1372    pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1373        let run = self
1374            .store
1375            .get_run(run_id)
1376            .await?
1377            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1378
1379        if let Some(root_run_id) = chain_root(&run) {
1380            return self.resume_chain(run, root_run_id).await;
1381        }
1382
1383        let _active = self.track_execution(run_id).await;
1384
1385        let handler = self
1386            .handlers
1387            .get(&run.workflow_name)
1388            .ok_or_else(|| {
1389                EngineError::InvalidWorkflow(format!(
1390                    "no handler registered: {}",
1391                    run.workflow_name
1392                ))
1393            })?
1394            .clone();
1395
1396        #[cfg(feature = "prometheus")]
1397        gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
1398
1399        let run_start = Instant::now();
1400        let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1401
1402        // Replay the steps already persisted for this run: on a retry, an approval
1403        // a human already granted must not be asked again, and a run requeued to
1404        // `Pending` after its approval, human input or escalation resolved
1405        // (`ExecutionMode::Workers`, `retry_count` unchanged) must not re-run
1406        // completed steps. Neither must a run requeued by the reaper after its
1407        // worker lost the lease, which stays in the same attempt too. A
1408        // brand-new run has no steps, so this is a no-op.
1409        //
1410        // The handler version is checked first: replaying an incompatible
1411        // handler's steps risks serving one step's cached output to another
1412        // (`EngineError::ReplayDivergence`), so no step is replayed at all
1413        // when the handler changed incompatibly since the run was created.
1414        let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1415            ctx.load_replay_steps().await?;
1416            self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1417                .await
1418        } else {
1419            Err(EngineError::HandlerVersionMismatch {
1420                run_id,
1421                workflow_name: run.workflow_name.clone(),
1422                run_version: run
1423                    .handler_version
1424                    .clone()
1425                    .unwrap_or_else(|| "unknown".to_string()),
1426                current_version: handler
1427                    .version()
1428                    .map(str::to_string)
1429                    .unwrap_or_else(|| "unknown".to_string()),
1430            })
1431        };
1432
1433        self.finalize_run(
1434            run_id,
1435            &run.workflow_name,
1436            result,
1437            &ctx,
1438            run_start,
1439            run.labels,
1440        )
1441        .await
1442    }
1443
1444    /// Execute a run by its ID (used by the worker after pick_next_pending).
1445    ///
1446    /// Delegates to [`execute_handler_run`](Self::execute_handler_run).
1447    ///
1448    /// # Errors
1449    ///
1450    /// Returns [`EngineError`] if the run is not found or execution fails.
1451    #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
1452    pub async fn execute_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1453        self.execute_handler_run(run_id).await
1454    }
1455
1456    /// Resume a run after human approval.
1457    ///
1458    /// Re-executes the handler with step replay: completed steps return
1459    /// cached output, approved approval steps are skipped, and execution
1460    /// continues from the first unexecuted step.
1461    ///
1462    /// Supports multiple approval gates -- each resume replays all prior
1463    /// steps and stops at the next approval (or completes the run).
1464    ///
1465    /// When `run_id` is a child run of a sub-workflow, the root run of its
1466    /// chain is moved back to `Running` and resumed instead: it replays, and
1467    /// its open `Workflow` step re-enters the same child run. The returned
1468    /// result is the root run's.
1469    ///
1470    /// Like [`execute_handler_run`](Self::execute_handler_run), the handler
1471    /// only starts once [`AgentProvider::release_run`] has stopped whatever
1472    /// a previous execution of the run left running.
1473    ///
1474    /// # Errors
1475    ///
1476    /// Returns [`EngineError::InvalidWorkflow`] if no handler matches, or when
1477    /// the root run of a child cannot be resumed (it is running or finished:
1478    /// the child is then failed).
1479    /// Returns [`EngineError`] if execution fails or hits another approval.
1480    /// A failed release fails the execution with [`EngineError::Operation`].
1481    /// Returns [`EngineError::HandlerVersionMismatch`] when the handler's
1482    /// current version is incompatible with the run's `handler_version` --
1483    /// checked before any step is replayed.
1484    #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
1485    pub async fn resume_run(&self, run_id: Uuid) -> Result<WorkflowResult, EngineError> {
1486        let run = self
1487            .store
1488            .get_run(run_id)
1489            .await?
1490            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1491
1492        if let Some(root_run_id) = chain_root(&run) {
1493            return self.resume_chain(run, root_run_id).await;
1494        }
1495
1496        self.resume_loaded_run(run).await
1497    }
1498
1499    /// Resume the root run of a suspended child run.
1500    ///
1501    /// The root waits without `scheduled_at` while its child is suspended, so
1502    /// nothing but this path ever wakes it. It is moved to `Running` and
1503    /// resumed; its replay re-enters the child. A root that is not suspended
1504    /// (already running, or finished) cannot take the child back: the child
1505    /// is failed.
1506    ///
1507    /// When the child holds a worker lease (it was picked by a worker), the
1508    /// lease is transferred to the root in the same update that moves it to
1509    /// `Running`, then released on the child: the worker keeps renewing the
1510    /// root it now executes, and the reaper recovers the root if that worker
1511    /// dies. A child without a lease (inline execution, API-side resume)
1512    /// leaves the root without one, as before.
1513    async fn resume_chain(
1514        &self,
1515        child: Run,
1516        root_run_id: Uuid,
1517    ) -> Result<WorkflowResult, EngineError> {
1518        let child_run_id = child.id;
1519        let lease = child.worker_id.zip(child.lease_expires_at);
1520        let root = self
1521            .store
1522            .get_run(root_run_id)
1523            .await?
1524            .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1525
1526        match root.status.state {
1527            RunStatus::AwaitingApproval | RunStatus::Pending => {
1528                self.move_root_to_running(root_run_id, lease.as_ref())
1529                    .await?;
1530            }
1531            RunStatus::Sleeping => {
1532                self.store
1533                    .update_run_status(root_run_id, RunStatus::Pending)
1534                    .await?;
1535                self.move_root_to_running(root_run_id, lease.as_ref())
1536                    .await?;
1537            }
1538            other => {
1539                let reason = format!(
1540                    "cannot resume child run {child_run_id}: root run {root_run_id} is {other}"
1541                );
1542                if let Err(err) = self
1543                    .fail_or_schedule_retry(child_run_id, &reason, false, None, None)
1544                    .await
1545                {
1546                    error!(
1547                        run_id = %child_run_id,
1548                        error = %err,
1549                        "failed to fail a child run whose root cannot resume"
1550                    );
1551                }
1552                return Err(EngineError::InvalidWorkflow(reason));
1553            }
1554        }
1555
1556        if lease.is_some() {
1557            self.store
1558                .update_run(
1559                    child_run_id,
1560                    RunUpdate {
1561                        lease: Some(LeaseUpdate::Release),
1562                        ..RunUpdate::default()
1563                    },
1564                )
1565                .await?;
1566        }
1567
1568        info!(
1569            run_id = %child_run_id,
1570            root_run_id = %root_run_id,
1571            lease_transferred = lease.is_some(),
1572            "child run resumed through its root run"
1573        );
1574
1575        let root = self
1576            .store
1577            .get_run(root_run_id)
1578            .await?
1579            .ok_or(EngineError::Store(StoreError::RunNotFound(root_run_id)))?;
1580        self.resume_loaded_run(root).await
1581    }
1582
1583    /// Move the suspended root run of a chain to `Running`.
1584    ///
1585    /// With `lease`, the root takes the worker lease of the child that
1586    /// resumes it, atomically with the transition, so it is never `Running`
1587    /// without an owner.
1588    async fn move_root_to_running(
1589        &self,
1590        root_run_id: Uuid,
1591        lease: Option<&(String, DateTime<Utc>)>,
1592    ) -> Result<(), EngineError> {
1593        match lease {
1594            Some((worker_id, expires_at)) => {
1595                self.store
1596                    .update_run(
1597                        root_run_id,
1598                        RunUpdate {
1599                            status: Some(RunStatus::Running),
1600                            lease: Some(LeaseUpdate::Set {
1601                                worker_id: worker_id.clone(),
1602                                expires_at: *expires_at,
1603                            }),
1604                            ..RunUpdate::default()
1605                        },
1606                    )
1607                    .await?;
1608            }
1609            None => {
1610                self.store
1611                    .update_run_status(root_run_id, RunStatus::Running)
1612                    .await?;
1613            }
1614        }
1615        Ok(())
1616    }
1617
1618    /// Resume `run`, already loaded and already `Running`.
1619    async fn resume_loaded_run(&self, run: Run) -> Result<WorkflowResult, EngineError> {
1620        let run_id = run.id;
1621        let _active = self.track_execution(run_id).await;
1622        let handler = self
1623            .handlers
1624            .get(&run.workflow_name)
1625            .ok_or_else(|| {
1626                EngineError::InvalidWorkflow(format!(
1627                    "no handler registered: {}",
1628                    run.workflow_name
1629                ))
1630            })?
1631            .clone();
1632
1633        info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
1634
1635        let run_start = Instant::now();
1636        let mut ctx = self.build_context_with_guard(&run, handler.as_ref());
1637
1638        let result = if handler.is_version_compatible(run.handler_version.as_deref()) {
1639            ctx.load_replay_steps().await?;
1640            self.release_then_execute(run_id, handler.as_ref(), &mut ctx)
1641                .await
1642        } else {
1643            Err(EngineError::HandlerVersionMismatch {
1644                run_id,
1645                workflow_name: run.workflow_name.clone(),
1646                run_version: run
1647                    .handler_version
1648                    .clone()
1649                    .unwrap_or_else(|| "unknown".to_string()),
1650                current_version: handler
1651                    .version()
1652                    .map(str::to_string)
1653                    .unwrap_or_else(|| "unknown".to_string()),
1654            })
1655        };
1656
1657        self.finalize_run(
1658            run_id,
1659            &run.workflow_name,
1660            result,
1661            &ctx,
1662            run_start,
1663            run.labels,
1664        )
1665        .await
1666    }
1667
1668    /// Store a signal and resolve every step waiting for its `(name, key)`.
1669    ///
1670    /// Each waiting step validates the payload against the JSON schema it
1671    /// stored when it opened. A matching step is completed with the payload
1672    /// and its run, if `Sleeping`, goes back to `Pending`: under
1673    /// [`ExecutionMode::Local`] it resumes in a background task, under
1674    /// [`ExecutionMode::Workers`] a worker picks it up. A step whose schema
1675    /// the payload does not match keeps waiting and is listed in
1676    /// [`SignalDelivery::rejected`].
1677    ///
1678    /// The signal is stored even when nobody waits for it: a run opening its
1679    /// wait step later still finds it. A signal whose `idempotency_id` was
1680    /// already used is not stored nor delivered again, and comes back with
1681    /// [`SignalDelivery::duplicate`] set.
1682    ///
1683    /// Publishes [`Event::SignalReceived`] for every stored signal.
1684    ///
1685    /// # Errors
1686    ///
1687    /// Returns [`EngineError::InvalidSignal`] when `name` or `key` is empty,
1688    /// and [`EngineError::Store`] when the signal cannot be stored or its
1689    /// waiters cannot be listed.
1690    ///
1691    /// # Examples
1692    ///
1693    /// ```no_run
1694    /// use std::sync::Arc;
1695    /// use ironflow_engine::engine::Engine;
1696    /// use ironflow_engine::error::EngineError;
1697    /// use ironflow_store::entities::NewSignal;
1698    /// use serde_json::json;
1699    ///
1700    /// # async fn example(engine: Arc<Engine>) -> Result<(), EngineError> {
1701    /// let delivery = engine
1702    ///     .deliver_signal(NewSignal {
1703    ///         name: "ci.pipeline_finished".to_string(),
1704    ///         key: "4f2a9c1".to_string(),
1705    ///         payload: json!({"status": "success"}),
1706    ///         idempotency_id: Some("delivery-42".to_string()),
1707    ///     })
1708    ///     .await?;
1709    /// println!("{} runs resumed", delivery.resumed.len());
1710    /// # Ok(())
1711    /// # }
1712    /// ```
1713    pub async fn deliver_signal(
1714        self: &Arc<Self>,
1715        signal: NewSignal,
1716    ) -> Result<SignalDelivery, EngineError> {
1717        if signal.name.trim().is_empty() {
1718            return Err(EngineError::InvalidSignal(
1719                "signal name must not be empty".to_string(),
1720            ));
1721        }
1722        if signal.key.trim().is_empty() {
1723            return Err(EngineError::InvalidSignal(
1724                "signal key must not be empty".to_string(),
1725            ));
1726        }
1727
1728        let stored = match self.store.insert_signal(signal).await? {
1729            SignalInsert::Created(stored) => stored,
1730            SignalInsert::Duplicate(existing) => {
1731                info!(
1732                    signal_id = %existing.id,
1733                    signal = %existing.name,
1734                    key = %existing.key,
1735                    "duplicate signal ignored"
1736                );
1737                return Ok(SignalDelivery {
1738                    signal_id: existing.id,
1739                    duplicate: true,
1740                    resumed: Vec::new(),
1741                    rejected: Vec::new(),
1742                });
1743            }
1744        };
1745
1746        let waiters = self
1747            .store
1748            .list_signal_waiters(&stored.name, &stored.key)
1749            .await?;
1750        let mut resumed = Vec::new();
1751        let mut rejected = Vec::new();
1752
1753        for step in waiters {
1754            if let Err(error) = validate_step_payload(step.input.as_ref(), &stored.payload) {
1755                rejected.push(SignalRejected {
1756                    run_id: step.run_id,
1757                    step_id: step.id,
1758                    error,
1759                });
1760                continue;
1761            }
1762
1763            match self
1764                .store
1765                .resolve_signal_step(step.id, received_output(&stored))
1766                .await
1767            {
1768                Ok(SignalStepResolution::Resolved {
1769                    run_id,
1770                    run_resumed,
1771                }) => {
1772                    resumed.push(SignalResumed {
1773                        run_id,
1774                        step_id: step.id,
1775                    });
1776                    if run_resumed && self.execution_mode == ExecutionMode::Local {
1777                        self.spawn_local_resume(run_id);
1778                    }
1779                }
1780                // A concurrent delivery or the timeout resolved it first.
1781                Ok(SignalStepResolution::NotWaiting { .. }) => {}
1782                Err(err) => {
1783                    error!(
1784                        run_id = %step.run_id,
1785                        step_id = %step.id,
1786                        error = %err,
1787                        "failed to resolve a waiting signal step"
1788                    );
1789                    rejected.push(SignalRejected {
1790                        run_id: step.run_id,
1791                        step_id: step.id,
1792                        error: err.to_string(),
1793                    });
1794                }
1795            }
1796        }
1797
1798        info!(
1799            signal_id = %stored.id,
1800            signal = %stored.name,
1801            key = %stored.key,
1802            resumed = resumed.len(),
1803            rejected = rejected.len(),
1804            "signal received"
1805        );
1806        self.event_publisher
1807            .publish(Event::SignalReceived(SignalReceivedEvent {
1808                signal_id: stored.id,
1809                name: stored.name.clone(),
1810                key: stored.key.clone(),
1811                resumed_runs: resumed.iter().map(|r| r.run_id).collect(),
1812                at: stored.received_at,
1813            }));
1814
1815        Ok(SignalDelivery {
1816            signal_id: stored.id,
1817            duplicate: false,
1818            resumed,
1819            rejected,
1820        })
1821    }
1822
1823    /// Send a typed signal: shorthand for [`deliver_signal`](Self::deliver_signal)
1824    /// with `S::NAME` as the name and `signal` as the payload.
1825    ///
1826    /// `key` identifies the occurrence (a commit SHA, an order ID).
1827    /// `idempotency_id`, when set, makes a redelivery of the same event (a
1828    /// webhook retried by its sender) a no-op.
1829    ///
1830    /// # Errors
1831    ///
1832    /// Returns [`EngineError::Serialization`] when `signal` cannot be
1833    /// serialized, and every error of [`deliver_signal`](Self::deliver_signal).
1834    ///
1835    /// # Examples
1836    ///
1837    /// ```no_run
1838    /// use std::sync::Arc;
1839    /// use ironflow_engine::engine::Engine;
1840    /// use ironflow_engine::error::EngineError;
1841    /// use ironflow_engine::signal::Signal;
1842    /// use schemars::JsonSchema;
1843    /// use serde::{Deserialize, Serialize};
1844    ///
1845    /// #[derive(Serialize, Deserialize, JsonSchema)]
1846    /// struct PipelineFinished {
1847    ///     status: String,
1848    /// }
1849    ///
1850    /// impl Signal for PipelineFinished {
1851    ///     const NAME: &'static str = "ci.pipeline_finished";
1852    /// }
1853    ///
1854    /// # async fn example(engine: Arc<Engine>) -> Result<(), EngineError> {
1855    /// let finished = PipelineFinished { status: "success".to_string() };
1856    /// engine.send_signal(&finished, "4f2a9c1", Some("delivery-42")).await?;
1857    /// # Ok(())
1858    /// # }
1859    /// ```
1860    pub async fn send_signal<S: Signal>(
1861        self: &Arc<Self>,
1862        signal: &S,
1863        key: &str,
1864        idempotency_id: Option<&str>,
1865    ) -> Result<SignalDelivery, EngineError> {
1866        let payload = to_value(signal)?;
1867        self.deliver_signal(NewSignal {
1868            name: S::NAME.to_string(),
1869            key: key.to_string(),
1870            payload,
1871            idempotency_id: idempotency_id.map(str::to_string),
1872        })
1873        .await
1874    }
1875
1876    /// Resume a run requeued to `Pending` in a background task.
1877    ///
1878    /// Used under [`ExecutionMode::Local`], where no worker would pick the
1879    /// run up. The state change already happened, so a failed resume is
1880    /// logged, not rolled back.
1881    ///
1882    /// The task first waits for the execution of the run still in flight in
1883    /// this process, if any: a run paused then resumed while a step runs
1884    /// would otherwise be replayed next to the execution that has not
1885    /// noticed the pause yet. Once it is gone, the run is read again:
1886    /// `Pending` is started, `Running` was left unfinished by the previous
1887    /// execution and is resumed as is, any other status means the previous
1888    /// execution finished the run or it was taken elsewhere, and nothing is
1889    /// done.
1890    pub(crate) fn spawn_local_resume(self: &Arc<Self>, run_id: Uuid) {
1891        let engine = Arc::clone(self);
1892        spawn(async move {
1893            engine.wait_until_idle(run_id).await;
1894            let run = match engine.store.get_run(run_id).await {
1895                Ok(Some(run)) => run,
1896                Ok(None) => {
1897                    error!(run_id = %run_id, "run to restart not found");
1898                    return;
1899                }
1900                Err(err) => {
1901                    error!(run_id = %run_id, error = %err, "failed to load a run to restart");
1902                    return;
1903                }
1904            };
1905            match run.status.state {
1906                RunStatus::Pending => {
1907                    if let Err(err) = engine
1908                        .store
1909                        .update_run_status(run_id, RunStatus::Running)
1910                        .await
1911                    {
1912                        error!(run_id = %run_id, error = %err, "failed to restart a woken run");
1913                        return;
1914                    }
1915                }
1916                RunStatus::Running => {}
1917                status => {
1918                    info!(
1919                        run_id = %run_id,
1920                        status = %status,
1921                        "run no longer waiting to restart, resume skipped"
1922                    );
1923                    return;
1924                }
1925            }
1926            if let Err(err) = engine.resume_run(run_id).await {
1927                error!(run_id = %run_id, error = %err, "failed to resume a woken run");
1928            }
1929        });
1930    }
1931
1932    /// Record a run failure, replaying the run later when retries remain.
1933    ///
1934    /// This is the single place where a failed run's fate is decided. When
1935    /// `retryable` is true and the run has not exhausted `max_retries`, the run
1936    /// moves to [`RunStatus::Retrying`] with `scheduled_at` set to
1937    /// `now + backoff` -- [`pick_next_pending`](ironflow_store::store::RunStore::pick_next_pending)
1938    /// picks it up again once that time has passed. Otherwise the run moves to
1939    /// [`RunStatus::Failed`].
1940    ///
1941    /// Either way, steps left non-terminal by the failed attempt are closed via
1942    /// [`fail_orphaned_steps`](Self::fail_orphaned_steps) so they are never
1943    /// confused with the next attempt's steps, and the sub-workflow runs it
1944    /// left non-terminal are cancelled by
1945    /// [`cancel_descendants`](Self::cancel_descendants) (a failure there is
1946    /// logged, not returned).
1947    ///
1948    /// Callers pass `retryable` explicitly rather than an error value, because
1949    /// the worker classifies failures it observes from the outside (a timeout, a
1950    /// panicked task) that never produce an [`EngineError`]. Use
1951    /// [`is_run_retryable`] to classify an
1952    /// [`EngineError`].
1953    ///
1954    /// `cost_usd` and `duration_ms` are `None` when the caller does not know the
1955    /// totals (a timeout or a panic observed from outside the handler); the
1956    /// values already stored on the run are then left untouched.
1957    ///
1958    /// A run an operator paused meanwhile is left untouched and
1959    /// [`RunStatus::Paused`] is returned: the resume replays the attempt.
1960    ///
1961    /// Returns the status the run was moved to.
1962    ///
1963    /// # Errors
1964    ///
1965    /// Returns [`EngineError::Store`] if the run does not exist or the update
1966    /// cannot be persisted.
1967    ///
1968    /// # Examples
1969    ///
1970    /// ```no_run
1971    /// use ironflow_engine::engine::Engine;
1972    /// use ironflow_engine::error::EngineError;
1973    /// use uuid::Uuid;
1974    ///
1975    /// # async fn example(engine: &Engine, run_id: Uuid) -> Result<(), EngineError> {
1976    /// let status = engine
1977    ///     .fail_or_schedule_retry(run_id, "run timed out after 300s", true, None, None)
1978    ///     .await?;
1979    /// # Ok(())
1980    /// # }
1981    /// ```
1982    pub async fn fail_or_schedule_retry(
1983        &self,
1984        run_id: Uuid,
1985        error: &str,
1986        retryable: bool,
1987        cost_usd: Option<Decimal>,
1988        duration_ms: Option<u64>,
1989    ) -> Result<RunStatus, EngineError> {
1990        let run = self
1991            .store
1992            .get_run(run_id)
1993            .await?
1994            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
1995
1996        // An operator paused the run: its fate is decided by the resume or
1997        // the cancellation, not by the attempt that stopped.
1998        if run.status.state == RunStatus::Paused {
1999            info!(run_id = %run_id, error = %error, "run paused, failure not recorded");
2000            return Ok(RunStatus::Paused);
2001        }
2002
2003        let has_attempts_left = run.retry_count < run.max_retries;
2004        let update = if retryable && has_attempts_left {
2005            let backoff = backoff_for_retry(run.retry_count);
2006            let scheduled_at = Utc::now() + TimeDelta::milliseconds(backoff.as_millis() as i64);
2007
2008            info!(
2009                run_id = %run_id,
2010                workflow = %run.workflow_name,
2011                attempt = run.retry_count + 1,
2012                max_retries = run.max_retries,
2013                backoff_secs = backoff.as_secs(),
2014                scheduled_at = %scheduled_at,
2015                "run failed, scheduling retry"
2016            );
2017
2018            RunUpdate {
2019                status: Some(RunStatus::Retrying),
2020                error: Some(error.to_string()),
2021                increment_retry: true,
2022                cost_usd,
2023                duration_ms,
2024                scheduled_at: Some(scheduled_at),
2025                ..RunUpdate::default()
2026            }
2027        } else {
2028            RunUpdate {
2029                status: Some(RunStatus::Failed),
2030                error: Some(error.to_string()),
2031                cost_usd,
2032                duration_ms,
2033                completed_at: Some(Utc::now()),
2034                ..RunUpdate::default()
2035            }
2036        };
2037
2038        let status = update.status.unwrap_or(RunStatus::Failed);
2039        self.store.update_run(run_id, update).await?;
2040        self.fail_orphaned_steps(run_id, error).await?;
2041        // The attempt is over: a retry starts new children, and nothing drives
2042        // those this attempt left running.
2043        self.cancel_descendants_of_stopped_run(run_id, error).await;
2044
2045        Ok(status)
2046    }
2047
2048    /// Mark the steps of a run requeued after a lost worker lease as interrupted.
2049    ///
2050    /// Every `Running` step is marked `Failed` with
2051    /// [`STEP_INTERRUPTED_ERROR`](ironflow_store::store::STEP_INTERRUPTED_ERROR).
2052    /// The run keeps its attempt number, so when it is picked up again its
2053    /// finished steps are replayed and each interrupted step is executed again
2054    /// at the same position; an interrupted `Workflow` step re-enters the same
2055    /// child run. `Pending` and `AwaitingApproval` steps are left as they are,
2056    /// unlike [`fail_orphaned_steps`](Self::fail_orphaned_steps), which ends
2057    /// the run's steps for good.
2058    ///
2059    /// Errors from individual step updates are logged but do not abort the cleanup.
2060    ///
2061    /// # Errors
2062    ///
2063    /// Returns [`EngineError`] if listing steps fails.
2064    ///
2065    /// # Examples
2066    ///
2067    /// ```no_run
2068    /// use ironflow_engine::engine::Engine;
2069    /// use ironflow_engine::error::EngineError;
2070    /// use uuid::Uuid;
2071    ///
2072    /// # async fn example(engine: &Engine, run_id: Uuid) -> Result<(), EngineError> {
2073    /// engine.interrupt_running_steps(run_id).await?;
2074    /// # Ok(())
2075    /// # }
2076    /// ```
2077    pub async fn interrupt_running_steps(&self, run_id: Uuid) -> Result<(), EngineError> {
2078        interrupt_running_steps(self.store.as_ref(), run_id).await
2079    }
2080
2081    /// Fail all non-terminal steps for a run.
2082    ///
2083    /// Called after a run is marked as failed (timeout, error, panic) to clean up
2084    /// orphaned steps that are still in `Running`, `Pending`, or `AwaitingApproval`.
2085    ///
2086    /// - `Running` / `AwaitingApproval` steps are marked `Failed`.
2087    /// - `Pending` steps are marked `Skipped` (FSM does not allow Pending -> Failed).
2088    ///
2089    /// Errors from individual step updates are logged but do not abort the cleanup.
2090    ///
2091    /// # Errors
2092    ///
2093    /// Returns [`EngineError`] if listing steps fails.
2094    pub async fn fail_orphaned_steps(
2095        &self,
2096        run_id: Uuid,
2097        error_message: &str,
2098    ) -> Result<(), EngineError> {
2099        let steps = self.store.list_steps(run_id).await?;
2100        let now = Utc::now();
2101
2102        for step in steps {
2103            if step.status.state.is_terminal() {
2104                continue;
2105            }
2106
2107            let (target_status, error) = match step.status.state {
2108                StepStatus::Running | StepStatus::AwaitingApproval => {
2109                    let err = if step.error.is_some() {
2110                        None
2111                    } else {
2112                        Some(error_message.to_string())
2113                    };
2114                    (StepStatus::Failed, err)
2115                }
2116                StepStatus::Pending => (StepStatus::Skipped, None),
2117                _ => continue,
2118            };
2119
2120            if let Err(e) = self
2121                .store
2122                .update_step(
2123                    step.id,
2124                    StepUpdate {
2125                        status: Some(target_status),
2126                        error,
2127                        completed_at: Some(now),
2128                        ..StepUpdate::default()
2129                    },
2130                )
2131                .await
2132            {
2133                warn!(
2134                    run_id = %run_id,
2135                    step_id = %step.id,
2136                    step_name = %step.name,
2137                    error = %e,
2138                    "failed to cleanup orphaned step"
2139                );
2140            } else {
2141                info!(
2142                    run_id = %run_id,
2143                    step_id = %step.id,
2144                    step_name = %step.name,
2145                    from = %step.status.state,
2146                    to = %target_status,
2147                    "cleaned up orphaned step"
2148                );
2149            }
2150        }
2151
2152        Ok(())
2153    }
2154
2155    /// Execute `handler` once [`AgentProvider::release_run`] has stopped
2156    /// whatever a previous execution of the run left running (an agent pod
2157    /// writing to a shared worktree). A failed release fails the execution
2158    /// before its first step, with [`EngineError::Operation`].
2159    async fn release_then_execute(
2160        &self,
2161        run_id: Uuid,
2162        handler: &dyn WorkflowHandler,
2163        ctx: &mut WorkflowContext,
2164    ) -> Result<(), EngineError> {
2165        match self.provider.release_run(&run_id.to_string()).await {
2166            Ok(()) => handler.execute(ctx).await,
2167            Err(e) => Err(EngineError::Operation(OperationError::Agent(e))),
2168        }
2169    }
2170
2171    /// Finalize a run with the given result and context.
2172    ///
2173    /// On success: updates run to Completed with cost, duration, and completed_at.
2174    /// On failure: updates run to Failed with error, cost, duration, and completed_at.
2175    /// Always: fetches and returns the final Run.
2176    async fn finalize_run(
2177        &self,
2178        run_id: Uuid,
2179        workflow_name: &str,
2180        result: Result<(), EngineError>,
2181        ctx: &WorkflowContext,
2182        run_start: Instant,
2183        run_labels: HashMap<String, String>,
2184    ) -> Result<WorkflowResult, EngineError> {
2185        // Covers the whole run: previous attempts plus this one, so a retried
2186        // run reports the time it really consumed.
2187        let total_duration = ctx.carried_duration_ms() + run_start.elapsed().as_millis() as u64;
2188        let completed_at = Utc::now();
2189
2190        // Paused while it ran: whatever the handler returned, the run stays
2191        // paused for the operator. Only what this execution spent is kept;
2192        // the resume replays from the first step that did not complete.
2193        if let Some(run) = self.store.get_run(run_id).await?
2194            && run.status.state == RunStatus::Paused
2195        {
2196            let run = self
2197                .store
2198                .update_run_returning(
2199                    run_id,
2200                    RunUpdate {
2201                        cost_usd: Some(ctx.total_cost_usd()),
2202                        duration_ms: Some(total_duration),
2203                        ..RunUpdate::default()
2204                    },
2205                )
2206                .await?;
2207            info!(
2208                run_id = %run_id,
2209                outcome = ?result.err().map(|err| err.to_string()),
2210                "run paused, execution stopped"
2211            );
2212            return Ok(WorkflowResult {
2213                run,
2214                steps: ctx.step_results().to_vec(),
2215            });
2216        }
2217
2218        // The pause stopped this execution but the run was resumed since: the
2219        // resume continues or restarts it, so nothing is recorded here.
2220        if matches!(result, Err(EngineError::RunPaused { .. })) {
2221            let run = self
2222                .store
2223                .get_run(run_id)
2224                .await?
2225                .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2226            info!(
2227                run_id = %run_id,
2228                status = %run.status.state,
2229                "run resumed after the pause, execution stopped"
2230            );
2231            return Ok(WorkflowResult {
2232                run,
2233                steps: ctx.step_results().to_vec(),
2234            });
2235        }
2236
2237        let final_status;
2238        let final_run;
2239
2240        match result {
2241            Ok(()) => {
2242                final_status = if ctx.has_allowed_failure() {
2243                    RunStatus::Warning
2244                } else {
2245                    RunStatus::Completed
2246                };
2247                final_run = self
2248                    .store
2249                    .update_run_returning(
2250                        run_id,
2251                        RunUpdate {
2252                            status: Some(final_status),
2253                            cost_usd: Some(ctx.total_cost_usd()),
2254                            duration_ms: Some(total_duration),
2255                            completed_at: Some(completed_at),
2256                            output: ctx.output().cloned(),
2257                            ..RunUpdate::default()
2258                        },
2259                    )
2260                    .await?;
2261
2262                info!(
2263                    run_id = %run_id,
2264                    status = %final_status,
2265                    cost_usd = %ctx.total_cost_usd(),
2266                    duration_ms = total_duration,
2267                    "run completed"
2268                );
2269            }
2270            Err(EngineError::ApprovalRequired {
2271                run_id: approval_run_id,
2272                step_id,
2273                ref message,
2274            }) => {
2275                final_status = RunStatus::AwaitingApproval;
2276                final_run = self
2277                    .store
2278                    .update_run_returning(
2279                        run_id,
2280                        RunUpdate {
2281                            status: Some(RunStatus::AwaitingApproval),
2282                            cost_usd: Some(ctx.total_cost_usd()),
2283                            duration_ms: Some(total_duration),
2284                            ..RunUpdate::default()
2285                        },
2286                    )
2287                    .await?;
2288
2289                info!(
2290                    run_id = %approval_run_id,
2291                    step_id = %step_id,
2292                    message = %message,
2293                    "run awaiting approval"
2294                );
2295
2296                self.publish_approval_requested(approval_run_id, step_id, message)
2297                    .await?;
2298            }
2299            Err(EngineError::ChildSuspended {
2300                run_id: child_run_id,
2301                ref cause,
2302            }) => {
2303                final_status = cause.suspension_status();
2304                // No `scheduled_at`: the suspended descendant owns the wake-up
2305                // and resumes this run through `resume_chain`. A wake-up armed
2306                // here too would resume the chain twice.
2307                final_run = self
2308                    .store
2309                    .update_run_returning(
2310                        run_id,
2311                        RunUpdate {
2312                            status: Some(final_status),
2313                            cost_usd: Some(ctx.total_cost_usd()),
2314                            duration_ms: Some(total_duration),
2315                            ..RunUpdate::default()
2316                        },
2317                    )
2318                    .await?;
2319
2320                let leaf = cause.suspension_leaf();
2321                info!(
2322                    run_id = %run_id,
2323                    child_run_id = %child_run_id,
2324                    status = %final_status,
2325                    cause = %leaf,
2326                    "run suspended with its child run"
2327                );
2328
2329                match leaf {
2330                    EngineError::ApprovalRequired {
2331                        run_id: approval_run_id,
2332                        step_id,
2333                        message,
2334                    } => {
2335                        self.publish_approval_requested(*approval_run_id, *step_id, message)
2336                            .await?;
2337                    }
2338                    EngineError::SignalWaiting {
2339                        run_id: wait_run_id,
2340                        step_id,
2341                        step_name,
2342                        name,
2343                        key,
2344                        deadline_at,
2345                    } => {
2346                        self.event_publisher
2347                            .publish(Event::SignalAwaited(SignalAwaitedEvent {
2348                                run_id: *wait_run_id,
2349                                step_id: *step_id,
2350                                step_name: step_name.clone(),
2351                                name: name.clone(),
2352                                key: key.clone(),
2353                                deadline_at: *deadline_at,
2354                                at: Utc::now(),
2355                            }));
2356                    }
2357                    // A human input or a delay publishes no suspension event,
2358                    // like on a top-level run.
2359                    _ => {}
2360                }
2361            }
2362            Err(EngineError::HumanInputRequired {
2363                run_id: input_run_id,
2364                step_id,
2365                ref message,
2366            }) => {
2367                final_status = RunStatus::AwaitingApproval;
2368                final_run = self
2369                    .store
2370                    .update_run_returning(
2371                        run_id,
2372                        RunUpdate {
2373                            status: Some(RunStatus::AwaitingApproval),
2374                            cost_usd: Some(ctx.total_cost_usd()),
2375                            duration_ms: Some(total_duration),
2376                            ..RunUpdate::default()
2377                        },
2378                    )
2379                    .await?;
2380
2381                // No `ApprovalRequested` event: a human input is not an approval.
2382                info!(
2383                    run_id = %input_run_id,
2384                    step_id = %step_id,
2385                    message = %message,
2386                    "run awaiting human input"
2387                );
2388            }
2389            Err(EngineError::DelaySleeping {
2390                run_id: delay_run_id,
2391                step_id,
2392                wake_at,
2393            }) => {
2394                final_status = RunStatus::Sleeping;
2395                final_run = self
2396                    .store
2397                    .update_run_returning(
2398                        run_id,
2399                        RunUpdate {
2400                            status: Some(RunStatus::Sleeping),
2401                            cost_usd: Some(ctx.total_cost_usd()),
2402                            duration_ms: Some(total_duration),
2403                            scheduled_at: Some(wake_at),
2404                            ..RunUpdate::default()
2405                        },
2406                    )
2407                    .await?;
2408
2409                info!(
2410                    run_id = %delay_run_id,
2411                    step_id = %step_id,
2412                    wake_at = %wake_at,
2413                    "run sleeping until delay elapses"
2414                );
2415            }
2416            Err(EngineError::CapacitySleeping {
2417                run_id: capacity_run_id,
2418                step_id,
2419                ref kind,
2420                wake_at,
2421            }) => {
2422                final_status = RunStatus::Sleeping;
2423                final_run = self
2424                    .store
2425                    .update_run_returning(
2426                        run_id,
2427                        RunUpdate {
2428                            status: Some(RunStatus::Sleeping),
2429                            cost_usd: Some(ctx.total_cost_usd()),
2430                            duration_ms: Some(total_duration),
2431                            scheduled_at: Some(wake_at),
2432                            capacity_wait_kind: Some(ProviderKind::new(kind.as_str())),
2433                            ..RunUpdate::default()
2434                        },
2435                    )
2436                    .await?;
2437
2438                info!(
2439                    run_id = %capacity_run_id,
2440                    step_id = %step_id,
2441                    kind = %kind,
2442                    wake_at = %wake_at,
2443                    "run sleeping until provider capacity returns"
2444                );
2445            }
2446            Err(EngineError::SignalWaiting {
2447                run_id: wait_run_id,
2448                step_id,
2449                ref step_name,
2450                ref name,
2451                ref key,
2452                deadline_at,
2453            }) => {
2454                final_status = RunStatus::Sleeping;
2455                // Atomic with the step lock: a signal delivered since the step
2456                // opened leaves the run due right away instead of until the
2457                // deadline.
2458                let waiting = self
2459                    .store
2460                    .suspend_run_on_signal(run_id, step_id, deadline_at)
2461                    .await?;
2462                final_run = self
2463                    .store
2464                    .update_run_returning(
2465                        run_id,
2466                        RunUpdate {
2467                            cost_usd: Some(ctx.total_cost_usd()),
2468                            duration_ms: Some(total_duration),
2469                            ..RunUpdate::default()
2470                        },
2471                    )
2472                    .await?;
2473
2474                if waiting {
2475                    self.event_publisher
2476                        .publish(Event::SignalAwaited(SignalAwaitedEvent {
2477                            run_id: wait_run_id,
2478                            step_id,
2479                            step_name: step_name.clone(),
2480                            name: name.clone(),
2481                            key: key.clone(),
2482                            deadline_at,
2483                            at: Utc::now(),
2484                        }));
2485                }
2486
2487                info!(
2488                    run_id = %wait_run_id,
2489                    step_id = %step_id,
2490                    signal = %name,
2491                    key = %key,
2492                    deadline_at = %deadline_at,
2493                    waiting,
2494                    "run sleeping until a signal arrives"
2495                );
2496            }
2497            Err(err) => {
2498                // A guardrail stop (budget or workflow guard) is deliberate,
2499                // not a breakage: the run is cancelled, never failed and
2500                // never replayed.
2501                let guardrail_stop = matches!(
2502                    err,
2503                    EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
2504                );
2505
2506                final_status = if guardrail_stop {
2507                    if let Err(store_err) = self
2508                        .store
2509                        .update_run(
2510                            run_id,
2511                            RunUpdate {
2512                                status: Some(RunStatus::Cancelled),
2513                                error: Some(err.to_string()),
2514                                cost_usd: Some(ctx.total_cost_usd()),
2515                                duration_ms: Some(total_duration),
2516                                completed_at: Some(completed_at),
2517                                output: ctx.output().cloned(),
2518                                ..RunUpdate::default()
2519                            },
2520                        )
2521                        .await
2522                    {
2523                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
2524                    }
2525                    if let Err(cleanup_err) = self
2526                        .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
2527                        .await
2528                    {
2529                        error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
2530                    }
2531                    RunStatus::Cancelled
2532                } else {
2533                    // Written before the failure so a failed run keeps the
2534                    // verdict the handler set before returning its error.
2535                    if let Some(output) = ctx.output()
2536                        && let Err(store_err) = self
2537                            .store
2538                            .update_run(
2539                                run_id,
2540                                RunUpdate {
2541                                    output: Some(output.clone()),
2542                                    ..RunUpdate::default()
2543                                },
2544                            )
2545                            .await
2546                    {
2547                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run output");
2548                    }
2549                    self.fail_or_schedule_retry(
2550                        run_id,
2551                        &err.to_string(),
2552                        is_run_retryable(&err),
2553                        Some(ctx.total_cost_usd()),
2554                        Some(total_duration),
2555                    )
2556                    .await
2557                    .unwrap_or_else(|store_err| {
2558                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
2559                        RunStatus::Failed
2560                    })
2561                };
2562
2563                if matches!(err, EngineError::RunBudgetExceeded { .. }) {
2564                    self.on_run_budget_exceeded(workflow_name, run_id, &err);
2565                }
2566
2567                error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
2568
2569                self.publish_run_status_changed(
2570                    workflow_name,
2571                    run_id,
2572                    final_status,
2573                    Some(err.to_string()),
2574                    ctx,
2575                    total_duration,
2576                    run_labels,
2577                );
2578
2579                #[cfg(feature = "prometheus")]
2580                self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2581
2582                return Err(err);
2583            }
2584        }
2585
2586        self.publish_run_status_changed(
2587            workflow_name,
2588            run_id,
2589            final_status,
2590            None,
2591            ctx,
2592            total_duration,
2593            run_labels,
2594        );
2595
2596        #[cfg(feature = "prometheus")]
2597        self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
2598
2599        Ok(WorkflowResult {
2600            run: final_run,
2601            steps: ctx.step_results().to_vec(),
2602        })
2603    }
2604
2605    /// Publish [`Event::ApprovalRequested`] for an approval gate that opened
2606    /// on `run_id`, with the requirement recorded when the gate opened.
2607    async fn publish_approval_requested(
2608        &self,
2609        run_id: Uuid,
2610        step_id: Uuid,
2611        message: &str,
2612    ) -> Result<(), EngineError> {
2613        let requirement = self
2614            .store
2615            .get_step(step_id)
2616            .await?
2617            .and_then(|s| s.approval_requirement);
2618        self.event_publisher
2619            .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
2620                run_id,
2621                step_id,
2622                message: message.to_string(),
2623                requirement,
2624                at: Utc::now(),
2625            }));
2626        Ok(())
2627    }
2628
2629    /// Fail every ancestor of a child run, closest first.
2630    ///
2631    /// Used when a gate inside a sub-workflow is rejected: the child run
2632    /// fails, and the runs suspended with it (its parent, up to the root)
2633    /// must not stay suspended on a child that will never resume. Each
2634    /// ancestor goes through [`fail_or_schedule_retry`](Self::fail_or_schedule_retry)
2635    /// without a retry, which also fails its open `Workflow` step. A run that
2636    /// is not a child of a sub-workflow has no ancestor: nothing happens.
2637    ///
2638    /// A paused ancestor is not failed: the root of the chain is set to resume
2639    /// to `Pending`, so its replay observes the failed child once an operator
2640    /// resumes it.
2641    ///
2642    /// # Errors
2643    ///
2644    /// Returns [`EngineError::Store`] if a run of the chain does not exist or
2645    /// its failure cannot be persisted.
2646    ///
2647    /// # Examples
2648    ///
2649    /// ```no_run
2650    /// use ironflow_engine::engine::Engine;
2651    /// use ironflow_engine::error::EngineError;
2652    /// use uuid::Uuid;
2653    ///
2654    /// # async fn example(engine: &Engine, child_run_id: Uuid) -> Result<(), EngineError> {
2655    /// engine
2656    ///     .fail_ancestors(child_run_id, "approval rejected in a sub-workflow")
2657    ///     .await?;
2658    /// # Ok(())
2659    /// # }
2660    /// ```
2661    pub async fn fail_ancestors(&self, run_id: Uuid, reason: &str) -> Result<(), EngineError> {
2662        let mut current = self
2663            .store
2664            .get_run(run_id)
2665            .await?
2666            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
2667        // Labels are data: a chain that loops back on itself stops there.
2668        let mut visited = HashSet::from([run_id]);
2669
2670        while let Some(parent_id) = chain_parent(&current) {
2671            if !visited.insert(parent_id) {
2672                break;
2673            }
2674            let status = self
2675                .fail_or_schedule_retry(parent_id, reason, false, None, None)
2676                .await?;
2677            // A paused chain waits for its operator: the root replays on
2678            // resume and observes the failed child, like a cancelled one.
2679            if status == RunStatus::Paused {
2680                self.requeue_paused_root(&current).await?;
2681                break;
2682            }
2683            info!(
2684                run_id = %run_id,
2685                ancestor_run_id = %parent_id,
2686                status = %status,
2687                "ancestor run failed with its child"
2688            );
2689            current = self
2690                .store
2691                .get_run(parent_id)
2692                .await?
2693                .ok_or(EngineError::Store(StoreError::RunNotFound(parent_id)))?;
2694        }
2695
2696        Ok(())
2697    }
2698
2699    /// Emit Prometheus metrics for a completed run.
2700    #[cfg(feature = "prometheus")]
2701    fn emit_run_metrics(
2702        &self,
2703        workflow_name: &str,
2704        status: RunStatus,
2705        duration_ms: u64,
2706        ctx: &WorkflowContext,
2707    ) {
2708        let status_str = status.to_string();
2709        let wf = workflow_name.to_string();
2710
2711        counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
2712        histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
2713            .record(duration_ms as f64 / 1000.0);
2714        histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
2715            ctx.total_cost_usd()
2716                .to_string()
2717                .parse::<f64>()
2718                .unwrap_or(0.0),
2719        );
2720        gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
2721    }
2722
2723    /// Record the metric and publish the audit event for a run that hit its
2724    /// cost cap.
2725    ///
2726    /// A non-[`RunBudgetExceeded`](EngineError::RunBudgetExceeded) error is
2727    /// ignored, so callers can pass the error unconditionally.
2728    fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
2729        let EngineError::RunBudgetExceeded {
2730            limit_usd,
2731            spent_usd,
2732            step_budget_usd,
2733            ..
2734        } = err
2735        else {
2736            return;
2737        };
2738
2739        #[cfg(feature = "prometheus")]
2740        counter!(
2741            RUN_BUDGET_EXCEEDED_TOTAL,
2742            "workflow" => workflow_name.to_string(),
2743            "scope" => "run",
2744        )
2745        .increment(1);
2746
2747        self.event_publisher
2748            .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
2749                run_id,
2750                workflow_name: workflow_name.to_string(),
2751                limit_usd: *limit_usd,
2752                spent_usd: *spent_usd,
2753                step_budget_usd: *step_budget_usd,
2754                at: Utc::now(),
2755            }));
2756    }
2757
2758    /// Publish a run status changed event to all registered subscribers.
2759    ///
2760    /// `from` is always `Running` because `finalize_run` is only called
2761    /// from a running state.
2762    #[allow(clippy::too_many_arguments)]
2763    fn publish_run_status_changed(
2764        &self,
2765        workflow_name: &str,
2766        run_id: Uuid,
2767        to: RunStatus,
2768        error: Option<String>,
2769        ctx: &WorkflowContext,
2770        duration_ms: u64,
2771        labels: HashMap<String, String>,
2772    ) {
2773        let now = Utc::now();
2774        let cost_usd = ctx.total_cost_usd();
2775        let wf = workflow_name.to_string();
2776
2777        self.event_publisher
2778            .publish(Event::RunStatusChanged(RunStatusChangedEvent {
2779                run_id,
2780                workflow_name: wf.clone(),
2781                from: RunStatus::Running,
2782                to,
2783                error: error.clone(),
2784                cost_usd,
2785                duration_ms,
2786                labels: labels.clone(),
2787                at: now,
2788            }));
2789
2790        if to == RunStatus::Failed {
2791            self.event_publisher
2792                .publish(Event::RunFailed(RunFailedEvent {
2793                    run_id,
2794                    workflow_name: wf,
2795                    error,
2796                    cost_usd,
2797                    duration_ms,
2798                    labels,
2799                    at: now,
2800                }));
2801        }
2802    }
2803}
2804
2805impl fmt::Debug for Engine {
2806    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
2807        f.debug_struct("Engine")
2808            .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
2809            .finish_non_exhaustive()
2810    }
2811}
2812
2813#[cfg(test)]
2814mod tests {
2815    use super::*;
2816    use crate::config::ShellConfig;
2817    use crate::handler::{HandlerFuture, WorkflowHandler};
2818    use ironflow_core::providers::claude::ClaudeCodeProvider;
2819    use ironflow_core::providers::record_replay::RecordReplayProvider;
2820    use ironflow_store::memory::InMemoryStore;
2821    use ironflow_store::models::{MAX_PRIORITY, MIN_PRIORITY, StepStatus};
2822    use serde_json::json;
2823
2824    // Test handler that echoes a message via shell
2825    struct EchoWorkflow;
2826
2827    impl WorkflowHandler for EchoWorkflow {
2828        fn name(&self) -> &str {
2829            "echo-workflow"
2830        }
2831
2832        fn describe(&self) -> WorkflowInfo {
2833            WorkflowInfo {
2834                description: "A simple workflow that echoes hello".to_string(),
2835                source_code: None,
2836                sub_workflows: Vec::new(),
2837                category: None,
2838                version: self.version().map(str::to_string),
2839                compatible_versions: Vec::new(),
2840                input_schema: None,
2841                default_labels: HashMap::new(),
2842                schedule: self.schedule().cloned(),
2843                default_max_cost_usd: self.default_max_cost_usd(),
2844                priority: self.priority(),
2845            }
2846        }
2847
2848        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2849            Box::pin(async move {
2850                ctx.shell("greet", ShellConfig::new("echo hello")).await?;
2851                Ok(())
2852            })
2853        }
2854    }
2855
2856    // Test handler that fails
2857    struct FailingWorkflow;
2858
2859    impl WorkflowHandler for FailingWorkflow {
2860        fn name(&self) -> &str {
2861            "failing-workflow"
2862        }
2863
2864        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2865            Box::pin(async move {
2866                ctx.shell("fail", ShellConfig::new("exit 1")).await?;
2867                Ok(())
2868            })
2869        }
2870    }
2871
2872    fn create_test_engine() -> Engine {
2873        let store = Arc::new(InMemoryStore::new());
2874        let inner = ClaudeCodeProvider::new();
2875        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2876            inner,
2877            "/tmp/ironflow-fixtures",
2878        ));
2879        Engine::new(store, provider)
2880    }
2881
2882    #[test]
2883    fn engine_new_creates_instance() {
2884        let engine = create_test_engine();
2885        assert_eq!(engine.handler_names().len(), 0);
2886    }
2887
2888    #[test]
2889    fn execution_mode_defaults_to_local() {
2890        let engine = create_test_engine();
2891        assert_eq!(engine.execution_mode(), ExecutionMode::Local);
2892    }
2893
2894    #[test]
2895    fn with_execution_mode_overrides_the_default() {
2896        let engine = create_test_engine().with_execution_mode(ExecutionMode::Workers);
2897        assert_eq!(engine.execution_mode(), ExecutionMode::Workers);
2898    }
2899
2900    #[test]
2901    fn engine_register_handler() {
2902        let mut engine = create_test_engine();
2903        let result = engine.register(EchoWorkflow);
2904        assert!(result.is_ok());
2905        assert_eq!(engine.handler_names().len(), 1);
2906        assert!(engine.handler_names().contains(&"echo-workflow"));
2907    }
2908
2909    #[test]
2910    fn engine_register_duplicate_returns_error() {
2911        let mut engine = create_test_engine();
2912        engine.register(EchoWorkflow).unwrap();
2913        let result = engine.register(EchoWorkflow);
2914        assert!(result.is_err());
2915    }
2916
2917    #[test]
2918    fn engine_get_handler_found() {
2919        let mut engine = create_test_engine();
2920        engine.register(EchoWorkflow).unwrap();
2921        let handler = engine.get_handler("echo-workflow");
2922        assert!(handler.is_some());
2923    }
2924
2925    #[test]
2926    fn engine_get_handler_not_found() {
2927        let engine = create_test_engine();
2928        let handler = engine.get_handler("nonexistent");
2929        assert!(handler.is_none());
2930    }
2931
2932    #[test]
2933    fn engine_handler_names_lists_all() {
2934        let mut engine = create_test_engine();
2935        engine.register(EchoWorkflow).unwrap();
2936        engine.register(FailingWorkflow).unwrap();
2937        let names = engine.handler_names();
2938        assert_eq!(names.len(), 2);
2939        assert!(names.contains(&"echo-workflow"));
2940        assert!(names.contains(&"failing-workflow"));
2941    }
2942
2943    #[test]
2944    fn engine_handler_info_returns_description() {
2945        let mut engine = create_test_engine();
2946        engine.register(EchoWorkflow).unwrap();
2947        let info = engine.handler_info("echo-workflow");
2948        assert!(info.is_some());
2949        let info = info.unwrap();
2950        assert_eq!(info.description, "A simple workflow that echoes hello");
2951    }
2952
2953    struct CategorizedWorkflow;
2954
2955    impl WorkflowHandler for CategorizedWorkflow {
2956        fn name(&self) -> &str {
2957            "categorized"
2958        }
2959        fn category(&self) -> Option<&str> {
2960            Some("data/etl")
2961        }
2962        fn execute<'a>(
2963            &'a self,
2964            _ctx: &'a mut WorkflowContext,
2965        ) -> crate::handler::HandlerFuture<'a> {
2966            Box::pin(async move { Ok(()) })
2967        }
2968    }
2969
2970    #[test]
2971    fn engine_default_describe_propagates_category() {
2972        let mut engine = create_test_engine();
2973        engine.register(CategorizedWorkflow).unwrap();
2974        let info = engine.handler_info("categorized").unwrap();
2975        assert_eq!(info.category.as_deref(), Some("data/etl"));
2976    }
2977
2978    #[test]
2979    fn engine_default_describe_without_category() {
2980        let mut engine = create_test_engine();
2981        engine.register(EchoWorkflow).unwrap();
2982        let info = engine.handler_info("echo-workflow").unwrap();
2983        assert!(info.category.is_none());
2984    }
2985
2986    // -----------------------------------------------------------------------
2987    // Schedule tests
2988    // -----------------------------------------------------------------------
2989
2990    struct ScheduledWorkflow {
2991        schedule: CronSchedule,
2992    }
2993
2994    impl ScheduledWorkflow {
2995        fn new() -> Self {
2996            Self {
2997                schedule: CronSchedule::new("0 0 * * * *").unwrap(),
2998            }
2999        }
3000    }
3001
3002    impl WorkflowHandler for ScheduledWorkflow {
3003        fn name(&self) -> &str {
3004            "scheduled"
3005        }
3006        fn schedule(&self) -> Option<&CronSchedule> {
3007            Some(&self.schedule)
3008        }
3009        fn execute<'a>(
3010            &'a self,
3011            _ctx: &'a mut WorkflowContext,
3012        ) -> crate::handler::HandlerFuture<'a> {
3013            Box::pin(async move { Ok(()) })
3014        }
3015    }
3016
3017    #[test]
3018    fn engine_default_describe_propagates_schedule() {
3019        let mut engine = create_test_engine();
3020        engine.register(ScheduledWorkflow::new()).unwrap();
3021        let info = engine.handler_info("scheduled").unwrap();
3022        assert_eq!(
3023            info.schedule.as_ref().map(|s| s.as_str()),
3024            Some("0 0 * * * *")
3025        );
3026    }
3027
3028    #[test]
3029    fn engine_default_describe_without_schedule() {
3030        let mut engine = create_test_engine();
3031        engine.register(EchoWorkflow).unwrap();
3032        let info = engine.handler_info("echo-workflow").unwrap();
3033        assert!(info.schedule.is_none());
3034    }
3035
3036    #[test]
3037    fn scheduled_handlers_returns_only_scheduled() {
3038        let mut engine = create_test_engine();
3039        engine.register(EchoWorkflow).unwrap();
3040        engine.register(ScheduledWorkflow::new()).unwrap();
3041        engine.register(FailingWorkflow).unwrap();
3042
3043        let scheduled = engine.scheduled_handlers();
3044        assert_eq!(scheduled.len(), 1);
3045        assert_eq!(scheduled[0].0, "scheduled");
3046        assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
3047    }
3048
3049    #[test]
3050    fn scheduled_handlers_empty_when_none_scheduled() {
3051        let mut engine = create_test_engine();
3052        engine.register(EchoWorkflow).unwrap();
3053        engine.register(FailingWorkflow).unwrap();
3054
3055        let scheduled = engine.scheduled_handlers();
3056        assert!(scheduled.is_empty());
3057    }
3058
3059    struct BadCategoryWorkflow(&'static str);
3060
3061    impl WorkflowHandler for BadCategoryWorkflow {
3062        fn name(&self) -> &str {
3063            "bad-category"
3064        }
3065        fn category(&self) -> Option<&str> {
3066            Some(self.0)
3067        }
3068        fn execute<'a>(
3069            &'a self,
3070            _ctx: &'a mut WorkflowContext,
3071        ) -> crate::handler::HandlerFuture<'a> {
3072            Box::pin(async move { Ok(()) })
3073        }
3074    }
3075
3076    #[test]
3077    fn engine_register_rejects_empty_category() {
3078        let mut engine = create_test_engine();
3079        let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
3080        match err {
3081            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
3082            other => panic!("expected InvalidWorkflow, got {other:?}"),
3083        }
3084    }
3085
3086    #[test]
3087    fn engine_register_rejects_leading_slash_category() {
3088        let mut engine = create_test_engine();
3089        let err = engine
3090            .register(BadCategoryWorkflow("/data/etl"))
3091            .unwrap_err();
3092        match err {
3093            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
3094            other => panic!("expected InvalidWorkflow, got {other:?}"),
3095        }
3096    }
3097
3098    #[test]
3099    fn engine_register_rejects_trailing_slash_category() {
3100        let mut engine = create_test_engine();
3101        let err = engine
3102            .register(BadCategoryWorkflow("data/etl/"))
3103            .unwrap_err();
3104        match err {
3105            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
3106            other => panic!("expected InvalidWorkflow, got {other:?}"),
3107        }
3108    }
3109
3110    #[test]
3111    fn engine_register_rejects_double_slash_category() {
3112        let mut engine = create_test_engine();
3113        let err = engine
3114            .register(BadCategoryWorkflow("data//etl"))
3115            .unwrap_err();
3116        match err {
3117            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
3118            other => panic!("expected InvalidWorkflow, got {other:?}"),
3119        }
3120    }
3121
3122    #[test]
3123    fn engine_register_rejects_whitespace_only_segment_category() {
3124        let mut engine = create_test_engine();
3125        let err = engine
3126            .register(BadCategoryWorkflow("data/ /etl"))
3127            .unwrap_err();
3128        match err {
3129            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
3130            other => panic!("expected InvalidWorkflow, got {other:?}"),
3131        }
3132    }
3133
3134    #[test]
3135    fn engine_register_accepts_valid_nested_category() {
3136        let mut engine = create_test_engine();
3137        assert!(engine.register(CategorizedWorkflow).is_ok());
3138    }
3139
3140    #[tokio::test]
3141    async fn engine_unknown_workflow_returns_error() {
3142        let engine = create_test_engine();
3143        let result = engine
3144            .run_handler("unknown", TriggerKind::Manual, json!({}))
3145            .await;
3146        assert!(result.is_err());
3147        match result {
3148            Err(EngineError::InvalidWorkflow(msg)) => {
3149                assert!(msg.contains("no handler registered"));
3150            }
3151            _ => panic!("expected InvalidWorkflow error"),
3152        }
3153    }
3154
3155    #[tokio::test]
3156    async fn engine_enqueue_handler_creates_pending_run() {
3157        let mut engine = create_test_engine();
3158        engine.register(EchoWorkflow).unwrap();
3159
3160        let run = engine
3161            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
3162            .await
3163            .unwrap();
3164        assert_eq!(run.status.state, RunStatus::Pending);
3165        assert_eq!(run.workflow_name, "echo-workflow");
3166    }
3167
3168    #[tokio::test]
3169    async fn enqueue_handler_leaves_the_run_unattributed() {
3170        let mut engine = create_test_engine();
3171        engine.register(EchoWorkflow).unwrap();
3172
3173        let run = engine
3174            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
3175            .await
3176            .unwrap();
3177
3178        assert!(run.created_by.is_none());
3179    }
3180
3181    #[tokio::test]
3182    async fn enqueue_handler_with_options_records_the_author() {
3183        let mut engine = create_test_engine();
3184        engine.register(EchoWorkflow).unwrap();
3185        let actor = RunActor::User {
3186            user_id: Uuid::now_v7(),
3187        };
3188
3189        let run = engine
3190            .enqueue_handler_with_options(
3191                "echo-workflow",
3192                TriggerKind::Api,
3193                json!({}),
3194                EnqueueOptions {
3195                    created_by: Some(actor.clone()),
3196                    ..Default::default()
3197                },
3198            )
3199            .await
3200            .unwrap()
3201            .into_run();
3202
3203        assert_eq!(run.created_by, Some(actor));
3204    }
3205
3206    #[tokio::test]
3207    async fn enqueue_handler_with_options_accepts_no_author() {
3208        let mut engine = create_test_engine();
3209        engine.register(EchoWorkflow).unwrap();
3210
3211        let run = engine
3212            .enqueue_handler_with_options(
3213                "echo-workflow",
3214                TriggerKind::Cron {
3215                    schedule: "0 * * * * *".to_string(),
3216                    schedule_id: None,
3217                    scheduled_for: None,
3218                },
3219                json!({}),
3220                EnqueueOptions::default(),
3221            )
3222            .await
3223            .unwrap()
3224            .into_run();
3225
3226        assert!(run.created_by.is_none());
3227    }
3228
3229    #[tokio::test]
3230    async fn enqueue_handler_with_options_stores_concurrency_limits() {
3231        let mut engine = create_test_engine();
3232        engine.register(EchoWorkflow).unwrap();
3233        let limits = vec![
3234            ConcurrencyLimit::new("repo:acme", 2),
3235            ConcurrencyLimit::new("tenant:42", 5),
3236        ];
3237
3238        let run = engine
3239            .enqueue_handler_with_options(
3240                "echo-workflow",
3241                TriggerKind::Api,
3242                json!({}),
3243                EnqueueOptions {
3244                    concurrency_limits: limits.clone(),
3245                    ..Default::default()
3246                },
3247            )
3248            .await
3249            .unwrap()
3250            .into_run();
3251
3252        assert_eq!(run.concurrency_limits, limits);
3253    }
3254
3255    #[tokio::test]
3256    async fn enqueue_rejects_invalid_concurrency_limits() {
3257        let mut engine = create_test_engine();
3258        engine.register(EchoWorkflow).unwrap();
3259
3260        let invalid = [
3261            vec![ConcurrencyLimit::new("repo:acme", 0)],
3262            vec![ConcurrencyLimit::new("", 1)],
3263            vec![
3264                ConcurrencyLimit::new("repo:acme", 1),
3265                ConcurrencyLimit::new("repo:acme", 2),
3266            ],
3267        ];
3268        for concurrency_limits in invalid {
3269            let err = engine
3270                .enqueue_handler_with_options(
3271                    "echo-workflow",
3272                    TriggerKind::Api,
3273                    json!({}),
3274                    EnqueueOptions {
3275                        concurrency_limits,
3276                        ..Default::default()
3277                    },
3278                )
3279                .await
3280                .unwrap_err();
3281            assert!(
3282                matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3283                "{err:?}"
3284            );
3285        }
3286
3287        // Validated before the handler lookup.
3288        let err = engine
3289            .enqueue_handler_with_options(
3290                "not-registered",
3291                TriggerKind::Api,
3292                json!({}),
3293                EnqueueOptions {
3294                    concurrency_limits: vec![ConcurrencyLimit::new("repo:acme", 0)],
3295                    ..Default::default()
3296                },
3297            )
3298            .await
3299            .unwrap_err();
3300        assert!(
3301            matches!(err, EngineError::InvalidConcurrencyLimit(_)),
3302            "{err:?}"
3303        );
3304
3305        let page = engine
3306            .store()
3307            .list_runs(RunFilter::default(), 1, 10)
3308            .await
3309            .unwrap();
3310        assert_eq!(page.total, 0, "no run may be created");
3311    }
3312
3313    struct UrgentWorkflow;
3314
3315    impl WorkflowHandler for UrgentWorkflow {
3316        fn name(&self) -> &str {
3317            "urgent-workflow"
3318        }
3319
3320        fn priority(&self) -> i16 {
3321            60
3322        }
3323
3324        fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3325            Box::pin(async { Ok(()) })
3326        }
3327    }
3328
3329    #[tokio::test]
3330    async fn enqueue_priority_defaults_to_the_handler_priority() {
3331        let mut engine = create_test_engine();
3332        engine.register(EchoWorkflow).unwrap();
3333        engine.register(UrgentWorkflow).unwrap();
3334
3335        let echo = engine
3336            .enqueue_handler_with_options(
3337                "echo-workflow",
3338                TriggerKind::Api,
3339                json!({}),
3340                EnqueueOptions::default(),
3341            )
3342            .await
3343            .unwrap()
3344            .into_run();
3345        assert_eq!(echo.priority, 0);
3346
3347        let urgent = engine
3348            .enqueue_handler_with_options(
3349                "urgent-workflow",
3350                TriggerKind::Api,
3351                json!({}),
3352                EnqueueOptions::default(),
3353            )
3354            .await
3355            .unwrap()
3356            .into_run();
3357        assert_eq!(urgent.priority, 60);
3358    }
3359
3360    #[tokio::test]
3361    async fn enqueue_priority_explicit_value_overrides_the_handler() {
3362        let mut engine = create_test_engine();
3363        engine.register(UrgentWorkflow).unwrap();
3364
3365        let run = engine
3366            .enqueue_handler_with_options(
3367                "urgent-workflow",
3368                TriggerKind::Api,
3369                json!({}),
3370                EnqueueOptions {
3371                    priority: Some(-20),
3372                    ..Default::default()
3373                },
3374            )
3375            .await
3376            .unwrap()
3377            .into_run();
3378        assert_eq!(run.priority, -20);
3379
3380        let stored = engine.store().get_run(run.id).await.unwrap().unwrap();
3381        assert_eq!(stored.priority, -20);
3382    }
3383
3384    #[tokio::test]
3385    async fn enqueue_priority_out_of_range_is_rejected() {
3386        let mut engine = create_test_engine();
3387        engine.register(EchoWorkflow).unwrap();
3388
3389        for priority in [MAX_PRIORITY + 1, MIN_PRIORITY - 1] {
3390            let err = engine
3391                .enqueue_handler_with_options(
3392                    "echo-workflow",
3393                    TriggerKind::Api,
3394                    json!({}),
3395                    EnqueueOptions {
3396                        priority: Some(priority),
3397                        ..Default::default()
3398                    },
3399                )
3400                .await
3401                .unwrap_err();
3402            assert!(matches!(err, EngineError::InvalidPriority(_)), "{err:?}");
3403        }
3404
3405        // Validated before the handler lookup.
3406        let err = engine
3407            .enqueue_handler_with_options(
3408                "not-registered",
3409                TriggerKind::Api,
3410                json!({}),
3411                EnqueueOptions {
3412                    priority: Some(MAX_PRIORITY + 1),
3413                    ..Default::default()
3414                },
3415            )
3416            .await
3417            .unwrap_err();
3418        assert!(matches!(err, EngineError::InvalidPriority(_)), "{err:?}");
3419
3420        let page = engine
3421            .store()
3422            .list_runs(RunFilter::default(), 1, 10)
3423            .await
3424            .unwrap();
3425        assert_eq!(page.total, 0, "no run may be created");
3426    }
3427
3428    #[tokio::test]
3429    async fn enqueue_priority_bounds_are_accepted() {
3430        let mut engine = create_test_engine();
3431        engine.register(EchoWorkflow).unwrap();
3432
3433        for priority in [MIN_PRIORITY, MAX_PRIORITY] {
3434            let run = engine
3435                .enqueue_handler_with_options(
3436                    "echo-workflow",
3437                    TriggerKind::Api,
3438                    json!({}),
3439                    EnqueueOptions {
3440                        priority: Some(priority),
3441                        ..Default::default()
3442                    },
3443                )
3444                .await
3445                .unwrap()
3446                .into_run();
3447            assert_eq!(run.priority, priority);
3448        }
3449    }
3450
3451    struct GpuWorkflow;
3452
3453    impl WorkflowHandler for GpuWorkflow {
3454        fn name(&self) -> &str {
3455            "gpu-workflow"
3456        }
3457
3458        fn required_worker_tags(&self) -> Vec<String> {
3459            vec!["gpu".to_string()]
3460        }
3461
3462        fn execute<'a>(&'a self, _ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3463            Box::pin(async move { Ok(()) })
3464        }
3465    }
3466
3467    #[tokio::test]
3468    async fn enqueue_merges_handler_and_request_worker_tags() {
3469        let mut engine = create_test_engine();
3470        engine.register(GpuWorkflow).unwrap();
3471
3472        let run = engine
3473            .enqueue_handler_with_options(
3474                "gpu-workflow",
3475                TriggerKind::Api,
3476                json!({}),
3477                EnqueueOptions {
3478                    worker_tags: vec!["region:eu".to_string(), "gpu".to_string()],
3479                    ..Default::default()
3480                },
3481            )
3482            .await
3483            .unwrap()
3484            .into_run();
3485
3486        assert_eq!(
3487            run.worker_tags,
3488            vec!["gpu".to_string(), "region:eu".to_string()]
3489        );
3490    }
3491
3492    #[tokio::test]
3493    async fn enqueue_without_worker_tags_keeps_handler_tags() {
3494        let mut engine = create_test_engine();
3495        engine.register(GpuWorkflow).unwrap();
3496        engine.register(EchoWorkflow).unwrap();
3497
3498        let gpu = engine
3499            .enqueue_handler("gpu-workflow", TriggerKind::Api, json!({}), 0)
3500            .await
3501            .unwrap();
3502        assert_eq!(gpu.worker_tags, vec!["gpu".to_string()]);
3503
3504        let echo = engine
3505            .enqueue_handler("echo-workflow", TriggerKind::Api, json!({}), 0)
3506            .await
3507            .unwrap();
3508        assert!(echo.worker_tags.is_empty());
3509    }
3510
3511    #[tokio::test]
3512    async fn enqueue_rejects_invalid_worker_tags() {
3513        let mut engine = create_test_engine();
3514        engine.register(EchoWorkflow).unwrap();
3515
3516        for worker_tags in [
3517            vec!["bad,tag".to_string()],
3518            vec![" ".to_string()],
3519            vec!["x".repeat(65)],
3520        ] {
3521            let err = engine
3522                .enqueue_handler_with_options(
3523                    "echo-workflow",
3524                    TriggerKind::Api,
3525                    json!({}),
3526                    EnqueueOptions {
3527                        worker_tags,
3528                        ..Default::default()
3529                    },
3530                )
3531                .await
3532                .unwrap_err();
3533            assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3534        }
3535
3536        // Validated before the handler lookup.
3537        let err = engine
3538            .enqueue_handler_with_options(
3539                "not-registered",
3540                TriggerKind::Api,
3541                json!({}),
3542                EnqueueOptions {
3543                    worker_tags: vec!["bad,tag".to_string()],
3544                    ..Default::default()
3545                },
3546            )
3547            .await
3548            .unwrap_err();
3549        assert!(matches!(err, EngineError::InvalidWorkerTag(_)), "{err:?}");
3550    }
3551
3552    #[test]
3553    fn worker_tags_are_unset_by_default() {
3554        let engine = create_test_engine();
3555        assert!(engine.worker_tags().is_none());
3556    }
3557
3558    #[test]
3559    fn set_worker_tags_stores_the_tags() {
3560        let mut engine = create_test_engine();
3561        engine.set_worker_tags(vec!["gpu".to_string()]);
3562        assert_eq!(engine.worker_tags(), Some(&["gpu".to_string()][..]));
3563
3564        engine.set_worker_tags(Vec::new());
3565        assert_eq!(engine.worker_tags(), Some(&[][..]));
3566    }
3567
3568    #[tokio::test]
3569    async fn run_handler_records_handler_worker_tags() {
3570        let mut engine = create_test_engine();
3571        engine.register(GpuWorkflow).unwrap();
3572
3573        let result = engine
3574            .run_handler("gpu-workflow", TriggerKind::Manual, json!({}))
3575            .await
3576            .unwrap();
3577        assert_eq!(result.run.worker_tags, vec!["gpu".to_string()]);
3578    }
3579
3580    #[tokio::test]
3581    async fn run_handler_leaves_the_run_unattributed() {
3582        let mut engine = create_test_engine();
3583        engine.register(EchoWorkflow).unwrap();
3584
3585        let run = engine
3586            .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
3587            .await
3588            .unwrap()
3589            .run;
3590
3591        assert!(run.created_by.is_none());
3592    }
3593
3594    #[tokio::test]
3595    async fn run_handler_priority_comes_from_the_handler() {
3596        let mut engine = create_test_engine();
3597        engine.register(UrgentWorkflow).unwrap();
3598
3599        let run = engine
3600            .run_handler("urgent-workflow", TriggerKind::Manual, json!({}))
3601            .await
3602            .unwrap()
3603            .run;
3604
3605        assert_eq!(run.priority, 60);
3606    }
3607
3608    #[tokio::test]
3609    async fn engine_register_boxed() {
3610        let mut engine = create_test_engine();
3611        let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
3612        let result = engine.register_boxed(handler);
3613        assert!(result.is_ok());
3614        assert_eq!(engine.handler_names().len(), 1);
3615    }
3616
3617    #[tokio::test]
3618    async fn engine_store_and_provider_accessors() {
3619        let store = Arc::new(InMemoryStore::new());
3620        let inner = ClaudeCodeProvider::new();
3621        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
3622            inner,
3623            "/tmp/ironflow-fixtures",
3624        ));
3625        let engine = Engine::new(store.clone(), provider.clone());
3626
3627        // Verify accessors return references
3628        let _ = engine.store();
3629        let _ = engine.provider();
3630    }
3631
3632    // -----------------------------------------------------------------------
3633    // Operation trait tests
3634    // -----------------------------------------------------------------------
3635
3636    use crate::operation::{Operation, OperationContext};
3637    use async_trait::async_trait;
3638    use ironflow_core::error::OperationError;
3639    use ironflow_store::models::StepKind;
3640
3641    struct FakeGitlabOp {
3642        project_id: u64,
3643        title: String,
3644    }
3645
3646    #[async_trait]
3647    impl Operation for FakeGitlabOp {
3648        fn kind(&self) -> &str {
3649            "gitlab"
3650        }
3651
3652        async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3653            Ok(json!({
3654                "issue_id": 42,
3655                "project_id": self.project_id,
3656                "title": self.title,
3657            }))
3658        }
3659
3660        fn input(&self) -> Option<Value> {
3661            Some(json!({
3662                "project_id": self.project_id,
3663                "title": self.title,
3664            }))
3665        }
3666    }
3667
3668    struct FailingOp;
3669
3670    #[async_trait]
3671    impl Operation for FailingOp {
3672        fn kind(&self) -> &str {
3673            "broken-service"
3674        }
3675
3676        async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
3677            Err(OperationError::Http {
3678                status: None,
3679                message: "service unavailable".to_string(),
3680            })
3681        }
3682    }
3683
3684    struct OperationWorkflow;
3685
3686    impl WorkflowHandler for OperationWorkflow {
3687        fn name(&self) -> &str {
3688            "operation-workflow"
3689        }
3690
3691        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3692            Box::pin(async move {
3693                let op = FakeGitlabOp {
3694                    project_id: 123,
3695                    title: "Bug report".to_string(),
3696                };
3697                ctx.operation("create-issue", &op).await?;
3698                Ok(())
3699            })
3700        }
3701    }
3702
3703    struct FailingOperationWorkflow;
3704
3705    impl WorkflowHandler for FailingOperationWorkflow {
3706        fn name(&self) -> &str {
3707            "failing-operation-workflow"
3708        }
3709
3710        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3711            Box::pin(async move {
3712                ctx.operation("broken-call", &FailingOp).await?;
3713                Ok(())
3714            })
3715        }
3716    }
3717
3718    struct MixedWorkflow;
3719
3720    impl WorkflowHandler for MixedWorkflow {
3721        fn name(&self) -> &str {
3722            "mixed-workflow"
3723        }
3724
3725        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3726            Box::pin(async move {
3727                ctx.shell("build", ShellConfig::new("echo built")).await?;
3728                let op = FakeGitlabOp {
3729                    project_id: 456,
3730                    title: "Deploy done".to_string(),
3731                };
3732                let result = ctx.operation("notify-gitlab", &op).await?;
3733                assert_eq!(result.output["issue_id"], 42);
3734                Ok(())
3735            })
3736        }
3737    }
3738
3739    #[tokio::test]
3740    async fn operation_step_happy_path() {
3741        let mut engine = create_test_engine();
3742        engine.register(OperationWorkflow).unwrap();
3743
3744        let run = engine
3745            .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
3746            .await
3747            .unwrap()
3748            .run;
3749
3750        assert_eq!(run.status.state, RunStatus::Completed);
3751
3752        let steps = engine.store().list_steps(run.id).await.unwrap();
3753
3754        assert_eq!(steps.len(), 1);
3755        assert_eq!(steps[0].name, "create-issue");
3756        assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
3757        assert_eq!(
3758            steps[0].status.state,
3759            ironflow_store::models::StepStatus::Completed
3760        );
3761
3762        let output = steps[0].output.as_ref().unwrap();
3763        assert_eq!(output["issue_id"], 42);
3764        assert_eq!(output["project_id"], 123);
3765
3766        let input = steps[0].input.as_ref().unwrap();
3767        assert_eq!(input["project_id"], 123);
3768        assert_eq!(input["title"], "Bug report");
3769    }
3770
3771    #[tokio::test]
3772    async fn operation_step_failure_marks_run_failed() {
3773        let mut engine = create_test_engine();
3774        engine.register(FailingOperationWorkflow).unwrap();
3775
3776        let result = engine
3777            .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
3778            .await;
3779
3780        assert!(result.is_err());
3781    }
3782
3783    #[tokio::test]
3784    async fn operation_mixed_with_shell_steps() {
3785        let mut engine = create_test_engine();
3786        engine.register(MixedWorkflow).unwrap();
3787
3788        let run = engine
3789            .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
3790            .await
3791            .unwrap()
3792            .run;
3793
3794        assert_eq!(run.status.state, RunStatus::Completed);
3795
3796        let steps = engine.store().list_steps(run.id).await.unwrap();
3797
3798        assert_eq!(steps.len(), 2);
3799        assert_eq!(steps[0].kind, StepKind::Shell);
3800        assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
3801        assert_eq!(steps[0].position, 0);
3802        assert_eq!(steps[1].position, 1);
3803    }
3804
3805    // -----------------------------------------------------------------------
3806    // Approval + resume tests
3807    // -----------------------------------------------------------------------
3808
3809    use crate::config::ApprovalConfig;
3810
3811    struct SingleApprovalWorkflow;
3812
3813    impl WorkflowHandler for SingleApprovalWorkflow {
3814        fn name(&self) -> &str {
3815            "single-approval"
3816        }
3817
3818        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3819            Box::pin(async move {
3820                ctx.shell("build", ShellConfig::new("echo built")).await?;
3821                ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
3822                ctx.shell("deploy", ShellConfig::new("echo deployed"))
3823                    .await?;
3824                Ok(())
3825            })
3826        }
3827    }
3828
3829    struct DoubleApprovalWorkflow;
3830
3831    impl WorkflowHandler for DoubleApprovalWorkflow {
3832        fn name(&self) -> &str {
3833            "double-approval"
3834        }
3835
3836        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
3837            Box::pin(async move {
3838                ctx.shell("build", ShellConfig::new("echo built")).await?;
3839                ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
3840                    .await?;
3841                ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
3842                    .await?;
3843                ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
3844                    .await?;
3845                ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
3846                    .await?;
3847                Ok(())
3848            })
3849        }
3850    }
3851
3852    #[tokio::test]
3853    async fn approval_pauses_run() {
3854        let mut engine = create_test_engine();
3855        engine.register(SingleApprovalWorkflow).unwrap();
3856
3857        let run = engine
3858            .run_handler("single-approval", TriggerKind::Manual, json!({}))
3859            .await
3860            .unwrap()
3861            .run;
3862
3863        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3864
3865        let steps = engine.store().list_steps(run.id).await.unwrap();
3866        assert_eq!(steps.len(), 2); // build + approval gate
3867        assert_eq!(steps[0].kind, StepKind::Shell);
3868        assert_eq!(steps[0].status.state, StepStatus::Completed);
3869        assert_eq!(steps[1].kind, StepKind::Approval);
3870        assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
3871    }
3872
3873    #[tokio::test]
3874    async fn approval_resume_completes_run() {
3875        let mut engine = create_test_engine();
3876        engine.register(SingleApprovalWorkflow).unwrap();
3877
3878        // First execution: pauses at approval
3879        let run = engine
3880            .run_handler("single-approval", TriggerKind::Manual, json!({}))
3881            .await
3882            .unwrap()
3883            .run;
3884        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3885
3886        // Simulate approval: transition to Running
3887        engine
3888            .store()
3889            .update_run_status(run.id, RunStatus::Running)
3890            .await
3891            .unwrap();
3892
3893        // Resume: replays build, skips approval, executes deploy
3894        let resumed = engine.resume_run(run.id).await.unwrap().run;
3895        assert_eq!(resumed.status.state, RunStatus::Completed);
3896
3897        let steps = engine.store().list_steps(run.id).await.unwrap();
3898        assert_eq!(steps.len(), 3); // build + approval + deploy
3899        assert_eq!(steps[0].name, "build");
3900        assert_eq!(steps[0].status.state, StepStatus::Completed);
3901        assert_eq!(steps[1].name, "gate");
3902        assert_eq!(steps[1].kind, StepKind::Approval);
3903        assert_eq!(steps[1].status.state, StepStatus::Completed);
3904        assert_eq!(steps[2].name, "deploy");
3905        assert_eq!(steps[2].status.state, StepStatus::Completed);
3906    }
3907
3908    #[tokio::test]
3909    async fn double_approval_two_resumes() {
3910        let mut engine = create_test_engine();
3911        engine.register(DoubleApprovalWorkflow).unwrap();
3912
3913        // First execution: pauses at staging-gate
3914        let run = engine
3915            .run_handler("double-approval", TriggerKind::Manual, json!({}))
3916            .await
3917            .unwrap()
3918            .run;
3919        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
3920
3921        let steps = engine.store().list_steps(run.id).await.unwrap();
3922        assert_eq!(steps.len(), 2); // build + staging-gate
3923
3924        // First approval
3925        engine
3926            .store()
3927            .update_run_status(run.id, RunStatus::Running)
3928            .await
3929            .unwrap();
3930
3931        let resumed = engine.resume_run(run.id).await.unwrap().run;
3932        assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
3933
3934        let steps = engine.store().list_steps(run.id).await.unwrap();
3935        assert_eq!(steps.len(), 4); // build + staging-gate + deploy-staging + prod-gate
3936
3937        // Second approval
3938        engine
3939            .store()
3940            .update_run_status(run.id, RunStatus::Running)
3941            .await
3942            .unwrap();
3943
3944        let final_run = engine.resume_run(run.id).await.unwrap().run;
3945        assert_eq!(final_run.status.state, RunStatus::Completed);
3946
3947        let steps = engine.store().list_steps(run.id).await.unwrap();
3948        assert_eq!(steps.len(), 5);
3949        assert_eq!(steps[0].name, "build");
3950        assert_eq!(steps[1].name, "staging-gate");
3951        assert_eq!(steps[2].name, "deploy-staging");
3952        assert_eq!(steps[3].name, "prod-gate");
3953        assert_eq!(steps[4].name, "deploy-prod");
3954
3955        for step in &steps {
3956            assert_eq!(step.status.state, StepStatus::Completed);
3957        }
3958    }
3959
3960    // -----------------------------------------------------------------------
3961    // fail_orphaned_steps tests
3962    // -----------------------------------------------------------------------
3963
3964    use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
3965
3966    async fn create_step_with_status(
3967        store: &Arc<dyn Store>,
3968        run_id: Uuid,
3969        name: &str,
3970        position: u32,
3971        status: StepStatus,
3972    ) -> ironflow_store::models::Step {
3973        let step = store
3974            .create_step(NewStep {
3975                run_id,
3976                trace_id: step_trace_id(run_id, name, position),
3977                name: name.to_string(),
3978                kind: StepKind::Shell,
3979                position,
3980                input: None,
3981                is_error_handler: false,
3982            })
3983            .await
3984            .unwrap();
3985
3986        match status {
3987            StepStatus::Pending => {}
3988            StepStatus::Running => {
3989                store
3990                    .update_step(
3991                        step.id,
3992                        StepUpdate {
3993                            status: Some(StepStatus::Running),
3994                            ..StepUpdate::default()
3995                        },
3996                    )
3997                    .await
3998                    .unwrap();
3999            }
4000            StepStatus::Completed => {
4001                store
4002                    .update_step(
4003                        step.id,
4004                        StepUpdate {
4005                            status: Some(StepStatus::Running),
4006                            ..StepUpdate::default()
4007                        },
4008                    )
4009                    .await
4010                    .unwrap();
4011                store
4012                    .update_step(
4013                        step.id,
4014                        StepUpdate {
4015                            status: Some(StepStatus::Completed),
4016                            ..StepUpdate::default()
4017                        },
4018                    )
4019                    .await
4020                    .unwrap();
4021            }
4022            StepStatus::AwaitingApproval => {
4023                store
4024                    .update_step(
4025                        step.id,
4026                        StepUpdate {
4027                            status: Some(StepStatus::Running),
4028                            ..StepUpdate::default()
4029                        },
4030                    )
4031                    .await
4032                    .unwrap();
4033                store
4034                    .update_step(
4035                        step.id,
4036                        StepUpdate {
4037                            status: Some(StepStatus::AwaitingApproval),
4038                            ..StepUpdate::default()
4039                        },
4040                    )
4041                    .await
4042                    .unwrap();
4043            }
4044            _ => panic!("unsupported status for test helper: {status}"),
4045        }
4046
4047        store.get_step(step.id).await.unwrap().unwrap()
4048    }
4049
4050    #[tokio::test]
4051    async fn fail_orphaned_steps_marks_running_as_failed() {
4052        let engine = create_test_engine();
4053        let run = engine
4054            .store()
4055            .create_run(NewRun {
4056                created_by: None,
4057                workflow_name: "test".to_string(),
4058                trigger: TriggerKind::Manual,
4059                payload: json!({}),
4060                max_retries: 0,
4061                handler_version: None,
4062                labels: HashMap::new(),
4063                scheduled_at: None,
4064                idempotency_key: None,
4065                concurrency_key: None,
4066                priority: 0,
4067                concurrency_limits: Vec::new(),
4068                max_cost_usd: None,
4069                worker_tags: Vec::new(),
4070            })
4071            .await
4072            .unwrap()
4073            .into_run();
4074
4075        let step = create_step_with_status(
4076            engine.store(),
4077            run.id,
4078            "running-step",
4079            0,
4080            StepStatus::Running,
4081        )
4082        .await;
4083
4084        engine
4085            .fail_orphaned_steps(run.id, "parent run timed out")
4086            .await
4087            .unwrap();
4088
4089        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
4090        assert_eq!(updated.status.state, StepStatus::Failed);
4091        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
4092        assert!(updated.completed_at.is_some());
4093    }
4094
4095    #[tokio::test]
4096    async fn fail_orphaned_steps_marks_pending_as_skipped() {
4097        let engine = create_test_engine();
4098        let run = engine
4099            .store()
4100            .create_run(NewRun {
4101                created_by: None,
4102                workflow_name: "test".to_string(),
4103                trigger: TriggerKind::Manual,
4104                payload: json!({}),
4105                max_retries: 0,
4106                handler_version: None,
4107                labels: HashMap::new(),
4108                scheduled_at: None,
4109                idempotency_key: None,
4110                concurrency_key: None,
4111                priority: 0,
4112                concurrency_limits: Vec::new(),
4113                max_cost_usd: None,
4114                worker_tags: Vec::new(),
4115            })
4116            .await
4117            .unwrap()
4118            .into_run();
4119
4120        let step = create_step_with_status(
4121            engine.store(),
4122            run.id,
4123            "pending-step",
4124            0,
4125            StepStatus::Pending,
4126        )
4127        .await;
4128
4129        engine
4130            .fail_orphaned_steps(run.id, "parent run timed out")
4131            .await
4132            .unwrap();
4133
4134        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
4135        assert_eq!(updated.status.state, StepStatus::Skipped);
4136        assert!(updated.error.is_none());
4137        assert!(updated.completed_at.is_some());
4138    }
4139
4140    #[tokio::test]
4141    async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
4142        let engine = create_test_engine();
4143        let run = engine
4144            .store()
4145            .create_run(NewRun {
4146                created_by: None,
4147                workflow_name: "test".to_string(),
4148                trigger: TriggerKind::Manual,
4149                payload: json!({}),
4150                max_retries: 0,
4151                handler_version: None,
4152                labels: HashMap::new(),
4153                scheduled_at: None,
4154                idempotency_key: None,
4155                concurrency_key: None,
4156                priority: 0,
4157                concurrency_limits: Vec::new(),
4158                max_cost_usd: None,
4159                worker_tags: Vec::new(),
4160            })
4161            .await
4162            .unwrap()
4163            .into_run();
4164
4165        let step = create_step_with_status(
4166            engine.store(),
4167            run.id,
4168            "approval-step",
4169            0,
4170            StepStatus::AwaitingApproval,
4171        )
4172        .await;
4173
4174        engine
4175            .fail_orphaned_steps(run.id, "parent run timed out")
4176            .await
4177            .unwrap();
4178
4179        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
4180        assert_eq!(updated.status.state, StepStatus::Failed);
4181        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
4182        assert!(updated.completed_at.is_some());
4183    }
4184
4185    #[tokio::test]
4186    async fn fail_orphaned_steps_skips_terminal_steps() {
4187        let engine = create_test_engine();
4188        let run = engine
4189            .store()
4190            .create_run(NewRun {
4191                created_by: None,
4192                workflow_name: "test".to_string(),
4193                trigger: TriggerKind::Manual,
4194                payload: json!({}),
4195                max_retries: 0,
4196                handler_version: None,
4197                labels: HashMap::new(),
4198                scheduled_at: None,
4199                idempotency_key: None,
4200                concurrency_key: None,
4201                priority: 0,
4202                concurrency_limits: Vec::new(),
4203                max_cost_usd: None,
4204                worker_tags: Vec::new(),
4205            })
4206            .await
4207            .unwrap()
4208            .into_run();
4209
4210        let completed_step =
4211            create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
4212        let running_step =
4213            create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
4214                .await;
4215
4216        engine
4217            .fail_orphaned_steps(run.id, "parent run timed out")
4218            .await
4219            .unwrap();
4220
4221        let completed = engine
4222            .store()
4223            .get_step(completed_step.id)
4224            .await
4225            .unwrap()
4226            .unwrap();
4227        assert_eq!(completed.status.state, StepStatus::Completed);
4228
4229        let failed = engine
4230            .store()
4231            .get_step(running_step.id)
4232            .await
4233            .unwrap()
4234            .unwrap();
4235        assert_eq!(failed.status.state, StepStatus::Failed);
4236    }
4237
4238    #[tokio::test]
4239    async fn fail_orphaned_steps_mixed_states() {
4240        let engine = create_test_engine();
4241        let run = engine
4242            .store()
4243            .create_run(NewRun {
4244                created_by: None,
4245                workflow_name: "test".to_string(),
4246                trigger: TriggerKind::Manual,
4247                payload: json!({}),
4248                max_retries: 0,
4249                handler_version: None,
4250                labels: HashMap::new(),
4251                scheduled_at: None,
4252                idempotency_key: None,
4253                concurrency_key: None,
4254                priority: 0,
4255                concurrency_limits: Vec::new(),
4256                max_cost_usd: None,
4257                worker_tags: Vec::new(),
4258            })
4259            .await
4260            .unwrap()
4261            .into_run();
4262
4263        let s_completed =
4264            create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
4265                .await;
4266        let s_running =
4267            create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
4268        let s_pending =
4269            create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
4270
4271        engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
4272
4273        let r_completed = engine
4274            .store()
4275            .get_step(s_completed.id)
4276            .await
4277            .unwrap()
4278            .unwrap();
4279        assert_eq!(r_completed.status.state, StepStatus::Completed);
4280
4281        let r_running = engine
4282            .store()
4283            .get_step(s_running.id)
4284            .await
4285            .unwrap()
4286            .unwrap();
4287        assert_eq!(r_running.status.state, StepStatus::Failed);
4288        assert_eq!(r_running.error.as_deref(), Some("timeout"));
4289
4290        let r_pending = engine
4291            .store()
4292            .get_step(s_pending.id)
4293            .await
4294            .unwrap()
4295            .unwrap();
4296        assert_eq!(r_pending.status.state, StepStatus::Skipped);
4297        assert!(r_pending.error.is_none());
4298    }
4299
4300    #[tokio::test]
4301    async fn fail_orphaned_steps_no_steps_is_noop() {
4302        let engine = create_test_engine();
4303        let run = engine
4304            .store()
4305            .create_run(NewRun {
4306                created_by: None,
4307                workflow_name: "test".to_string(),
4308                trigger: TriggerKind::Manual,
4309                payload: json!({}),
4310                max_retries: 0,
4311                handler_version: None,
4312                labels: HashMap::new(),
4313                scheduled_at: None,
4314                idempotency_key: None,
4315                concurrency_key: None,
4316                priority: 0,
4317                concurrency_limits: Vec::new(),
4318                max_cost_usd: None,
4319                worker_tags: Vec::new(),
4320            })
4321            .await
4322            .unwrap()
4323            .into_run();
4324
4325        let result = engine.fail_orphaned_steps(run.id, "timeout").await;
4326        assert!(result.is_ok());
4327    }
4328
4329    #[tokio::test]
4330    async fn fail_orphaned_steps_preserves_existing_error() {
4331        let engine = create_test_engine();
4332        let run = engine
4333            .store()
4334            .create_run(NewRun {
4335                created_by: None,
4336                workflow_name: "test".to_string(),
4337                trigger: TriggerKind::Manual,
4338                payload: json!({}),
4339                max_retries: 0,
4340                handler_version: None,
4341                labels: HashMap::new(),
4342                scheduled_at: None,
4343                idempotency_key: None,
4344                concurrency_key: None,
4345                priority: 0,
4346                concurrency_limits: Vec::new(),
4347                max_cost_usd: None,
4348                worker_tags: Vec::new(),
4349            })
4350            .await
4351            .unwrap()
4352            .into_run();
4353
4354        let step_with_error = create_step_with_status(
4355            engine.store(),
4356            run.id,
4357            "already-errored",
4358            0,
4359            StepStatus::Running,
4360        )
4361        .await;
4362
4363        engine
4364            .store()
4365            .update_step(
4366                step_with_error.id,
4367                StepUpdate {
4368                    error: Some("real error from provider".to_string()),
4369                    ..StepUpdate::default()
4370                },
4371            )
4372            .await
4373            .unwrap();
4374
4375        let step_no_error = create_step_with_status(
4376            engine.store(),
4377            run.id,
4378            "no-error-yet",
4379            1,
4380            StepStatus::Running,
4381        )
4382        .await;
4383
4384        engine
4385            .fail_orphaned_steps(run.id, "parent run failed")
4386            .await
4387            .unwrap();
4388
4389        let updated_with = engine
4390            .store()
4391            .get_step(step_with_error.id)
4392            .await
4393            .unwrap()
4394            .unwrap();
4395        assert_eq!(updated_with.status.state, StepStatus::Failed);
4396        assert_eq!(
4397            updated_with.error.as_deref(),
4398            Some("real error from provider"),
4399        );
4400
4401        let updated_without = engine
4402            .store()
4403            .get_step(step_no_error.id)
4404            .await
4405            .unwrap()
4406            .unwrap();
4407        assert_eq!(updated_without.status.state, StepStatus::Failed);
4408        assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
4409    }
4410}