Skip to main content

ironflow_engine/testing/
result.rs

1//! What a [`TestEngine`](crate::testing::TestEngine) run leaves behind.
2//!
3//! [`TestResult`] and [`TestStep`] read back the run and the steps the engine
4//! persisted, so an assertion sees exactly what the API and the dashboard would
5//! serve for that run.
6
7use std::time::Duration;
8
9use rust_decimal::Decimal;
10use serde_json::Value;
11use uuid::Uuid;
12
13use ironflow_store::models::{Run, RunStatus, Step, StepKind, StepStatus};
14
15use crate::executor::{StepOutput, StepResult};
16
17/// Stand-in for a step that recorded no input, and for a run with no output.
18static NULL: Value = Value::Null;
19
20/// One persisted step, with assertion-friendly accessors.
21///
22/// # Examples
23///
24/// ```no_run
25/// use ironflow_engine::testing::TestResult;
26/// use ironflow_store::models::StepStatus;
27///
28/// # fn example(result: &TestResult) {
29/// let build = result.step("build");
30/// assert_eq!(build.status(), StepStatus::Completed);
31/// assert_eq!(build.step_output().exit_code(), Some(0));
32/// # }
33/// ```
34#[derive(Debug, Clone)]
35pub struct TestStep {
36    step: Step,
37    output: Value,
38}
39
40impl TestStep {
41    /// Wrap a persisted step, defaulting a missing output to [`Value::Null`].
42    pub(crate) fn new(step: Step) -> Self {
43        let output = step.output.clone().unwrap_or(Value::Null);
44        Self { step, output }
45    }
46
47    /// The step name, as passed to the context method that created it.
48    pub fn name(&self) -> &str {
49        &self.step.name
50    }
51
52    /// The kind of operation this step ran.
53    pub fn kind(&self) -> &StepKind {
54        &self.step.kind
55    }
56
57    /// The terminal status the step reached.
58    pub fn status(&self) -> StepStatus {
59        self.step.status.state
60    }
61
62    /// Persisted output, [`Value::Null`] when the step produced none.
63    pub fn output(&self) -> &Value {
64        &self.output
65    }
66
67    /// Persisted input, the serialized step config.
68    pub fn input(&self) -> &Value {
69        self.step.input.as_ref().unwrap_or(&NULL)
70    }
71
72    /// Error message recorded on the step, if it failed or was rejected.
73    pub fn error(&self) -> Option<&str> {
74        self.step.error.as_deref()
75    }
76
77    /// Wall-clock duration of the step.
78    pub fn duration(&self) -> Duration {
79        Duration::from_millis(self.step.duration_ms)
80    }
81
82    /// Cost charged by the step, in USD. Zero for everything but agent steps.
83    pub fn cost_usd(&self) -> Decimal {
84        self.step.cost_usd
85    }
86
87    /// Whether the step reached [`StepStatus::Completed`].
88    pub fn is_completed(&self) -> bool {
89        self.status() == StepStatus::Completed
90    }
91
92    /// Whether the step came from an
93    /// [`on_error`](crate::context::WorkflowContext::on_error) handler.
94    pub fn is_error_handler(&self) -> bool {
95        self.step.is_error_handler
96    }
97
98    /// The underlying store record, for assertions the accessors do not cover.
99    pub fn raw(&self) -> &Step {
100        &self.step
101    }
102
103    /// The persisted output read through the typed [`StepOutput`] accessors:
104    /// `stdout()`, `status()`, `body()`, `json::<T>()`...
105    ///
106    /// # Examples
107    ///
108    /// ```no_run
109    /// use ironflow_engine::executor::SubWorkflowOutput;
110    /// use ironflow_engine::testing::TestResult;
111    ///
112    /// # fn example(result: &TestResult) -> Result<(), ironflow_engine::error::EngineError> {
113    /// assert_eq!(result.step("build").step_output().stdout(), "compiled");
114    /// let child: SubWorkflowOutput = result.step("collect").step_output().json()?;
115    /// println!("child run {}", child.run_id());
116    /// # Ok(())
117    /// # }
118    /// ```
119    pub fn step_output(&self) -> StepOutput {
120        StepOutput::from(&self.step)
121    }
122}
123
124/// Outcome of a [`TestEngine`](crate::testing::TestEngine) run.
125///
126/// A handler that fails is *not* an error here: the run is read back with
127/// [`RunStatus::Failed`] and [`error`](Self::error) set. Only wiring failures
128/// (unknown handler, duplicate handler name, store error) surface as an `Err`
129/// from [`TestEngine::run`](crate::testing::TestEngine::run).
130///
131/// # Examples
132///
133/// ```no_run
134/// use ironflow_engine::testing::TestResult;
135/// use ironflow_store::models::RunStatus;
136///
137/// # fn example(result: &TestResult) {
138/// assert_eq!(result.status(), RunStatus::Completed);
139/// assert_eq!(result.step_names(), vec!["build", "test", "deploy"]);
140/// # }
141/// ```
142#[derive(Debug, Clone)]
143pub struct TestResult {
144    run: Run,
145    steps: Vec<TestStep>,
146    step_results: Vec<StepResult>,
147    error: Option<String>,
148}
149
150impl TestResult {
151    /// Assemble a result from what the store holds after execution.
152    pub(crate) fn new(
153        run: Run,
154        steps: Vec<Step>,
155        step_results: Vec<StepResult>,
156        error: Option<String>,
157    ) -> Self {
158        Self {
159            run,
160            steps: steps.into_iter().map(TestStep::new).collect(),
161            step_results,
162            error,
163        }
164    }
165
166    /// The status the run finished in.
167    pub fn status(&self) -> RunStatus {
168        self.run.status.state
169    }
170
171    /// The persisted run record.
172    pub fn run(&self) -> &Run {
173        &self.run
174    }
175
176    /// The run identifier, for [`TestEngine::resume`](crate::testing::TestEngine::resume).
177    pub fn run_id(&self) -> Uuid {
178        self.run.id
179    }
180
181    /// Every persisted step, ordered by position.
182    pub fn steps(&self) -> &[TestStep] {
183        &self.steps
184    }
185
186    /// The step names, ordered by position.
187    pub fn step_names(&self) -> Vec<&str> {
188        self.steps.iter().map(TestStep::name).collect()
189    }
190
191    /// The first step carrying `name`.
192    ///
193    /// Steps of a parallel wave share a position, and a handler may reuse a
194    /// name: disambiguate those with [`steps`](Self::steps).
195    ///
196    /// # Panics
197    ///
198    /// Panics when no step carries that name; the message lists the names that
199    /// exist.
200    ///
201    /// # Examples
202    ///
203    /// ```no_run
204    /// use ironflow_engine::testing::TestResult;
205    ///
206    /// # fn example(result: &TestResult) {
207    /// assert!(result.step("deploy").is_completed());
208    /// # }
209    /// ```
210    pub fn step(&self, name: &str) -> &TestStep {
211        self.try_step(name).unwrap_or_else(|| {
212            panic!(
213                "no step named {name:?} in this run; steps are {:?}",
214                self.step_names()
215            )
216        })
217    }
218
219    /// The first step carrying `name`, or `None`.
220    pub fn try_step(&self, name: &str) -> Option<&TestStep> {
221        self.steps.iter().find(|step| step.name() == name)
222    }
223
224    /// Output the handler set with
225    /// [`set_output`](crate::context::WorkflowContext::set_output),
226    /// [`Value::Null`] when it set none.
227    ///
228    /// The output of the last step is `steps().last()`, or
229    /// [`step(name).output()`](TestStep::output) by name.
230    ///
231    /// # Examples
232    ///
233    /// ```no_run
234    /// use ironflow_engine::testing::TestResult;
235    /// use serde_json::json;
236    ///
237    /// # fn example(result: &TestResult) {
238    /// assert_eq!(result.output(), &json!({"approved": true}));
239    /// # }
240    /// ```
241    pub fn output(&self) -> &Value {
242        self.run.output.as_ref().unwrap_or(&NULL)
243    }
244
245    /// Wall-clock duration recorded on the run.
246    pub fn duration(&self) -> Duration {
247        Duration::from_millis(self.run.duration_ms)
248    }
249
250    /// Total cost of the run, in USD.
251    pub fn cost_usd(&self) -> Decimal {
252        self.run.cost_usd
253    }
254
255    /// Why the run stopped, if it did not complete.
256    pub fn error(&self) -> Option<&str> {
257        self.error.as_deref()
258    }
259
260    /// Per-step metrics as the engine reported them.
261    ///
262    /// Empty when the run failed: the engine returns the error instead of a
263    /// [`WorkflowResult`](crate::engine::WorkflowResult) on that path. Use
264    /// [`steps`](Self::steps), which always reflects the store.
265    pub fn step_results(&self) -> &[StepResult] {
266        &self.step_results
267    }
268
269    /// Whether the run reached [`RunStatus::Completed`].
270    pub fn is_completed(&self) -> bool {
271        self.status() == RunStatus::Completed
272    }
273}
274
275#[cfg(test)]
276mod tests {
277    use std::collections::HashMap;
278    use std::sync::Arc;
279
280    use chrono::Utc;
281    use serde_json::json;
282
283    use ironflow_store::memory::InMemoryStore;
284    use ironflow_store::models::{
285        NewRun, NewStep, RunUpdate, StepUpdate, TriggerKind, step_trace_id,
286    };
287    use ironflow_store::store::RunStore;
288
289    use super::*;
290
291    /// Build a two-step run in a store: `build` completed, `deploy` failed.
292    ///
293    /// `Run` and `Step` are `#[non_exhaustive]`, so they can only be obtained
294    /// from a store.
295    async fn persisted_run() -> (Run, Vec<Step>) {
296        persisted_run_with_output(None).await
297    }
298
299    /// Same as [`persisted_run`], with `output` recorded on the run.
300    async fn persisted_run_with_output(output: Option<Value>) -> (Run, Vec<Step>) {
301        let store = Arc::new(InMemoryStore::new());
302        let run = store
303            .create_run(NewRun {
304                created_by: None,
305                workflow_name: "deploy".to_string(),
306                trigger: TriggerKind::Manual,
307                payload: json!({}),
308                max_retries: 0,
309                handler_version: None,
310                labels: HashMap::new(),
311                scheduled_at: None,
312                idempotency_key: None,
313                concurrency_key: None,
314                concurrency_limits: Vec::new(),
315                max_cost_usd: None,
316                worker_tags: Vec::new(),
317            })
318            .await
319            .expect("create run")
320            .into_run();
321
322        for (position, name) in [(0u32, "build"), (1, "deploy")] {
323            let step = store
324                .create_step(NewStep {
325                    run_id: run.id,
326                    trace_id: step_trace_id(run.id, name, position),
327                    name: name.to_string(),
328                    kind: StepKind::Shell,
329                    position,
330                    input: Some(json!({"type": "shell", "command": name})),
331                    is_error_handler: false,
332                })
333                .await
334                .expect("create step");
335
336            store
337                .update_step(
338                    step.id,
339                    StepUpdate {
340                        status: Some(StepStatus::Running),
341                        started_at: Some(Utc::now()),
342                        ..StepUpdate::default()
343                    },
344                )
345                .await
346                .expect("start step");
347
348            let update = if name == "build" {
349                StepUpdate {
350                    status: Some(StepStatus::Completed),
351                    output: Some(json!({"stdout": "built", "stderr": "", "exit_code": 0})),
352                    duration_ms: Some(12),
353                    completed_at: Some(Utc::now()),
354                    ..StepUpdate::default()
355                }
356            } else {
357                StepUpdate {
358                    status: Some(StepStatus::Failed),
359                    error: Some("boom".to_string()),
360                    completed_at: Some(Utc::now()),
361                    ..StepUpdate::default()
362                }
363            };
364            store
365                .update_step(step.id, update)
366                .await
367                .expect("finish the step");
368        }
369
370        store
371            .update_run(
372                run.id,
373                RunUpdate {
374                    output,
375                    ..RunUpdate::default()
376                },
377            )
378            .await
379            .expect("record the run output");
380
381        let steps = store.list_steps(run.id).await.expect("list steps");
382        let run = store
383            .get_run(run.id)
384            .await
385            .expect("get run")
386            .expect("the run exists");
387        (run, steps)
388    }
389
390    #[tokio::test]
391    async fn accessors_read_back_what_was_persisted() {
392        let (run, steps) = persisted_run().await;
393        let result = TestResult::new(run, steps, Vec::new(), Some("boom".to_string()));
394
395        assert_eq!(result.step_names(), vec!["build", "deploy"]);
396        assert_eq!(result.error(), Some("boom"));
397        assert!(result.step_results().is_empty());
398        assert!(!result.is_completed());
399
400        let build = result.step("build");
401        assert!(build.is_completed());
402        assert_eq!(build.output()["stdout"], "built");
403        assert_eq!(build.input()["command"], "build");
404        assert_eq!(build.duration(), Duration::from_millis(12));
405        assert_eq!(build.cost_usd(), Decimal::ZERO);
406        assert_eq!(build.kind(), &StepKind::Shell);
407        assert!(!build.is_error_handler());
408        assert_eq!(build.raw().name, "build");
409        assert!(build.error().is_none());
410    }
411
412    #[tokio::test]
413    async fn set_output_is_the_test_result_output() {
414        let (run, steps) = persisted_run_with_output(Some(json!({"approved": true}))).await;
415        let result = TestResult::new(run, steps, Vec::new(), None);
416
417        assert_eq!(result.output(), &json!({"approved": true}));
418        // The last step output stays reachable through the steps.
419        assert_eq!(result.step("deploy").output(), &Value::Null);
420        assert_eq!(result.steps().last().map(TestStep::name), Some("deploy"));
421    }
422
423    #[tokio::test]
424    async fn output_is_null_when_the_handler_set_none() {
425        let (run, steps) = persisted_run().await;
426        let result = TestResult::new(run, steps, Vec::new(), None);
427
428        // `build` has an output, but the run itself has none.
429        assert_eq!(result.output(), &Value::Null);
430        assert_eq!(result.step("build").output()["stdout"], "built");
431        // `deploy` failed without an output.
432        assert_eq!(result.step("deploy").status(), StepStatus::Failed);
433        assert_eq!(result.step("deploy").error(), Some("boom"));
434    }
435
436    #[tokio::test]
437    async fn run_level_accessors_reflect_the_store() {
438        let (run, steps) = persisted_run().await;
439        let run_id = run.id;
440        let result = TestResult::new(run, steps, Vec::new(), None);
441
442        assert_eq!(result.run_id(), run_id);
443        assert_eq!(result.run().workflow_name, "deploy");
444        assert_eq!(result.status(), RunStatus::Pending);
445        assert_eq!(result.duration(), Duration::ZERO);
446        assert_eq!(result.cost_usd(), Decimal::ZERO);
447        assert_eq!(result.steps().len(), 2);
448    }
449
450    #[tokio::test]
451    async fn try_step_returns_none_for_an_unknown_name() {
452        let (run, steps) = persisted_run().await;
453        let result = TestResult::new(run, steps, Vec::new(), None);
454
455        assert!(result.try_step("nope").is_none());
456    }
457
458    #[tokio::test]
459    #[should_panic(expected = "no step named \"nope\"")]
460    async fn step_panics_for_an_unknown_name() {
461        let (run, steps) = persisted_run().await;
462        let result = TestResult::new(run, steps, Vec::new(), None);
463
464        result.step("nope");
465    }
466
467    #[tokio::test]
468    async fn a_run_without_steps_has_a_null_output() {
469        let (run, _) = persisted_run().await;
470        let result = TestResult::new(run, Vec::new(), Vec::new(), None);
471
472        assert_eq!(result.output(), &Value::Null);
473        assert!(result.step_names().is_empty());
474    }
475}