ironflow_engine/executor/
stored.rs1use serde_json::Value;
4
5use ironflow_store::entities::Step;
6
7use super::{StepArtifacts, StepOutput};
8
9impl From<&Step> for StepOutput {
31 fn from(step: &Step) -> Self {
32 Self {
33 output: step.output.clone().unwrap_or(Value::Null),
34 duration_ms: step.duration_ms,
35 cost_usd: step.cost_usd,
36 input_tokens: step.input_tokens,
37 cache_read_input_tokens: step.cache_read_input_tokens,
38 cache_creation_input_tokens: step.cache_creation_input_tokens,
39 output_tokens: step.output_tokens,
40 model: None,
41 debug_messages: None,
42 artifacts: StepArtifacts::default(),
43 account_id: step.account_id,
44 }
45 }
46}
47
48#[cfg(test)]
49mod tests {
50 use std::collections::HashMap;
51
52 use ironflow_store::entities::{
53 NewRun, NewStep, StepKind, StepStatus, StepUpdate, TriggerKind,
54 };
55 use ironflow_store::memory::InMemoryStore;
56 use ironflow_store::store::RunStore;
57 use rust_decimal::Decimal;
58 use serde_json::json;
59 use uuid::Uuid;
60
61 use super::*;
62
63 async fn stored_step(output: Option<Value>) -> Step {
65 let store = InMemoryStore::new();
66 let run = store
67 .create_run(NewRun {
68 created_by: None,
69 workflow_name: "collect".to_string(),
70 trigger: TriggerKind::Manual,
71 payload: json!({}),
72 max_retries: 0,
73 handler_version: None,
74 labels: HashMap::new(),
75 scheduled_at: None,
76 idempotency_key: None,
77 concurrency_key: None,
78 max_cost_usd: None,
79 })
80 .await
81 .expect("create run")
82 .into_run();
83 let step = store
84 .create_step(NewStep {
85 run_id: run.id,
86 trace_id: Uuid::now_v7(),
87 name: "disk".to_string(),
88 kind: StepKind::Shell,
89 position: 0,
90 input: None,
91 is_error_handler: false,
92 })
93 .await
94 .expect("create step");
95 store
96 .update_step(
97 step.id,
98 StepUpdate {
99 status: Some(StepStatus::Running),
100 ..StepUpdate::default()
101 },
102 )
103 .await
104 .expect("to running");
105 if output.is_some() {
106 store
107 .update_step(
108 step.id,
109 StepUpdate {
110 status: Some(StepStatus::Completed),
111 output,
112 duration_ms: Some(12),
113 cost_usd: Some(Decimal::new(3, 2)),
114 ..StepUpdate::default()
115 },
116 )
117 .await
118 .expect("to completed");
119 }
120 store.get_step(step.id).await.expect("get").expect("exists")
121 }
122
123 #[tokio::test]
124 async fn a_stored_shell_step_reads_through_the_accessors() {
125 let step = stored_step(Some(
126 json!({"stdout": "42%\n", "stderr": "", "exit_code": 0}),
127 ))
128 .await;
129
130 let output = StepOutput::from(&step);
131
132 assert_eq!(output.stdout(), "42%\n");
133 assert_eq!(output.exit_code(), Some(0));
134 assert!(output.is_success());
135 assert_eq!(output.duration_ms, 12);
136 assert_eq!(output.cost_usd, Decimal::new(3, 2));
137 }
138
139 #[tokio::test]
140 async fn step_output_error_reads_from_stored_step() {
141 let step = stored_step(Some(json!({"error": "x"}))).await;
142 assert_eq!(StepOutput::from(&step).error(), Some("x"));
143
144 let step = stored_step(None).await;
145 assert_eq!(StepOutput::from(&step).error(), None);
146 }
147
148 #[tokio::test]
149 async fn a_step_without_output_reads_as_empty() {
150 let step = stored_step(None).await;
151
152 let output = StepOutput::from(&step);
153
154 assert_eq!(output.output, Value::Null);
155 assert_eq!(output.stdout(), "");
156 assert_eq!(output.exit_code(), None);
157 assert!(!output.is_success());
158 }
159}