Skip to main content

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}