use std::num::{NonZeroU32, NonZeroUsize};
use std::time::Duration;
use crate::policies::{BackoffPolicy, RestartPolicy};
use crate::tasks::{TaskRef, TaskSpec};
#[derive(Clone, Debug)]
pub struct SupervisorConfig {
pub grace: Duration,
pub max_concurrent: Option<NonZeroUsize>,
pub bus_capacity: usize,
pub registry_queue_capacity: usize,
pub restart: RestartPolicy,
pub backoff: BackoffPolicy,
pub timeout: Option<Duration>,
pub max_retries: Option<NonZeroU32>,
}
impl SupervisorConfig {
#[inline]
#[must_use]
pub fn bus_capacity_clamped(&self) -> usize {
self.bus_capacity.max(1)
}
#[inline]
#[must_use]
pub fn registry_queue_capacity_clamped(&self) -> usize {
self.registry_queue_capacity.max(1)
}
pub fn task_spec(&self, task: TaskRef) -> TaskSpec {
TaskSpec::new(task, self.restart, self.backoff, self.timeout)
.with_max_retries(self.max_retries)
}
}
impl Default for SupervisorConfig {
fn default() -> Self {
Self {
grace: Duration::from_secs(60),
max_concurrent: None,
bus_capacity: 1024,
registry_queue_capacity: 1024,
timeout: None,
restart: RestartPolicy::default(),
backoff: BackoffPolicy::default(),
max_retries: None,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn default_has_no_concurrency_limit_or_timeout() {
let cfg = SupervisorConfig::default();
assert_eq!(cfg.max_concurrent, None, "default is unlimited concurrency");
assert_eq!(cfg.timeout, None, "default has no per-task timeout");
assert_eq!(cfg.registry_queue_capacity, 1024);
}
#[test]
fn optional_fields_round_trip() {
let cfg = SupervisorConfig {
max_concurrent: NonZeroUsize::new(4),
timeout: Some(Duration::from_secs(30)),
max_retries: NonZeroU32::new(5),
..Default::default()
};
assert_eq!(cfg.max_concurrent.map(NonZeroUsize::get), Some(4));
assert_eq!(cfg.timeout, Some(Duration::from_secs(30)));
assert_eq!(cfg.max_retries.map(NonZeroU32::get), Some(5));
}
#[test]
fn bus_capacity_clamped_never_zero() {
let cfg = SupervisorConfig {
bus_capacity: 0,
..Default::default()
};
assert_eq!(cfg.bus_capacity_clamped(), 1);
}
#[test]
fn registry_queue_capacity_clamped_never_zero() {
let cfg = SupervisorConfig {
registry_queue_capacity: 0,
..Default::default()
};
assert_eq!(cfg.registry_queue_capacity_clamped(), 1);
}
#[test]
fn task_spec_carries_config_defaults() {
use crate::{TaskContext, TaskFn, TaskRef};
let task: TaskRef = TaskFn::arc("bridge", |_ctx: TaskContext| async { Ok(()) });
let cfg = SupervisorConfig {
max_retries: NonZeroU32::new(7),
timeout: Some(Duration::from_secs(3)),
..Default::default()
};
let spec = cfg.task_spec(task);
assert_eq!(
spec.max_retries(),
NonZeroU32::new(7),
"task_spec must bridge max_retries from config into the spec"
);
assert_eq!(
spec.timeout(),
Some(Duration::from_secs(3)),
"task_spec must bridge timeout from config into the spec"
);
}
}