Skip to main content

millipede_core/autoscale/
system_status.rs

1use super::Snapshotter;
2use tokio::time::Instant;
3
4/// The concurrency adjustment recommended by current system load.
5#[non_exhaustive]
6#[derive(Debug, Clone, Copy, PartialEq, Eq)]
7pub enum ScaleDecision {
8    /// Increase desired concurrency.
9    ScaleUp,
10    /// Decrease desired concurrency.
11    ScaleDown,
12    /// Keep desired concurrency unchanged.
13    Hold,
14}
15
16/// Options controlling load-history evaluation.
17#[derive(Debug, Clone, Default)]
18#[non_exhaustive]
19#[must_use = "system status options do nothing unless passed to SystemStatus::new"]
20pub struct SystemStatusOptions {
21    /// Minimum observations a signal needs before it contributes to scale-up.
22    pub min_samples: usize,
23}
24
25/// Evaluates load-signal histories into scaling decisions.
26pub struct SystemStatus {
27    options: SystemStatusOptions,
28}
29
30impl SystemStatus {
31    /// Creates a status evaluator from the supplied options.
32    pub fn new(options: SystemStatusOptions) -> Self {
33        Self { options }
34    }
35
36    /// Evaluates all configured signals at the current Tokio-clock instant.
37    pub fn evaluate(
38        &self,
39        snapshotter: &Snapshotter,
40        desired_utilization_ratio: f32,
41        _now: Instant,
42    ) -> ScaleDecision {
43        let min_samples = self.options.min_samples.max(1);
44        let mut ratios = Vec::new();
45
46        for signal in snapshotter.signals() {
47            let samples = signal.sample(snapshotter.window());
48
49            if samples
50                .iter()
51                .max_by_key(|sample| sample.at)
52                .is_some_and(|sample| sample.overloaded)
53            {
54                return ScaleDecision::ScaleDown;
55            }
56
57            if samples.len() >= min_samples {
58                let healthy = samples.iter().filter(|sample| !sample.overloaded).count();
59                ratios.push(healthy as f32 / samples.len() as f32);
60            }
61        }
62
63        if ratios.is_empty() {
64            return ScaleDecision::Hold;
65        }
66
67        let mean = ratios.iter().sum::<f32>() / ratios.len() as f32;
68        if mean >= desired_utilization_ratio {
69            ScaleDecision::ScaleUp
70        } else {
71            ScaleDecision::Hold
72        }
73    }
74}