use std::time::Duration;
use rust_decimal::Decimal;
use serde_json::Value;
use uuid::Uuid;
use ironflow_store::models::{Run, RunStatus, Step, StepKind, StepStatus};
use crate::executor::{StepOutput, StepResult};
static NULL: Value = Value::Null;
#[derive(Debug, Clone)]
pub struct TestStep {
step: Step,
output: Value,
}
impl TestStep {
pub(crate) fn new(step: Step) -> Self {
let output = step.output.clone().unwrap_or(Value::Null);
Self { step, output }
}
pub fn name(&self) -> &str {
&self.step.name
}
pub fn kind(&self) -> &StepKind {
&self.step.kind
}
pub fn status(&self) -> StepStatus {
self.step.status.state
}
pub fn output(&self) -> &Value {
&self.output
}
pub fn input(&self) -> &Value {
self.step.input.as_ref().unwrap_or(&NULL)
}
pub fn error(&self) -> Option<&str> {
self.step.error.as_deref()
}
pub fn duration(&self) -> Duration {
Duration::from_millis(self.step.duration_ms)
}
pub fn cost_usd(&self) -> Decimal {
self.step.cost_usd
}
pub fn is_completed(&self) -> bool {
self.status() == StepStatus::Completed
}
pub fn is_error_handler(&self) -> bool {
self.step.is_error_handler
}
pub fn raw(&self) -> &Step {
&self.step
}
pub fn step_output(&self) -> StepOutput {
StepOutput::from(&self.step)
}
}
#[derive(Debug, Clone)]
pub struct TestResult {
run: Run,
steps: Vec<TestStep>,
step_results: Vec<StepResult>,
error: Option<String>,
}
impl TestResult {
pub(crate) fn new(
run: Run,
steps: Vec<Step>,
step_results: Vec<StepResult>,
error: Option<String>,
) -> Self {
Self {
run,
steps: steps.into_iter().map(TestStep::new).collect(),
step_results,
error,
}
}
pub fn status(&self) -> RunStatus {
self.run.status.state
}
pub fn run(&self) -> &Run {
&self.run
}
pub fn run_id(&self) -> Uuid {
self.run.id
}
pub fn steps(&self) -> &[TestStep] {
&self.steps
}
pub fn step_names(&self) -> Vec<&str> {
self.steps.iter().map(TestStep::name).collect()
}
pub fn step(&self, name: &str) -> &TestStep {
self.try_step(name).unwrap_or_else(|| {
panic!(
"no step named {name:?} in this run; steps are {:?}",
self.step_names()
)
})
}
pub fn try_step(&self, name: &str) -> Option<&TestStep> {
self.steps.iter().find(|step| step.name() == name)
}
pub fn output(&self) -> &Value {
self.steps.last().map_or(&NULL, TestStep::output)
}
pub fn duration(&self) -> Duration {
Duration::from_millis(self.run.duration_ms)
}
pub fn cost_usd(&self) -> Decimal {
self.run.cost_usd
}
pub fn error(&self) -> Option<&str> {
self.error.as_deref()
}
pub fn step_results(&self) -> &[StepResult] {
&self.step_results
}
pub fn is_completed(&self) -> bool {
self.status() == RunStatus::Completed
}
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use std::sync::Arc;
use chrono::Utc;
use serde_json::json;
use ironflow_store::memory::InMemoryStore;
use ironflow_store::models::{NewRun, NewStep, StepUpdate, TriggerKind, step_trace_id};
use ironflow_store::store::RunStore;
use super::*;
async fn persisted_run() -> (Run, Vec<Step>) {
let store = Arc::new(InMemoryStore::new());
let run = store
.create_run(NewRun {
created_by: None,
workflow_name: "deploy".to_string(),
trigger: TriggerKind::Manual,
payload: json!({}),
max_retries: 0,
handler_version: None,
labels: HashMap::new(),
scheduled_at: None,
idempotency_key: None,
max_cost_usd: None,
})
.await
.expect("create run")
.into_run();
for (position, name) in [(0u32, "build"), (1, "deploy")] {
let step = store
.create_step(NewStep {
run_id: run.id,
trace_id: step_trace_id(run.id, name, position),
name: name.to_string(),
kind: StepKind::Shell,
position,
input: Some(json!({"type": "shell", "command": name})),
is_error_handler: false,
})
.await
.expect("create step");
store
.update_step(
step.id,
StepUpdate {
status: Some(StepStatus::Running),
started_at: Some(Utc::now()),
..StepUpdate::default()
},
)
.await
.expect("start step");
let update = if name == "build" {
StepUpdate {
status: Some(StepStatus::Completed),
output: Some(json!({"stdout": "built", "stderr": "", "exit_code": 0})),
duration_ms: Some(12),
completed_at: Some(Utc::now()),
..StepUpdate::default()
}
} else {
StepUpdate {
status: Some(StepStatus::Failed),
error: Some("boom".to_string()),
completed_at: Some(Utc::now()),
..StepUpdate::default()
}
};
store
.update_step(step.id, update)
.await
.expect("finish the step");
}
let steps = store.list_steps(run.id).await.expect("list steps");
let run = store
.get_run(run.id)
.await
.expect("get run")
.expect("the run exists");
(run, steps)
}
#[tokio::test]
async fn accessors_read_back_what_was_persisted() {
let (run, steps) = persisted_run().await;
let result = TestResult::new(run, steps, Vec::new(), Some("boom".to_string()));
assert_eq!(result.step_names(), vec!["build", "deploy"]);
assert_eq!(result.error(), Some("boom"));
assert!(result.step_results().is_empty());
assert!(!result.is_completed());
let build = result.step("build");
assert!(build.is_completed());
assert_eq!(build.output()["stdout"], "built");
assert_eq!(build.input()["command"], "build");
assert_eq!(build.duration(), Duration::from_millis(12));
assert_eq!(build.cost_usd(), Decimal::ZERO);
assert_eq!(build.kind(), &StepKind::Shell);
assert!(!build.is_error_handler());
assert_eq!(build.raw().name, "build");
assert!(build.error().is_none());
}
#[tokio::test]
async fn output_is_the_last_step_output() {
let (run, steps) = persisted_run().await;
let result = TestResult::new(run, steps, Vec::new(), None);
assert_eq!(result.output(), &Value::Null);
assert_eq!(result.step("deploy").status(), StepStatus::Failed);
assert_eq!(result.step("deploy").error(), Some("boom"));
}
#[tokio::test]
async fn run_level_accessors_reflect_the_store() {
let (run, steps) = persisted_run().await;
let run_id = run.id;
let result = TestResult::new(run, steps, Vec::new(), None);
assert_eq!(result.run_id(), run_id);
assert_eq!(result.run().workflow_name, "deploy");
assert_eq!(result.status(), RunStatus::Pending);
assert_eq!(result.duration(), Duration::ZERO);
assert_eq!(result.cost_usd(), Decimal::ZERO);
assert_eq!(result.steps().len(), 2);
}
#[tokio::test]
async fn try_step_returns_none_for_an_unknown_name() {
let (run, steps) = persisted_run().await;
let result = TestResult::new(run, steps, Vec::new(), None);
assert!(result.try_step("nope").is_none());
}
#[tokio::test]
#[should_panic(expected = "no step named \"nope\"")]
async fn step_panics_for_an_unknown_name() {
let (run, steps) = persisted_run().await;
let result = TestResult::new(run, steps, Vec::new(), None);
result.step("nope");
}
#[tokio::test]
async fn a_run_without_steps_has_a_null_output() {
let (run, _) = persisted_run().await;
let result = TestResult::new(run, Vec::new(), Vec::new(), None);
assert_eq!(result.output(), &Value::Null);
assert!(result.step_names().is_empty());
}
}