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, validate_priority};
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    /// Queue priority of the child run.
40    ///
41    /// Set through [`WorkflowOptions::priority`] and recorded in the step
42    /// input. `None` gives the child the priority of its own handler.
43    #[serde(default, skip_serializing_if = "Option::is_none")]
44    pub priority: Option<i16>,
45    /// Tolerate a failed child: the step completes with the child's failure
46    /// in its output instead of failing the parent.
47    #[serde(default, skip_serializing_if = "Not::not")]
48    pub allow_failure: bool,
49}
50
51impl WorkflowStepConfig {
52    /// Create a new workflow step config.
53    ///
54    /// # Examples
55    ///
56    /// ```
57    /// use ironflow_engine::config::WorkflowStepConfig;
58    /// use serde_json::json;
59    ///
60    /// let config = WorkflowStepConfig::new("deploy", json!({}));
61    /// assert_eq!(config.workflow_name, "deploy");
62    /// ```
63    pub fn new(workflow_name: &str, payload: Value) -> Self {
64        Self {
65            workflow_name: workflow_name.to_string(),
66            payload,
67            retry: None,
68            concurrency_key: None,
69            priority: None,
70            allow_failure: false,
71        }
72    }
73
74    /// Set a step-level retry policy.
75    ///
76    /// # Examples
77    ///
78    /// ```
79    /// use ironflow_core::retry::RetryPolicy;
80    /// use ironflow_engine::config::WorkflowStepConfig;
81    /// use serde_json::json;
82    ///
83    /// let config = WorkflowStepConfig::new("build", json!({}))
84    ///     .retry_policy(RetryPolicy::new(3));
85    /// assert!(config.retry.is_some());
86    /// ```
87    pub fn retry_policy(mut self, policy: RetryPolicy) -> Self {
88        self.retry = Some(policy);
89        self
90    }
91
92    /// Tolerate a failed child run: the step completes with a
93    /// [`SubWorkflowOutput`](crate::executor::SubWorkflowOutput) reporting the
94    /// failure and the parent run ends as `Warning`.
95    ///
96    /// # Examples
97    ///
98    /// ```
99    /// use ironflow_engine::config::WorkflowStepConfig;
100    /// use serde_json::json;
101    ///
102    /// let config = WorkflowStepConfig::new("build", json!({})).allow_failure();
103    /// assert!(config.allow_failure);
104    /// ```
105    pub fn allow_failure(mut self) -> Self {
106        self.allow_failure = true;
107        self
108    }
109}
110
111/// Options of a sub-workflow step started with
112/// [`workflow_with`](crate::context::WorkflowContext::workflow_with).
113///
114/// # Examples
115///
116/// ```
117/// use ironflow_engine::config::WorkflowOptions;
118///
119/// let options = WorkflowOptions::new().allow_failure();
120/// assert!(options.allow_failure);
121/// assert!(!WorkflowOptions::default().allow_failure);
122/// ```
123#[derive(Debug, Clone, Default, PartialEq, Eq)]
124pub struct WorkflowOptions {
125    /// Tolerate a failed child run instead of failing the parent.
126    pub allow_failure: bool,
127    concurrency_key: Option<String>,
128    priority: Option<i16>,
129}
130
131impl WorkflowOptions {
132    /// Create options with every flag off.
133    ///
134    /// # Examples
135    ///
136    /// ```
137    /// use ironflow_engine::config::WorkflowOptions;
138    ///
139    /// assert_eq!(WorkflowOptions::new(), WorkflowOptions::default());
140    /// ```
141    pub fn new() -> Self {
142        Self::default()
143    }
144
145    /// Tolerate a failed child run: the step completes with the failure in its
146    /// output and the parent run ends as `Warning`. A suspension is never
147    /// tolerated.
148    ///
149    /// # Examples
150    ///
151    /// ```
152    /// use ironflow_engine::config::WorkflowOptions;
153    ///
154    /// assert!(WorkflowOptions::new().allow_failure().allow_failure);
155    /// ```
156    pub fn allow_failure(mut self) -> Self {
157        self.allow_failure = true;
158        self
159    }
160
161    /// Make the child run exclusive on `key`.
162    ///
163    /// While another non-terminal run (Pending, Running, Retrying, Sleeping,
164    /// AwaitingApproval) holds the same key, no child run is created and the
165    /// step completes with a
166    /// [`SubWorkflowOutcome::Conflict`](crate::executor::SubWorkflowOutcome::Conflict).
167    ///
168    /// # Panics
169    ///
170    /// Panics if `key` is empty or only whitespace, or longer than
171    /// [`MAX_CONCURRENCY_KEY_LEN`] bytes.
172    ///
173    /// # Examples
174    ///
175    /// ```
176    /// use ironflow_engine::config::WorkflowOptions;
177    ///
178    /// let options = WorkflowOptions::new().concurrency_key("issue:12");
179    /// assert_eq!(options.concurrency_key_ref(), Some("issue:12"));
180    /// ```
181    ///
182    /// ```should_panic
183    /// use ironflow_engine::config::WorkflowOptions;
184    ///
185    /// WorkflowOptions::new().concurrency_key("  ");
186    /// ```
187    pub fn concurrency_key(mut self, key: impl Into<String>) -> Self {
188        let key = key.into();
189        assert!(!key.trim().is_empty(), "concurrency key must not be empty");
190        assert!(
191            key.len() <= MAX_CONCURRENCY_KEY_LEN,
192            "concurrency key must be at most {MAX_CONCURRENCY_KEY_LEN} bytes"
193        );
194        self.concurrency_key = Some(key);
195        self
196    }
197
198    /// The concurrency key, if one was set.
199    ///
200    /// # Examples
201    ///
202    /// ```
203    /// use ironflow_engine::config::WorkflowOptions;
204    ///
205    /// let options = WorkflowOptions::new().concurrency_key("deploy:prod");
206    /// assert_eq!(options.concurrency_key_ref(), Some("deploy:prod"));
207    /// ```
208    pub fn concurrency_key_ref(&self) -> Option<&str> {
209        self.concurrency_key.as_deref()
210    }
211
212    /// Set the queue priority of the child run, overriding the priority of
213    /// the child handler.
214    ///
215    /// The child executes inside its parent's worker; the priority is
216    /// recorded on the child run and orders it whenever it waits in the queue.
217    ///
218    /// # Panics
219    ///
220    /// Panics if `priority` is outside
221    /// [`MIN_PRIORITY`](ironflow_store::entities::MIN_PRIORITY)`..=`[`MAX_PRIORITY`](ironflow_store::entities::MAX_PRIORITY).
222    ///
223    /// # Examples
224    ///
225    /// ```
226    /// use ironflow_engine::config::WorkflowOptions;
227    ///
228    /// let options = WorkflowOptions::new().priority(20);
229    /// assert_eq!(options.priority_ref(), Some(20));
230    /// ```
231    ///
232    /// ```should_panic
233    /// use ironflow_engine::config::WorkflowOptions;
234    ///
235    /// WorkflowOptions::new().priority(101);
236    /// ```
237    pub fn priority(mut self, priority: i16) -> Self {
238        if let Err(message) = validate_priority(priority) {
239            panic!("{message}");
240        }
241        self.priority = Some(priority);
242        self
243    }
244
245    /// The priority of the child run, if one was set.
246    ///
247    /// # Examples
248    ///
249    /// ```
250    /// use ironflow_engine::config::WorkflowOptions;
251    ///
252    /// assert_eq!(WorkflowOptions::new().priority_ref(), None);
253    /// assert_eq!(WorkflowOptions::new().priority(-5).priority_ref(), Some(-5));
254    /// ```
255    pub fn priority_ref(&self) -> Option<i16> {
256        self.priority
257    }
258
259    /// Consume the options and return the concurrency key and the priority.
260    pub(crate) fn into_parts(self) -> (Option<String>, Option<i16>) {
261        (self.concurrency_key, self.priority)
262    }
263}
264
265#[cfg(test)]
266mod tests {
267    use super::*;
268    use ironflow_store::entities::{MAX_PRIORITY, MIN_PRIORITY};
269    use serde_json::{from_str, json, to_string, to_value};
270
271    #[test]
272    fn new_sets_fields() {
273        let config = WorkflowStepConfig::new("build", json!({"key": "val"}));
274        assert_eq!(config.workflow_name, "build");
275        assert_eq!(config.payload["key"], "val");
276    }
277
278    #[test]
279    fn serde_roundtrip() {
280        let config = WorkflowStepConfig::new("deploy", json!({"env": "prod"}));
281        let json = serde_json::to_string(&config).unwrap();
282        let back: WorkflowStepConfig = serde_json::from_str(&json).unwrap();
283        assert_eq!(back.workflow_name, "deploy");
284        assert_eq!(back.payload["env"], "prod");
285    }
286
287    #[test]
288    fn a_config_predating_retry_still_deserializes() {
289        let config: WorkflowStepConfig =
290            serde_json::from_str(r#"{"workflow_name":"build","payload":{"key":"val"}}"#)
291                .expect("deserialize");
292        assert!(config.retry.is_none());
293    }
294
295    #[test]
296    fn a_config_without_concurrency_key_omits_it() {
297        let config = WorkflowStepConfig::new("build", json!({}));
298        let value = to_value(&config).expect("serialize");
299        assert!(value.get("concurrency_key").is_none());
300    }
301
302    #[test]
303    fn concurrency_key_roundtrip() {
304        let mut config = WorkflowStepConfig::new("build", json!({}));
305        config.concurrency_key = Some("issue:12".to_string());
306        let json = to_string(&config).expect("serialize");
307        let back: WorkflowStepConfig = from_str(&json).expect("deserialize");
308        assert_eq!(back.concurrency_key.as_deref(), Some("issue:12"));
309    }
310
311    #[test]
312    fn options_carry_the_concurrency_key() {
313        let options = WorkflowOptions::new().concurrency_key("issue:12");
314        assert_eq!(options.concurrency_key_ref(), Some("issue:12"));
315        assert_eq!(options.into_parts().0.as_deref(), Some("issue:12"));
316        assert!(WorkflowOptions::new().concurrency_key_ref().is_none());
317    }
318
319    #[test]
320    fn options_accept_a_key_at_the_length_limit() {
321        let key = "a".repeat(MAX_CONCURRENCY_KEY_LEN);
322        let options = WorkflowOptions::new().concurrency_key(key.clone());
323        assert_eq!(options.concurrency_key_ref(), Some(key.as_str()));
324    }
325
326    #[test]
327    #[should_panic(expected = "concurrency key must not be empty")]
328    fn options_reject_an_empty_key() {
329        let _ = WorkflowOptions::new().concurrency_key("");
330    }
331
332    #[test]
333    #[should_panic(expected = "concurrency key must not be empty")]
334    fn options_reject_a_blank_key() {
335        let _ = WorkflowOptions::new().concurrency_key(" \t ");
336    }
337
338    #[test]
339    #[should_panic(expected = "concurrency key must be at most")]
340    fn options_reject_a_key_over_the_limit() {
341        let _ = WorkflowOptions::new().concurrency_key("a".repeat(MAX_CONCURRENCY_KEY_LEN + 1));
342    }
343
344    #[test]
345    fn options_priority_is_carried_and_split() {
346        let options = WorkflowOptions::new()
347            .concurrency_key("issue:12")
348            .priority(MAX_PRIORITY);
349        assert_eq!(options.priority_ref(), Some(MAX_PRIORITY));
350        let (key, priority) = options.into_parts();
351        assert_eq!(key.as_deref(), Some("issue:12"));
352        assert_eq!(priority, Some(MAX_PRIORITY));
353        assert_eq!(
354            WorkflowOptions::new().priority(MIN_PRIORITY).priority_ref(),
355            Some(MIN_PRIORITY)
356        );
357        assert!(WorkflowOptions::new().priority_ref().is_none());
358    }
359
360    #[test]
361    #[should_panic(expected = "priority must be between -100 and 100")]
362    fn options_priority_rejects_above_the_maximum() {
363        let _ = WorkflowOptions::new().priority(MAX_PRIORITY + 1);
364    }
365
366    #[test]
367    #[should_panic(expected = "priority must be between -100 and 100")]
368    fn options_priority_rejects_below_the_minimum() {
369        let _ = WorkflowOptions::new().priority(MIN_PRIORITY - 1);
370    }
371
372    #[test]
373    fn step_config_priority_is_omitted_when_unset_and_round_trips() {
374        let config = WorkflowStepConfig::new("build", json!({}));
375        assert!(
376            to_value(&config)
377                .expect("serialize")
378                .get("priority")
379                .is_none()
380        );
381
382        let mut config = WorkflowStepConfig::new("build", json!({}));
383        config.priority = Some(-30);
384        let back: WorkflowStepConfig =
385            from_str(&to_string(&config).expect("serialize")).expect("deserialize");
386        assert_eq!(back.priority, Some(-30));
387    }
388
389    #[test]
390    fn allow_failure_defaults_to_false() {
391        let config = WorkflowStepConfig::new("build", json!({}));
392        assert!(!config.allow_failure);
393    }
394
395    #[test]
396    fn allow_failure_is_omitted_from_json_when_false() {
397        let config = WorkflowStepConfig::new("build", json!({}));
398        let value = serde_json::to_value(&config).expect("serialize");
399        assert!(value.get("allow_failure").is_none());
400    }
401
402    #[test]
403    fn allow_failure_roundtrip() {
404        let config = WorkflowStepConfig::new("build", json!({})).allow_failure();
405        let json = serde_json::to_string(&config).expect("serialize");
406        let back: WorkflowStepConfig = serde_json::from_str(&json).expect("deserialize");
407        assert!(back.allow_failure);
408    }
409
410    #[test]
411    fn options_builder_sets_allow_failure() {
412        assert!(!WorkflowOptions::new().allow_failure);
413        assert!(WorkflowOptions::new().allow_failure().allow_failure);
414    }
415
416    #[test]
417    fn retry_policy_roundtrip() {
418        let config = WorkflowStepConfig::new("deploy", json!({})).retry_policy(RetryPolicy::new(3));
419        let json = serde_json::to_string(&config).expect("serialize");
420        let back: WorkflowStepConfig = serde_json::from_str(&json).expect("deserialize");
421        assert_eq!(back.retry.as_ref().unwrap().max_retries(), 3);
422    }
423}