Skip to main content

ironflow_engine/
engine.rs

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