ironflow_engine/executor/workflow_output.rs
1//! [`SubWorkflowOutput`] -- what a parent gets back from a sub-workflow.
2//!
3//! [`SubWorkflowOutcome`] adds the case of a sub-workflow started with a
4//! concurrency key that another active run already holds.
5
6use std::fmt;
7
8use rust_decimal::Decimal;
9use serde::{Deserialize, Serialize};
10use uuid::Uuid;
11
12use ironflow_store::entities::RunStatus;
13
14/// Result of a [`workflow`](crate::context::WorkflowContext::workflow) step.
15///
16/// While planning, no child run is created: [`run_id`](Self::run_id) is
17/// [`Uuid::nil`] and the metrics are zero.
18///
19/// It is also the persisted output of the step, so a stored workflow step
20/// reads back with [`StepOutput::json`](super::StepOutput::json).
21///
22/// # Examples
23///
24/// ```
25/// use ironflow_engine::executor::SubWorkflowOutput;
26/// use ironflow_store::entities::RunStatus;
27/// use rust_decimal::Decimal;
28/// use uuid::Uuid;
29///
30/// let run_id = Uuid::now_v7();
31/// let output = SubWorkflowOutput::new(run_id, "collect", RunStatus::Completed, Decimal::ZERO, 1200);
32/// assert_eq!(output.run_id(), run_id);
33/// assert_eq!(output.workflow_name(), "collect");
34/// assert_eq!(output.status(), RunStatus::Completed);
35/// assert_eq!(output.duration_ms(), 1200);
36/// ```
37#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
38pub struct SubWorkflowOutput {
39 run_id: Uuid,
40 workflow_name: String,
41 status: RunStatus,
42 cost_usd: Decimal,
43 duration_ms: u64,
44 #[serde(default, skip_serializing_if = "Option::is_none")]
45 error: Option<String>,
46}
47
48impl SubWorkflowOutput {
49 /// Assemble the result of a child run.
50 ///
51 /// # Examples
52 ///
53 /// ```
54 /// use ironflow_engine::executor::SubWorkflowOutput;
55 /// use ironflow_store::entities::RunStatus;
56 /// use rust_decimal::Decimal;
57 /// use uuid::Uuid;
58 ///
59 /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Warning, Decimal::ONE, 0);
60 /// assert_eq!(output.cost_usd(), Decimal::ONE);
61 /// ```
62 pub fn new(
63 run_id: Uuid,
64 workflow_name: &str,
65 status: RunStatus,
66 cost_usd: Decimal,
67 duration_ms: u64,
68 ) -> Self {
69 Self {
70 run_id,
71 workflow_name: workflow_name.to_string(),
72 status,
73 cost_usd,
74 duration_ms,
75 error: None,
76 }
77 }
78
79 /// Attach the error of a child run tolerated by `allow_failure`.
80 ///
81 /// # Examples
82 ///
83 /// ```
84 /// use ironflow_engine::executor::SubWorkflowOutput;
85 /// use ironflow_store::entities::RunStatus;
86 /// use rust_decimal::Decimal;
87 /// use uuid::Uuid;
88 ///
89 /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Failed, Decimal::ZERO, 0)
90 /// .with_error("boom");
91 /// assert_eq!(output.error(), Some("boom"));
92 /// ```
93 pub fn with_error(mut self, error: impl Into<String>) -> Self {
94 self.error = Some(error.into());
95 self
96 }
97
98 /// The child run, to read its steps from the store. [`Uuid::nil`] while
99 /// planning.
100 ///
101 /// # Examples
102 ///
103 /// ```
104 /// use ironflow_engine::executor::SubWorkflowOutput;
105 /// use ironflow_store::entities::RunStatus;
106 /// use rust_decimal::Decimal;
107 /// use uuid::Uuid;
108 ///
109 /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Completed, Decimal::ZERO, 0);
110 /// assert!(output.run_id().is_nil());
111 /// ```
112 pub fn run_id(&self) -> Uuid {
113 self.run_id
114 }
115
116 /// Name of the child workflow.
117 pub fn workflow_name(&self) -> &str {
118 &self.workflow_name
119 }
120
121 /// Final status of the child run: `Completed`, `Warning` when one of its
122 /// `allow_failure` steps failed, or (only when the step was started with
123 /// `allow_failure`) `Failed` / `Cancelled`.
124 pub fn status(&self) -> RunStatus {
125 self.status
126 }
127
128 /// Cost of the child run, in USD. Already included in the parent's cost.
129 pub fn cost_usd(&self) -> Decimal {
130 self.cost_usd
131 }
132
133 /// Wall-clock duration of the child run, in milliseconds.
134 pub fn duration_ms(&self) -> u64 {
135 self.duration_ms
136 }
137
138 /// Error of a failed or cancelled child run tolerated by `allow_failure`,
139 /// `None` otherwise.
140 ///
141 /// # Examples
142 ///
143 /// ```
144 /// use ironflow_engine::executor::SubWorkflowOutput;
145 /// use ironflow_store::entities::RunStatus;
146 /// use rust_decimal::Decimal;
147 /// use uuid::Uuid;
148 ///
149 /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Completed, Decimal::ZERO, 0);
150 /// assert_eq!(output.error(), None);
151 /// ```
152 pub fn error(&self) -> Option<&str> {
153 self.error.as_deref()
154 }
155}
156
157/// A sub-workflow skipped because its concurrency key is held by another
158/// active run.
159///
160/// Persisted as the step output `{"concurrency_conflict": {"key": .., "run_id": ..}}`
161/// and replayed as-is on resume.
162///
163/// # Examples
164///
165/// ```
166/// use ironflow_engine::executor::ConcurrencyConflict;
167/// use uuid::Uuid;
168///
169/// let holder = Uuid::now_v7();
170/// let conflict = ConcurrencyConflict::new("issue:12", holder);
171/// assert_eq!(conflict.key(), "issue:12");
172/// assert_eq!(conflict.run_id(), holder);
173/// assert!(conflict.to_string().contains("issue:12"));
174/// ```
175#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
176pub struct ConcurrencyConflict {
177 key: String,
178 run_id: Uuid,
179}
180
181impl ConcurrencyConflict {
182 /// Record that `run_id` holds `key`.
183 ///
184 /// # Examples
185 ///
186 /// ```
187 /// use ironflow_engine::executor::ConcurrencyConflict;
188 /// use uuid::Uuid;
189 ///
190 /// let conflict = ConcurrencyConflict::new("deploy:prod", Uuid::nil());
191 /// assert_eq!(conflict.key(), "deploy:prod");
192 /// ```
193 pub fn new(key: impl Into<String>, run_id: Uuid) -> Self {
194 Self {
195 key: key.into(),
196 run_id,
197 }
198 }
199
200 /// The contested concurrency key.
201 ///
202 /// # Examples
203 ///
204 /// ```
205 /// use ironflow_engine::executor::ConcurrencyConflict;
206 /// use uuid::Uuid;
207 ///
208 /// assert_eq!(ConcurrencyConflict::new("k", Uuid::nil()).key(), "k");
209 /// ```
210 pub fn key(&self) -> &str {
211 &self.key
212 }
213
214 /// The active run holding the key, at the time of the conflict.
215 ///
216 /// # Examples
217 ///
218 /// ```
219 /// use ironflow_engine::executor::ConcurrencyConflict;
220 /// use uuid::Uuid;
221 ///
222 /// let holder = Uuid::now_v7();
223 /// assert_eq!(ConcurrencyConflict::new("k", holder).run_id(), holder);
224 /// ```
225 pub fn run_id(&self) -> Uuid {
226 self.run_id
227 }
228}
229
230impl fmt::Display for ConcurrencyConflict {
231 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
232 write!(
233 f,
234 "conflict on concurrency key {:?} (run {})",
235 self.key, self.run_id
236 )
237 }
238}
239
240/// Result of a
241/// [`workflow_with`](crate::context::WorkflowContext::workflow_with) step.
242///
243/// # Examples
244///
245/// ```
246/// use ironflow_engine::executor::{ConcurrencyConflict, SubWorkflowOutcome};
247/// use uuid::Uuid;
248///
249/// let outcome = SubWorkflowOutcome::Conflict(ConcurrencyConflict::new("issue:12", Uuid::nil()));
250/// match &outcome {
251/// SubWorkflowOutcome::Completed(output) => println!("child run {}", output.run_id()),
252/// SubWorkflowOutcome::Conflict(conflict) => println!("held by {}", conflict.run_id()),
253/// }
254/// assert!(outcome.output().is_none());
255/// ```
256#[derive(Debug, Clone, PartialEq)]
257pub enum SubWorkflowOutcome {
258 /// The child run was created and finished.
259 Completed(SubWorkflowOutput),
260 /// No child run was created: another active run holds the concurrency key.
261 Conflict(ConcurrencyConflict),
262}
263
264impl SubWorkflowOutcome {
265 /// The conflict, when no child run was created.
266 ///
267 /// # Examples
268 ///
269 /// ```
270 /// use ironflow_engine::executor::{ConcurrencyConflict, SubWorkflowOutcome};
271 /// use uuid::Uuid;
272 ///
273 /// let outcome = SubWorkflowOutcome::Conflict(ConcurrencyConflict::new("k", Uuid::nil()));
274 /// assert_eq!(outcome.conflict().map(|c| c.key()), Some("k"));
275 /// ```
276 pub fn conflict(&self) -> Option<&ConcurrencyConflict> {
277 match self {
278 SubWorkflowOutcome::Conflict(conflict) => Some(conflict),
279 SubWorkflowOutcome::Completed(_) => None,
280 }
281 }
282
283 /// The child result, when the child run was created and finished.
284 ///
285 /// # Examples
286 ///
287 /// ```
288 /// use ironflow_engine::executor::{SubWorkflowOutcome, SubWorkflowOutput};
289 /// use ironflow_store::entities::RunStatus;
290 /// use rust_decimal::Decimal;
291 /// use uuid::Uuid;
292 ///
293 /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Completed, Decimal::ZERO, 0);
294 /// let outcome = SubWorkflowOutcome::Completed(output.clone());
295 /// assert_eq!(outcome.output(), Some(&output));
296 /// assert!(outcome.conflict().is_none());
297 /// ```
298 pub fn output(&self) -> Option<&SubWorkflowOutput> {
299 match self {
300 SubWorkflowOutcome::Completed(output) => Some(output),
301 SubWorkflowOutcome::Conflict(_) => None,
302 }
303 }
304}
305
306/// The persisted output of a completed `Workflow` step, as read back on replay.
307///
308/// `Conflict` is tried first: a flat [`SubWorkflowOutput`] has no
309/// `concurrency_conflict` field, so outputs recorded before concurrency keys
310/// existed still read back as `Completed`.
311#[derive(Debug, Deserialize)]
312#[serde(untagged)]
313pub(crate) enum RecordedWorkflowStep {
314 /// The step was skipped on a concurrency conflict.
315 Conflict {
316 /// The recorded conflict.
317 concurrency_conflict: ConcurrencyConflict,
318 },
319 /// The child run finished.
320 Completed(SubWorkflowOutput),
321}
322
323impl From<RecordedWorkflowStep> for SubWorkflowOutcome {
324 fn from(recorded: RecordedWorkflowStep) -> Self {
325 match recorded {
326 RecordedWorkflowStep::Conflict {
327 concurrency_conflict,
328 } => SubWorkflowOutcome::Conflict(concurrency_conflict),
329 RecordedWorkflowStep::Completed(output) => SubWorkflowOutcome::Completed(output),
330 }
331 }
332}
333
334#[cfg(test)]
335mod tests {
336 use serde_json::{from_value, json, to_value};
337
338 use super::*;
339
340 #[test]
341 fn serializes_to_the_persisted_step_output() {
342 let run_id = Uuid::now_v7();
343 let output = SubWorkflowOutput::new(
344 run_id,
345 "collect",
346 RunStatus::Warning,
347 Decimal::new(25, 2),
348 1200,
349 );
350
351 assert_eq!(
352 to_value(&output).expect("serialize"),
353 json!({
354 "run_id": run_id,
355 "workflow_name": "collect",
356 "status": "warning",
357 "cost_usd": 0.25,
358 "duration_ms": 1200,
359 })
360 );
361 }
362
363 #[test]
364 fn a_persisted_step_output_reads_back() {
365 let run_id = Uuid::now_v7();
366 let stored = json!({
367 "run_id": run_id,
368 "workflow_name": "collect",
369 "status": "completed",
370 "cost_usd": 0,
371 "duration_ms": 7,
372 });
373
374 let output: SubWorkflowOutput = from_value(stored).expect("deserialize");
375
376 assert_eq!(output.run_id(), run_id);
377 assert_eq!(output.status(), RunStatus::Completed);
378 assert_eq!(output.cost_usd(), Decimal::ZERO);
379 assert_eq!(output.duration_ms(), 7);
380 }
381
382 #[test]
383 fn a_recorded_conflict_round_trips() {
384 let holder = Uuid::now_v7();
385 let conflict = ConcurrencyConflict::new("issue:12", holder);
386 let stored = json!({ "concurrency_conflict": conflict });
387
388 assert_eq!(
389 stored,
390 json!({ "concurrency_conflict": { "key": "issue:12", "run_id": holder } })
391 );
392
393 let recorded: RecordedWorkflowStep = from_value(stored).expect("deserialize");
394 let outcome = SubWorkflowOutcome::from(recorded);
395 assert_eq!(outcome, SubWorkflowOutcome::Conflict(conflict));
396 assert_eq!(
397 outcome.conflict().map(ConcurrencyConflict::run_id),
398 Some(holder)
399 );
400 assert!(outcome.output().is_none());
401 }
402
403 #[test]
404 fn a_flat_output_still_reads_back_as_completed() {
405 let run_id = Uuid::now_v7();
406 let stored = json!({
407 "run_id": run_id,
408 "workflow_name": "collect",
409 "status": "completed",
410 "cost_usd": 0,
411 "duration_ms": 7,
412 });
413
414 let recorded: RecordedWorkflowStep = from_value(stored).expect("deserialize");
415 let outcome = SubWorkflowOutcome::from(recorded);
416 let output = outcome.output().expect("completed outcome");
417 assert_eq!(output.run_id(), run_id);
418 assert_eq!(output.duration_ms(), 7);
419 assert!(outcome.conflict().is_none());
420 }
421
422 #[test]
423 fn conflict_display_names_the_key_and_the_holder() {
424 let holder = Uuid::now_v7();
425 let text = ConcurrencyConflict::new("issue:12", holder).to_string();
426 assert!(text.contains("\"issue:12\""));
427 assert!(text.contains(&holder.to_string()));
428 }
429
430 #[test]
431 fn an_error_roundtrips() {
432 let output = SubWorkflowOutput::new(
433 Uuid::now_v7(),
434 "collect",
435 RunStatus::Failed,
436 Decimal::ZERO,
437 3,
438 )
439 .with_error("boom");
440
441 let back: SubWorkflowOutput =
442 from_value(to_value(&output).expect("serialize")).expect("deserialize");
443
444 assert_eq!(back, output);
445 assert_eq!(back.error(), Some("boom"));
446 }
447
448 #[test]
449 fn a_legacy_output_without_error_reads_back_as_none() {
450 let output: SubWorkflowOutput = from_value(json!({
451 "run_id": Uuid::now_v7(),
452 "workflow_name": "collect",
453 "status": "completed",
454 "cost_usd": 0,
455 "duration_ms": 7,
456 }))
457 .expect("deserialize");
458
459 assert_eq!(output.error(), None);
460 }
461
462 #[test]
463 fn an_output_without_run_id_is_refused() {
464 let result = from_value::<SubWorkflowOutput>(json!({
465 "workflow_name": "collect",
466 "status": "completed",
467 "cost_usd": 0,
468 "duration_ms": 0,
469 }));
470 assert!(result.is_err());
471 }
472}