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::de::DeserializeOwned;
10use serde::{Deserialize, Serialize};
11use serde_json::{Value, from_value};
12use uuid::Uuid;
13
14use ironflow_store::entities::RunStatus;
15
16use crate::error::EngineError;
17
18/// Result of a [`workflow`](crate::context::WorkflowContext::workflow) step.
19///
20/// While planning, no child run is created: [`run_id`](Self::run_id) is
21/// [`Uuid::nil`] and the metrics are zero.
22///
23/// It is also the persisted output of the step, so a stored workflow step
24/// reads back with [`StepOutput::json`](super::StepOutput::json).
25///
26/// # Examples
27///
28/// ```
29/// use ironflow_engine::executor::SubWorkflowOutput;
30/// use ironflow_store::entities::RunStatus;
31/// use rust_decimal::Decimal;
32/// use uuid::Uuid;
33///
34/// let run_id = Uuid::now_v7();
35/// let output = SubWorkflowOutput::new(run_id, "collect", RunStatus::Completed, Decimal::ZERO, 1200);
36/// assert_eq!(output.run_id(), run_id);
37/// assert_eq!(output.workflow_name(), "collect");
38/// assert_eq!(output.status(), RunStatus::Completed);
39/// assert_eq!(output.duration_ms(), 1200);
40/// ```
41#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
42pub struct SubWorkflowOutput {
43 run_id: Uuid,
44 workflow_name: String,
45 status: RunStatus,
46 cost_usd: Decimal,
47 duration_ms: u64,
48 #[serde(default, skip_serializing_if = "Option::is_none")]
49 error: Option<String>,
50 #[serde(default, skip_serializing_if = "Option::is_none")]
51 output: Option<Value>,
52}
53
54impl SubWorkflowOutput {
55 /// Assemble the result of a child run.
56 ///
57 /// # Examples
58 ///
59 /// ```
60 /// use ironflow_engine::executor::SubWorkflowOutput;
61 /// use ironflow_store::entities::RunStatus;
62 /// use rust_decimal::Decimal;
63 /// use uuid::Uuid;
64 ///
65 /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Warning, Decimal::ONE, 0);
66 /// assert_eq!(output.cost_usd(), Decimal::ONE);
67 /// ```
68 pub fn new(
69 run_id: Uuid,
70 workflow_name: &str,
71 status: RunStatus,
72 cost_usd: Decimal,
73 duration_ms: u64,
74 ) -> Self {
75 Self {
76 run_id,
77 workflow_name: workflow_name.to_string(),
78 status,
79 cost_usd,
80 duration_ms,
81 error: None,
82 output: None,
83 }
84 }
85
86 /// Attach the error of a child run tolerated by `allow_failure`.
87 ///
88 /// # Examples
89 ///
90 /// ```
91 /// use ironflow_engine::executor::SubWorkflowOutput;
92 /// use ironflow_store::entities::RunStatus;
93 /// use rust_decimal::Decimal;
94 /// use uuid::Uuid;
95 ///
96 /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Failed, Decimal::ZERO, 0)
97 /// .with_error("boom");
98 /// assert_eq!(output.error(), Some("boom"));
99 /// ```
100 pub fn with_error(mut self, error: impl Into<String>) -> Self {
101 self.error = Some(error.into());
102 self
103 }
104
105 /// Attach the output the child handler set with
106 /// [`set_output`](crate::context::WorkflowContext::set_output).
107 ///
108 /// # Examples
109 ///
110 /// ```
111 /// use ironflow_engine::executor::SubWorkflowOutput;
112 /// use ironflow_store::entities::RunStatus;
113 /// use rust_decimal::Decimal;
114 /// use serde_json::json;
115 /// use uuid::Uuid;
116 ///
117 /// # use ironflow_engine::error::EngineError;
118 /// # fn example() -> Result<(), EngineError> {
119 /// let output = SubWorkflowOutput::new(Uuid::nil(), "review", RunStatus::Completed, Decimal::ZERO, 0)
120 /// .with_output(Some(json!(true)));
121 /// assert_eq!(output.output::<bool>()?, Some(true));
122 /// # Ok(())
123 /// # }
124 /// ```
125 pub fn with_output(mut self, output: Option<Value>) -> Self {
126 self.output = output;
127 self
128 }
129
130 /// The child run, to read its steps from the store. [`Uuid::nil`] while
131 /// planning.
132 ///
133 /// # Examples
134 ///
135 /// ```
136 /// use ironflow_engine::executor::SubWorkflowOutput;
137 /// use ironflow_store::entities::RunStatus;
138 /// use rust_decimal::Decimal;
139 /// use uuid::Uuid;
140 ///
141 /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Completed, Decimal::ZERO, 0);
142 /// assert!(output.run_id().is_nil());
143 /// ```
144 pub fn run_id(&self) -> Uuid {
145 self.run_id
146 }
147
148 /// Name of the child workflow.
149 pub fn workflow_name(&self) -> &str {
150 &self.workflow_name
151 }
152
153 /// Final status of the child run: `Completed`, `Warning` when one of its
154 /// `allow_failure` steps failed, or (only when the step was started with
155 /// `allow_failure`) `Failed` / `Cancelled`.
156 pub fn status(&self) -> RunStatus {
157 self.status
158 }
159
160 /// Cost of the child run, in USD. Already included in the parent's cost.
161 pub fn cost_usd(&self) -> Decimal {
162 self.cost_usd
163 }
164
165 /// Wall-clock duration of the child run, in milliseconds.
166 pub fn duration_ms(&self) -> u64 {
167 self.duration_ms
168 }
169
170 /// Error of a failed or cancelled child run tolerated by `allow_failure`,
171 /// `None` otherwise.
172 ///
173 /// # Examples
174 ///
175 /// ```
176 /// use ironflow_engine::executor::SubWorkflowOutput;
177 /// use ironflow_store::entities::RunStatus;
178 /// use rust_decimal::Decimal;
179 /// use uuid::Uuid;
180 ///
181 /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Completed, Decimal::ZERO, 0);
182 /// assert_eq!(output.error(), None);
183 /// ```
184 pub fn error(&self) -> Option<&str> {
185 self.error.as_deref()
186 }
187
188 /// The typed output the child handler set with
189 /// [`set_output`](crate::context::WorkflowContext::set_output).
190 ///
191 /// `Ok(None)` when the child set no output. A failed child tolerated by
192 /// `allow_failure` keeps the output it set before failing. On replay the
193 /// value is read from the recorded step, not from the child run.
194 ///
195 /// # Errors
196 ///
197 /// Returns [`EngineError::Serialization`] if the output does not
198 /// deserialize into `T`.
199 ///
200 /// # Examples
201 ///
202 /// ```
203 /// use ironflow_engine::executor::SubWorkflowOutput;
204 /// use ironflow_store::entities::RunStatus;
205 /// use rust_decimal::Decimal;
206 /// use serde::Deserialize;
207 /// use serde_json::json;
208 /// use uuid::Uuid;
209 ///
210 /// #[derive(Deserialize)]
211 /// struct Review {
212 /// approved: bool,
213 /// }
214 ///
215 /// # use ironflow_engine::error::EngineError;
216 /// # fn example() -> Result<(), EngineError> {
217 /// let child = SubWorkflowOutput::new(Uuid::nil(), "review", RunStatus::Completed, Decimal::ZERO, 0)
218 /// .with_output(Some(json!({"approved": true})));
219 /// let review: Option<Review> = child.output()?;
220 /// assert!(review.is_some_and(|r| r.approved));
221 /// # Ok(())
222 /// # }
223 /// ```
224 pub fn output<T: DeserializeOwned>(&self) -> Result<Option<T>, EngineError> {
225 self.output
226 .as_ref()
227 .map(|value| from_value(value.clone()))
228 .transpose()
229 .map_err(EngineError::Serialization)
230 }
231}
232
233/// A sub-workflow skipped because its concurrency key is held by another
234/// active run.
235///
236/// Persisted as the step output `{"concurrency_conflict": {"key": .., "run_id": ..}}`
237/// and replayed as-is on resume.
238///
239/// # Examples
240///
241/// ```
242/// use ironflow_engine::executor::ConcurrencyConflict;
243/// use uuid::Uuid;
244///
245/// let holder = Uuid::now_v7();
246/// let conflict = ConcurrencyConflict::new("issue:12", holder);
247/// assert_eq!(conflict.key(), "issue:12");
248/// assert_eq!(conflict.run_id(), holder);
249/// assert!(conflict.to_string().contains("issue:12"));
250/// ```
251#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
252pub struct ConcurrencyConflict {
253 key: String,
254 run_id: Uuid,
255}
256
257impl ConcurrencyConflict {
258 /// Record that `run_id` holds `key`.
259 ///
260 /// # Examples
261 ///
262 /// ```
263 /// use ironflow_engine::executor::ConcurrencyConflict;
264 /// use uuid::Uuid;
265 ///
266 /// let conflict = ConcurrencyConflict::new("deploy:prod", Uuid::nil());
267 /// assert_eq!(conflict.key(), "deploy:prod");
268 /// ```
269 pub fn new(key: impl Into<String>, run_id: Uuid) -> Self {
270 Self {
271 key: key.into(),
272 run_id,
273 }
274 }
275
276 /// The contested concurrency key.
277 ///
278 /// # Examples
279 ///
280 /// ```
281 /// use ironflow_engine::executor::ConcurrencyConflict;
282 /// use uuid::Uuid;
283 ///
284 /// assert_eq!(ConcurrencyConflict::new("k", Uuid::nil()).key(), "k");
285 /// ```
286 pub fn key(&self) -> &str {
287 &self.key
288 }
289
290 /// The active run holding the key, at the time of the conflict.
291 ///
292 /// # Examples
293 ///
294 /// ```
295 /// use ironflow_engine::executor::ConcurrencyConflict;
296 /// use uuid::Uuid;
297 ///
298 /// let holder = Uuid::now_v7();
299 /// assert_eq!(ConcurrencyConflict::new("k", holder).run_id(), holder);
300 /// ```
301 pub fn run_id(&self) -> Uuid {
302 self.run_id
303 }
304}
305
306impl fmt::Display for ConcurrencyConflict {
307 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
308 write!(
309 f,
310 "conflict on concurrency key {:?} (run {})",
311 self.key, self.run_id
312 )
313 }
314}
315
316/// Result of a
317/// [`workflow_with`](crate::context::WorkflowContext::workflow_with) step.
318///
319/// # Examples
320///
321/// ```
322/// use ironflow_engine::executor::{ConcurrencyConflict, SubWorkflowOutcome};
323/// use uuid::Uuid;
324///
325/// let outcome = SubWorkflowOutcome::Conflict(ConcurrencyConflict::new("issue:12", Uuid::nil()));
326/// match &outcome {
327/// SubWorkflowOutcome::Completed(output) => println!("child run {}", output.run_id()),
328/// SubWorkflowOutcome::Conflict(conflict) => println!("held by {}", conflict.run_id()),
329/// }
330/// assert!(outcome.output().is_none());
331/// ```
332#[derive(Debug, Clone, PartialEq)]
333pub enum SubWorkflowOutcome {
334 /// The child run was created and finished.
335 Completed(SubWorkflowOutput),
336 /// No child run was created: another active run holds the concurrency key.
337 Conflict(ConcurrencyConflict),
338}
339
340impl SubWorkflowOutcome {
341 /// The conflict, when no child run was created.
342 ///
343 /// # Examples
344 ///
345 /// ```
346 /// use ironflow_engine::executor::{ConcurrencyConflict, SubWorkflowOutcome};
347 /// use uuid::Uuid;
348 ///
349 /// let outcome = SubWorkflowOutcome::Conflict(ConcurrencyConflict::new("k", Uuid::nil()));
350 /// assert_eq!(outcome.conflict().map(|c| c.key()), Some("k"));
351 /// ```
352 pub fn conflict(&self) -> Option<&ConcurrencyConflict> {
353 match self {
354 SubWorkflowOutcome::Conflict(conflict) => Some(conflict),
355 SubWorkflowOutcome::Completed(_) => None,
356 }
357 }
358
359 /// The child result, when the child run was created and finished.
360 ///
361 /// # Examples
362 ///
363 /// ```
364 /// use ironflow_engine::executor::{SubWorkflowOutcome, SubWorkflowOutput};
365 /// use ironflow_store::entities::RunStatus;
366 /// use rust_decimal::Decimal;
367 /// use uuid::Uuid;
368 ///
369 /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Completed, Decimal::ZERO, 0);
370 /// let outcome = SubWorkflowOutcome::Completed(output.clone());
371 /// assert_eq!(outcome.output(), Some(&output));
372 /// assert!(outcome.conflict().is_none());
373 /// ```
374 pub fn output(&self) -> Option<&SubWorkflowOutput> {
375 match self {
376 SubWorkflowOutcome::Completed(output) => Some(output),
377 SubWorkflowOutcome::Conflict(_) => None,
378 }
379 }
380}
381
382/// The persisted output of a completed `Workflow` step, as read back on replay.
383///
384/// `Conflict` is tried first: a flat [`SubWorkflowOutput`] has no
385/// `concurrency_conflict` field, so outputs recorded before concurrency keys
386/// existed still read back as `Completed`.
387#[derive(Debug, Deserialize)]
388#[serde(untagged)]
389pub(crate) enum RecordedWorkflowStep {
390 /// The step was skipped on a concurrency conflict.
391 Conflict {
392 /// The recorded conflict.
393 concurrency_conflict: ConcurrencyConflict,
394 },
395 /// The child run finished.
396 Completed(SubWorkflowOutput),
397}
398
399impl From<RecordedWorkflowStep> for SubWorkflowOutcome {
400 fn from(recorded: RecordedWorkflowStep) -> Self {
401 match recorded {
402 RecordedWorkflowStep::Conflict {
403 concurrency_conflict,
404 } => SubWorkflowOutcome::Conflict(concurrency_conflict),
405 RecordedWorkflowStep::Completed(output) => SubWorkflowOutcome::Completed(output),
406 }
407 }
408}
409
410#[cfg(test)]
411mod tests {
412 use serde_json::{from_value, json, to_value};
413
414 use super::*;
415
416 #[test]
417 fn serializes_to_the_persisted_step_output() {
418 let run_id = Uuid::now_v7();
419 let output = SubWorkflowOutput::new(
420 run_id,
421 "collect",
422 RunStatus::Warning,
423 Decimal::new(25, 2),
424 1200,
425 );
426
427 assert_eq!(
428 to_value(&output).expect("serialize"),
429 json!({
430 "run_id": run_id,
431 "workflow_name": "collect",
432 "status": "warning",
433 "cost_usd": 0.25,
434 "duration_ms": 1200,
435 })
436 );
437 }
438
439 #[test]
440 fn a_persisted_step_output_reads_back() {
441 let run_id = Uuid::now_v7();
442 let stored = json!({
443 "run_id": run_id,
444 "workflow_name": "collect",
445 "status": "completed",
446 "cost_usd": 0,
447 "duration_ms": 7,
448 });
449
450 let output: SubWorkflowOutput = from_value(stored).expect("deserialize");
451
452 assert_eq!(output.run_id(), run_id);
453 assert_eq!(output.status(), RunStatus::Completed);
454 assert_eq!(output.cost_usd(), Decimal::ZERO);
455 assert_eq!(output.duration_ms(), 7);
456 }
457
458 #[test]
459 fn a_recorded_conflict_round_trips() {
460 let holder = Uuid::now_v7();
461 let conflict = ConcurrencyConflict::new("issue:12", holder);
462 let stored = json!({ "concurrency_conflict": conflict });
463
464 assert_eq!(
465 stored,
466 json!({ "concurrency_conflict": { "key": "issue:12", "run_id": holder } })
467 );
468
469 let recorded: RecordedWorkflowStep = from_value(stored).expect("deserialize");
470 let outcome = SubWorkflowOutcome::from(recorded);
471 assert_eq!(outcome, SubWorkflowOutcome::Conflict(conflict));
472 assert_eq!(
473 outcome.conflict().map(ConcurrencyConflict::run_id),
474 Some(holder)
475 );
476 assert!(outcome.output().is_none());
477 }
478
479 #[test]
480 fn a_flat_output_still_reads_back_as_completed() {
481 let run_id = Uuid::now_v7();
482 let stored = json!({
483 "run_id": run_id,
484 "workflow_name": "collect",
485 "status": "completed",
486 "cost_usd": 0,
487 "duration_ms": 7,
488 });
489
490 let recorded: RecordedWorkflowStep = from_value(stored).expect("deserialize");
491 let outcome = SubWorkflowOutcome::from(recorded);
492 let output = outcome.output().expect("completed outcome");
493 assert_eq!(output.run_id(), run_id);
494 assert_eq!(output.duration_ms(), 7);
495 assert!(outcome.conflict().is_none());
496 }
497
498 #[test]
499 fn conflict_display_names_the_key_and_the_holder() {
500 let holder = Uuid::now_v7();
501 let text = ConcurrencyConflict::new("issue:12", holder).to_string();
502 assert!(text.contains("\"issue:12\""));
503 assert!(text.contains(&holder.to_string()));
504 }
505
506 #[test]
507 fn an_error_roundtrips() {
508 let output = SubWorkflowOutput::new(
509 Uuid::now_v7(),
510 "collect",
511 RunStatus::Failed,
512 Decimal::ZERO,
513 3,
514 )
515 .with_error("boom");
516
517 let back: SubWorkflowOutput =
518 from_value(to_value(&output).expect("serialize")).expect("deserialize");
519
520 assert_eq!(back, output);
521 assert_eq!(back.error(), Some("boom"));
522 }
523
524 #[test]
525 fn a_legacy_output_without_error_reads_back_as_none() {
526 let output: SubWorkflowOutput = from_value(json!({
527 "run_id": Uuid::now_v7(),
528 "workflow_name": "collect",
529 "status": "completed",
530 "cost_usd": 0,
531 "duration_ms": 7,
532 }))
533 .expect("deserialize");
534
535 assert_eq!(output.error(), None);
536 }
537
538 #[derive(Debug, PartialEq, Serialize, Deserialize)]
539 struct Verdict {
540 approved: bool,
541 score: u8,
542 }
543
544 fn child() -> SubWorkflowOutput {
545 SubWorkflowOutput::new(
546 Uuid::now_v7(),
547 "review",
548 RunStatus::Completed,
549 Decimal::ZERO,
550 5,
551 )
552 }
553
554 #[test]
555 fn set_output_none_reads_back_as_none() {
556 let output: Option<Verdict> = child().output().expect("no output");
557 assert_eq!(output, None);
558 }
559
560 #[test]
561 fn set_output_typed_value_reads_back() {
562 let output = child().with_output(Some(json!({"approved": true, "score": 7})));
563
564 let verdict: Option<Verdict> = output.output().expect("deserialize");
565 assert_eq!(
566 verdict,
567 Some(Verdict {
568 approved: true,
569 score: 7
570 })
571 );
572 }
573
574 #[test]
575 fn set_output_of_the_wrong_type_is_a_serialization_error() {
576 let output = child().with_output(Some(json!({"approved": "yes"})));
577
578 let result = output.output::<Verdict>();
579 assert!(matches!(result, Err(EngineError::Serialization(_))));
580 }
581
582 #[test]
583 fn set_output_round_trips_through_the_step_output() {
584 let output = child().with_output(Some(json!({"approved": false, "score": 1})));
585
586 let stored = to_value(&output).expect("serialize");
587 assert_eq!(stored["output"], json!({"approved": false, "score": 1}));
588
589 let back: SubWorkflowOutput = from_value(stored.clone()).expect("deserialize");
590 assert_eq!(back, output);
591
592 let recorded: RecordedWorkflowStep = from_value(stored).expect("deserialize");
593 let outcome = SubWorkflowOutcome::from(recorded);
594 let replayed = outcome.output().expect("completed outcome");
595 assert_eq!(
596 replayed.output::<Verdict>().expect("deserialize"),
597 Some(Verdict {
598 approved: false,
599 score: 1
600 })
601 );
602 }
603
604 #[test]
605 fn set_output_absent_is_not_serialized() {
606 let stored = to_value(child()).expect("serialize");
607 assert!(stored.get("output").is_none());
608 }
609
610 #[test]
611 fn set_output_legacy_step_output_without_the_field_reads_back_as_none() {
612 let output: SubWorkflowOutput = from_value(json!({
613 "run_id": Uuid::now_v7(),
614 "workflow_name": "collect",
615 "status": "completed",
616 "cost_usd": 0,
617 "duration_ms": 7,
618 }))
619 .expect("deserialize");
620
621 assert_eq!(output.output::<Verdict>().expect("no output"), None);
622 }
623
624 #[test]
625 fn an_output_without_run_id_is_refused() {
626 let result = from_value::<SubWorkflowOutput>(json!({
627 "workflow_name": "collect",
628 "status": "completed",
629 "cost_usd": 0,
630 "duration_ms": 0,
631 }));
632 assert!(result.is_err());
633 }
634}