ironflow_engine/config/
workflow.rs1use 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#[derive(Debug, Clone, Serialize, Deserialize)]
25pub struct WorkflowStepConfig {
26 pub workflow_name: String,
28 pub payload: Value,
30 #[serde(default, skip_serializing_if = "Option::is_none")]
32 pub retry: Option<RetryPolicy>,
33 #[serde(default, skip_serializing_if = "Option::is_none")]
38 pub concurrency_key: Option<String>,
39 #[serde(default, skip_serializing_if = "Not::not")]
42 pub allow_failure: bool,
43}
44
45impl WorkflowStepConfig {
46 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 pub fn retry_policy(mut self, policy: RetryPolicy) -> Self {
81 self.retry = Some(policy);
82 self
83 }
84
85 pub fn allow_failure(mut self) -> Self {
99 self.allow_failure = true;
100 self
101 }
102}
103
104#[derive(Debug, Clone, Default, PartialEq, Eq)]
117pub struct WorkflowOptions {
118 pub allow_failure: bool,
120 concurrency_key: Option<String>,
121}
122
123impl WorkflowOptions {
124 pub fn new() -> Self {
134 Self::default()
135 }
136
137 pub fn allow_failure(mut self) -> Self {
149 self.allow_failure = true;
150 self
151 }
152
153 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 pub fn concurrency_key_ref(&self) -> Option<&str> {
201 self.concurrency_key.as_deref()
202 }
203
204 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}