1pub const METRIC_PARTITION_BYTES: &str = "mkit_server_partition_bytes";
6pub const ALERT_INTERVAL_MS: u64 = 10 * 60 * 1000;
8
9#[derive(Debug, Clone, Copy, PartialEq, Eq)]
11pub enum PressureLevel {
12 Warn,
14 Critical,
16}
17
18impl PressureLevel {
19 #[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#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
32pub struct PressureState {
33 active: [bool; 2],
34 last_emitted: [Option<u64>; 2],
35}
36
37#[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 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#[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
85pub 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 assert_eq!(
166 observe(state, 90, 100, 101).1,
167 vec![PressureLevel::Critical]
168 );
169 }
170}