Skip to main content

ironflow_engine/config/
workflow.rs

1//! Configuration for workflow (sub-workflow) steps.
2
3use std::ops::Not;
4
5use ironflow_core::retry::RetryPolicy;
6use ironflow_store::entities::MAX_CONCURRENCY_KEY_LEN;
7use serde::{Deserialize, Serialize};
8use serde_json::Value;
9
10/// Configuration for invoking a registered workflow as a sub-step.
11///
12/// The engine will look up the handler by [`workflow_name`](WorkflowStepConfig::workflow_name)
13/// and execute it as a child run with its own steps and lifecycle.
14///
15/// # Examples
16///
17/// ```
18/// use ironflow_engine::config::WorkflowStepConfig;
19/// use serde_json::json;
20///
21/// let config = WorkflowStepConfig::new("build", json!({"branch": "main"}));
22/// assert_eq!(config.workflow_name, "build");
23/// ```
24#[derive(Debug, Clone, Serialize, Deserialize)]
25pub struct WorkflowStepConfig {
26    /// Name of the registered workflow handler to invoke.
27    pub workflow_name: String,
28    /// Payload to pass to the child workflow run.
29    pub payload: Value,
30    /// Optional step-level retry policy.
31    #[serde(default, skip_serializing_if = "Option::is_none")]
32    pub retry: Option<RetryPolicy>,
33    /// Concurrency key held by the child run while it is not terminal.
34    ///
35    /// Set through [`WorkflowOptions::concurrency_key`] and recorded in the
36    /// step input.
37    #[serde(default, skip_serializing_if = "Option::is_none")]
38    pub concurrency_key: Option<String>,
39    /// Tolerate a failed child: the step completes with the child's failure
40    /// in its output instead of failing the parent.
41    #[serde(default, skip_serializing_if = "Not::not")]
42    pub allow_failure: bool,
43}
44
45impl WorkflowStepConfig {
46    /// Create a new workflow step config.
47    ///
48    /// # Examples
49    ///
50    /// ```
51    /// use ironflow_engine::config::WorkflowStepConfig;
52    /// use serde_json::json;
53    ///
54    /// let config = WorkflowStepConfig::new("deploy", json!({}));
55    /// assert_eq!(config.workflow_name, "deploy");
56    /// ```
57    pub fn new(workflow_name: &str, payload: Value) -> Self {
58        Self {
59            workflow_name: workflow_name.to_string(),
60            payload,
61            retry: None,
62            concurrency_key: None,
63            allow_failure: false,
64        }
65    }
66
67    /// Set a step-level retry policy.
68    ///
69    /// # Examples
70    ///
71    /// ```
72    /// use ironflow_core::retry::RetryPolicy;
73    /// use ironflow_engine::config::WorkflowStepConfig;
74    /// use serde_json::json;
75    ///
76    /// let config = WorkflowStepConfig::new("build", json!({}))
77    ///     .retry_policy(RetryPolicy::new(3));
78    /// assert!(config.retry.is_some());
79    /// ```
80    pub fn retry_policy(mut self, policy: RetryPolicy) -> Self {
81        self.retry = Some(policy);
82        self
83    }
84
85    /// Tolerate a failed child run: the step completes with a
86    /// [`SubWorkflowOutput`](crate::executor::SubWorkflowOutput) reporting the
87    /// failure and the parent run ends as `Warning`.
88    ///
89    /// # Examples
90    ///
91    /// ```
92    /// use ironflow_engine::config::WorkflowStepConfig;
93    /// use serde_json::json;
94    ///
95    /// let config = WorkflowStepConfig::new("build", json!({})).allow_failure();
96    /// assert!(config.allow_failure);
97    /// ```
98    pub fn allow_failure(mut self) -> Self {
99        self.allow_failure = true;
100        self
101    }
102}
103
104/// Options of a sub-workflow step started with
105/// [`workflow_with`](crate::context::WorkflowContext::workflow_with).
106///
107/// # Examples
108///
109/// ```
110/// use ironflow_engine::config::WorkflowOptions;
111///
112/// let options = WorkflowOptions::new().allow_failure();
113/// assert!(options.allow_failure);
114/// assert!(!WorkflowOptions::default().allow_failure);
115/// ```
116#[derive(Debug, Clone, Default, PartialEq, Eq)]
117pub struct WorkflowOptions {
118    /// Tolerate a failed child run instead of failing the parent.
119    pub allow_failure: bool,
120    concurrency_key: Option<String>,
121}
122
123impl WorkflowOptions {
124    /// Create options with every flag off.
125    ///
126    /// # Examples
127    ///
128    /// ```
129    /// use ironflow_engine::config::WorkflowOptions;
130    ///
131    /// assert_eq!(WorkflowOptions::new(), WorkflowOptions::default());
132    /// ```
133    pub fn new() -> Self {
134        Self::default()
135    }
136
137    /// Tolerate a failed child run: the step completes with the failure in its
138    /// output and the parent run ends as `Warning`. A suspension is never
139    /// tolerated.
140    ///
141    /// # Examples
142    ///
143    /// ```
144    /// use ironflow_engine::config::WorkflowOptions;
145    ///
146    /// assert!(WorkflowOptions::new().allow_failure().allow_failure);
147    /// ```
148    pub fn allow_failure(mut self) -> Self {
149        self.allow_failure = true;
150        self
151    }
152
153    /// Make the child run exclusive on `key`.
154    ///
155    /// While another non-terminal run (Pending, Running, Retrying, Sleeping,
156    /// AwaitingApproval) holds the same key, no child run is created and the
157    /// step completes with a
158    /// [`SubWorkflowOutcome::Conflict`](crate::executor::SubWorkflowOutcome::Conflict).
159    ///
160    /// # Panics
161    ///
162    /// Panics if `key` is empty or only whitespace, or longer than
163    /// [`MAX_CONCURRENCY_KEY_LEN`] bytes.
164    ///
165    /// # Examples
166    ///
167    /// ```
168    /// use ironflow_engine::config::WorkflowOptions;
169    ///
170    /// let options = WorkflowOptions::new().concurrency_key("issue:12");
171    /// assert_eq!(options.concurrency_key_ref(), Some("issue:12"));
172    /// ```
173    ///
174    /// ```should_panic
175    /// use ironflow_engine::config::WorkflowOptions;
176    ///
177    /// WorkflowOptions::new().concurrency_key("  ");
178    /// ```
179    pub fn concurrency_key(mut self, key: impl Into<String>) -> Self {
180        let key = key.into();
181        assert!(!key.trim().is_empty(), "concurrency key must not be empty");
182        assert!(
183            key.len() <= MAX_CONCURRENCY_KEY_LEN,
184            "concurrency key must be at most {MAX_CONCURRENCY_KEY_LEN} bytes"
185        );
186        self.concurrency_key = Some(key);
187        self
188    }
189
190    /// The concurrency key, if one was set.
191    ///
192    /// # Examples
193    ///
194    /// ```
195    /// use ironflow_engine::config::WorkflowOptions;
196    ///
197    /// let options = WorkflowOptions::new().concurrency_key("deploy:prod");
198    /// assert_eq!(options.concurrency_key_ref(), Some("deploy:prod"));
199    /// ```
200    pub fn concurrency_key_ref(&self) -> Option<&str> {
201        self.concurrency_key.as_deref()
202    }
203
204    /// Consume the options and return the concurrency key.
205    pub(crate) fn into_concurrency_key(self) -> Option<String> {
206        self.concurrency_key
207    }
208}
209
210#[cfg(test)]
211mod tests {
212    use super::*;
213    use serde_json::{from_str, json, to_string, to_value};
214
215    #[test]
216    fn new_sets_fields() {
217        let config = WorkflowStepConfig::new("build", json!({"key": "val"}));
218        assert_eq!(config.workflow_name, "build");
219        assert_eq!(config.payload["key"], "val");
220    }
221
222    #[test]
223    fn serde_roundtrip() {
224        let config = WorkflowStepConfig::new("deploy", json!({"env": "prod"}));
225        let json = serde_json::to_string(&config).unwrap();
226        let back: WorkflowStepConfig = serde_json::from_str(&json).unwrap();
227        assert_eq!(back.workflow_name, "deploy");
228        assert_eq!(back.payload["env"], "prod");
229    }
230
231    #[test]
232    fn a_config_predating_retry_still_deserializes() {
233        let config: WorkflowStepConfig =
234            serde_json::from_str(r#"{"workflow_name":"build","payload":{"key":"val"}}"#)
235                .expect("deserialize");
236        assert!(config.retry.is_none());
237    }
238
239    #[test]
240    fn a_config_without_concurrency_key_omits_it() {
241        let config = WorkflowStepConfig::new("build", json!({}));
242        let value = to_value(&config).expect("serialize");
243        assert!(value.get("concurrency_key").is_none());
244    }
245
246    #[test]
247    fn concurrency_key_roundtrip() {
248        let mut config = WorkflowStepConfig::new("build", json!({}));
249        config.concurrency_key = Some("issue:12".to_string());
250        let json = to_string(&config).expect("serialize");
251        let back: WorkflowStepConfig = from_str(&json).expect("deserialize");
252        assert_eq!(back.concurrency_key.as_deref(), Some("issue:12"));
253    }
254
255    #[test]
256    fn options_carry_the_concurrency_key() {
257        let options = WorkflowOptions::new().concurrency_key("issue:12");
258        assert_eq!(options.concurrency_key_ref(), Some("issue:12"));
259        assert_eq!(options.into_concurrency_key().as_deref(), Some("issue:12"));
260        assert!(WorkflowOptions::new().concurrency_key_ref().is_none());
261    }
262
263    #[test]
264    fn options_accept_a_key_at_the_length_limit() {
265        let key = "a".repeat(MAX_CONCURRENCY_KEY_LEN);
266        let options = WorkflowOptions::new().concurrency_key(key.clone());
267        assert_eq!(options.concurrency_key_ref(), Some(key.as_str()));
268    }
269
270    #[test]
271    #[should_panic(expected = "concurrency key must not be empty")]
272    fn options_reject_an_empty_key() {
273        let _ = WorkflowOptions::new().concurrency_key("");
274    }
275
276    #[test]
277    #[should_panic(expected = "concurrency key must not be empty")]
278    fn options_reject_a_blank_key() {
279        let _ = WorkflowOptions::new().concurrency_key(" \t ");
280    }
281
282    #[test]
283    #[should_panic(expected = "concurrency key must be at most")]
284    fn options_reject_a_key_over_the_limit() {
285        let _ = WorkflowOptions::new().concurrency_key("a".repeat(MAX_CONCURRENCY_KEY_LEN + 1));
286    }
287
288    #[test]
289    fn allow_failure_defaults_to_false() {
290        let config = WorkflowStepConfig::new("build", json!({}));
291        assert!(!config.allow_failure);
292    }
293
294    #[test]
295    fn allow_failure_is_omitted_from_json_when_false() {
296        let config = WorkflowStepConfig::new("build", json!({}));
297        let value = serde_json::to_value(&config).expect("serialize");
298        assert!(value.get("allow_failure").is_none());
299    }
300
301    #[test]
302    fn allow_failure_roundtrip() {
303        let config = WorkflowStepConfig::new("build", json!({})).allow_failure();
304        let json = serde_json::to_string(&config).expect("serialize");
305        let back: WorkflowStepConfig = serde_json::from_str(&json).expect("deserialize");
306        assert!(back.allow_failure);
307    }
308
309    #[test]
310    fn options_builder_sets_allow_failure() {
311        assert!(!WorkflowOptions::new().allow_failure);
312        assert!(WorkflowOptions::new().allow_failure().allow_failure);
313    }
314
315    #[test]
316    fn retry_policy_roundtrip() {
317        let config = WorkflowStepConfig::new("deploy", json!({})).retry_policy(RetryPolicy::new(3));
318        let json = serde_json::to_string(&config).expect("serialize");
319        let back: WorkflowStepConfig = serde_json::from_str(&json).expect("deserialize");
320        assert_eq!(back.retry.as_ref().unwrap().max_retries(), 3);
321    }
322}