use std::ops::Not;
use ironflow_core::retry::RetryPolicy;
use ironflow_store::entities::{MAX_CONCURRENCY_KEY_LEN, validate_priority};
use serde::{Deserialize, Serialize};
use serde_json::Value;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkflowStepConfig {
pub workflow_name: String,
pub payload: Value,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub retry: Option<RetryPolicy>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub concurrency_key: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub priority: Option<i16>,
#[serde(default, skip_serializing_if = "Not::not")]
pub allow_failure: bool,
}
impl WorkflowStepConfig {
pub fn new(workflow_name: &str, payload: Value) -> Self {
Self {
workflow_name: workflow_name.to_string(),
payload,
retry: None,
concurrency_key: None,
priority: None,
allow_failure: false,
}
}
pub fn retry_policy(mut self, policy: RetryPolicy) -> Self {
self.retry = Some(policy);
self
}
pub fn allow_failure(mut self) -> Self {
self.allow_failure = true;
self
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct WorkflowOptions {
pub allow_failure: bool,
concurrency_key: Option<String>,
priority: Option<i16>,
}
impl WorkflowOptions {
pub fn new() -> Self {
Self::default()
}
pub fn allow_failure(mut self) -> Self {
self.allow_failure = true;
self
}
pub fn concurrency_key(mut self, key: impl Into<String>) -> Self {
let key = key.into();
assert!(!key.trim().is_empty(), "concurrency key must not be empty");
assert!(
key.len() <= MAX_CONCURRENCY_KEY_LEN,
"concurrency key must be at most {MAX_CONCURRENCY_KEY_LEN} bytes"
);
self.concurrency_key = Some(key);
self
}
pub fn concurrency_key_ref(&self) -> Option<&str> {
self.concurrency_key.as_deref()
}
pub fn priority(mut self, priority: i16) -> Self {
if let Err(message) = validate_priority(priority) {
panic!("{message}");
}
self.priority = Some(priority);
self
}
pub fn priority_ref(&self) -> Option<i16> {
self.priority
}
pub(crate) fn into_parts(self) -> (Option<String>, Option<i16>) {
(self.concurrency_key, self.priority)
}
}
#[cfg(test)]
mod tests {
use super::*;
use ironflow_store::entities::{MAX_PRIORITY, MIN_PRIORITY};
use serde_json::{from_str, json, to_string, to_value};
#[test]
fn new_sets_fields() {
let config = WorkflowStepConfig::new("build", json!({"key": "val"}));
assert_eq!(config.workflow_name, "build");
assert_eq!(config.payload["key"], "val");
}
#[test]
fn serde_roundtrip() {
let config = WorkflowStepConfig::new("deploy", json!({"env": "prod"}));
let json = serde_json::to_string(&config).unwrap();
let back: WorkflowStepConfig = serde_json::from_str(&json).unwrap();
assert_eq!(back.workflow_name, "deploy");
assert_eq!(back.payload["env"], "prod");
}
#[test]
fn a_config_predating_retry_still_deserializes() {
let config: WorkflowStepConfig =
serde_json::from_str(r#"{"workflow_name":"build","payload":{"key":"val"}}"#)
.expect("deserialize");
assert!(config.retry.is_none());
}
#[test]
fn a_config_without_concurrency_key_omits_it() {
let config = WorkflowStepConfig::new("build", json!({}));
let value = to_value(&config).expect("serialize");
assert!(value.get("concurrency_key").is_none());
}
#[test]
fn concurrency_key_roundtrip() {
let mut config = WorkflowStepConfig::new("build", json!({}));
config.concurrency_key = Some("issue:12".to_string());
let json = to_string(&config).expect("serialize");
let back: WorkflowStepConfig = from_str(&json).expect("deserialize");
assert_eq!(back.concurrency_key.as_deref(), Some("issue:12"));
}
#[test]
fn options_carry_the_concurrency_key() {
let options = WorkflowOptions::new().concurrency_key("issue:12");
assert_eq!(options.concurrency_key_ref(), Some("issue:12"));
assert_eq!(options.into_parts().0.as_deref(), Some("issue:12"));
assert!(WorkflowOptions::new().concurrency_key_ref().is_none());
}
#[test]
fn options_accept_a_key_at_the_length_limit() {
let key = "a".repeat(MAX_CONCURRENCY_KEY_LEN);
let options = WorkflowOptions::new().concurrency_key(key.clone());
assert_eq!(options.concurrency_key_ref(), Some(key.as_str()));
}
#[test]
#[should_panic(expected = "concurrency key must not be empty")]
fn options_reject_an_empty_key() {
let _ = WorkflowOptions::new().concurrency_key("");
}
#[test]
#[should_panic(expected = "concurrency key must not be empty")]
fn options_reject_a_blank_key() {
let _ = WorkflowOptions::new().concurrency_key(" \t ");
}
#[test]
#[should_panic(expected = "concurrency key must be at most")]
fn options_reject_a_key_over_the_limit() {
let _ = WorkflowOptions::new().concurrency_key("a".repeat(MAX_CONCURRENCY_KEY_LEN + 1));
}
#[test]
fn options_priority_is_carried_and_split() {
let options = WorkflowOptions::new()
.concurrency_key("issue:12")
.priority(MAX_PRIORITY);
assert_eq!(options.priority_ref(), Some(MAX_PRIORITY));
let (key, priority) = options.into_parts();
assert_eq!(key.as_deref(), Some("issue:12"));
assert_eq!(priority, Some(MAX_PRIORITY));
assert_eq!(
WorkflowOptions::new().priority(MIN_PRIORITY).priority_ref(),
Some(MIN_PRIORITY)
);
assert!(WorkflowOptions::new().priority_ref().is_none());
}
#[test]
#[should_panic(expected = "priority must be between -100 and 100")]
fn options_priority_rejects_above_the_maximum() {
let _ = WorkflowOptions::new().priority(MAX_PRIORITY + 1);
}
#[test]
#[should_panic(expected = "priority must be between -100 and 100")]
fn options_priority_rejects_below_the_minimum() {
let _ = WorkflowOptions::new().priority(MIN_PRIORITY - 1);
}
#[test]
fn step_config_priority_is_omitted_when_unset_and_round_trips() {
let config = WorkflowStepConfig::new("build", json!({}));
assert!(
to_value(&config)
.expect("serialize")
.get("priority")
.is_none()
);
let mut config = WorkflowStepConfig::new("build", json!({}));
config.priority = Some(-30);
let back: WorkflowStepConfig =
from_str(&to_string(&config).expect("serialize")).expect("deserialize");
assert_eq!(back.priority, Some(-30));
}
#[test]
fn allow_failure_defaults_to_false() {
let config = WorkflowStepConfig::new("build", json!({}));
assert!(!config.allow_failure);
}
#[test]
fn allow_failure_is_omitted_from_json_when_false() {
let config = WorkflowStepConfig::new("build", json!({}));
let value = serde_json::to_value(&config).expect("serialize");
assert!(value.get("allow_failure").is_none());
}
#[test]
fn allow_failure_roundtrip() {
let config = WorkflowStepConfig::new("build", json!({})).allow_failure();
let json = serde_json::to_string(&config).expect("serialize");
let back: WorkflowStepConfig = serde_json::from_str(&json).expect("deserialize");
assert!(back.allow_failure);
}
#[test]
fn options_builder_sets_allow_failure() {
assert!(!WorkflowOptions::new().allow_failure);
assert!(WorkflowOptions::new().allow_failure().allow_failure);
}
#[test]
fn retry_policy_roundtrip() {
let config = WorkflowStepConfig::new("deploy", json!({})).retry_policy(RetryPolicy::new(3));
let json = serde_json::to_string(&config).expect("serialize");
let back: WorkflowStepConfig = serde_json::from_str(&json).expect("deserialize");
assert_eq!(back.retry.as_ref().unwrap().max_retries(), 3);
}
}