use crate::error::fatal;
use crate::records::validate_instance_id;
use serde::Deserialize;
use spate_core::coordination::CoordinationError;
use std::time::Duration;
#[derive(Clone, Debug, Deserialize)]
#[serde(deny_unknown_fields, default)]
pub struct CoordinationConfig {
#[serde(with = "humantime_serde")]
pub lease_duration: Duration,
#[serde(with = "humantime_serde")]
pub op_timeout: Duration,
pub instance_id: Option<String>,
pub max_attempts: u32,
pub max_in_flight: u32,
#[serde(with = "humantime_serde")]
pub replan_interval: Duration,
#[serde(with = "humantime_serde")]
pub reconcile_interval: Duration,
pub startup_max_attempts: u32,
#[serde(with = "humantime_serde")]
pub rebalance_delay: Duration,
#[serde(with = "humantime_serde")]
pub drain_deadline: Duration,
}
impl Default for CoordinationConfig {
fn default() -> CoordinationConfig {
CoordinationConfig {
lease_duration: Duration::from_secs(30),
op_timeout: Duration::from_secs(10),
instance_id: None,
max_attempts: 4,
max_in_flight: 8,
replan_interval: Duration::from_secs(60),
reconcile_interval: Duration::from_secs(30),
startup_max_attempts: 8,
rebalance_delay: Duration::from_secs(20),
drain_deadline: Duration::from_secs(10),
}
}
}
impl CoordinationConfig {
pub fn validate(&self) -> Result<(), CoordinationError> {
if self.op_timeout < Duration::from_millis(50) {
return Err(fatal(format!(
"op_timeout must be >= 50ms, got {:?}",
self.op_timeout
)));
}
if self.lease_duration < Duration::from_millis(300) {
return Err(fatal(format!(
"lease_duration must be >= 300ms, got {:?}",
self.lease_duration
)));
}
if self.lease_duration < self.op_timeout * 2 {
return Err(fatal(format!(
"lease_duration ({:?}) must be >= 2 x op_timeout ({:?}): a single slow \
store write must not outlive the lease it renews",
self.lease_duration, self.op_timeout
)));
}
if self.max_attempts == 0 {
return Err(fatal("max_attempts must be >= 1"));
}
if self.max_in_flight == 0 {
return Err(fatal("max_in_flight must be >= 1"));
}
if self.replan_interval < self.lease_duration {
return Err(fatal(format!(
"replan_interval ({:?}) must be >= lease_duration ({:?}): replanning \
faster than leadership can be observed to fail is churn, not freshness",
self.replan_interval, self.lease_duration
)));
}
if let Some(id) = &self.instance_id {
validate_instance_id(id)?;
}
if self.startup_max_attempts == 0 {
return Err(fatal("startup_max_attempts must be >= 1"));
}
if self.drain_deadline < self.op_timeout {
return Err(fatal(format!(
"drain_deadline ({:?}) must be >= op_timeout ({:?}): a deadline shorter \
than one store round-trip forces every revocation before its final \
commit can land, so no drain could ever finish cooperatively",
self.drain_deadline, self.op_timeout
)));
}
Ok(())
}
#[must_use]
pub fn renew_interval(&self) -> Duration {
self.lease_duration / 3
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn defaults_validate_and_follow_the_documented_ratios() {
let config = CoordinationConfig::default();
config.validate().unwrap();
assert_eq!(config.lease_duration, Duration::from_secs(30));
assert_eq!(config.renew_interval(), Duration::from_secs(10));
assert_eq!(config.max_attempts, 4);
assert_eq!(config.max_in_flight, 8);
assert_eq!(
config.rebalance_delay,
Duration::from_secs(20),
"sized for a pod bounce, not for expensive task startup"
);
assert!(
config.drain_deadline >= config.op_timeout,
"a drain deadline below one store round-trip forces every revocation"
);
}
#[test]
fn a_zero_rebalance_delay_is_legal() {
CoordinationConfig {
rebalance_delay: Duration::ZERO,
..Default::default()
}
.validate()
.expect("zero delay is a supported configuration");
}
#[test]
fn floors_reject_with_the_rule_spelled_out() {
let cases: Vec<(CoordinationConfig, &str)> = vec![
(
CoordinationConfig {
op_timeout: Duration::from_millis(10),
..Default::default()
},
"op_timeout",
),
(
CoordinationConfig {
lease_duration: Duration::from_millis(100),
op_timeout: Duration::from_millis(50),
..Default::default()
},
"lease_duration",
),
(
CoordinationConfig {
lease_duration: Duration::from_secs(15),
op_timeout: Duration::from_secs(10),
..Default::default()
},
"2 x op_timeout",
),
(
CoordinationConfig {
max_attempts: 0,
..Default::default()
},
"max_attempts",
),
(
CoordinationConfig {
max_in_flight: 0,
..Default::default()
},
"max_in_flight",
),
(
CoordinationConfig {
replan_interval: Duration::from_secs(5),
..Default::default()
},
"replan_interval",
),
(
CoordinationConfig {
instance_id: Some("a.b".into()),
..Default::default()
},
"instance_id",
),
(
CoordinationConfig {
drain_deadline: Duration::from_millis(1),
..Default::default()
},
"drain_deadline",
),
];
for (config, needle) in cases {
let err = config.validate().unwrap_err();
assert!(err.to_string().contains(needle), "{err}");
}
}
#[test]
fn yaml_round_trip_with_humantime_and_unknown_field_rejection() {
let config: CoordinationConfig =
serde_yaml::from_str("lease_duration: 45s\nmax_in_flight: 4\ninstance_id: pod-3\n")
.unwrap();
assert_eq!(config.lease_duration, Duration::from_secs(45));
assert_eq!(config.max_in_flight, 4);
assert_eq!(config.instance_id.as_deref(), Some("pod-3"));
config.validate().unwrap();
let err = serde_yaml::from_str::<CoordinationConfig>("lease_secs: 45\n").unwrap_err();
assert!(err.to_string().contains("lease_secs"), "{err}");
}
}