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