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