use serde::{Deserialize, Serialize};
use crate::config::validation::{require_nonempty, require_nonzero};
use crate::errors::OrionError;
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct ClusterConfig {
pub enabled: bool,
pub redis_url: String,
pub epoch_poll_interval_ms: u64,
pub instance_id: String,
}
impl Default for ClusterConfig {
fn default() -> Self {
Self {
enabled: false,
redis_url: String::new(),
epoch_poll_interval_ms: 2000,
instance_id: String::new(),
}
}
}
impl ClusterConfig {
pub fn effective_instance_id(&self) -> String {
if self.instance_id.is_empty() {
uuid::Uuid::new_v4().to_string()
} else {
self.instance_id.clone()
}
}
pub(crate) fn validate(&self) -> Result<(), OrionError> {
require_nonzero(
self.epoch_poll_interval_ms,
"cluster.epoch_poll_interval_ms",
)?;
if self.enabled {
require_nonempty(&self.redis_url, "cluster.redis_url")?;
}
if self.instance_id.len() > 64 {
return Err(OrionError::Config {
message: "cluster.instance_id must be at most 64 characters \
(it doubles as the Kafka group.instance.id)"
.to_string(),
});
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_default_is_disabled() {
let config = ClusterConfig::default();
assert!(!config.enabled);
assert_eq!(config.epoch_poll_interval_ms, 2000);
assert!(config.validate().is_ok());
}
#[test]
fn test_enabled_requires_redis_url() {
let config = ClusterConfig {
enabled: true,
..Default::default()
};
assert!(config.validate().is_err());
let config = ClusterConfig {
enabled: true,
redis_url: "redis://localhost:6379".to_string(),
..Default::default()
};
assert!(config.validate().is_ok());
}
#[test]
fn test_rejects_oversized_instance_id() {
let config = ClusterConfig {
instance_id: "x".repeat(65),
..Default::default()
};
assert!(config.validate().is_err());
}
}