ironflow_engine/config/
workflow.rs1use 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#[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 = "Option::is_none")]
44 pub priority: Option<i16>,
45 #[serde(default, skip_serializing_if = "Not::not")]
48 pub allow_failure: bool,
49}
50
51impl WorkflowStepConfig {
52 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 pub fn retry_policy(mut self, policy: RetryPolicy) -> Self {
88 self.retry = Some(policy);
89 self
90 }
91
92 pub fn allow_failure(mut self) -> Self {
106 self.allow_failure = true;
107 self
108 }
109}
110
111#[derive(Debug, Clone, Default, PartialEq, Eq)]
124pub struct WorkflowOptions {
125 pub allow_failure: bool,
127 concurrency_key: Option<String>,
128 priority: Option<i16>,
129}
130
131impl WorkflowOptions {
132 pub fn new() -> Self {
142 Self::default()
143 }
144
145 pub fn allow_failure(mut self) -> Self {
157 self.allow_failure = true;
158 self
159 }
160
161 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 pub fn concurrency_key_ref(&self) -> Option<&str> {
209 self.concurrency_key.as_deref()
210 }
211
212 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 pub fn priority_ref(&self) -> Option<i16> {
256 self.priority
257 }
258
259 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}