use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct WorkerPoolConfig {
#[serde(default = "default_min_threads")]
pub min_threads: usize,
#[serde(default)]
pub max_threads: usize,
#[serde(default = "default_grow_below")]
pub grow_below: f64,
#[serde(default = "default_shrink_above")]
pub shrink_above: f64,
#[serde(default = "default_emergency_above")]
pub emergency_above: f64,
#[serde(default = "default_memory_pressure_cap")]
pub memory_pressure_cap: f64,
#[serde(default = "default_scale_interval_secs")]
pub scale_interval_secs: u64,
#[serde(default = "default_async_concurrency")]
pub async_concurrency: usize,
#[serde(default = "default_health_saturation_timeout_secs")]
pub health_saturation_timeout_secs: u64,
}
fn default_min_threads() -> usize {
2
}
fn default_grow_below() -> f64 {
0.60
}
fn default_shrink_above() -> f64 {
0.85
}
fn default_emergency_above() -> f64 {
0.95
}
fn default_memory_pressure_cap() -> f64 {
0.80
}
fn default_scale_interval_secs() -> u64 {
5
}
fn default_async_concurrency() -> usize {
32
}
fn default_health_saturation_timeout_secs() -> u64 {
30
}
impl Default for WorkerPoolConfig {
fn default() -> Self {
Self {
min_threads: default_min_threads(),
max_threads: 0,
grow_below: default_grow_below(),
shrink_above: default_shrink_above(),
emergency_above: default_emergency_above(),
memory_pressure_cap: default_memory_pressure_cap(),
scale_interval_secs: default_scale_interval_secs(),
async_concurrency: default_async_concurrency(),
health_saturation_timeout_secs: default_health_saturation_timeout_secs(),
}
}
}
impl WorkerPoolConfig {
pub fn from_cascade(key: &str) -> Result<Self, crate::config::ConfigError> {
let (mut pool_cfg, min_explicit) = if let Some(cfg) = crate::config::try_get() {
let parsed: Self = cfg.unmarshal_key(key).unwrap_or_default();
let explicit = cfg.contains(&format!("{key}.min_threads"));
(parsed, explicit)
} else {
tracing::debug!("Config cascade not initialised, using default WorkerPoolConfig");
(Self::default(), false)
};
if !min_explicit {
pool_cfg.clamp_derived_min(detected_parallelism());
}
pool_cfg.validate()?;
Ok(pool_cfg)
}
pub fn validate(&self) -> Result<(), crate::config::ConfigError> {
if self.max_threads != 0 && self.min_threads > self.max_threads {
return Err(crate::config::ConfigError::InvalidValue {
key: "worker_pool.min_threads".into(),
reason: format!(
"min_threads ({}) > max_threads ({})",
self.min_threads, self.max_threads
),
});
}
if self.grow_below >= self.shrink_above {
return Err(crate::config::ConfigError::InvalidValue {
key: "worker_pool.grow_below".into(),
reason: format!(
"grow_below ({}) >= shrink_above ({})",
self.grow_below, self.shrink_above
),
});
}
if self.shrink_above >= self.emergency_above {
return Err(crate::config::ConfigError::InvalidValue {
key: "worker_pool.shrink_above".into(),
reason: format!(
"shrink_above ({}) >= emergency_above ({})",
self.shrink_above, self.emergency_above
),
});
}
if self.async_concurrency == 0 {
return Err(crate::config::ConfigError::InvalidValue {
key: "worker_pool.async_concurrency".into(),
reason: "must be >= 1 (fan_out_async iterates via step_by)".into(),
});
}
if self.min_threads == 0 {
return Err(crate::config::ConfigError::InvalidValue {
key: "worker_pool.min_threads".into(),
reason: "must be >= 1 (zero permits busy-spin the pool)".into(),
});
}
Ok(())
}
pub fn resolve_max_threads(&mut self) {
self.max_threads = self.effective_max(detected_parallelism());
}
fn effective_max(&self, available: usize) -> usize {
if self.max_threads == 0 {
available
} else {
self.max_threads.min(available)
}
}
fn clamp_derived_min(&mut self, available: usize) {
let ceiling = self.effective_max(available).max(1);
if self.min_threads > ceiling {
tracing::info!(
derived_min = self.min_threads,
clamped_min = ceiling,
available,
"worker_pool min_threads default exceeds the CPU-derived max_threads; clamping min down to max"
);
self.min_threads = ceiling;
}
}
}
fn detected_parallelism() -> usize {
std::thread::available_parallelism().map_or(4, std::num::NonZero::get)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn validate_rejects_zero_async_concurrency() {
let cfg = WorkerPoolConfig {
async_concurrency: 0,
..Default::default()
};
let err = cfg.validate().unwrap_err();
assert!(matches!(
err,
crate::config::ConfigError::InvalidValue { .. }
));
}
#[test]
fn validate_accepts_one_async_concurrency() {
let cfg = WorkerPoolConfig {
async_concurrency: 1,
..Default::default()
};
assert!(cfg.validate().is_ok());
}
#[test]
fn one_cpu_derivation_clamps_default_min() {
let mut cfg = WorkerPoolConfig::default();
cfg.clamp_derived_min(1);
assert_eq!(cfg.min_threads, 1);
cfg.max_threads = cfg.effective_max(1);
assert_eq!(cfg.max_threads, 1);
assert!(cfg.validate().is_ok());
}
#[test]
fn one_cpu_explicit_min_still_fails_validation() {
let mut cfg = WorkerPoolConfig {
min_threads: 4,
..Default::default()
};
cfg.max_threads = cfg.effective_max(1);
let err = cfg.validate().unwrap_err();
assert!(matches!(
err,
crate::config::ConfigError::InvalidValue { .. }
));
}
#[test]
fn clamp_derived_min_noop_when_ceiling_suffices() {
let mut cfg = WorkerPoolConfig::default();
cfg.clamp_derived_min(8);
assert_eq!(cfg.min_threads, 2);
}
#[test]
fn validate_rejects_zero_min_threads() {
let cfg = WorkerPoolConfig {
min_threads: 0,
..Default::default()
};
let err = cfg.validate().unwrap_err();
assert!(matches!(
err,
crate::config::ConfigError::InvalidValue { .. }
));
}
}