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