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 environment_id: step.environment_id.clone(),
45 }
46 }
47}
48
49#[cfg(test)]
50mod tests {
51 use std::collections::HashMap;
52
53 use ironflow_store::entities::{
54 NewRun, NewStep, StepKind, StepStatus, StepUpdate, TriggerKind,
55 };
56 use ironflow_store::memory::InMemoryStore;
57 use ironflow_store::store::RunStore;
58 use rust_decimal::Decimal;
59 use serde_json::json;
60 use uuid::Uuid;
61
62 use super::*;
63
64 async fn stored_step(output: Option<Value>) -> Step {
66 let store = InMemoryStore::new();
67 let run = store
68 .create_run(NewRun {
69 created_by: None,
70 workflow_name: "collect".to_string(),
71 trigger: TriggerKind::Manual,
72 payload: json!({}),
73 max_retries: 0,
74 handler_version: None,
75 labels: HashMap::new(),
76 scheduled_at: None,
77 idempotency_key: None,
78 concurrency_key: None,
79 priority: 0,
80 concurrency_limits: Vec::new(),
81 max_cost_usd: None,
82 worker_tags: Vec::new(),
83 })
84 .await
85 .expect("create run")
86 .into_run();
87 let step = store
88 .create_step(NewStep {
89 run_id: run.id,
90 trace_id: Uuid::now_v7(),
91 name: "disk".to_string(),
92 kind: StepKind::Shell,
93 position: 0,
94 input: None,
95 is_error_handler: false,
96 })
97 .await
98 .expect("create step");
99 store
100 .update_step(
101 step.id,
102 StepUpdate {
103 status: Some(StepStatus::Running),
104 ..StepUpdate::default()
105 },
106 )
107 .await
108 .expect("to running");
109 if output.is_some() {
110 store
111 .update_step(
112 step.id,
113 StepUpdate {
114 status: Some(StepStatus::Completed),
115 output,
116 duration_ms: Some(12),
117 cost_usd: Some(Decimal::new(3, 2)),
118 ..StepUpdate::default()
119 },
120 )
121 .await
122 .expect("to completed");
123 }
124 store.get_step(step.id).await.expect("get").expect("exists")
125 }
126
127 #[tokio::test]
128 async fn a_stored_shell_step_reads_through_the_accessors() {
129 let step = stored_step(Some(
130 json!({"stdout": "42%\n", "stderr": "", "exit_code": 0}),
131 ))
132 .await;
133
134 let output = StepOutput::from(&step);
135
136 assert_eq!(output.stdout(), "42%\n");
137 assert_eq!(output.exit_code(), Some(0));
138 assert!(output.is_success());
139 assert_eq!(output.duration_ms, 12);
140 assert_eq!(output.cost_usd, Decimal::new(3, 2));
141 }
142
143 #[tokio::test]
144 async fn step_output_error_reads_from_stored_step() {
145 let step = stored_step(Some(json!({"error": "x"}))).await;
146 assert_eq!(StepOutput::from(&step).error(), Some("x"));
147
148 let step = stored_step(None).await;
149 assert_eq!(StepOutput::from(&step).error(), None);
150 }
151
152 #[tokio::test]
153 async fn a_step_without_output_reads_as_empty() {
154 let step = stored_step(None).await;
155
156 let output = StepOutput::from(&step);
157
158 assert_eq!(output.output, Value::Null);
159 assert_eq!(output.stdout(), "");
160 assert_eq!(output.exit_code(), None);
161 assert!(!output.is_success());
162 }
163
164 #[tokio::test]
165 async fn a_stored_step_carries_its_environment_id() {
166 let mut step = stored_step(None).await;
167 assert_eq!(StepOutput::from(&step).environment_id, None);
168
169 step.environment_id = Some("ironflow-env-0a1b2c".to_string());
170 assert_eq!(
171 StepOutput::from(&step).environment_id.as_deref(),
172 Some("ironflow-env-0a1b2c")
173 );
174 }
175}