Skip to main content

ironflow_engine/executor/
workflow_output.rs

1//! [`SubWorkflowOutput`] -- what a parent gets back from a sub-workflow.
2
3use rust_decimal::Decimal;
4use serde::{Deserialize, Serialize};
5use uuid::Uuid;
6
7use ironflow_store::entities::RunStatus;
8
9/// Result of a [`workflow`](crate::context::WorkflowContext::workflow) step.
10///
11/// While planning, no child run is created: [`run_id`](Self::run_id) is
12/// [`Uuid::nil`] and the metrics are zero.
13///
14/// It is also the persisted output of the step, so a stored workflow step
15/// reads back with [`StepOutput::json`](super::StepOutput::json).
16///
17/// # Examples
18///
19/// ```
20/// use ironflow_engine::executor::SubWorkflowOutput;
21/// use ironflow_store::entities::RunStatus;
22/// use rust_decimal::Decimal;
23/// use uuid::Uuid;
24///
25/// let run_id = Uuid::now_v7();
26/// let output = SubWorkflowOutput::new(run_id, "collect", RunStatus::Completed, Decimal::ZERO, 1200);
27/// assert_eq!(output.run_id(), run_id);
28/// assert_eq!(output.workflow_name(), "collect");
29/// assert_eq!(output.status(), RunStatus::Completed);
30/// assert_eq!(output.duration_ms(), 1200);
31/// ```
32#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
33pub struct SubWorkflowOutput {
34    run_id: Uuid,
35    workflow_name: String,
36    status: RunStatus,
37    cost_usd: Decimal,
38    duration_ms: u64,
39    #[serde(default, skip_serializing_if = "Option::is_none")]
40    error: Option<String>,
41}
42
43impl SubWorkflowOutput {
44    /// Assemble the result of a child run.
45    ///
46    /// # Examples
47    ///
48    /// ```
49    /// use ironflow_engine::executor::SubWorkflowOutput;
50    /// use ironflow_store::entities::RunStatus;
51    /// use rust_decimal::Decimal;
52    /// use uuid::Uuid;
53    ///
54    /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Warning, Decimal::ONE, 0);
55    /// assert_eq!(output.cost_usd(), Decimal::ONE);
56    /// ```
57    pub fn new(
58        run_id: Uuid,
59        workflow_name: &str,
60        status: RunStatus,
61        cost_usd: Decimal,
62        duration_ms: u64,
63    ) -> Self {
64        Self {
65            run_id,
66            workflow_name: workflow_name.to_string(),
67            status,
68            cost_usd,
69            duration_ms,
70            error: None,
71        }
72    }
73
74    /// Attach the error of a child run tolerated by `allow_failure`.
75    ///
76    /// # Examples
77    ///
78    /// ```
79    /// use ironflow_engine::executor::SubWorkflowOutput;
80    /// use ironflow_store::entities::RunStatus;
81    /// use rust_decimal::Decimal;
82    /// use uuid::Uuid;
83    ///
84    /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Failed, Decimal::ZERO, 0)
85    ///     .with_error("boom");
86    /// assert_eq!(output.error(), Some("boom"));
87    /// ```
88    pub fn with_error(mut self, error: impl Into<String>) -> Self {
89        self.error = Some(error.into());
90        self
91    }
92
93    /// The child run, to read its steps from the store. [`Uuid::nil`] while
94    /// planning.
95    ///
96    /// # Examples
97    ///
98    /// ```
99    /// use ironflow_engine::executor::SubWorkflowOutput;
100    /// use ironflow_store::entities::RunStatus;
101    /// use rust_decimal::Decimal;
102    /// use uuid::Uuid;
103    ///
104    /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Completed, Decimal::ZERO, 0);
105    /// assert!(output.run_id().is_nil());
106    /// ```
107    pub fn run_id(&self) -> Uuid {
108        self.run_id
109    }
110
111    /// Name of the child workflow.
112    pub fn workflow_name(&self) -> &str {
113        &self.workflow_name
114    }
115
116    /// Final status of the child run: `Completed`, `Warning` when one of its
117    /// `allow_failure` steps failed, or (only when the step was started with
118    /// `allow_failure`) `Failed` / `Cancelled`.
119    pub fn status(&self) -> RunStatus {
120        self.status
121    }
122
123    /// Cost of the child run, in USD. Already included in the parent's cost.
124    pub fn cost_usd(&self) -> Decimal {
125        self.cost_usd
126    }
127
128    /// Wall-clock duration of the child run, in milliseconds.
129    pub fn duration_ms(&self) -> u64 {
130        self.duration_ms
131    }
132
133    /// Error of a failed or cancelled child run tolerated by `allow_failure`,
134    /// `None` otherwise.
135    ///
136    /// # Examples
137    ///
138    /// ```
139    /// use ironflow_engine::executor::SubWorkflowOutput;
140    /// use ironflow_store::entities::RunStatus;
141    /// use rust_decimal::Decimal;
142    /// use uuid::Uuid;
143    ///
144    /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Completed, Decimal::ZERO, 0);
145    /// assert_eq!(output.error(), None);
146    /// ```
147    pub fn error(&self) -> Option<&str> {
148        self.error.as_deref()
149    }
150}
151
152#[cfg(test)]
153mod tests {
154    use serde_json::{from_value, json, to_value};
155
156    use super::*;
157
158    #[test]
159    fn serializes_to_the_persisted_step_output() {
160        let run_id = Uuid::now_v7();
161        let output = SubWorkflowOutput::new(
162            run_id,
163            "collect",
164            RunStatus::Warning,
165            Decimal::new(25, 2),
166            1200,
167        );
168
169        assert_eq!(
170            to_value(&output).expect("serialize"),
171            json!({
172                "run_id": run_id,
173                "workflow_name": "collect",
174                "status": "warning",
175                "cost_usd": 0.25,
176                "duration_ms": 1200,
177            })
178        );
179    }
180
181    #[test]
182    fn a_persisted_step_output_reads_back() {
183        let run_id = Uuid::now_v7();
184        let stored = json!({
185            "run_id": run_id,
186            "workflow_name": "collect",
187            "status": "completed",
188            "cost_usd": 0,
189            "duration_ms": 7,
190        });
191
192        let output: SubWorkflowOutput = from_value(stored).expect("deserialize");
193
194        assert_eq!(output.run_id(), run_id);
195        assert_eq!(output.status(), RunStatus::Completed);
196        assert_eq!(output.cost_usd(), Decimal::ZERO);
197        assert_eq!(output.duration_ms(), 7);
198    }
199
200    #[test]
201    fn an_error_roundtrips() {
202        let output = SubWorkflowOutput::new(
203            Uuid::now_v7(),
204            "collect",
205            RunStatus::Failed,
206            Decimal::ZERO,
207            3,
208        )
209        .with_error("boom");
210
211        let back: SubWorkflowOutput =
212            from_value(to_value(&output).expect("serialize")).expect("deserialize");
213
214        assert_eq!(back, output);
215        assert_eq!(back.error(), Some("boom"));
216    }
217
218    #[test]
219    fn a_legacy_output_without_error_reads_back_as_none() {
220        let output: SubWorkflowOutput = from_value(json!({
221            "run_id": Uuid::now_v7(),
222            "workflow_name": "collect",
223            "status": "completed",
224            "cost_usd": 0,
225            "duration_ms": 7,
226        }))
227        .expect("deserialize");
228
229        assert_eq!(output.error(), None);
230    }
231
232    #[test]
233    fn an_output_without_run_id_is_refused() {
234        let result = from_value::<SubWorkflowOutput>(json!({
235            "workflow_name": "collect",
236            "status": "completed",
237            "cost_usd": 0,
238            "duration_ms": 0,
239        }));
240        assert!(result.is_err());
241    }
242}