Skip to main content

ironflow_engine/executor/
stored.rs

1//! Read a persisted step back as a [`StepOutput`].
2
3use serde_json::Value;
4
5use ironflow_store::entities::Step;
6
7use super::{StepArtifacts, StepOutput};
8
9/// View a persisted step through the typed [`StepOutput`] accessors.
10///
11/// Useful to read the steps of a sub-workflow run, listed from the store by
12/// [`SubWorkflowOutput::run_id`](super::SubWorkflowOutput::run_id). A step
13/// without output (still running, failed, skipped) reads as an empty output.
14/// The model and the debug conversation are not carried over.
15///
16/// # Examples
17///
18/// ```no_run
19/// use ironflow_engine::context::WorkflowContext;
20/// use ironflow_engine::error::EngineError;
21/// use ironflow_engine::executor::{StepOutput, SubWorkflowOutput};
22///
23/// # async fn example(ctx: &WorkflowContext, child: &SubWorkflowOutput) -> Result<(), EngineError> {
24/// for step in ctx.store().list_steps(child.run_id()).await? {
25///     println!("{}: {}", step.name, StepOutput::from(&step).stdout());
26/// }
27/// # Ok(())
28/// # }
29/// ```
30impl 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    /// A shell step of a fresh run, completed with `output` when given.
65    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}