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