Skip to main content

ironflow_engine/
engine.rs

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