Skip to main content

mkit_server/telemetry/
pressure.rs

1//! Alert on physical bytes against the put soft limit. No storage I/O or
2//! wall clock is hidden here: callers supply measurements and time.
3
4/// Gauge of physical bytes, labelled by partition kind (or `database`).
5pub const METRIC_PARTITION_BYTES: &str = "mkit_server_partition_bytes";
6/// Minimum time between alerts of the same level, even after re-entry.
7pub const ALERT_INTERVAL_MS: u64 = 10 * 60 * 1000;
8
9/// Storage-pressure severity.
10#[derive(Debug, Clone, Copy, PartialEq, Eq)]
11pub enum PressureLevel {
12    /// At least 70% of the physical soft limit, clearing at 65%.
13    Warn,
14    /// At least 90% of the physical soft limit, clearing at 85%.
15    Critical,
16}
17
18impl PressureLevel {
19    /// Structured log label.
20    #[must_use]
21    pub const fn label(self) -> &'static str {
22        match self {
23            Self::Warn => "warn",
24            Self::Critical => "critical",
25        }
26    }
27}
28
29/// Per-instance state. Clearing a level preserves its last emission time:
30/// bouncing around a threshold cannot bypass the ten-minute limit.
31#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
32pub struct PressureState {
33    active: [bool; 2],
34    last_emitted: [Option<u64>; 2],
35}
36
37/// Pure pressure transition with an injected clock. Emits only the highest
38/// active severity; hysteresis keeps it active until five points below its
39/// threshold. Clock rollback cannot bypass the limiter. A zero soft limit
40/// means no writable capacity and is reported as 100%.
41#[must_use]
42pub fn observe(
43    mut state: PressureState,
44    bytes: u64,
45    limit_bytes: u64,
46    now_ms: u64,
47) -> (PressureState, Vec<PressureLevel>) {
48    // Integer cross-products preserve exact boundaries even above 2^53.
49    let used = u128::from(bytes) * 100;
50    for (i, threshold) in [70_u128, 90].into_iter().enumerate() {
51        if limit_bytes == 0 || used >= u128::from(limit_bytes) * threshold {
52            state.active[i] = true;
53        } else if used <= u128::from(limit_bytes) * (threshold - 5) {
54            state.active[i] = false;
55        }
56    }
57    let highest = if state.active[1] {
58        Some((1, PressureLevel::Critical))
59    } else if state.active[0] {
60        Some((0, PressureLevel::Warn))
61    } else {
62        None
63    };
64    let mut emitted = Vec::new();
65    if let Some((i, level)) = highest
66        && state.last_emitted[i].is_none_or(|last| now_ms.saturating_sub(last) >= ALERT_INTERVAL_MS)
67    {
68        state.last_emitted[i] = Some(now_ms);
69        emitted.push(level);
70    }
71    (state, emitted)
72}
73
74/// Percentage for display only; the decision uses integer cross-products.
75#[must_use]
76#[allow(clippy::cast_precision_loss)]
77pub fn percentage(bytes: u64, limit_bytes: u64) -> f64 {
78    if limit_bytes == 0 {
79        100.0
80    } else {
81        bytes as f64 / limit_bytes as f64 * 100.0
82    }
83}
84
85/// Emit the same structured fields on both runtimes. Critical pressure is
86/// an error event; warning pressure is a warn event.
87pub fn emit(level: PressureLevel, kind: &str, bytes: u64, limit_bytes: u64) {
88    let pct = percentage(bytes, limit_bytes);
89    match level {
90        PressureLevel::Warn => tracing::warn!(
91            event = "storage_pressure",
92            level = level.label(),
93            kind,
94            bytes,
95            limit_bytes,
96            pct
97        ),
98        PressureLevel::Critical => tracing::error!(
99            event = "storage_pressure",
100            level = level.label(),
101            kind,
102            bytes,
103            limit_bytes,
104            pct
105        ),
106    }
107}
108
109#[cfg(test)]
110mod tests {
111    use super::*;
112
113    #[test]
114    fn thresholds_and_highest_severity() {
115        for (bytes, expected) in [
116            (69, vec![]),
117            (70, vec![PressureLevel::Warn]),
118            (89, vec![PressureLevel::Warn]),
119            (90, vec![PressureLevel::Critical]),
120            (120, vec![PressureLevel::Critical]),
121        ] {
122            assert_eq!(observe(PressureState::default(), bytes, 100, 0).1, expected);
123        }
124        assert_eq!(
125            observe(PressureState::default(), 0, 0, 0).1,
126            vec![PressureLevel::Critical]
127        );
128        assert_eq!(
129            observe(PressureState::default(), u64::MAX, u64::MAX, 0).1,
130            vec![PressureLevel::Critical]
131        );
132    }
133
134    #[test]
135    fn hysteresis_clears_at_five_points() {
136        let (warn, _) = observe(PressureState::default(), 70, 100, 0);
137        assert_eq!(
138            observe(warn, 66, 100, ALERT_INTERVAL_MS).1,
139            vec![PressureLevel::Warn]
140        );
141        assert!(observe(warn, 65, 100, ALERT_INTERVAL_MS).1.is_empty());
142        let (critical, _) = observe(PressureState::default(), 90, 100, 0);
143        assert_eq!(
144            observe(critical, 86, 100, ALERT_INTERVAL_MS).1,
145            vec![PressureLevel::Critical]
146        );
147        assert_eq!(
148            observe(critical, 85, 100, ALERT_INTERVAL_MS).1,
149            vec![PressureLevel::Warn]
150        );
151    }
152
153    #[test]
154    fn rate_limit_survives_clear_reentry_and_clock_rollback() {
155        let (state, _) = observe(PressureState::default(), 70, 100, 100);
156        let (cleared, _) = observe(state, 0, 100, 101);
157        for at in [0, 102, ALERT_INTERVAL_MS + 99] {
158            assert!(observe(cleared, 70, 100, at).1.is_empty());
159        }
160        assert_eq!(
161            observe(cleared, 70, 100, ALERT_INTERVAL_MS + 100).1,
162            vec![PressureLevel::Warn]
163        );
164        // Critical has its own clock, independent of an earlier warning.
165        assert_eq!(
166            observe(state, 90, 100, 101).1,
167            vec![PressureLevel::Critical]
168        );
169    }
170}