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::{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}