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 max_cost_usd: None,
78 })
79 .await
80 .expect("create run")
81 .into_run();
82 let step = store
83 .create_step(NewStep {
84 run_id: run.id,
85 trace_id: Uuid::now_v7(),
86 name: "disk".to_string(),
87 kind: StepKind::Shell,
88 position: 0,
89 input: None,
90 is_error_handler: false,
91 })
92 .await
93 .expect("create step");
94 store
95 .update_step(
96 step.id,
97 StepUpdate {
98 status: Some(StepStatus::Running),
99 ..StepUpdate::default()
100 },
101 )
102 .await
103 .expect("to running");
104 if output.is_some() {
105 store
106 .update_step(
107 step.id,
108 StepUpdate {
109 status: Some(StepStatus::Completed),
110 output,
111 duration_ms: Some(12),
112 cost_usd: Some(Decimal::new(3, 2)),
113 ..StepUpdate::default()
114 },
115 )
116 .await
117 .expect("to completed");
118 }
119 store.get_step(step.id).await.expect("get").expect("exists")
120 }
121
122 #[tokio::test]
123 async fn a_stored_shell_step_reads_through_the_accessors() {
124 let step = stored_step(Some(
125 json!({"stdout": "42%\n", "stderr": "", "exit_code": 0}),
126 ))
127 .await;
128
129 let output = StepOutput::from(&step);
130
131 assert_eq!(output.stdout(), "42%\n");
132 assert_eq!(output.exit_code(), Some(0));
133 assert!(output.is_success());
134 assert_eq!(output.duration_ms, 12);
135 assert_eq!(output.cost_usd, Decimal::new(3, 2));
136 }
137
138 #[tokio::test]
139 async fn a_step_without_output_reads_as_empty() {
140 let step = stored_step(None).await;
141
142 let output = StepOutput::from(&step);
143
144 assert_eq!(output.output, Value::Null);
145 assert_eq!(output.stdout(), "");
146 assert_eq!(output.exit_code(), None);
147 assert!(!output.is_success());
148 }
149}