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    ApprovalRequestedEvent, Event, EventPublisher, EventSubscriber, RunBudgetExceededEvent,
44    RunFailedEvent, 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                // The requirement was recorded when the gate opened.
1401                let requirement = self
1402                    .store
1403                    .get_step(step_id)
1404                    .await?
1405                    .and_then(|s| s.approval_requirement);
1406                self.event_publisher
1407                    .publish(Event::ApprovalRequested(ApprovalRequestedEvent {
1408                        run_id: approval_run_id,
1409                        step_id,
1410                        message: message.clone(),
1411                        requirement,
1412                        at: Utc::now(),
1413                    }));
1414            }
1415            Err(EngineError::DelaySleeping {
1416                run_id: delay_run_id,
1417                step_id,
1418                wake_at,
1419            }) => {
1420                final_status = RunStatus::Sleeping;
1421                final_run = self
1422                    .store
1423                    .update_run_returning(
1424                        run_id,
1425                        RunUpdate {
1426                            status: Some(RunStatus::Sleeping),
1427                            cost_usd: Some(ctx.total_cost_usd()),
1428                            duration_ms: Some(total_duration),
1429                            scheduled_at: Some(wake_at),
1430                            ..RunUpdate::default()
1431                        },
1432                    )
1433                    .await?;
1434
1435                info!(
1436                    run_id = %delay_run_id,
1437                    step_id = %step_id,
1438                    wake_at = %wake_at,
1439                    "run sleeping until delay elapses"
1440                );
1441            }
1442            Err(err) => {
1443                // A guardrail stop (budget or workflow guard) is deliberate,
1444                // not a breakage: the run is cancelled, never failed and
1445                // never replayed.
1446                let guardrail_stop = matches!(
1447                    err,
1448                    EngineError::RunBudgetExceeded { .. } | EngineError::WorkflowGuardRejected(_)
1449                );
1450
1451                final_status = if guardrail_stop {
1452                    if let Err(store_err) = self
1453                        .store
1454                        .update_run(
1455                            run_id,
1456                            RunUpdate {
1457                                status: Some(RunStatus::Cancelled),
1458                                error: Some(err.to_string()),
1459                                cost_usd: Some(ctx.total_cost_usd()),
1460                                duration_ms: Some(total_duration),
1461                                completed_at: Some(completed_at),
1462                                ..RunUpdate::default()
1463                            },
1464                        )
1465                        .await
1466                    {
1467                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run cancellation");
1468                    }
1469                    if let Err(cleanup_err) = self
1470                        .fail_orphaned_steps(run_id, "run stopped: guardrail limit reached")
1471                        .await
1472                    {
1473                        error!(run_id = %run_id, store_error = %cleanup_err, "failed to cleanup orphaned steps");
1474                    }
1475                    RunStatus::Cancelled
1476                } else {
1477                    self.fail_or_schedule_retry(
1478                        run_id,
1479                        &err.to_string(),
1480                        is_run_retryable(&err),
1481                        Some(ctx.total_cost_usd()),
1482                        Some(total_duration),
1483                    )
1484                    .await
1485                    .unwrap_or_else(|store_err| {
1486                        error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
1487                        RunStatus::Failed
1488                    })
1489                };
1490
1491                if matches!(err, EngineError::RunBudgetExceeded { .. }) {
1492                    self.on_run_budget_exceeded(workflow_name, run_id, &err);
1493                }
1494
1495                error!(run_id = %run_id, status = %final_status, error = %err, "run stopped");
1496
1497                self.publish_run_status_changed(
1498                    workflow_name,
1499                    run_id,
1500                    final_status,
1501                    Some(err.to_string()),
1502                    ctx,
1503                    total_duration,
1504                    run_labels,
1505                );
1506
1507                #[cfg(feature = "prometheus")]
1508                self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1509
1510                return Err(err);
1511            }
1512        }
1513
1514        self.publish_run_status_changed(
1515            workflow_name,
1516            run_id,
1517            final_status,
1518            None,
1519            ctx,
1520            total_duration,
1521            run_labels,
1522        );
1523
1524        #[cfg(feature = "prometheus")]
1525        self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
1526
1527        Ok(WorkflowResult {
1528            run: final_run,
1529            steps: ctx.step_results().to_vec(),
1530        })
1531    }
1532
1533    /// Emit Prometheus metrics for a completed run.
1534    #[cfg(feature = "prometheus")]
1535    fn emit_run_metrics(
1536        &self,
1537        workflow_name: &str,
1538        status: RunStatus,
1539        duration_ms: u64,
1540        ctx: &WorkflowContext,
1541    ) {
1542        let status_str = status.to_string();
1543        let wf = workflow_name.to_string();
1544
1545        counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
1546        histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
1547            .record(duration_ms as f64 / 1000.0);
1548        histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
1549            ctx.total_cost_usd()
1550                .to_string()
1551                .parse::<f64>()
1552                .unwrap_or(0.0),
1553        );
1554        gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
1555    }
1556
1557    /// Record the metric and publish the audit event for a run that hit its
1558    /// cost cap.
1559    ///
1560    /// A non-[`RunBudgetExceeded`](EngineError::RunBudgetExceeded) error is
1561    /// ignored, so callers can pass the error unconditionally.
1562    fn on_run_budget_exceeded(&self, workflow_name: &str, run_id: Uuid, err: &EngineError) {
1563        let EngineError::RunBudgetExceeded {
1564            limit_usd,
1565            spent_usd,
1566            step_budget_usd,
1567            ..
1568        } = err
1569        else {
1570            return;
1571        };
1572
1573        #[cfg(feature = "prometheus")]
1574        counter!(
1575            RUN_BUDGET_EXCEEDED_TOTAL,
1576            "workflow" => workflow_name.to_string(),
1577            "scope" => "run",
1578        )
1579        .increment(1);
1580
1581        self.event_publisher
1582            .publish(Event::RunBudgetExceeded(RunBudgetExceededEvent {
1583                run_id,
1584                workflow_name: workflow_name.to_string(),
1585                limit_usd: *limit_usd,
1586                spent_usd: *spent_usd,
1587                step_budget_usd: *step_budget_usd,
1588                at: Utc::now(),
1589            }));
1590    }
1591
1592    /// Publish a run status changed event to all registered subscribers.
1593    ///
1594    /// `from` is always `Running` because `finalize_run` is only called
1595    /// from a running state.
1596    #[allow(clippy::too_many_arguments)]
1597    fn publish_run_status_changed(
1598        &self,
1599        workflow_name: &str,
1600        run_id: Uuid,
1601        to: RunStatus,
1602        error: Option<String>,
1603        ctx: &WorkflowContext,
1604        duration_ms: u64,
1605        labels: HashMap<String, String>,
1606    ) {
1607        let now = Utc::now();
1608        let cost_usd = ctx.total_cost_usd();
1609        let wf = workflow_name.to_string();
1610
1611        self.event_publisher
1612            .publish(Event::RunStatusChanged(RunStatusChangedEvent {
1613                run_id,
1614                workflow_name: wf.clone(),
1615                from: RunStatus::Running,
1616                to,
1617                error: error.clone(),
1618                cost_usd,
1619                duration_ms,
1620                labels: labels.clone(),
1621                at: now,
1622            }));
1623
1624        if to == RunStatus::Failed {
1625            self.event_publisher
1626                .publish(Event::RunFailed(RunFailedEvent {
1627                    run_id,
1628                    workflow_name: wf,
1629                    error,
1630                    cost_usd,
1631                    duration_ms,
1632                    labels,
1633                    at: now,
1634                }));
1635        }
1636    }
1637}
1638
1639impl fmt::Debug for Engine {
1640    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
1641        f.debug_struct("Engine")
1642            .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
1643            .finish_non_exhaustive()
1644    }
1645}
1646
1647#[cfg(test)]
1648mod tests {
1649    use super::*;
1650    use crate::config::ShellConfig;
1651    use crate::handler::{HandlerFuture, WorkflowHandler};
1652    use ironflow_core::providers::claude::ClaudeCodeProvider;
1653    use ironflow_core::providers::record_replay::RecordReplayProvider;
1654    use ironflow_store::memory::InMemoryStore;
1655    use ironflow_store::models::StepStatus;
1656    use serde_json::json;
1657
1658    // Test handler that echoes a message via shell
1659    struct EchoWorkflow;
1660
1661    impl WorkflowHandler for EchoWorkflow {
1662        fn name(&self) -> &str {
1663            "echo-workflow"
1664        }
1665
1666        fn describe(&self) -> WorkflowInfo {
1667            WorkflowInfo {
1668                description: "A simple workflow that echoes hello".to_string(),
1669                source_code: None,
1670                sub_workflows: Vec::new(),
1671                category: None,
1672                version: self.version().map(str::to_string),
1673                compatible_versions: Vec::new(),
1674                input_schema: None,
1675                default_labels: HashMap::new(),
1676                schedule: self.schedule().cloned(),
1677                default_max_cost_usd: self.default_max_cost_usd(),
1678            }
1679        }
1680
1681        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1682            Box::pin(async move {
1683                ctx.shell("greet", ShellConfig::new("echo hello")).await?;
1684                Ok(())
1685            })
1686        }
1687    }
1688
1689    // Test handler that fails
1690    struct FailingWorkflow;
1691
1692    impl WorkflowHandler for FailingWorkflow {
1693        fn name(&self) -> &str {
1694            "failing-workflow"
1695        }
1696
1697        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1698            Box::pin(async move {
1699                ctx.shell("fail", ShellConfig::new("exit 1")).await?;
1700                Ok(())
1701            })
1702        }
1703    }
1704
1705    fn create_test_engine() -> Engine {
1706        let store = Arc::new(InMemoryStore::new());
1707        let inner = ClaudeCodeProvider::new();
1708        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
1709            inner,
1710            "/tmp/ironflow-fixtures",
1711        ));
1712        Engine::new(store, provider)
1713    }
1714
1715    #[test]
1716    fn engine_new_creates_instance() {
1717        let engine = create_test_engine();
1718        assert_eq!(engine.handler_names().len(), 0);
1719    }
1720
1721    #[test]
1722    fn engine_register_handler() {
1723        let mut engine = create_test_engine();
1724        let result = engine.register(EchoWorkflow);
1725        assert!(result.is_ok());
1726        assert_eq!(engine.handler_names().len(), 1);
1727        assert!(engine.handler_names().contains(&"echo-workflow"));
1728    }
1729
1730    #[test]
1731    fn engine_register_duplicate_returns_error() {
1732        let mut engine = create_test_engine();
1733        engine.register(EchoWorkflow).unwrap();
1734        let result = engine.register(EchoWorkflow);
1735        assert!(result.is_err());
1736    }
1737
1738    #[test]
1739    fn engine_get_handler_found() {
1740        let mut engine = create_test_engine();
1741        engine.register(EchoWorkflow).unwrap();
1742        let handler = engine.get_handler("echo-workflow");
1743        assert!(handler.is_some());
1744    }
1745
1746    #[test]
1747    fn engine_get_handler_not_found() {
1748        let engine = create_test_engine();
1749        let handler = engine.get_handler("nonexistent");
1750        assert!(handler.is_none());
1751    }
1752
1753    #[test]
1754    fn engine_handler_names_lists_all() {
1755        let mut engine = create_test_engine();
1756        engine.register(EchoWorkflow).unwrap();
1757        engine.register(FailingWorkflow).unwrap();
1758        let names = engine.handler_names();
1759        assert_eq!(names.len(), 2);
1760        assert!(names.contains(&"echo-workflow"));
1761        assert!(names.contains(&"failing-workflow"));
1762    }
1763
1764    #[test]
1765    fn engine_handler_info_returns_description() {
1766        let mut engine = create_test_engine();
1767        engine.register(EchoWorkflow).unwrap();
1768        let info = engine.handler_info("echo-workflow");
1769        assert!(info.is_some());
1770        let info = info.unwrap();
1771        assert_eq!(info.description, "A simple workflow that echoes hello");
1772    }
1773
1774    struct CategorizedWorkflow;
1775
1776    impl WorkflowHandler for CategorizedWorkflow {
1777        fn name(&self) -> &str {
1778            "categorized"
1779        }
1780        fn category(&self) -> Option<&str> {
1781            Some("data/etl")
1782        }
1783        fn execute<'a>(
1784            &'a self,
1785            _ctx: &'a mut WorkflowContext,
1786        ) -> crate::handler::HandlerFuture<'a> {
1787            Box::pin(async move { Ok(()) })
1788        }
1789    }
1790
1791    #[test]
1792    fn engine_default_describe_propagates_category() {
1793        let mut engine = create_test_engine();
1794        engine.register(CategorizedWorkflow).unwrap();
1795        let info = engine.handler_info("categorized").unwrap();
1796        assert_eq!(info.category.as_deref(), Some("data/etl"));
1797    }
1798
1799    #[test]
1800    fn engine_default_describe_without_category() {
1801        let mut engine = create_test_engine();
1802        engine.register(EchoWorkflow).unwrap();
1803        let info = engine.handler_info("echo-workflow").unwrap();
1804        assert!(info.category.is_none());
1805    }
1806
1807    // -----------------------------------------------------------------------
1808    // Schedule tests
1809    // -----------------------------------------------------------------------
1810
1811    struct ScheduledWorkflow {
1812        schedule: CronSchedule,
1813    }
1814
1815    impl ScheduledWorkflow {
1816        fn new() -> Self {
1817            Self {
1818                schedule: CronSchedule::new("0 0 * * * *").unwrap(),
1819            }
1820        }
1821    }
1822
1823    impl WorkflowHandler for ScheduledWorkflow {
1824        fn name(&self) -> &str {
1825            "scheduled"
1826        }
1827        fn schedule(&self) -> Option<&CronSchedule> {
1828            Some(&self.schedule)
1829        }
1830        fn execute<'a>(
1831            &'a self,
1832            _ctx: &'a mut WorkflowContext,
1833        ) -> crate::handler::HandlerFuture<'a> {
1834            Box::pin(async move { Ok(()) })
1835        }
1836    }
1837
1838    #[test]
1839    fn engine_default_describe_propagates_schedule() {
1840        let mut engine = create_test_engine();
1841        engine.register(ScheduledWorkflow::new()).unwrap();
1842        let info = engine.handler_info("scheduled").unwrap();
1843        assert_eq!(
1844            info.schedule.as_ref().map(|s| s.as_str()),
1845            Some("0 0 * * * *")
1846        );
1847    }
1848
1849    #[test]
1850    fn engine_default_describe_without_schedule() {
1851        let mut engine = create_test_engine();
1852        engine.register(EchoWorkflow).unwrap();
1853        let info = engine.handler_info("echo-workflow").unwrap();
1854        assert!(info.schedule.is_none());
1855    }
1856
1857    #[test]
1858    fn scheduled_handlers_returns_only_scheduled() {
1859        let mut engine = create_test_engine();
1860        engine.register(EchoWorkflow).unwrap();
1861        engine.register(ScheduledWorkflow::new()).unwrap();
1862        engine.register(FailingWorkflow).unwrap();
1863
1864        let scheduled = engine.scheduled_handlers();
1865        assert_eq!(scheduled.len(), 1);
1866        assert_eq!(scheduled[0].0, "scheduled");
1867        assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
1868    }
1869
1870    #[test]
1871    fn scheduled_handlers_empty_when_none_scheduled() {
1872        let mut engine = create_test_engine();
1873        engine.register(EchoWorkflow).unwrap();
1874        engine.register(FailingWorkflow).unwrap();
1875
1876        let scheduled = engine.scheduled_handlers();
1877        assert!(scheduled.is_empty());
1878    }
1879
1880    struct BadCategoryWorkflow(&'static str);
1881
1882    impl WorkflowHandler for BadCategoryWorkflow {
1883        fn name(&self) -> &str {
1884            "bad-category"
1885        }
1886        fn category(&self) -> Option<&str> {
1887            Some(self.0)
1888        }
1889        fn execute<'a>(
1890            &'a self,
1891            _ctx: &'a mut WorkflowContext,
1892        ) -> crate::handler::HandlerFuture<'a> {
1893            Box::pin(async move { Ok(()) })
1894        }
1895    }
1896
1897    #[test]
1898    fn engine_register_rejects_empty_category() {
1899        let mut engine = create_test_engine();
1900        let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
1901        match err {
1902            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
1903            other => panic!("expected InvalidWorkflow, got {other:?}"),
1904        }
1905    }
1906
1907    #[test]
1908    fn engine_register_rejects_leading_slash_category() {
1909        let mut engine = create_test_engine();
1910        let err = engine
1911            .register(BadCategoryWorkflow("/data/etl"))
1912            .unwrap_err();
1913        match err {
1914            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
1915            other => panic!("expected InvalidWorkflow, got {other:?}"),
1916        }
1917    }
1918
1919    #[test]
1920    fn engine_register_rejects_trailing_slash_category() {
1921        let mut engine = create_test_engine();
1922        let err = engine
1923            .register(BadCategoryWorkflow("data/etl/"))
1924            .unwrap_err();
1925        match err {
1926            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
1927            other => panic!("expected InvalidWorkflow, got {other:?}"),
1928        }
1929    }
1930
1931    #[test]
1932    fn engine_register_rejects_double_slash_category() {
1933        let mut engine = create_test_engine();
1934        let err = engine
1935            .register(BadCategoryWorkflow("data//etl"))
1936            .unwrap_err();
1937        match err {
1938            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
1939            other => panic!("expected InvalidWorkflow, got {other:?}"),
1940        }
1941    }
1942
1943    #[test]
1944    fn engine_register_rejects_whitespace_only_segment_category() {
1945        let mut engine = create_test_engine();
1946        let err = engine
1947            .register(BadCategoryWorkflow("data/ /etl"))
1948            .unwrap_err();
1949        match err {
1950            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
1951            other => panic!("expected InvalidWorkflow, got {other:?}"),
1952        }
1953    }
1954
1955    #[test]
1956    fn engine_register_accepts_valid_nested_category() {
1957        let mut engine = create_test_engine();
1958        assert!(engine.register(CategorizedWorkflow).is_ok());
1959    }
1960
1961    #[tokio::test]
1962    async fn engine_unknown_workflow_returns_error() {
1963        let engine = create_test_engine();
1964        let result = engine
1965            .run_handler("unknown", TriggerKind::Manual, json!({}))
1966            .await;
1967        assert!(result.is_err());
1968        match result {
1969            Err(EngineError::InvalidWorkflow(msg)) => {
1970                assert!(msg.contains("no handler registered"));
1971            }
1972            _ => panic!("expected InvalidWorkflow error"),
1973        }
1974    }
1975
1976    #[tokio::test]
1977    async fn engine_enqueue_handler_creates_pending_run() {
1978        let mut engine = create_test_engine();
1979        engine.register(EchoWorkflow).unwrap();
1980
1981        let run = engine
1982            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1983            .await
1984            .unwrap();
1985        assert_eq!(run.status.state, RunStatus::Pending);
1986        assert_eq!(run.workflow_name, "echo-workflow");
1987    }
1988
1989    #[tokio::test]
1990    async fn enqueue_handler_leaves_the_run_unattributed() {
1991        let mut engine = create_test_engine();
1992        engine.register(EchoWorkflow).unwrap();
1993
1994        let run = engine
1995            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1996            .await
1997            .unwrap();
1998
1999        assert!(run.created_by.is_none());
2000    }
2001
2002    #[tokio::test]
2003    async fn enqueue_handler_with_options_records_the_author() {
2004        let mut engine = create_test_engine();
2005        engine.register(EchoWorkflow).unwrap();
2006        let actor = RunActor::User {
2007            user_id: Uuid::now_v7(),
2008        };
2009
2010        let run = engine
2011            .enqueue_handler_with_options(
2012                "echo-workflow",
2013                TriggerKind::Api,
2014                json!({}),
2015                EnqueueOptions {
2016                    created_by: Some(actor.clone()),
2017                    ..Default::default()
2018                },
2019            )
2020            .await
2021            .unwrap()
2022            .into_run();
2023
2024        assert_eq!(run.created_by, Some(actor));
2025    }
2026
2027    #[tokio::test]
2028    async fn enqueue_handler_with_options_accepts_no_author() {
2029        let mut engine = create_test_engine();
2030        engine.register(EchoWorkflow).unwrap();
2031
2032        let run = engine
2033            .enqueue_handler_with_options(
2034                "echo-workflow",
2035                TriggerKind::Cron {
2036                    schedule: "0 * * * * *".to_string(),
2037                },
2038                json!({}),
2039                EnqueueOptions::default(),
2040            )
2041            .await
2042            .unwrap()
2043            .into_run();
2044
2045        assert!(run.created_by.is_none());
2046    }
2047
2048    #[tokio::test]
2049    async fn run_handler_leaves_the_run_unattributed() {
2050        let mut engine = create_test_engine();
2051        engine.register(EchoWorkflow).unwrap();
2052
2053        let run = engine
2054            .run_handler("echo-workflow", TriggerKind::Manual, json!({}))
2055            .await
2056            .unwrap()
2057            .run;
2058
2059        assert!(run.created_by.is_none());
2060    }
2061
2062    #[tokio::test]
2063    async fn engine_register_boxed() {
2064        let mut engine = create_test_engine();
2065        let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
2066        let result = engine.register_boxed(handler);
2067        assert!(result.is_ok());
2068        assert_eq!(engine.handler_names().len(), 1);
2069    }
2070
2071    #[tokio::test]
2072    async fn engine_store_and_provider_accessors() {
2073        let store = Arc::new(InMemoryStore::new());
2074        let inner = ClaudeCodeProvider::new();
2075        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
2076            inner,
2077            "/tmp/ironflow-fixtures",
2078        ));
2079        let engine = Engine::new(store.clone(), provider.clone());
2080
2081        // Verify accessors return references
2082        let _ = engine.store();
2083        let _ = engine.provider();
2084    }
2085
2086    // -----------------------------------------------------------------------
2087    // Operation trait tests
2088    // -----------------------------------------------------------------------
2089
2090    use crate::operation::{Operation, OperationContext};
2091    use async_trait::async_trait;
2092    use ironflow_core::error::OperationError;
2093    use ironflow_store::models::StepKind;
2094
2095    struct FakeGitlabOp {
2096        project_id: u64,
2097        title: String,
2098    }
2099
2100    #[async_trait]
2101    impl Operation for FakeGitlabOp {
2102        fn kind(&self) -> &str {
2103            "gitlab"
2104        }
2105
2106        async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2107            Ok(json!({
2108                "issue_id": 42,
2109                "project_id": self.project_id,
2110                "title": self.title,
2111            }))
2112        }
2113
2114        fn input(&self) -> Option<Value> {
2115            Some(json!({
2116                "project_id": self.project_id,
2117                "title": self.title,
2118            }))
2119        }
2120    }
2121
2122    struct FailingOp;
2123
2124    #[async_trait]
2125    impl Operation for FailingOp {
2126        fn kind(&self) -> &str {
2127            "broken-service"
2128        }
2129
2130        async fn execute(&self, _ctx: &OperationContext) -> Result<Value, OperationError> {
2131            Err(OperationError::Http {
2132                status: None,
2133                message: "service unavailable".to_string(),
2134            })
2135        }
2136    }
2137
2138    struct OperationWorkflow;
2139
2140    impl WorkflowHandler for OperationWorkflow {
2141        fn name(&self) -> &str {
2142            "operation-workflow"
2143        }
2144
2145        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2146            Box::pin(async move {
2147                let op = FakeGitlabOp {
2148                    project_id: 123,
2149                    title: "Bug report".to_string(),
2150                };
2151                ctx.operation("create-issue", &op).await?;
2152                Ok(())
2153            })
2154        }
2155    }
2156
2157    struct FailingOperationWorkflow;
2158
2159    impl WorkflowHandler for FailingOperationWorkflow {
2160        fn name(&self) -> &str {
2161            "failing-operation-workflow"
2162        }
2163
2164        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2165            Box::pin(async move {
2166                ctx.operation("broken-call", &FailingOp).await?;
2167                Ok(())
2168            })
2169        }
2170    }
2171
2172    struct MixedWorkflow;
2173
2174    impl WorkflowHandler for MixedWorkflow {
2175        fn name(&self) -> &str {
2176            "mixed-workflow"
2177        }
2178
2179        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2180            Box::pin(async move {
2181                ctx.shell("build", ShellConfig::new("echo built")).await?;
2182                let op = FakeGitlabOp {
2183                    project_id: 456,
2184                    title: "Deploy done".to_string(),
2185                };
2186                let result = ctx.operation("notify-gitlab", &op).await?;
2187                assert_eq!(result.output["issue_id"], 42);
2188                Ok(())
2189            })
2190        }
2191    }
2192
2193    #[tokio::test]
2194    async fn operation_step_happy_path() {
2195        let mut engine = create_test_engine();
2196        engine.register(OperationWorkflow).unwrap();
2197
2198        let run = engine
2199            .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
2200            .await
2201            .unwrap()
2202            .run;
2203
2204        assert_eq!(run.status.state, RunStatus::Completed);
2205
2206        let steps = engine.store().list_steps(run.id).await.unwrap();
2207
2208        assert_eq!(steps.len(), 1);
2209        assert_eq!(steps[0].name, "create-issue");
2210        assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
2211        assert_eq!(
2212            steps[0].status.state,
2213            ironflow_store::models::StepStatus::Completed
2214        );
2215
2216        let output = steps[0].output.as_ref().unwrap();
2217        assert_eq!(output["issue_id"], 42);
2218        assert_eq!(output["project_id"], 123);
2219
2220        let input = steps[0].input.as_ref().unwrap();
2221        assert_eq!(input["project_id"], 123);
2222        assert_eq!(input["title"], "Bug report");
2223    }
2224
2225    #[tokio::test]
2226    async fn operation_step_failure_marks_run_failed() {
2227        let mut engine = create_test_engine();
2228        engine.register(FailingOperationWorkflow).unwrap();
2229
2230        let result = engine
2231            .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
2232            .await;
2233
2234        assert!(result.is_err());
2235    }
2236
2237    #[tokio::test]
2238    async fn operation_mixed_with_shell_steps() {
2239        let mut engine = create_test_engine();
2240        engine.register(MixedWorkflow).unwrap();
2241
2242        let run = engine
2243            .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
2244            .await
2245            .unwrap()
2246            .run;
2247
2248        assert_eq!(run.status.state, RunStatus::Completed);
2249
2250        let steps = engine.store().list_steps(run.id).await.unwrap();
2251
2252        assert_eq!(steps.len(), 2);
2253        assert_eq!(steps[0].kind, StepKind::Shell);
2254        assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
2255        assert_eq!(steps[0].position, 0);
2256        assert_eq!(steps[1].position, 1);
2257    }
2258
2259    // -----------------------------------------------------------------------
2260    // Approval + resume tests
2261    // -----------------------------------------------------------------------
2262
2263    use crate::config::ApprovalConfig;
2264
2265    struct SingleApprovalWorkflow;
2266
2267    impl WorkflowHandler for SingleApprovalWorkflow {
2268        fn name(&self) -> &str {
2269            "single-approval"
2270        }
2271
2272        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2273            Box::pin(async move {
2274                ctx.shell("build", ShellConfig::new("echo built")).await?;
2275                ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
2276                ctx.shell("deploy", ShellConfig::new("echo deployed"))
2277                    .await?;
2278                Ok(())
2279            })
2280        }
2281    }
2282
2283    struct DoubleApprovalWorkflow;
2284
2285    impl WorkflowHandler for DoubleApprovalWorkflow {
2286        fn name(&self) -> &str {
2287            "double-approval"
2288        }
2289
2290        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
2291            Box::pin(async move {
2292                ctx.shell("build", ShellConfig::new("echo built")).await?;
2293                ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
2294                    .await?;
2295                ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
2296                    .await?;
2297                ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
2298                    .await?;
2299                ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
2300                    .await?;
2301                Ok(())
2302            })
2303        }
2304    }
2305
2306    #[tokio::test]
2307    async fn approval_pauses_run() {
2308        let mut engine = create_test_engine();
2309        engine.register(SingleApprovalWorkflow).unwrap();
2310
2311        let run = engine
2312            .run_handler("single-approval", TriggerKind::Manual, json!({}))
2313            .await
2314            .unwrap()
2315            .run;
2316
2317        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2318
2319        let steps = engine.store().list_steps(run.id).await.unwrap();
2320        assert_eq!(steps.len(), 2); // build + approval gate
2321        assert_eq!(steps[0].kind, StepKind::Shell);
2322        assert_eq!(steps[0].status.state, StepStatus::Completed);
2323        assert_eq!(steps[1].kind, StepKind::Approval);
2324        assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
2325    }
2326
2327    #[tokio::test]
2328    async fn approval_resume_completes_run() {
2329        let mut engine = create_test_engine();
2330        engine.register(SingleApprovalWorkflow).unwrap();
2331
2332        // First execution: pauses at approval
2333        let run = engine
2334            .run_handler("single-approval", TriggerKind::Manual, json!({}))
2335            .await
2336            .unwrap()
2337            .run;
2338        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2339
2340        // Simulate approval: transition to Running
2341        engine
2342            .store()
2343            .update_run_status(run.id, RunStatus::Running)
2344            .await
2345            .unwrap();
2346
2347        // Resume: replays build, skips approval, executes deploy
2348        let resumed = engine.resume_run(run.id).await.unwrap().run;
2349        assert_eq!(resumed.status.state, RunStatus::Completed);
2350
2351        let steps = engine.store().list_steps(run.id).await.unwrap();
2352        assert_eq!(steps.len(), 3); // build + approval + deploy
2353        assert_eq!(steps[0].name, "build");
2354        assert_eq!(steps[0].status.state, StepStatus::Completed);
2355        assert_eq!(steps[1].name, "gate");
2356        assert_eq!(steps[1].kind, StepKind::Approval);
2357        assert_eq!(steps[1].status.state, StepStatus::Completed);
2358        assert_eq!(steps[2].name, "deploy");
2359        assert_eq!(steps[2].status.state, StepStatus::Completed);
2360    }
2361
2362    #[tokio::test]
2363    async fn double_approval_two_resumes() {
2364        let mut engine = create_test_engine();
2365        engine.register(DoubleApprovalWorkflow).unwrap();
2366
2367        // First execution: pauses at staging-gate
2368        let run = engine
2369            .run_handler("double-approval", TriggerKind::Manual, json!({}))
2370            .await
2371            .unwrap()
2372            .run;
2373        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
2374
2375        let steps = engine.store().list_steps(run.id).await.unwrap();
2376        assert_eq!(steps.len(), 2); // build + staging-gate
2377
2378        // First approval
2379        engine
2380            .store()
2381            .update_run_status(run.id, RunStatus::Running)
2382            .await
2383            .unwrap();
2384
2385        let resumed = engine.resume_run(run.id).await.unwrap().run;
2386        assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
2387
2388        let steps = engine.store().list_steps(run.id).await.unwrap();
2389        assert_eq!(steps.len(), 4); // build + staging-gate + deploy-staging + prod-gate
2390
2391        // Second approval
2392        engine
2393            .store()
2394            .update_run_status(run.id, RunStatus::Running)
2395            .await
2396            .unwrap();
2397
2398        let final_run = engine.resume_run(run.id).await.unwrap().run;
2399        assert_eq!(final_run.status.state, RunStatus::Completed);
2400
2401        let steps = engine.store().list_steps(run.id).await.unwrap();
2402        assert_eq!(steps.len(), 5);
2403        assert_eq!(steps[0].name, "build");
2404        assert_eq!(steps[1].name, "staging-gate");
2405        assert_eq!(steps[2].name, "deploy-staging");
2406        assert_eq!(steps[3].name, "prod-gate");
2407        assert_eq!(steps[4].name, "deploy-prod");
2408
2409        for step in &steps {
2410            assert_eq!(step.status.state, StepStatus::Completed);
2411        }
2412    }
2413
2414    // -----------------------------------------------------------------------
2415    // fail_orphaned_steps tests
2416    // -----------------------------------------------------------------------
2417
2418    use ironflow_store::models::{NewStep, StepUpdate, step_trace_id};
2419
2420    async fn create_step_with_status(
2421        store: &Arc<dyn Store>,
2422        run_id: Uuid,
2423        name: &str,
2424        position: u32,
2425        status: StepStatus,
2426    ) -> ironflow_store::models::Step {
2427        let step = store
2428            .create_step(NewStep {
2429                run_id,
2430                trace_id: step_trace_id(run_id, name, position),
2431                name: name.to_string(),
2432                kind: StepKind::Shell,
2433                position,
2434                input: None,
2435                is_error_handler: false,
2436            })
2437            .await
2438            .unwrap();
2439
2440        match status {
2441            StepStatus::Pending => {}
2442            StepStatus::Running => {
2443                store
2444                    .update_step(
2445                        step.id,
2446                        StepUpdate {
2447                            status: Some(StepStatus::Running),
2448                            ..StepUpdate::default()
2449                        },
2450                    )
2451                    .await
2452                    .unwrap();
2453            }
2454            StepStatus::Completed => {
2455                store
2456                    .update_step(
2457                        step.id,
2458                        StepUpdate {
2459                            status: Some(StepStatus::Running),
2460                            ..StepUpdate::default()
2461                        },
2462                    )
2463                    .await
2464                    .unwrap();
2465                store
2466                    .update_step(
2467                        step.id,
2468                        StepUpdate {
2469                            status: Some(StepStatus::Completed),
2470                            ..StepUpdate::default()
2471                        },
2472                    )
2473                    .await
2474                    .unwrap();
2475            }
2476            StepStatus::AwaitingApproval => {
2477                store
2478                    .update_step(
2479                        step.id,
2480                        StepUpdate {
2481                            status: Some(StepStatus::Running),
2482                            ..StepUpdate::default()
2483                        },
2484                    )
2485                    .await
2486                    .unwrap();
2487                store
2488                    .update_step(
2489                        step.id,
2490                        StepUpdate {
2491                            status: Some(StepStatus::AwaitingApproval),
2492                            ..StepUpdate::default()
2493                        },
2494                    )
2495                    .await
2496                    .unwrap();
2497            }
2498            _ => panic!("unsupported status for test helper: {status}"),
2499        }
2500
2501        store.get_step(step.id).await.unwrap().unwrap()
2502    }
2503
2504    #[tokio::test]
2505    async fn fail_orphaned_steps_marks_running_as_failed() {
2506        let engine = create_test_engine();
2507        let run = engine
2508            .store()
2509            .create_run(NewRun {
2510                created_by: None,
2511                workflow_name: "test".to_string(),
2512                trigger: TriggerKind::Manual,
2513                payload: json!({}),
2514                max_retries: 0,
2515                handler_version: None,
2516                labels: HashMap::new(),
2517                scheduled_at: None,
2518                idempotency_key: None,
2519                max_cost_usd: None,
2520            })
2521            .await
2522            .unwrap()
2523            .into_run();
2524
2525        let step = create_step_with_status(
2526            engine.store(),
2527            run.id,
2528            "running-step",
2529            0,
2530            StepStatus::Running,
2531        )
2532        .await;
2533
2534        engine
2535            .fail_orphaned_steps(run.id, "parent run timed out")
2536            .await
2537            .unwrap();
2538
2539        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2540        assert_eq!(updated.status.state, StepStatus::Failed);
2541        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2542        assert!(updated.completed_at.is_some());
2543    }
2544
2545    #[tokio::test]
2546    async fn fail_orphaned_steps_marks_pending_as_skipped() {
2547        let engine = create_test_engine();
2548        let run = engine
2549            .store()
2550            .create_run(NewRun {
2551                created_by: None,
2552                workflow_name: "test".to_string(),
2553                trigger: TriggerKind::Manual,
2554                payload: json!({}),
2555                max_retries: 0,
2556                handler_version: None,
2557                labels: HashMap::new(),
2558                scheduled_at: None,
2559                idempotency_key: None,
2560                max_cost_usd: None,
2561            })
2562            .await
2563            .unwrap()
2564            .into_run();
2565
2566        let step = create_step_with_status(
2567            engine.store(),
2568            run.id,
2569            "pending-step",
2570            0,
2571            StepStatus::Pending,
2572        )
2573        .await;
2574
2575        engine
2576            .fail_orphaned_steps(run.id, "parent run timed out")
2577            .await
2578            .unwrap();
2579
2580        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2581        assert_eq!(updated.status.state, StepStatus::Skipped);
2582        assert!(updated.error.is_none());
2583        assert!(updated.completed_at.is_some());
2584    }
2585
2586    #[tokio::test]
2587    async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
2588        let engine = create_test_engine();
2589        let run = engine
2590            .store()
2591            .create_run(NewRun {
2592                created_by: None,
2593                workflow_name: "test".to_string(),
2594                trigger: TriggerKind::Manual,
2595                payload: json!({}),
2596                max_retries: 0,
2597                handler_version: None,
2598                labels: HashMap::new(),
2599                scheduled_at: None,
2600                idempotency_key: None,
2601                max_cost_usd: None,
2602            })
2603            .await
2604            .unwrap()
2605            .into_run();
2606
2607        let step = create_step_with_status(
2608            engine.store(),
2609            run.id,
2610            "approval-step",
2611            0,
2612            StepStatus::AwaitingApproval,
2613        )
2614        .await;
2615
2616        engine
2617            .fail_orphaned_steps(run.id, "parent run timed out")
2618            .await
2619            .unwrap();
2620
2621        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
2622        assert_eq!(updated.status.state, StepStatus::Failed);
2623        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
2624        assert!(updated.completed_at.is_some());
2625    }
2626
2627    #[tokio::test]
2628    async fn fail_orphaned_steps_skips_terminal_steps() {
2629        let engine = create_test_engine();
2630        let run = engine
2631            .store()
2632            .create_run(NewRun {
2633                created_by: None,
2634                workflow_name: "test".to_string(),
2635                trigger: TriggerKind::Manual,
2636                payload: json!({}),
2637                max_retries: 0,
2638                handler_version: None,
2639                labels: HashMap::new(),
2640                scheduled_at: None,
2641                idempotency_key: None,
2642                max_cost_usd: None,
2643            })
2644            .await
2645            .unwrap()
2646            .into_run();
2647
2648        let completed_step =
2649            create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
2650        let running_step =
2651            create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
2652                .await;
2653
2654        engine
2655            .fail_orphaned_steps(run.id, "parent run timed out")
2656            .await
2657            .unwrap();
2658
2659        let completed = engine
2660            .store()
2661            .get_step(completed_step.id)
2662            .await
2663            .unwrap()
2664            .unwrap();
2665        assert_eq!(completed.status.state, StepStatus::Completed);
2666
2667        let failed = engine
2668            .store()
2669            .get_step(running_step.id)
2670            .await
2671            .unwrap()
2672            .unwrap();
2673        assert_eq!(failed.status.state, StepStatus::Failed);
2674    }
2675
2676    #[tokio::test]
2677    async fn fail_orphaned_steps_mixed_states() {
2678        let engine = create_test_engine();
2679        let run = engine
2680            .store()
2681            .create_run(NewRun {
2682                created_by: None,
2683                workflow_name: "test".to_string(),
2684                trigger: TriggerKind::Manual,
2685                payload: json!({}),
2686                max_retries: 0,
2687                handler_version: None,
2688                labels: HashMap::new(),
2689                scheduled_at: None,
2690                idempotency_key: None,
2691                max_cost_usd: None,
2692            })
2693            .await
2694            .unwrap()
2695            .into_run();
2696
2697        let s_completed =
2698            create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
2699                .await;
2700        let s_running =
2701            create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
2702        let s_pending =
2703            create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
2704
2705        engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
2706
2707        let r_completed = engine
2708            .store()
2709            .get_step(s_completed.id)
2710            .await
2711            .unwrap()
2712            .unwrap();
2713        assert_eq!(r_completed.status.state, StepStatus::Completed);
2714
2715        let r_running = engine
2716            .store()
2717            .get_step(s_running.id)
2718            .await
2719            .unwrap()
2720            .unwrap();
2721        assert_eq!(r_running.status.state, StepStatus::Failed);
2722        assert_eq!(r_running.error.as_deref(), Some("timeout"));
2723
2724        let r_pending = engine
2725            .store()
2726            .get_step(s_pending.id)
2727            .await
2728            .unwrap()
2729            .unwrap();
2730        assert_eq!(r_pending.status.state, StepStatus::Skipped);
2731        assert!(r_pending.error.is_none());
2732    }
2733
2734    #[tokio::test]
2735    async fn fail_orphaned_steps_no_steps_is_noop() {
2736        let engine = create_test_engine();
2737        let run = engine
2738            .store()
2739            .create_run(NewRun {
2740                created_by: None,
2741                workflow_name: "test".to_string(),
2742                trigger: TriggerKind::Manual,
2743                payload: json!({}),
2744                max_retries: 0,
2745                handler_version: None,
2746                labels: HashMap::new(),
2747                scheduled_at: None,
2748                idempotency_key: None,
2749                max_cost_usd: None,
2750            })
2751            .await
2752            .unwrap()
2753            .into_run();
2754
2755        let result = engine.fail_orphaned_steps(run.id, "timeout").await;
2756        assert!(result.is_ok());
2757    }
2758
2759    #[tokio::test]
2760    async fn fail_orphaned_steps_preserves_existing_error() {
2761        let engine = create_test_engine();
2762        let run = engine
2763            .store()
2764            .create_run(NewRun {
2765                created_by: None,
2766                workflow_name: "test".to_string(),
2767                trigger: TriggerKind::Manual,
2768                payload: json!({}),
2769                max_retries: 0,
2770                handler_version: None,
2771                labels: HashMap::new(),
2772                scheduled_at: None,
2773                idempotency_key: None,
2774                max_cost_usd: None,
2775            })
2776            .await
2777            .unwrap()
2778            .into_run();
2779
2780        let step_with_error = create_step_with_status(
2781            engine.store(),
2782            run.id,
2783            "already-errored",
2784            0,
2785            StepStatus::Running,
2786        )
2787        .await;
2788
2789        engine
2790            .store()
2791            .update_step(
2792                step_with_error.id,
2793                StepUpdate {
2794                    error: Some("real error from provider".to_string()),
2795                    ..StepUpdate::default()
2796                },
2797            )
2798            .await
2799            .unwrap();
2800
2801        let step_no_error = create_step_with_status(
2802            engine.store(),
2803            run.id,
2804            "no-error-yet",
2805            1,
2806            StepStatus::Running,
2807        )
2808        .await;
2809
2810        engine
2811            .fail_orphaned_steps(run.id, "parent run failed")
2812            .await
2813            .unwrap();
2814
2815        let updated_with = engine
2816            .store()
2817            .get_step(step_with_error.id)
2818            .await
2819            .unwrap()
2820            .unwrap();
2821        assert_eq!(updated_with.status.state, StepStatus::Failed);
2822        assert_eq!(
2823            updated_with.error.as_deref(),
2824            Some("real error from provider"),
2825        );
2826
2827        let updated_without = engine
2828            .store()
2829            .get_step(step_no_error.id)
2830            .await
2831            .unwrap()
2832            .unwrap();
2833        assert_eq!(updated_without.status.state, StepStatus::Failed);
2834        assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
2835    }
2836}