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