Skip to main content

ironflow_engine/
engine.rs

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