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}