1use 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::StepResult;
16
17static NULL: Value = Value::Null;
19
20#[derive(Debug, Clone)]
35pub struct TestStep {
36 step: Step,
37 output: Value,
38}
39
40impl TestStep {
41 pub(crate) fn new(step: Step) -> Self {
43 let output = step.output.clone().unwrap_or(Value::Null);
44 Self { step, output }
45 }
46
47 pub fn name(&self) -> &str {
49 &self.step.name
50 }
51
52 pub fn kind(&self) -> &StepKind {
54 &self.step.kind
55 }
56
57 pub fn status(&self) -> StepStatus {
59 self.step.status.state
60 }
61
62 pub fn output(&self) -> &Value {
64 &self.output
65 }
66
67 pub fn input(&self) -> &Value {
69 self.step.input.as_ref().unwrap_or(&NULL)
70 }
71
72 pub fn error(&self) -> Option<&str> {
74 self.step.error.as_deref()
75 }
76
77 pub fn duration(&self) -> Duration {
79 Duration::from_millis(self.step.duration_ms)
80 }
81
82 pub fn cost_usd(&self) -> Decimal {
84 self.step.cost_usd
85 }
86
87 pub fn is_completed(&self) -> bool {
89 self.status() == StepStatus::Completed
90 }
91
92 pub fn is_error_handler(&self) -> bool {
95 self.step.is_error_handler
96 }
97
98 pub fn raw(&self) -> &Step {
100 &self.step
101 }
102}
103
104#[derive(Debug, Clone)]
123pub struct TestResult {
124 run: Run,
125 steps: Vec<TestStep>,
126 step_results: Vec<StepResult>,
127 error: Option<String>,
128}
129
130impl TestResult {
131 pub(crate) fn new(
133 run: Run,
134 steps: Vec<Step>,
135 step_results: Vec<StepResult>,
136 error: Option<String>,
137 ) -> Self {
138 Self {
139 run,
140 steps: steps.into_iter().map(TestStep::new).collect(),
141 step_results,
142 error,
143 }
144 }
145
146 pub fn status(&self) -> RunStatus {
148 self.run.status.state
149 }
150
151 pub fn run(&self) -> &Run {
153 &self.run
154 }
155
156 pub fn run_id(&self) -> Uuid {
158 self.run.id
159 }
160
161 pub fn steps(&self) -> &[TestStep] {
163 &self.steps
164 }
165
166 pub fn step_names(&self) -> Vec<&str> {
168 self.steps.iter().map(TestStep::name).collect()
169 }
170
171 pub fn step(&self, name: &str) -> &TestStep {
191 self.try_step(name).unwrap_or_else(|| {
192 panic!(
193 "no step named {name:?} in this run; steps are {:?}",
194 self.step_names()
195 )
196 })
197 }
198
199 pub fn try_step(&self, name: &str) -> Option<&TestStep> {
201 self.steps.iter().find(|step| step.name() == name)
202 }
203
204 pub fn output(&self) -> &Value {
206 self.steps.last().map_or(&NULL, TestStep::output)
207 }
208
209 pub fn duration(&self) -> Duration {
211 Duration::from_millis(self.run.duration_ms)
212 }
213
214 pub fn cost_usd(&self) -> Decimal {
216 self.run.cost_usd
217 }
218
219 pub fn error(&self) -> Option<&str> {
221 self.error.as_deref()
222 }
223
224 pub fn step_results(&self) -> &[StepResult] {
230 &self.step_results
231 }
232
233 pub fn is_completed(&self) -> bool {
235 self.status() == RunStatus::Completed
236 }
237}
238
239#[cfg(test)]
240mod tests {
241 use std::collections::HashMap;
242 use std::sync::Arc;
243
244 use chrono::Utc;
245 use serde_json::json;
246
247 use ironflow_store::memory::InMemoryStore;
248 use ironflow_store::models::{NewRun, NewStep, StepUpdate, TriggerKind, step_trace_id};
249 use ironflow_store::store::RunStore;
250
251 use super::*;
252
253 async fn persisted_run() -> (Run, Vec<Step>) {
258 let store = Arc::new(InMemoryStore::new());
259 let run = store
260 .create_run(NewRun {
261 created_by: None,
262 workflow_name: "deploy".to_string(),
263 trigger: TriggerKind::Manual,
264 payload: json!({}),
265 max_retries: 0,
266 handler_version: None,
267 labels: HashMap::new(),
268 scheduled_at: None,
269 idempotency_key: None,
270 max_cost_usd: None,
271 })
272 .await
273 .expect("create run")
274 .into_run();
275
276 for (position, name) in [(0u32, "build"), (1, "deploy")] {
277 let step = store
278 .create_step(NewStep {
279 run_id: run.id,
280 trace_id: step_trace_id(run.id, name, position),
281 name: name.to_string(),
282 kind: StepKind::Shell,
283 position,
284 input: Some(json!({"type": "shell", "command": name})),
285 is_error_handler: false,
286 })
287 .await
288 .expect("create step");
289
290 store
291 .update_step(
292 step.id,
293 StepUpdate {
294 status: Some(StepStatus::Running),
295 started_at: Some(Utc::now()),
296 ..StepUpdate::default()
297 },
298 )
299 .await
300 .expect("start step");
301
302 let update = if name == "build" {
303 StepUpdate {
304 status: Some(StepStatus::Completed),
305 output: Some(json!({"stdout": "built", "stderr": "", "exit_code": 0})),
306 duration_ms: Some(12),
307 completed_at: Some(Utc::now()),
308 ..StepUpdate::default()
309 }
310 } else {
311 StepUpdate {
312 status: Some(StepStatus::Failed),
313 error: Some("boom".to_string()),
314 completed_at: Some(Utc::now()),
315 ..StepUpdate::default()
316 }
317 };
318 store
319 .update_step(step.id, update)
320 .await
321 .expect("finish the step");
322 }
323
324 let steps = store.list_steps(run.id).await.expect("list steps");
325 let run = store
326 .get_run(run.id)
327 .await
328 .expect("get run")
329 .expect("the run exists");
330 (run, steps)
331 }
332
333 #[tokio::test]
334 async fn accessors_read_back_what_was_persisted() {
335 let (run, steps) = persisted_run().await;
336 let result = TestResult::new(run, steps, Vec::new(), Some("boom".to_string()));
337
338 assert_eq!(result.step_names(), vec!["build", "deploy"]);
339 assert_eq!(result.error(), Some("boom"));
340 assert!(result.step_results().is_empty());
341 assert!(!result.is_completed());
342
343 let build = result.step("build");
344 assert!(build.is_completed());
345 assert_eq!(build.output()["stdout"], "built");
346 assert_eq!(build.input()["command"], "build");
347 assert_eq!(build.duration(), Duration::from_millis(12));
348 assert_eq!(build.cost_usd(), Decimal::ZERO);
349 assert_eq!(build.kind(), &StepKind::Shell);
350 assert!(!build.is_error_handler());
351 assert_eq!(build.raw().name, "build");
352 assert!(build.error().is_none());
353 }
354
355 #[tokio::test]
356 async fn output_is_the_last_step_output() {
357 let (run, steps) = persisted_run().await;
358 let result = TestResult::new(run, steps, Vec::new(), None);
359
360 assert_eq!(result.output(), &Value::Null);
362 assert_eq!(result.step("deploy").status(), StepStatus::Failed);
363 assert_eq!(result.step("deploy").error(), Some("boom"));
364 }
365
366 #[tokio::test]
367 async fn run_level_accessors_reflect_the_store() {
368 let (run, steps) = persisted_run().await;
369 let run_id = run.id;
370 let result = TestResult::new(run, steps, Vec::new(), None);
371
372 assert_eq!(result.run_id(), run_id);
373 assert_eq!(result.run().workflow_name, "deploy");
374 assert_eq!(result.status(), RunStatus::Pending);
375 assert_eq!(result.duration(), Duration::ZERO);
376 assert_eq!(result.cost_usd(), Decimal::ZERO);
377 assert_eq!(result.steps().len(), 2);
378 }
379
380 #[tokio::test]
381 async fn try_step_returns_none_for_an_unknown_name() {
382 let (run, steps) = persisted_run().await;
383 let result = TestResult::new(run, steps, Vec::new(), None);
384
385 assert!(result.try_step("nope").is_none());
386 }
387
388 #[tokio::test]
389 #[should_panic(expected = "no step named \"nope\"")]
390 async fn step_panics_for_an_unknown_name() {
391 let (run, steps) = persisted_run().await;
392 let result = TestResult::new(run, steps, Vec::new(), None);
393
394 result.step("nope");
395 }
396
397 #[tokio::test]
398 async fn a_run_without_steps_has_a_null_output() {
399 let (run, _) = persisted_run().await;
400 let result = TestResult::new(run, Vec::new(), Vec::new(), None);
401
402 assert_eq!(result.output(), &Value::Null);
403 assert!(result.step_names().is_empty());
404 }
405}