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, Utc};
15use serde_json::Value;
16use tracing::{error, info, warn};
17use uuid::Uuid;
18
19#[cfg(feature = "prometheus")]
20use ironflow_core::metric_names::{RUN_COST_USD, RUN_DURATION_SECONDS, RUNS_ACTIVE, RUNS_TOTAL};
21use ironflow_core::provider::AgentProvider;
22use ironflow_store::error::StoreError;
23use ironflow_store::models::{
24    NewRun, Run, RunStatus, RunUpdate, StepStatus, StepUpdate, TriggerKind,
25};
26use ironflow_store::store::Store;
27#[cfg(feature = "prometheus")]
28use metrics::{counter, gauge, histogram};
29
30use crate::context::WorkflowContext;
31use crate::error::EngineError;
32use crate::handler::{WorkflowHandler, WorkflowInfo};
33use crate::log_sender::LogSender;
34use crate::notify::{Event, EventPublisher, EventSubscriber};
35use crate::schedule::CronSchedule;
36
37/// The workflow orchestration engine.
38///
39/// Holds references to the store, agent provider, and a registry of
40/// [`WorkflowHandler`]s.
41///
42/// # Examples
43///
44/// ```no_run
45/// use std::sync::Arc;
46/// use ironflow_engine::engine::Engine;
47/// use ironflow_engine::config::ShellConfig;
48/// use ironflow_engine::handler::{WorkflowHandler, HandlerFuture, WorkflowInfo};
49/// use ironflow_engine::context::WorkflowContext;
50/// use ironflow_store::memory::InMemoryStore;
51/// use ironflow_store::models::TriggerKind;
52/// use ironflow_core::providers::claude::ClaudeCodeProvider;
53/// use serde_json::json;
54///
55/// struct CiWorkflow;
56/// impl WorkflowHandler for CiWorkflow {
57///     fn name(&self) -> &str { "ci" }
58///     fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
59///         Box::pin(async move {
60///             ctx.shell("test", ShellConfig::new("cargo test")).await?;
61///             Ok(())
62///         })
63///     }
64/// }
65///
66/// # async fn example() -> Result<(), ironflow_engine::error::EngineError> {
67/// let store = Arc::new(InMemoryStore::new());
68/// let provider = Arc::new(ClaudeCodeProvider::new());
69/// let mut engine = Engine::new(store, provider);
70/// engine.register(CiWorkflow)?;
71///
72/// let run = engine.run_handler("ci", TriggerKind::Manual, json!({})).await?;
73/// tracing::info!(run_id = %run.id, status = ?run.status, "run completed");
74/// # Ok(())
75/// # }
76/// ```
77pub struct Engine {
78    store: Arc<dyn Store>,
79    provider: Arc<dyn AgentProvider>,
80    handlers: HashMap<String, Arc<dyn WorkflowHandler>>,
81    event_publisher: EventPublisher,
82    log_sender: Option<LogSender>,
83}
84
85/// Validate a workflow category path.
86///
87/// A category is a `/`-separated list of non-empty segments. This function
88/// rejects empty paths, leading or trailing `/`, consecutive `/`, and
89/// segments containing only whitespace.
90///
91/// # Errors
92///
93/// Returns [`EngineError::InvalidWorkflow`] when the category is malformed.
94fn validate_category(handler_name: &str, category: &str) -> Result<(), EngineError> {
95    let reject = |reason: &str| {
96        Err(EngineError::InvalidWorkflow(format!(
97            "handler '{handler_name}' has invalid category '{category}': {reason}"
98        )))
99    };
100
101    if category.is_empty() {
102        return reject("empty category");
103    }
104    if category.starts_with('/') {
105        return reject("leading '/'");
106    }
107    if category.ends_with('/') {
108        return reject("trailing '/'");
109    }
110    for segment in category.split('/') {
111        if segment.is_empty() {
112            return reject("empty segment (double '/')");
113        }
114        if segment.trim().is_empty() {
115            return reject("whitespace-only segment");
116        }
117    }
118    Ok(())
119}
120
121impl Engine {
122    /// Create a new engine with the given store and agent provider.
123    ///
124    /// # Examples
125    ///
126    /// ```no_run
127    /// use std::sync::Arc;
128    /// use ironflow_engine::engine::Engine;
129    /// use ironflow_store::memory::InMemoryStore;
130    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
131    ///
132    /// let engine = Engine::new(
133    ///     Arc::new(InMemoryStore::new()),
134    ///     Arc::new(ClaudeCodeProvider::new()),
135    /// );
136    /// ```
137    pub fn new(store: Arc<dyn Store>, provider: Arc<dyn AgentProvider>) -> Self {
138        Self {
139            store,
140            provider,
141            handlers: HashMap::new(),
142            event_publisher: EventPublisher::new(),
143            log_sender: None,
144        }
145    }
146
147    /// Attach a log sender for real-time step output streaming.
148    ///
149    /// When set, all workflow contexts created by this engine will forward
150    /// step output (shell stdout/stderr, agent system messages) to the
151    /// given sender.
152    pub fn set_log_sender(&mut self, sender: LogSender) {
153        self.log_sender = Some(sender);
154    }
155
156    /// Returns a reference to the backing store.
157    pub fn store(&self) -> &Arc<dyn Store> {
158        &self.store
159    }
160
161    /// Returns a reference to the agent provider.
162    pub fn provider(&self) -> &Arc<dyn AgentProvider> {
163        &self.provider
164    }
165
166    /// Build a [`WorkflowContext`] with access to the handler registry.
167    fn build_context(&self, run_id: Uuid) -> WorkflowContext {
168        let handlers = self.handlers.clone();
169        let resolver: crate::context::HandlerResolver =
170            Arc::new(move |name: &str| handlers.get(name).cloned());
171        let mut ctx = WorkflowContext::with_handler_resolver(
172            run_id,
173            self.store.clone(),
174            self.provider.clone(),
175            resolver,
176        );
177        if let Some(ref sender) = self.log_sender {
178            ctx.set_log_sender(sender.clone());
179        }
180        ctx
181    }
182
183    // -----------------------------------------------------------------------
184    // Handler registration
185    // -----------------------------------------------------------------------
186
187    /// Register a [`WorkflowHandler`] for dynamic workflow execution.
188    ///
189    /// The handler is looked up by [`WorkflowHandler::name`] when executing
190    /// or enqueuing.
191    ///
192    /// # Errors
193    ///
194    /// Returns [`EngineError::InvalidWorkflow`] if a handler with the same
195    /// name is already registered.
196    ///
197    /// # Examples
198    ///
199    /// ```no_run
200    /// use std::sync::Arc;
201    /// use ironflow_engine::engine::Engine;
202    /// use ironflow_engine::handler::{WorkflowHandler, HandlerFuture};
203    /// use ironflow_engine::context::WorkflowContext;
204    /// use ironflow_engine::config::ShellConfig;
205    /// use ironflow_store::memory::InMemoryStore;
206    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
207    ///
208    /// struct MyWorkflow;
209    /// impl WorkflowHandler for MyWorkflow {
210    ///     fn name(&self) -> &str { "my-workflow" }
211    ///     fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
212    ///         Box::pin(async move {
213    ///             ctx.shell("step1", ShellConfig::new("echo done")).await?;
214    ///             Ok(())
215    ///         })
216    ///     }
217    /// }
218    ///
219    /// let mut engine = Engine::new(
220    ///     Arc::new(InMemoryStore::new()),
221    ///     Arc::new(ClaudeCodeProvider::new()),
222    /// );
223    /// engine.register(MyWorkflow)?;
224    /// # Ok::<(), ironflow_engine::error::EngineError>(())
225    /// ```
226    pub fn register(&mut self, handler: impl WorkflowHandler + 'static) -> Result<(), EngineError> {
227        let name = handler.name().to_string();
228        if self.handlers.contains_key(&name) {
229            return Err(EngineError::InvalidWorkflow(format!(
230                "handler '{}' already registered",
231                name
232            )));
233        }
234        if let Some(category) = handler.category() {
235            validate_category(&name, category)?;
236        }
237        self.handlers.insert(name, Arc::new(handler));
238        Ok(())
239    }
240
241    /// Register a pre-boxed workflow handler.
242    ///
243    /// # Errors
244    ///
245    /// Returns [`EngineError::InvalidWorkflow`] if a handler with the same
246    /// name is already registered or if its category is invalid.
247    pub fn register_boxed(&mut self, handler: Box<dyn WorkflowHandler>) -> Result<(), EngineError> {
248        let name = handler.name().to_string();
249        if self.handlers.contains_key(&name) {
250            return Err(EngineError::InvalidWorkflow(format!(
251                "handler '{}' already registered",
252                name
253            )));
254        }
255        if let Some(category) = handler.category() {
256            validate_category(&name, category)?;
257        }
258        self.handlers.insert(name, Arc::from(handler));
259        Ok(())
260    }
261
262    /// Get a registered handler by name.
263    pub fn get_handler(&self, name: &str) -> Option<&Arc<dyn WorkflowHandler>> {
264        self.handlers.get(name)
265    }
266
267    /// List registered handler names.
268    pub fn handler_names(&self) -> Vec<&str> {
269        self.handlers.keys().map(|s| s.as_str()).collect()
270    }
271
272    /// Get detailed info about a registered workflow handler.
273    pub fn handler_info(&self, name: &str) -> Option<WorkflowInfo> {
274        self.handlers.get(name).map(|h| h.describe())
275    }
276
277    /// List handlers that have a cron schedule configured.
278    ///
279    /// Returns pairs of `(workflow_name, cron_expression)` for all handlers
280    /// where [`WorkflowHandler::schedule`] returns `Some`.
281    ///
282    /// Use this to wire scheduled handlers into a cron scheduler
283    /// (e.g. `ironflow_runtime::Runtime::cron`).
284    ///
285    /// # Examples
286    ///
287    /// ```no_run
288    /// # use std::sync::Arc;
289    /// # use ironflow_engine::engine::Engine;
290    /// # use ironflow_store::memory::InMemoryStore;
291    /// # use ironflow_core::providers::claude::ClaudeCodeProvider;
292    /// let engine = Engine::new(
293    ///     Arc::new(InMemoryStore::new()),
294    ///     Arc::new(ClaudeCodeProvider::new()),
295    /// );
296    /// for (name, schedule) in engine.scheduled_handlers() {
297    ///     tracing::info!("{name} runs on schedule: {schedule}");
298    /// }
299    /// ```
300    pub fn scheduled_handlers(&self) -> Vec<(&str, &CronSchedule)> {
301        self.handlers
302            .iter()
303            .filter_map(|(name, handler)| handler.schedule().map(|sched| (name.as_str(), sched)))
304            .collect()
305    }
306
307    /// Register an event subscriber for domain events.
308    ///
309    /// The subscriber is called only for events whose type is in
310    /// `event_types`. Pass [`Event::ALL`] to receive every event.
311    ///
312    /// # Examples
313    ///
314    /// ```no_run
315    /// use ironflow_engine::engine::Engine;
316    /// use ironflow_engine::notify::{Event, WebhookSubscriber};
317    /// use ironflow_store::memory::InMemoryStore;
318    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
319    /// use std::sync::Arc;
320    ///
321    /// let mut engine = Engine::new(
322    ///     Arc::new(InMemoryStore::new()),
323    ///     Arc::new(ClaudeCodeProvider::new()),
324    /// );
325    ///
326    /// engine.subscribe(
327    ///     WebhookSubscriber::new("https://hooks.example.com/events"),
328    ///     &[Event::RUN_STATUS_CHANGED, Event::STEP_FAILED],
329    /// );
330    /// ```
331    pub fn subscribe(
332        &mut self,
333        subscriber: impl EventSubscriber + 'static,
334        event_types: &[&'static str],
335    ) {
336        self.event_publisher.subscribe(subscriber, event_types);
337    }
338
339    /// Returns a reference to the event publisher.
340    ///
341    /// Useful for publishing events from outside the engine (e.g. auth
342    /// routes in the API layer).
343    pub fn event_publisher(&self) -> &EventPublisher {
344        &self.event_publisher
345    }
346
347    // -----------------------------------------------------------------------
348    // Dynamic workflow execution (WorkflowHandler)
349    // -----------------------------------------------------------------------
350
351    /// Execute a registered handler inline.
352    ///
353    /// Creates a run, builds a [`WorkflowContext`], calls the handler's
354    /// [`execute`](WorkflowHandler::execute), and finalizes the run.
355    ///
356    /// # Errors
357    ///
358    /// Returns [`EngineError::InvalidWorkflow`] if no handler is registered
359    /// with that name. Returns [`EngineError`] if execution fails.
360    ///
361    /// # Examples
362    ///
363    /// ```no_run
364    /// use std::sync::Arc;
365    /// use ironflow_engine::engine::Engine;
366    /// use ironflow_store::memory::InMemoryStore;
367    /// use ironflow_store::models::TriggerKind;
368    /// use ironflow_core::providers::claude::ClaudeCodeProvider;
369    /// use serde_json::json;
370    ///
371    /// # async fn example(engine: &Engine) -> Result<(), ironflow_engine::error::EngineError> {
372    /// let run = engine.run_handler("deploy", TriggerKind::Manual, json!({})).await?;
373    /// # Ok(())
374    /// # }
375    /// ```
376    #[tracing::instrument(name = "engine.run_handler", skip_all, fields(workflow = %handler_name))]
377    pub async fn run_handler(
378        &self,
379        handler_name: &str,
380        trigger: TriggerKind,
381        payload: Value,
382    ) -> Result<Run, EngineError> {
383        let handler = self
384            .handlers
385            .get(handler_name)
386            .ok_or_else(|| {
387                EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
388            })?
389            .clone();
390
391        let handler_version = handler.version().map(str::to_string);
392        let run = self
393            .store
394            .create_run(NewRun {
395                workflow_name: handler_name.to_string(),
396                trigger,
397                payload,
398                max_retries: 0,
399                handler_version,
400                labels: handler.default_labels(),
401                scheduled_at: None,
402            })
403            .await?;
404
405        let run_id = run.id;
406        info!(run_id = %run_id, handler_version = run.handler_version.as_deref().unwrap_or(""), "run created");
407
408        self.store
409            .update_run_status(run_id, RunStatus::Running)
410            .await?;
411
412        #[cfg(feature = "prometheus")]
413        gauge!(RUNS_ACTIVE, "workflow" => handler_name.to_string()).increment(1.0);
414
415        let run_start = Instant::now();
416        let mut ctx = self.build_context(run_id);
417
418        let result = handler.execute(&mut ctx).await;
419        self.finalize_run(run_id, handler_name, result, &ctx, run_start)
420            .await
421    }
422
423    /// Enqueue a handler-based workflow for worker execution.
424    ///
425    /// The workflow name is stored in the run. The worker looks up the
426    /// handler by name when executing.
427    ///
428    /// # Errors
429    ///
430    /// Returns [`EngineError::InvalidWorkflow`] if no handler is registered.
431    #[tracing::instrument(name = "engine.enqueue_handler", skip_all, fields(workflow = %handler_name))]
432    pub async fn enqueue_handler(
433        &self,
434        handler_name: &str,
435        trigger: TriggerKind,
436        payload: Value,
437        max_retries: u32,
438    ) -> Result<Run, EngineError> {
439        self.enqueue_handler_with_options(
440            handler_name,
441            trigger,
442            payload,
443            max_retries,
444            HashMap::new(),
445            None,
446        )
447        .await
448    }
449
450    /// Enqueue a handler-based workflow with labels and optional deferred scheduling.
451    ///
452    /// # Errors
453    ///
454    /// Returns [`EngineError::InvalidWorkflow`] if no handler is registered.
455    #[tracing::instrument(name = "engine.enqueue_handler_with_options", skip_all, fields(workflow = %handler_name))]
456    pub async fn enqueue_handler_with_options(
457        &self,
458        handler_name: &str,
459        trigger: TriggerKind,
460        payload: Value,
461        max_retries: u32,
462        labels: HashMap<String, String>,
463        scheduled_at: Option<DateTime<Utc>>,
464    ) -> Result<Run, EngineError> {
465        let handler = self.handlers.get(handler_name).ok_or_else(|| {
466            EngineError::InvalidWorkflow(format!("no handler registered: {handler_name}"))
467        })?;
468
469        let handler_version = handler.version().map(str::to_string);
470        let mut merged_labels = handler.default_labels();
471        merged_labels.extend(labels);
472
473        let run = self
474            .store
475            .create_run(NewRun {
476                workflow_name: handler_name.to_string(),
477                trigger,
478                payload,
479                max_retries,
480                handler_version,
481                labels: merged_labels,
482                scheduled_at,
483            })
484            .await?;
485
486        info!(run_id = %run.id, workflow = %handler_name, "handler run enqueued");
487        Ok(run)
488    }
489
490    /// Execute a handler-based run (used by the worker after pick_next_pending).
491    ///
492    /// Looks up the handler by the run's `workflow_name` and executes it
493    /// with a fresh [`WorkflowContext`].
494    ///
495    /// # Errors
496    ///
497    /// Returns [`EngineError::InvalidWorkflow`] if no handler matches.
498    #[tracing::instrument(name = "engine.execute_handler_run", skip_all, fields(run_id = %run_id))]
499    pub async fn execute_handler_run(&self, run_id: Uuid) -> Result<Run, EngineError> {
500        let run = self
501            .store
502            .get_run(run_id)
503            .await?
504            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
505
506        let handler = self
507            .handlers
508            .get(&run.workflow_name)
509            .ok_or_else(|| {
510                EngineError::InvalidWorkflow(format!(
511                    "no handler registered: {}",
512                    run.workflow_name
513                ))
514            })?
515            .clone();
516
517        #[cfg(feature = "prometheus")]
518        gauge!(RUNS_ACTIVE, "workflow" => run.workflow_name.clone()).increment(1.0);
519
520        let run_start = Instant::now();
521        let mut ctx = self.build_context(run_id);
522
523        let result = handler.execute(&mut ctx).await;
524        self.finalize_run(run_id, &run.workflow_name, result, &ctx, run_start)
525            .await
526    }
527
528    /// Execute a run by its ID (used by the worker after pick_next_pending).
529    ///
530    /// Delegates to [`execute_handler_run`](Self::execute_handler_run).
531    ///
532    /// # Errors
533    ///
534    /// Returns [`EngineError`] if the run is not found or execution fails.
535    #[tracing::instrument(name = "engine.execute_run", skip_all, fields(run_id = %run_id))]
536    pub async fn execute_run(&self, run_id: Uuid) -> Result<Run, EngineError> {
537        self.execute_handler_run(run_id).await
538    }
539
540    /// Resume a run after human approval.
541    ///
542    /// Re-executes the handler with step replay: completed steps return
543    /// cached output, approved approval steps are skipped, and execution
544    /// continues from the first unexecuted step.
545    ///
546    /// Supports multiple approval gates -- each resume replays all prior
547    /// steps and stops at the next approval (or completes the run).
548    ///
549    /// # Errors
550    ///
551    /// Returns [`EngineError::InvalidWorkflow`] if no handler matches.
552    /// Returns [`EngineError`] if execution fails or hits another approval.
553    #[tracing::instrument(name = "engine.resume_run", skip_all, fields(run_id = %run_id))]
554    pub async fn resume_run(&self, run_id: Uuid) -> Result<Run, EngineError> {
555        let run = self
556            .store
557            .get_run(run_id)
558            .await?
559            .ok_or(EngineError::Store(StoreError::RunNotFound(run_id)))?;
560
561        let handler = self
562            .handlers
563            .get(&run.workflow_name)
564            .ok_or_else(|| {
565                EngineError::InvalidWorkflow(format!(
566                    "no handler registered: {}",
567                    run.workflow_name
568                ))
569            })?
570            .clone();
571
572        info!(run_id = %run_id, workflow = %run.workflow_name, "resuming run after approval");
573
574        let run_start = Instant::now();
575        let mut ctx = self.build_context(run_id);
576        ctx.load_replay_steps().await?;
577
578        let result = handler.execute(&mut ctx).await;
579        self.finalize_run(run_id, &run.workflow_name, result, &ctx, run_start)
580            .await
581    }
582
583    /// Fail all non-terminal steps for a run.
584    ///
585    /// Called after a run is marked as failed (timeout, error, panic) to clean up
586    /// orphaned steps that are still in `Running`, `Pending`, or `AwaitingApproval`.
587    ///
588    /// - `Running` / `AwaitingApproval` steps are marked `Failed`.
589    /// - `Pending` steps are marked `Skipped` (FSM does not allow Pending -> Failed).
590    ///
591    /// Errors from individual step updates are logged but do not abort the cleanup.
592    ///
593    /// # Errors
594    ///
595    /// Returns [`EngineError`] if listing steps fails.
596    pub async fn fail_orphaned_steps(
597        &self,
598        run_id: Uuid,
599        error_message: &str,
600    ) -> Result<(), EngineError> {
601        let steps = self.store.list_steps(run_id).await?;
602        let now = Utc::now();
603
604        for step in steps {
605            if step.status.state.is_terminal() {
606                continue;
607            }
608
609            let (target_status, error) = match step.status.state {
610                StepStatus::Running | StepStatus::AwaitingApproval => {
611                    let err = if step.error.is_some() {
612                        None
613                    } else {
614                        Some(error_message.to_string())
615                    };
616                    (StepStatus::Failed, err)
617                }
618                StepStatus::Pending => (StepStatus::Skipped, None),
619                _ => continue,
620            };
621
622            if let Err(e) = self
623                .store
624                .update_step(
625                    step.id,
626                    StepUpdate {
627                        status: Some(target_status),
628                        error,
629                        completed_at: Some(now),
630                        ..StepUpdate::default()
631                    },
632                )
633                .await
634            {
635                warn!(
636                    run_id = %run_id,
637                    step_id = %step.id,
638                    step_name = %step.name,
639                    error = %e,
640                    "failed to cleanup orphaned step"
641                );
642            } else {
643                info!(
644                    run_id = %run_id,
645                    step_id = %step.id,
646                    step_name = %step.name,
647                    from = %step.status.state,
648                    to = %target_status,
649                    "cleaned up orphaned step"
650                );
651            }
652        }
653
654        Ok(())
655    }
656
657    /// Finalize a run with the given result and context.
658    ///
659    /// On success: updates run to Completed with cost, duration, and completed_at.
660    /// On failure: updates run to Failed with error, cost, duration, and completed_at.
661    /// Always: fetches and returns the final Run.
662    async fn finalize_run(
663        &self,
664        run_id: Uuid,
665        workflow_name: &str,
666        result: Result<(), EngineError>,
667        ctx: &WorkflowContext,
668        run_start: Instant,
669    ) -> Result<Run, EngineError> {
670        let total_duration = run_start.elapsed().as_millis() as u64;
671        let completed_at = Utc::now();
672
673        let final_status;
674        let final_run;
675
676        match result {
677            Ok(()) => {
678                final_status = RunStatus::Completed;
679                final_run = self
680                    .store
681                    .update_run_returning(
682                        run_id,
683                        RunUpdate {
684                            status: Some(RunStatus::Completed),
685                            cost_usd: Some(ctx.total_cost_usd()),
686                            duration_ms: Some(total_duration),
687                            completed_at: Some(completed_at),
688                            ..RunUpdate::default()
689                        },
690                    )
691                    .await?;
692
693                info!(
694                    run_id = %run_id,
695                    cost_usd = %ctx.total_cost_usd(),
696                    duration_ms = total_duration,
697                    "run completed"
698                );
699            }
700            Err(EngineError::ApprovalRequired {
701                run_id: approval_run_id,
702                step_id,
703                ref message,
704            }) => {
705                final_status = RunStatus::AwaitingApproval;
706                final_run = self
707                    .store
708                    .update_run_returning(
709                        run_id,
710                        RunUpdate {
711                            status: Some(RunStatus::AwaitingApproval),
712                            cost_usd: Some(ctx.total_cost_usd()),
713                            duration_ms: Some(total_duration),
714                            ..RunUpdate::default()
715                        },
716                    )
717                    .await?;
718
719                info!(
720                    run_id = %approval_run_id,
721                    step_id = %step_id,
722                    message = %message,
723                    "run awaiting approval"
724                );
725            }
726            Err(err) => {
727                final_status = RunStatus::Failed;
728                if let Err(store_err) = self
729                    .store
730                    .update_run(
731                        run_id,
732                        RunUpdate {
733                            status: Some(RunStatus::Failed),
734                            error: Some(err.to_string()),
735                            cost_usd: Some(ctx.total_cost_usd()),
736                            duration_ms: Some(total_duration),
737                            completed_at: Some(completed_at),
738                            ..RunUpdate::default()
739                        },
740                    )
741                    .await
742                {
743                    error!(run_id = %run_id, store_error = %store_err, "failed to persist run failure");
744                }
745
746                error!(run_id = %run_id, error = %err, "run failed");
747
748                self.publish_run_status_changed(
749                    workflow_name,
750                    run_id,
751                    final_status,
752                    Some(err.to_string()),
753                    ctx,
754                    total_duration,
755                );
756
757                #[cfg(feature = "prometheus")]
758                self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
759
760                return Err(err);
761            }
762        }
763
764        self.publish_run_status_changed(
765            workflow_name,
766            run_id,
767            final_status,
768            None,
769            ctx,
770            total_duration,
771        );
772
773        #[cfg(feature = "prometheus")]
774        self.emit_run_metrics(workflow_name, final_status, total_duration, ctx);
775
776        Ok(final_run)
777    }
778
779    /// Emit Prometheus metrics for a completed run.
780    #[cfg(feature = "prometheus")]
781    fn emit_run_metrics(
782        &self,
783        workflow_name: &str,
784        status: RunStatus,
785        duration_ms: u64,
786        ctx: &WorkflowContext,
787    ) {
788        let status_str = status.to_string();
789        let wf = workflow_name.to_string();
790
791        counter!(RUNS_TOTAL, "workflow" => wf.clone(), "status" => status_str.clone()).increment(1);
792        histogram!(RUN_DURATION_SECONDS, "workflow" => wf.clone(), "status" => status_str)
793            .record(duration_ms as f64 / 1000.0);
794        histogram!(RUN_COST_USD, "workflow" => wf.clone()).record(
795            ctx.total_cost_usd()
796                .to_string()
797                .parse::<f64>()
798                .unwrap_or(0.0),
799        );
800        gauge!(RUNS_ACTIVE, "workflow" => wf).decrement(1.0);
801    }
802
803    /// Publish a run status changed event to all registered subscribers.
804    ///
805    /// `from` is always `Running` because `finalize_run` is only called
806    /// from a running state.
807    fn publish_run_status_changed(
808        &self,
809        workflow_name: &str,
810        run_id: Uuid,
811        to: RunStatus,
812        error: Option<String>,
813        ctx: &WorkflowContext,
814        duration_ms: u64,
815    ) {
816        let now = Utc::now();
817        let cost_usd = ctx.total_cost_usd();
818        let wf = workflow_name.to_string();
819
820        self.event_publisher.publish(Event::RunStatusChanged {
821            run_id,
822            workflow_name: wf.clone(),
823            from: RunStatus::Running,
824            to,
825            error: error.clone(),
826            cost_usd,
827            duration_ms,
828            at: now,
829        });
830
831        if to == RunStatus::Failed {
832            self.event_publisher.publish(Event::RunFailed {
833                run_id,
834                workflow_name: wf,
835                error,
836                cost_usd,
837                duration_ms,
838                at: now,
839            });
840        }
841    }
842}
843
844impl fmt::Debug for Engine {
845    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
846        f.debug_struct("Engine")
847            .field("handlers", &self.handlers.keys().collect::<Vec<_>>())
848            .finish_non_exhaustive()
849    }
850}
851
852#[cfg(test)]
853mod tests {
854    use super::*;
855    use crate::config::ShellConfig;
856    use crate::handler::{HandlerFuture, WorkflowHandler};
857    use ironflow_core::providers::claude::ClaudeCodeProvider;
858    use ironflow_core::providers::record_replay::RecordReplayProvider;
859    use ironflow_store::memory::InMemoryStore;
860    use ironflow_store::models::StepStatus;
861    use serde_json::json;
862
863    // Test handler that echoes a message via shell
864    struct EchoWorkflow;
865
866    impl WorkflowHandler for EchoWorkflow {
867        fn name(&self) -> &str {
868            "echo-workflow"
869        }
870
871        fn describe(&self) -> WorkflowInfo {
872            WorkflowInfo {
873                description: "A simple workflow that echoes hello".to_string(),
874                source_code: None,
875                sub_workflows: Vec::new(),
876                category: None,
877                version: self.version().map(str::to_string),
878                input_schema: None,
879                default_labels: HashMap::new(),
880                schedule: self.schedule().cloned(),
881            }
882        }
883
884        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
885            Box::pin(async move {
886                ctx.shell("greet", ShellConfig::new("echo hello")).await?;
887                Ok(())
888            })
889        }
890    }
891
892    // Test handler that fails
893    struct FailingWorkflow;
894
895    impl WorkflowHandler for FailingWorkflow {
896        fn name(&self) -> &str {
897            "failing-workflow"
898        }
899
900        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
901            Box::pin(async move {
902                ctx.shell("fail", ShellConfig::new("exit 1")).await?;
903                Ok(())
904            })
905        }
906    }
907
908    fn create_test_engine() -> Engine {
909        let store = Arc::new(InMemoryStore::new());
910        let inner = ClaudeCodeProvider::new();
911        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
912            inner,
913            "/tmp/ironflow-fixtures",
914        ));
915        Engine::new(store, provider)
916    }
917
918    #[test]
919    fn engine_new_creates_instance() {
920        let engine = create_test_engine();
921        assert_eq!(engine.handler_names().len(), 0);
922    }
923
924    #[test]
925    fn engine_register_handler() {
926        let mut engine = create_test_engine();
927        let result = engine.register(EchoWorkflow);
928        assert!(result.is_ok());
929        assert_eq!(engine.handler_names().len(), 1);
930        assert!(engine.handler_names().contains(&"echo-workflow"));
931    }
932
933    #[test]
934    fn engine_register_duplicate_returns_error() {
935        let mut engine = create_test_engine();
936        engine.register(EchoWorkflow).unwrap();
937        let result = engine.register(EchoWorkflow);
938        assert!(result.is_err());
939    }
940
941    #[test]
942    fn engine_get_handler_found() {
943        let mut engine = create_test_engine();
944        engine.register(EchoWorkflow).unwrap();
945        let handler = engine.get_handler("echo-workflow");
946        assert!(handler.is_some());
947    }
948
949    #[test]
950    fn engine_get_handler_not_found() {
951        let engine = create_test_engine();
952        let handler = engine.get_handler("nonexistent");
953        assert!(handler.is_none());
954    }
955
956    #[test]
957    fn engine_handler_names_lists_all() {
958        let mut engine = create_test_engine();
959        engine.register(EchoWorkflow).unwrap();
960        engine.register(FailingWorkflow).unwrap();
961        let names = engine.handler_names();
962        assert_eq!(names.len(), 2);
963        assert!(names.contains(&"echo-workflow"));
964        assert!(names.contains(&"failing-workflow"));
965    }
966
967    #[test]
968    fn engine_handler_info_returns_description() {
969        let mut engine = create_test_engine();
970        engine.register(EchoWorkflow).unwrap();
971        let info = engine.handler_info("echo-workflow");
972        assert!(info.is_some());
973        let info = info.unwrap();
974        assert_eq!(info.description, "A simple workflow that echoes hello");
975    }
976
977    struct CategorizedWorkflow;
978
979    impl WorkflowHandler for CategorizedWorkflow {
980        fn name(&self) -> &str {
981            "categorized"
982        }
983        fn category(&self) -> Option<&str> {
984            Some("data/etl")
985        }
986        fn execute<'a>(
987            &'a self,
988            _ctx: &'a mut WorkflowContext,
989        ) -> crate::handler::HandlerFuture<'a> {
990            Box::pin(async move { Ok(()) })
991        }
992    }
993
994    #[test]
995    fn engine_default_describe_propagates_category() {
996        let mut engine = create_test_engine();
997        engine.register(CategorizedWorkflow).unwrap();
998        let info = engine.handler_info("categorized").unwrap();
999        assert_eq!(info.category.as_deref(), Some("data/etl"));
1000    }
1001
1002    #[test]
1003    fn engine_default_describe_without_category() {
1004        let mut engine = create_test_engine();
1005        engine.register(EchoWorkflow).unwrap();
1006        let info = engine.handler_info("echo-workflow").unwrap();
1007        assert!(info.category.is_none());
1008    }
1009
1010    // -----------------------------------------------------------------------
1011    // Schedule tests
1012    // -----------------------------------------------------------------------
1013
1014    struct ScheduledWorkflow {
1015        schedule: CronSchedule,
1016    }
1017
1018    impl ScheduledWorkflow {
1019        fn new() -> Self {
1020            Self {
1021                schedule: CronSchedule::new("0 0 * * * *").unwrap(),
1022            }
1023        }
1024    }
1025
1026    impl WorkflowHandler for ScheduledWorkflow {
1027        fn name(&self) -> &str {
1028            "scheduled"
1029        }
1030        fn schedule(&self) -> Option<&CronSchedule> {
1031            Some(&self.schedule)
1032        }
1033        fn execute<'a>(
1034            &'a self,
1035            _ctx: &'a mut WorkflowContext,
1036        ) -> crate::handler::HandlerFuture<'a> {
1037            Box::pin(async move { Ok(()) })
1038        }
1039    }
1040
1041    #[test]
1042    fn engine_default_describe_propagates_schedule() {
1043        let mut engine = create_test_engine();
1044        engine.register(ScheduledWorkflow::new()).unwrap();
1045        let info = engine.handler_info("scheduled").unwrap();
1046        assert_eq!(
1047            info.schedule.as_ref().map(|s| s.as_str()),
1048            Some("0 0 * * * *")
1049        );
1050    }
1051
1052    #[test]
1053    fn engine_default_describe_without_schedule() {
1054        let mut engine = create_test_engine();
1055        engine.register(EchoWorkflow).unwrap();
1056        let info = engine.handler_info("echo-workflow").unwrap();
1057        assert!(info.schedule.is_none());
1058    }
1059
1060    #[test]
1061    fn scheduled_handlers_returns_only_scheduled() {
1062        let mut engine = create_test_engine();
1063        engine.register(EchoWorkflow).unwrap();
1064        engine.register(ScheduledWorkflow::new()).unwrap();
1065        engine.register(FailingWorkflow).unwrap();
1066
1067        let scheduled = engine.scheduled_handlers();
1068        assert_eq!(scheduled.len(), 1);
1069        assert_eq!(scheduled[0].0, "scheduled");
1070        assert_eq!(scheduled[0].1.as_str(), "0 0 * * * *");
1071    }
1072
1073    #[test]
1074    fn scheduled_handlers_empty_when_none_scheduled() {
1075        let mut engine = create_test_engine();
1076        engine.register(EchoWorkflow).unwrap();
1077        engine.register(FailingWorkflow).unwrap();
1078
1079        let scheduled = engine.scheduled_handlers();
1080        assert!(scheduled.is_empty());
1081    }
1082
1083    struct BadCategoryWorkflow(&'static str);
1084
1085    impl WorkflowHandler for BadCategoryWorkflow {
1086        fn name(&self) -> &str {
1087            "bad-category"
1088        }
1089        fn category(&self) -> Option<&str> {
1090            Some(self.0)
1091        }
1092        fn execute<'a>(
1093            &'a self,
1094            _ctx: &'a mut WorkflowContext,
1095        ) -> crate::handler::HandlerFuture<'a> {
1096            Box::pin(async move { Ok(()) })
1097        }
1098    }
1099
1100    #[test]
1101    fn engine_register_rejects_empty_category() {
1102        let mut engine = create_test_engine();
1103        let err = engine.register(BadCategoryWorkflow("")).unwrap_err();
1104        match err {
1105            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty category")),
1106            other => panic!("expected InvalidWorkflow, got {other:?}"),
1107        }
1108    }
1109
1110    #[test]
1111    fn engine_register_rejects_leading_slash_category() {
1112        let mut engine = create_test_engine();
1113        let err = engine
1114            .register(BadCategoryWorkflow("/data/etl"))
1115            .unwrap_err();
1116        match err {
1117            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("leading '/'")),
1118            other => panic!("expected InvalidWorkflow, got {other:?}"),
1119        }
1120    }
1121
1122    #[test]
1123    fn engine_register_rejects_trailing_slash_category() {
1124        let mut engine = create_test_engine();
1125        let err = engine
1126            .register(BadCategoryWorkflow("data/etl/"))
1127            .unwrap_err();
1128        match err {
1129            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("trailing '/'")),
1130            other => panic!("expected InvalidWorkflow, got {other:?}"),
1131        }
1132    }
1133
1134    #[test]
1135    fn engine_register_rejects_double_slash_category() {
1136        let mut engine = create_test_engine();
1137        let err = engine
1138            .register(BadCategoryWorkflow("data//etl"))
1139            .unwrap_err();
1140        match err {
1141            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("empty segment")),
1142            other => panic!("expected InvalidWorkflow, got {other:?}"),
1143        }
1144    }
1145
1146    #[test]
1147    fn engine_register_rejects_whitespace_only_segment_category() {
1148        let mut engine = create_test_engine();
1149        let err = engine
1150            .register(BadCategoryWorkflow("data/ /etl"))
1151            .unwrap_err();
1152        match err {
1153            EngineError::InvalidWorkflow(msg) => assert!(msg.contains("whitespace-only segment")),
1154            other => panic!("expected InvalidWorkflow, got {other:?}"),
1155        }
1156    }
1157
1158    #[test]
1159    fn engine_register_accepts_valid_nested_category() {
1160        let mut engine = create_test_engine();
1161        assert!(engine.register(CategorizedWorkflow).is_ok());
1162    }
1163
1164    #[tokio::test]
1165    async fn engine_unknown_workflow_returns_error() {
1166        let engine = create_test_engine();
1167        let result = engine
1168            .run_handler("unknown", TriggerKind::Manual, json!({}))
1169            .await;
1170        assert!(result.is_err());
1171        match result {
1172            Err(EngineError::InvalidWorkflow(msg)) => {
1173                assert!(msg.contains("no handler registered"));
1174            }
1175            _ => panic!("expected InvalidWorkflow error"),
1176        }
1177    }
1178
1179    #[tokio::test]
1180    async fn engine_enqueue_handler_creates_pending_run() {
1181        let mut engine = create_test_engine();
1182        engine.register(EchoWorkflow).unwrap();
1183
1184        let run = engine
1185            .enqueue_handler("echo-workflow", TriggerKind::Manual, json!({}), 0)
1186            .await
1187            .unwrap();
1188        assert_eq!(run.status.state, RunStatus::Pending);
1189        assert_eq!(run.workflow_name, "echo-workflow");
1190    }
1191
1192    #[tokio::test]
1193    async fn engine_register_boxed() {
1194        let mut engine = create_test_engine();
1195        let handler: Box<dyn WorkflowHandler> = Box::new(EchoWorkflow);
1196        let result = engine.register_boxed(handler);
1197        assert!(result.is_ok());
1198        assert_eq!(engine.handler_names().len(), 1);
1199    }
1200
1201    #[tokio::test]
1202    async fn engine_store_and_provider_accessors() {
1203        let store = Arc::new(InMemoryStore::new());
1204        let inner = ClaudeCodeProvider::new();
1205        let provider: Arc<dyn AgentProvider> = Arc::new(RecordReplayProvider::replay(
1206            inner,
1207            "/tmp/ironflow-fixtures",
1208        ));
1209        let engine = Engine::new(store.clone(), provider.clone());
1210
1211        // Verify accessors return references
1212        let _ = engine.store();
1213        let _ = engine.provider();
1214    }
1215
1216    // -----------------------------------------------------------------------
1217    // Operation trait tests
1218    // -----------------------------------------------------------------------
1219
1220    use crate::operation::Operation;
1221    use ironflow_store::models::StepKind;
1222    use std::future::Future;
1223    use std::pin::Pin;
1224
1225    struct FakeGitlabOp {
1226        project_id: u64,
1227        title: String,
1228    }
1229
1230    impl Operation for FakeGitlabOp {
1231        fn kind(&self) -> &str {
1232            "gitlab"
1233        }
1234
1235        fn execute(&self) -> Pin<Box<dyn Future<Output = Result<Value, EngineError>> + Send + '_>> {
1236            Box::pin(async move {
1237                Ok(json!({
1238                    "issue_id": 42,
1239                    "project_id": self.project_id,
1240                    "title": self.title,
1241                }))
1242            })
1243        }
1244
1245        fn input(&self) -> Option<Value> {
1246            Some(json!({
1247                "project_id": self.project_id,
1248                "title": self.title,
1249            }))
1250        }
1251    }
1252
1253    struct FailingOp;
1254
1255    impl Operation for FailingOp {
1256        fn kind(&self) -> &str {
1257            "broken-service"
1258        }
1259
1260        fn execute(&self) -> Pin<Box<dyn Future<Output = Result<Value, EngineError>> + Send + '_>> {
1261            Box::pin(async move { Err(EngineError::StepConfig("service unavailable".to_string())) })
1262        }
1263    }
1264
1265    struct OperationWorkflow;
1266
1267    impl WorkflowHandler for OperationWorkflow {
1268        fn name(&self) -> &str {
1269            "operation-workflow"
1270        }
1271
1272        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1273            Box::pin(async move {
1274                let op = FakeGitlabOp {
1275                    project_id: 123,
1276                    title: "Bug report".to_string(),
1277                };
1278                ctx.operation("create-issue", &op).await?;
1279                Ok(())
1280            })
1281        }
1282    }
1283
1284    struct FailingOperationWorkflow;
1285
1286    impl WorkflowHandler for FailingOperationWorkflow {
1287        fn name(&self) -> &str {
1288            "failing-operation-workflow"
1289        }
1290
1291        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1292            Box::pin(async move {
1293                ctx.operation("broken-call", &FailingOp).await?;
1294                Ok(())
1295            })
1296        }
1297    }
1298
1299    struct MixedWorkflow;
1300
1301    impl WorkflowHandler for MixedWorkflow {
1302        fn name(&self) -> &str {
1303            "mixed-workflow"
1304        }
1305
1306        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1307            Box::pin(async move {
1308                ctx.shell("build", ShellConfig::new("echo built")).await?;
1309                let op = FakeGitlabOp {
1310                    project_id: 456,
1311                    title: "Deploy done".to_string(),
1312                };
1313                let result = ctx.operation("notify-gitlab", &op).await?;
1314                assert_eq!(result.output["issue_id"], 42);
1315                Ok(())
1316            })
1317        }
1318    }
1319
1320    #[tokio::test]
1321    async fn operation_step_happy_path() {
1322        let mut engine = create_test_engine();
1323        engine.register(OperationWorkflow).unwrap();
1324
1325        let run = engine
1326            .run_handler("operation-workflow", TriggerKind::Manual, json!({}))
1327            .await
1328            .unwrap();
1329
1330        assert_eq!(run.status.state, RunStatus::Completed);
1331
1332        let steps = engine.store().list_steps(run.id).await.unwrap();
1333
1334        assert_eq!(steps.len(), 1);
1335        assert_eq!(steps[0].name, "create-issue");
1336        assert_eq!(steps[0].kind, StepKind::Custom("gitlab".to_string()));
1337        assert_eq!(
1338            steps[0].status.state,
1339            ironflow_store::models::StepStatus::Completed
1340        );
1341
1342        let output = steps[0].output.as_ref().unwrap();
1343        assert_eq!(output["issue_id"], 42);
1344        assert_eq!(output["project_id"], 123);
1345
1346        let input = steps[0].input.as_ref().unwrap();
1347        assert_eq!(input["project_id"], 123);
1348        assert_eq!(input["title"], "Bug report");
1349    }
1350
1351    #[tokio::test]
1352    async fn operation_step_failure_marks_run_failed() {
1353        let mut engine = create_test_engine();
1354        engine.register(FailingOperationWorkflow).unwrap();
1355
1356        let result = engine
1357            .run_handler("failing-operation-workflow", TriggerKind::Manual, json!({}))
1358            .await;
1359
1360        assert!(result.is_err());
1361    }
1362
1363    #[tokio::test]
1364    async fn operation_mixed_with_shell_steps() {
1365        let mut engine = create_test_engine();
1366        engine.register(MixedWorkflow).unwrap();
1367
1368        let run = engine
1369            .run_handler("mixed-workflow", TriggerKind::Manual, json!({}))
1370            .await
1371            .unwrap();
1372
1373        assert_eq!(run.status.state, RunStatus::Completed);
1374
1375        let steps = engine.store().list_steps(run.id).await.unwrap();
1376
1377        assert_eq!(steps.len(), 2);
1378        assert_eq!(steps[0].kind, StepKind::Shell);
1379        assert_eq!(steps[1].kind, StepKind::Custom("gitlab".to_string()));
1380        assert_eq!(steps[0].position, 0);
1381        assert_eq!(steps[1].position, 1);
1382    }
1383
1384    // -----------------------------------------------------------------------
1385    // Approval + resume tests
1386    // -----------------------------------------------------------------------
1387
1388    use crate::config::ApprovalConfig;
1389
1390    struct SingleApprovalWorkflow;
1391
1392    impl WorkflowHandler for SingleApprovalWorkflow {
1393        fn name(&self) -> &str {
1394            "single-approval"
1395        }
1396
1397        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1398            Box::pin(async move {
1399                ctx.shell("build", ShellConfig::new("echo built")).await?;
1400                ctx.approval("gate", ApprovalConfig::new("OK?")).await?;
1401                ctx.shell("deploy", ShellConfig::new("echo deployed"))
1402                    .await?;
1403                Ok(())
1404            })
1405        }
1406    }
1407
1408    struct DoubleApprovalWorkflow;
1409
1410    impl WorkflowHandler for DoubleApprovalWorkflow {
1411        fn name(&self) -> &str {
1412            "double-approval"
1413        }
1414
1415        fn execute<'a>(&'a self, ctx: &'a mut WorkflowContext) -> HandlerFuture<'a> {
1416            Box::pin(async move {
1417                ctx.shell("build", ShellConfig::new("echo built")).await?;
1418                ctx.approval("staging-gate", ApprovalConfig::new("Deploy staging?"))
1419                    .await?;
1420                ctx.shell("deploy-staging", ShellConfig::new("echo staging"))
1421                    .await?;
1422                ctx.approval("prod-gate", ApprovalConfig::new("Deploy prod?"))
1423                    .await?;
1424                ctx.shell("deploy-prod", ShellConfig::new("echo prod"))
1425                    .await?;
1426                Ok(())
1427            })
1428        }
1429    }
1430
1431    #[tokio::test]
1432    async fn approval_pauses_run() {
1433        let mut engine = create_test_engine();
1434        engine.register(SingleApprovalWorkflow).unwrap();
1435
1436        let run = engine
1437            .run_handler("single-approval", TriggerKind::Manual, json!({}))
1438            .await
1439            .unwrap();
1440
1441        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
1442
1443        let steps = engine.store().list_steps(run.id).await.unwrap();
1444        assert_eq!(steps.len(), 2); // build + approval gate
1445        assert_eq!(steps[0].kind, StepKind::Shell);
1446        assert_eq!(steps[0].status.state, StepStatus::Completed);
1447        assert_eq!(steps[1].kind, StepKind::Approval);
1448        assert_eq!(steps[1].status.state, StepStatus::AwaitingApproval);
1449    }
1450
1451    #[tokio::test]
1452    async fn approval_resume_completes_run() {
1453        let mut engine = create_test_engine();
1454        engine.register(SingleApprovalWorkflow).unwrap();
1455
1456        // First execution: pauses at approval
1457        let run = engine
1458            .run_handler("single-approval", TriggerKind::Manual, json!({}))
1459            .await
1460            .unwrap();
1461        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
1462
1463        // Simulate approval: transition to Running
1464        engine
1465            .store()
1466            .update_run_status(run.id, RunStatus::Running)
1467            .await
1468            .unwrap();
1469
1470        // Resume: replays build, skips approval, executes deploy
1471        let resumed = engine.resume_run(run.id).await.unwrap();
1472        assert_eq!(resumed.status.state, RunStatus::Completed);
1473
1474        let steps = engine.store().list_steps(run.id).await.unwrap();
1475        assert_eq!(steps.len(), 3); // build + approval + deploy
1476        assert_eq!(steps[0].name, "build");
1477        assert_eq!(steps[0].status.state, StepStatus::Completed);
1478        assert_eq!(steps[1].name, "gate");
1479        assert_eq!(steps[1].kind, StepKind::Approval);
1480        assert_eq!(steps[1].status.state, StepStatus::Completed);
1481        assert_eq!(steps[2].name, "deploy");
1482        assert_eq!(steps[2].status.state, StepStatus::Completed);
1483    }
1484
1485    #[tokio::test]
1486    async fn double_approval_two_resumes() {
1487        let mut engine = create_test_engine();
1488        engine.register(DoubleApprovalWorkflow).unwrap();
1489
1490        // First execution: pauses at staging-gate
1491        let run = engine
1492            .run_handler("double-approval", TriggerKind::Manual, json!({}))
1493            .await
1494            .unwrap();
1495        assert_eq!(run.status.state, RunStatus::AwaitingApproval);
1496
1497        let steps = engine.store().list_steps(run.id).await.unwrap();
1498        assert_eq!(steps.len(), 2); // build + staging-gate
1499
1500        // First approval
1501        engine
1502            .store()
1503            .update_run_status(run.id, RunStatus::Running)
1504            .await
1505            .unwrap();
1506
1507        let resumed = engine.resume_run(run.id).await.unwrap();
1508        assert_eq!(resumed.status.state, RunStatus::AwaitingApproval);
1509
1510        let steps = engine.store().list_steps(run.id).await.unwrap();
1511        assert_eq!(steps.len(), 4); // build + staging-gate + deploy-staging + prod-gate
1512
1513        // Second approval
1514        engine
1515            .store()
1516            .update_run_status(run.id, RunStatus::Running)
1517            .await
1518            .unwrap();
1519
1520        let final_run = engine.resume_run(run.id).await.unwrap();
1521        assert_eq!(final_run.status.state, RunStatus::Completed);
1522
1523        let steps = engine.store().list_steps(run.id).await.unwrap();
1524        assert_eq!(steps.len(), 5);
1525        assert_eq!(steps[0].name, "build");
1526        assert_eq!(steps[1].name, "staging-gate");
1527        assert_eq!(steps[2].name, "deploy-staging");
1528        assert_eq!(steps[3].name, "prod-gate");
1529        assert_eq!(steps[4].name, "deploy-prod");
1530
1531        for step in &steps {
1532            assert_eq!(step.status.state, StepStatus::Completed);
1533        }
1534    }
1535
1536    // -----------------------------------------------------------------------
1537    // fail_orphaned_steps tests
1538    // -----------------------------------------------------------------------
1539
1540    use ironflow_store::models::{NewStep, StepUpdate};
1541
1542    async fn create_step_with_status(
1543        store: &Arc<dyn Store>,
1544        run_id: Uuid,
1545        name: &str,
1546        position: u32,
1547        status: StepStatus,
1548    ) -> ironflow_store::models::Step {
1549        let step = store
1550            .create_step(NewStep {
1551                run_id,
1552                name: name.to_string(),
1553                kind: StepKind::Shell,
1554                position,
1555                input: None,
1556            })
1557            .await
1558            .unwrap();
1559
1560        match status {
1561            StepStatus::Pending => {}
1562            StepStatus::Running => {
1563                store
1564                    .update_step(
1565                        step.id,
1566                        StepUpdate {
1567                            status: Some(StepStatus::Running),
1568                            ..StepUpdate::default()
1569                        },
1570                    )
1571                    .await
1572                    .unwrap();
1573            }
1574            StepStatus::Completed => {
1575                store
1576                    .update_step(
1577                        step.id,
1578                        StepUpdate {
1579                            status: Some(StepStatus::Running),
1580                            ..StepUpdate::default()
1581                        },
1582                    )
1583                    .await
1584                    .unwrap();
1585                store
1586                    .update_step(
1587                        step.id,
1588                        StepUpdate {
1589                            status: Some(StepStatus::Completed),
1590                            ..StepUpdate::default()
1591                        },
1592                    )
1593                    .await
1594                    .unwrap();
1595            }
1596            StepStatus::AwaitingApproval => {
1597                store
1598                    .update_step(
1599                        step.id,
1600                        StepUpdate {
1601                            status: Some(StepStatus::Running),
1602                            ..StepUpdate::default()
1603                        },
1604                    )
1605                    .await
1606                    .unwrap();
1607                store
1608                    .update_step(
1609                        step.id,
1610                        StepUpdate {
1611                            status: Some(StepStatus::AwaitingApproval),
1612                            ..StepUpdate::default()
1613                        },
1614                    )
1615                    .await
1616                    .unwrap();
1617            }
1618            _ => panic!("unsupported status for test helper: {status}"),
1619        }
1620
1621        store.get_step(step.id).await.unwrap().unwrap()
1622    }
1623
1624    #[tokio::test]
1625    async fn fail_orphaned_steps_marks_running_as_failed() {
1626        let engine = create_test_engine();
1627        let run = engine
1628            .store()
1629            .create_run(NewRun {
1630                workflow_name: "test".to_string(),
1631                trigger: TriggerKind::Manual,
1632                payload: json!({}),
1633                max_retries: 0,
1634                handler_version: None,
1635                labels: HashMap::new(),
1636                scheduled_at: None,
1637            })
1638            .await
1639            .unwrap();
1640
1641        let step = create_step_with_status(
1642            engine.store(),
1643            run.id,
1644            "running-step",
1645            0,
1646            StepStatus::Running,
1647        )
1648        .await;
1649
1650        engine
1651            .fail_orphaned_steps(run.id, "parent run timed out")
1652            .await
1653            .unwrap();
1654
1655        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
1656        assert_eq!(updated.status.state, StepStatus::Failed);
1657        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
1658        assert!(updated.completed_at.is_some());
1659    }
1660
1661    #[tokio::test]
1662    async fn fail_orphaned_steps_marks_pending_as_skipped() {
1663        let engine = create_test_engine();
1664        let run = engine
1665            .store()
1666            .create_run(NewRun {
1667                workflow_name: "test".to_string(),
1668                trigger: TriggerKind::Manual,
1669                payload: json!({}),
1670                max_retries: 0,
1671                handler_version: None,
1672                labels: HashMap::new(),
1673                scheduled_at: None,
1674            })
1675            .await
1676            .unwrap();
1677
1678        let step = create_step_with_status(
1679            engine.store(),
1680            run.id,
1681            "pending-step",
1682            0,
1683            StepStatus::Pending,
1684        )
1685        .await;
1686
1687        engine
1688            .fail_orphaned_steps(run.id, "parent run timed out")
1689            .await
1690            .unwrap();
1691
1692        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
1693        assert_eq!(updated.status.state, StepStatus::Skipped);
1694        assert!(updated.error.is_none());
1695        assert!(updated.completed_at.is_some());
1696    }
1697
1698    #[tokio::test]
1699    async fn fail_orphaned_steps_marks_awaiting_approval_as_failed() {
1700        let engine = create_test_engine();
1701        let run = engine
1702            .store()
1703            .create_run(NewRun {
1704                workflow_name: "test".to_string(),
1705                trigger: TriggerKind::Manual,
1706                payload: json!({}),
1707                max_retries: 0,
1708                handler_version: None,
1709                labels: HashMap::new(),
1710                scheduled_at: None,
1711            })
1712            .await
1713            .unwrap();
1714
1715        let step = create_step_with_status(
1716            engine.store(),
1717            run.id,
1718            "approval-step",
1719            0,
1720            StepStatus::AwaitingApproval,
1721        )
1722        .await;
1723
1724        engine
1725            .fail_orphaned_steps(run.id, "parent run timed out")
1726            .await
1727            .unwrap();
1728
1729        let updated = engine.store().get_step(step.id).await.unwrap().unwrap();
1730        assert_eq!(updated.status.state, StepStatus::Failed);
1731        assert_eq!(updated.error.as_deref(), Some("parent run timed out"));
1732        assert!(updated.completed_at.is_some());
1733    }
1734
1735    #[tokio::test]
1736    async fn fail_orphaned_steps_skips_terminal_steps() {
1737        let engine = create_test_engine();
1738        let run = engine
1739            .store()
1740            .create_run(NewRun {
1741                workflow_name: "test".to_string(),
1742                trigger: TriggerKind::Manual,
1743                payload: json!({}),
1744                max_retries: 0,
1745                handler_version: None,
1746                labels: HashMap::new(),
1747                scheduled_at: None,
1748            })
1749            .await
1750            .unwrap();
1751
1752        let completed_step =
1753            create_step_with_status(engine.store(), run.id, "done", 0, StepStatus::Completed).await;
1754        let running_step =
1755            create_step_with_status(engine.store(), run.id, "in-flight", 1, StepStatus::Running)
1756                .await;
1757
1758        engine
1759            .fail_orphaned_steps(run.id, "parent run timed out")
1760            .await
1761            .unwrap();
1762
1763        let completed = engine
1764            .store()
1765            .get_step(completed_step.id)
1766            .await
1767            .unwrap()
1768            .unwrap();
1769        assert_eq!(completed.status.state, StepStatus::Completed);
1770
1771        let failed = engine
1772            .store()
1773            .get_step(running_step.id)
1774            .await
1775            .unwrap()
1776            .unwrap();
1777        assert_eq!(failed.status.state, StepStatus::Failed);
1778    }
1779
1780    #[tokio::test]
1781    async fn fail_orphaned_steps_mixed_states() {
1782        let engine = create_test_engine();
1783        let run = engine
1784            .store()
1785            .create_run(NewRun {
1786                workflow_name: "test".to_string(),
1787                trigger: TriggerKind::Manual,
1788                payload: json!({}),
1789                max_retries: 0,
1790                handler_version: None,
1791                labels: HashMap::new(),
1792                scheduled_at: None,
1793            })
1794            .await
1795            .unwrap();
1796
1797        let s_completed =
1798            create_step_with_status(engine.store(), run.id, "step-1", 0, StepStatus::Completed)
1799                .await;
1800        let s_running =
1801            create_step_with_status(engine.store(), run.id, "step-2", 1, StepStatus::Running).await;
1802        let s_pending =
1803            create_step_with_status(engine.store(), run.id, "step-3", 2, StepStatus::Pending).await;
1804
1805        engine.fail_orphaned_steps(run.id, "timeout").await.unwrap();
1806
1807        let r_completed = engine
1808            .store()
1809            .get_step(s_completed.id)
1810            .await
1811            .unwrap()
1812            .unwrap();
1813        assert_eq!(r_completed.status.state, StepStatus::Completed);
1814
1815        let r_running = engine
1816            .store()
1817            .get_step(s_running.id)
1818            .await
1819            .unwrap()
1820            .unwrap();
1821        assert_eq!(r_running.status.state, StepStatus::Failed);
1822        assert_eq!(r_running.error.as_deref(), Some("timeout"));
1823
1824        let r_pending = engine
1825            .store()
1826            .get_step(s_pending.id)
1827            .await
1828            .unwrap()
1829            .unwrap();
1830        assert_eq!(r_pending.status.state, StepStatus::Skipped);
1831        assert!(r_pending.error.is_none());
1832    }
1833
1834    #[tokio::test]
1835    async fn fail_orphaned_steps_no_steps_is_noop() {
1836        let engine = create_test_engine();
1837        let run = engine
1838            .store()
1839            .create_run(NewRun {
1840                workflow_name: "test".to_string(),
1841                trigger: TriggerKind::Manual,
1842                payload: json!({}),
1843                max_retries: 0,
1844                handler_version: None,
1845                labels: HashMap::new(),
1846                scheduled_at: None,
1847            })
1848            .await
1849            .unwrap();
1850
1851        let result = engine.fail_orphaned_steps(run.id, "timeout").await;
1852        assert!(result.is_ok());
1853    }
1854
1855    #[tokio::test]
1856    async fn fail_orphaned_steps_preserves_existing_error() {
1857        let engine = create_test_engine();
1858        let run = engine
1859            .store()
1860            .create_run(NewRun {
1861                workflow_name: "test".to_string(),
1862                trigger: TriggerKind::Manual,
1863                payload: json!({}),
1864                max_retries: 0,
1865                handler_version: None,
1866                labels: HashMap::new(),
1867                scheduled_at: None,
1868            })
1869            .await
1870            .unwrap();
1871
1872        let step_with_error = create_step_with_status(
1873            engine.store(),
1874            run.id,
1875            "already-errored",
1876            0,
1877            StepStatus::Running,
1878        )
1879        .await;
1880
1881        engine
1882            .store()
1883            .update_step(
1884                step_with_error.id,
1885                StepUpdate {
1886                    error: Some("real error from provider".to_string()),
1887                    ..StepUpdate::default()
1888                },
1889            )
1890            .await
1891            .unwrap();
1892
1893        let step_no_error = create_step_with_status(
1894            engine.store(),
1895            run.id,
1896            "no-error-yet",
1897            1,
1898            StepStatus::Running,
1899        )
1900        .await;
1901
1902        engine
1903            .fail_orphaned_steps(run.id, "parent run failed")
1904            .await
1905            .unwrap();
1906
1907        let updated_with = engine
1908            .store()
1909            .get_step(step_with_error.id)
1910            .await
1911            .unwrap()
1912            .unwrap();
1913        assert_eq!(updated_with.status.state, StepStatus::Failed);
1914        assert_eq!(
1915            updated_with.error.as_deref(),
1916            Some("real error from provider"),
1917        );
1918
1919        let updated_without = engine
1920            .store()
1921            .get_step(step_no_error.id)
1922            .await
1923            .unwrap()
1924            .unwrap();
1925        assert_eq!(updated_without.status.state, StepStatus::Failed);
1926        assert_eq!(updated_without.error.as_deref(), Some("parent run failed"),);
1927    }
1928}