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::{StepOutput, 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 pub fn step_output(&self) -> StepOutput {
120 StepOutput::from(&self.step)
121 }
122}
123
124#[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 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 pub fn status(&self) -> RunStatus {
168 self.run.status.state
169 }
170
171 pub fn run(&self) -> &Run {
173 &self.run
174 }
175
176 pub fn run_id(&self) -> Uuid {
178 self.run.id
179 }
180
181 pub fn steps(&self) -> &[TestStep] {
183 &self.steps
184 }
185
186 pub fn step_names(&self) -> Vec<&str> {
188 self.steps.iter().map(TestStep::name).collect()
189 }
190
191 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 pub fn try_step(&self, name: &str) -> Option<&TestStep> {
221 self.steps.iter().find(|step| step.name() == name)
222 }
223
224 pub fn output(&self) -> &Value {
242 self.run.output.as_ref().unwrap_or(&NULL)
243 }
244
245 pub fn duration(&self) -> Duration {
247 Duration::from_millis(self.run.duration_ms)
248 }
249
250 pub fn cost_usd(&self) -> Decimal {
252 self.run.cost_usd
253 }
254
255 pub fn error(&self) -> Option<&str> {
257 self.error.as_deref()
258 }
259
260 pub fn step_results(&self) -> &[StepResult] {
266 &self.step_results
267 }
268
269 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 async fn persisted_run() -> (Run, Vec<Step>) {
296 persisted_run_with_output(None).await
297 }
298
299 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 priority: 0,
315 concurrency_limits: Vec::new(),
316 max_cost_usd: None,
317 worker_tags: Vec::new(),
318 })
319 .await
320 .expect("create run")
321 .into_run();
322
323 for (position, name) in [(0u32, "build"), (1, "deploy")] {
324 let step = store
325 .create_step(NewStep {
326 run_id: run.id,
327 trace_id: step_trace_id(run.id, name, position),
328 name: name.to_string(),
329 kind: StepKind::Shell,
330 position,
331 input: Some(json!({"type": "shell", "command": name})),
332 is_error_handler: false,
333 })
334 .await
335 .expect("create step");
336
337 store
338 .update_step(
339 step.id,
340 StepUpdate {
341 status: Some(StepStatus::Running),
342 started_at: Some(Utc::now()),
343 ..StepUpdate::default()
344 },
345 )
346 .await
347 .expect("start step");
348
349 let update = if name == "build" {
350 StepUpdate {
351 status: Some(StepStatus::Completed),
352 output: Some(json!({"stdout": "built", "stderr": "", "exit_code": 0})),
353 duration_ms: Some(12),
354 completed_at: Some(Utc::now()),
355 ..StepUpdate::default()
356 }
357 } else {
358 StepUpdate {
359 status: Some(StepStatus::Failed),
360 error: Some("boom".to_string()),
361 completed_at: Some(Utc::now()),
362 ..StepUpdate::default()
363 }
364 };
365 store
366 .update_step(step.id, update)
367 .await
368 .expect("finish the step");
369 }
370
371 store
372 .update_run(
373 run.id,
374 RunUpdate {
375 output,
376 ..RunUpdate::default()
377 },
378 )
379 .await
380 .expect("record the run output");
381
382 let steps = store.list_steps(run.id).await.expect("list steps");
383 let run = store
384 .get_run(run.id)
385 .await
386 .expect("get run")
387 .expect("the run exists");
388 (run, steps)
389 }
390
391 #[tokio::test]
392 async fn accessors_read_back_what_was_persisted() {
393 let (run, steps) = persisted_run().await;
394 let result = TestResult::new(run, steps, Vec::new(), Some("boom".to_string()));
395
396 assert_eq!(result.step_names(), vec!["build", "deploy"]);
397 assert_eq!(result.error(), Some("boom"));
398 assert!(result.step_results().is_empty());
399 assert!(!result.is_completed());
400
401 let build = result.step("build");
402 assert!(build.is_completed());
403 assert_eq!(build.output()["stdout"], "built");
404 assert_eq!(build.input()["command"], "build");
405 assert_eq!(build.duration(), Duration::from_millis(12));
406 assert_eq!(build.cost_usd(), Decimal::ZERO);
407 assert_eq!(build.kind(), &StepKind::Shell);
408 assert!(!build.is_error_handler());
409 assert_eq!(build.raw().name, "build");
410 assert!(build.error().is_none());
411 }
412
413 #[tokio::test]
414 async fn set_output_is_the_test_result_output() {
415 let (run, steps) = persisted_run_with_output(Some(json!({"approved": true}))).await;
416 let result = TestResult::new(run, steps, Vec::new(), None);
417
418 assert_eq!(result.output(), &json!({"approved": true}));
419 assert_eq!(result.step("deploy").output(), &Value::Null);
421 assert_eq!(result.steps().last().map(TestStep::name), Some("deploy"));
422 }
423
424 #[tokio::test]
425 async fn output_is_null_when_the_handler_set_none() {
426 let (run, steps) = persisted_run().await;
427 let result = TestResult::new(run, steps, Vec::new(), None);
428
429 assert_eq!(result.output(), &Value::Null);
431 assert_eq!(result.step("build").output()["stdout"], "built");
432 assert_eq!(result.step("deploy").status(), StepStatus::Failed);
434 assert_eq!(result.step("deploy").error(), Some("boom"));
435 }
436
437 #[tokio::test]
438 async fn run_level_accessors_reflect_the_store() {
439 let (run, steps) = persisted_run().await;
440 let run_id = run.id;
441 let result = TestResult::new(run, steps, Vec::new(), None);
442
443 assert_eq!(result.run_id(), run_id);
444 assert_eq!(result.run().workflow_name, "deploy");
445 assert_eq!(result.status(), RunStatus::Pending);
446 assert_eq!(result.duration(), Duration::ZERO);
447 assert_eq!(result.cost_usd(), Decimal::ZERO);
448 assert_eq!(result.steps().len(), 2);
449 }
450
451 #[tokio::test]
452 async fn try_step_returns_none_for_an_unknown_name() {
453 let (run, steps) = persisted_run().await;
454 let result = TestResult::new(run, steps, Vec::new(), None);
455
456 assert!(result.try_step("nope").is_none());
457 }
458
459 #[tokio::test]
460 #[should_panic(expected = "no step named \"nope\"")]
461 async fn step_panics_for_an_unknown_name() {
462 let (run, steps) = persisted_run().await;
463 let result = TestResult::new(run, steps, Vec::new(), None);
464
465 result.step("nope");
466 }
467
468 #[tokio::test]
469 async fn a_run_without_steps_has_a_null_output() {
470 let (run, _) = persisted_run().await;
471 let result = TestResult::new(run, Vec::new(), Vec::new(), None);
472
473 assert_eq!(result.output(), &Value::Null);
474 assert!(result.step_names().is_empty());
475 }
476}