use std::sync::Once;
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(default, deny_unknown_fields)]
pub struct ConcurrencyConfig {
pub max_threads: Option<usize>,
}
static POOL_INIT: Once = Once::new();
pub(crate) fn resolve_thread_budget(config: Option<&ConcurrencyConfig>) -> usize {
if let Some(n) = config.and_then(|c| c.max_threads) {
return n.max(1);
}
num_cpus::get().min(8)
}
pub(crate) fn init_thread_pools(budget: usize) {
POOL_INIT.call_once(|| {
if let Err(_err) = rayon::ThreadPoolBuilder::new().num_threads(budget).build_global() {
tracing::debug!(budget, "global rayon pool already initialized; reusing existing pool");
}
});
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn budget_prefers_user_value() {
let cfg = ConcurrencyConfig { max_threads: Some(3) };
assert_eq!(resolve_thread_budget(Some(&cfg)), 3);
}
#[test]
fn budget_auto_is_sane() {
let budget = resolve_thread_budget(None);
assert!((1..=8).contains(&budget));
}
#[test]
fn budget_clamps_to_one() {
let cfg = ConcurrencyConfig { max_threads: Some(0) };
assert_eq!(resolve_thread_budget(Some(&cfg)), 1);
}
}