Skip to main content

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}
40
41impl SubWorkflowOutput {
42    /// Assemble the result of a child run.
43    ///
44    /// # Examples
45    ///
46    /// ```
47    /// use ironflow_engine::executor::SubWorkflowOutput;
48    /// use ironflow_store::entities::RunStatus;
49    /// use rust_decimal::Decimal;
50    /// use uuid::Uuid;
51    ///
52    /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Warning, Decimal::ONE, 0);
53    /// assert_eq!(output.cost_usd(), Decimal::ONE);
54    /// ```
55    pub fn new(
56        run_id: Uuid,
57        workflow_name: &str,
58        status: RunStatus,
59        cost_usd: Decimal,
60        duration_ms: u64,
61    ) -> Self {
62        Self {
63            run_id,
64            workflow_name: workflow_name.to_string(),
65            status,
66            cost_usd,
67            duration_ms,
68        }
69    }
70
71    /// The child run, to read its steps from the store. [`Uuid::nil`] while
72    /// planning.
73    ///
74    /// # Examples
75    ///
76    /// ```
77    /// use ironflow_engine::executor::SubWorkflowOutput;
78    /// use ironflow_store::entities::RunStatus;
79    /// use rust_decimal::Decimal;
80    /// use uuid::Uuid;
81    ///
82    /// let output = SubWorkflowOutput::new(Uuid::nil(), "collect", RunStatus::Completed, Decimal::ZERO, 0);
83    /// assert!(output.run_id().is_nil());
84    /// ```
85    pub fn run_id(&self) -> Uuid {
86        self.run_id
87    }
88
89    /// Name of the child workflow.
90    pub fn workflow_name(&self) -> &str {
91        &self.workflow_name
92    }
93
94    /// Final status of the child run: `Completed`, or `Warning` when one of its
95    /// `allow_failure` steps failed.
96    pub fn status(&self) -> RunStatus {
97        self.status
98    }
99
100    /// Cost of the child run, in USD. Already included in the parent's cost.
101    pub fn cost_usd(&self) -> Decimal {
102        self.cost_usd
103    }
104
105    /// Wall-clock duration of the child run, in milliseconds.
106    pub fn duration_ms(&self) -> u64 {
107        self.duration_ms
108    }
109}
110
111#[cfg(test)]
112mod tests {
113    use serde_json::{from_value, json, to_value};
114
115    use super::*;
116
117    #[test]
118    fn serializes_to_the_persisted_step_output() {
119        let run_id = Uuid::now_v7();
120        let output = SubWorkflowOutput::new(
121            run_id,
122            "collect",
123            RunStatus::Warning,
124            Decimal::new(25, 2),
125            1200,
126        );
127
128        assert_eq!(
129            to_value(&output).expect("serialize"),
130            json!({
131                "run_id": run_id,
132                "workflow_name": "collect",
133                "status": "warning",
134                "cost_usd": 0.25,
135                "duration_ms": 1200,
136            })
137        );
138    }
139
140    #[test]
141    fn a_persisted_step_output_reads_back() {
142        let run_id = Uuid::now_v7();
143        let stored = json!({
144            "run_id": run_id,
145            "workflow_name": "collect",
146            "status": "completed",
147            "cost_usd": 0,
148            "duration_ms": 7,
149        });
150
151        let output: SubWorkflowOutput = from_value(stored).expect("deserialize");
152
153        assert_eq!(output.run_id(), run_id);
154        assert_eq!(output.status(), RunStatus::Completed);
155        assert_eq!(output.cost_usd(), Decimal::ZERO);
156        assert_eq!(output.duration_ms(), 7);
157    }
158
159    #[test]
160    fn an_output_without_run_id_is_refused() {
161        let result = from_value::<SubWorkflowOutput>(json!({
162            "workflow_name": "collect",
163            "status": "completed",
164            "cost_usd": 0,
165            "duration_ms": 0,
166        }));
167        assert!(result.is_err());
168    }
169}