pub const METRIC_PARTITION_BYTES: &str = "mkit_server_partition_bytes";
pub const ALERT_INTERVAL_MS: u64 = 10 * 60 * 1000;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PressureLevel {
Warn,
Critical,
}
impl PressureLevel {
#[must_use]
pub const fn label(self) -> &'static str {
match self {
Self::Warn => "warn",
Self::Critical => "critical",
}
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct PressureState {
active: [bool; 2],
last_emitted: [Option<u64>; 2],
}
#[must_use]
pub fn observe(
mut state: PressureState,
bytes: u64,
limit_bytes: u64,
now_ms: u64,
) -> (PressureState, Vec<PressureLevel>) {
let used = u128::from(bytes) * 100;
for (i, threshold) in [70_u128, 90].into_iter().enumerate() {
if limit_bytes == 0 || used >= u128::from(limit_bytes) * threshold {
state.active[i] = true;
} else if used <= u128::from(limit_bytes) * (threshold - 5) {
state.active[i] = false;
}
}
let highest = if state.active[1] {
Some((1, PressureLevel::Critical))
} else if state.active[0] {
Some((0, PressureLevel::Warn))
} else {
None
};
let mut emitted = Vec::new();
if let Some((i, level)) = highest
&& state.last_emitted[i].is_none_or(|last| now_ms.saturating_sub(last) >= ALERT_INTERVAL_MS)
{
state.last_emitted[i] = Some(now_ms);
emitted.push(level);
}
(state, emitted)
}
#[must_use]
#[allow(clippy::cast_precision_loss)]
pub fn percentage(bytes: u64, limit_bytes: u64) -> f64 {
if limit_bytes == 0 {
100.0
} else {
bytes as f64 / limit_bytes as f64 * 100.0
}
}
pub fn emit(level: PressureLevel, kind: &str, bytes: u64, limit_bytes: u64) {
let pct = percentage(bytes, limit_bytes);
match level {
PressureLevel::Warn => tracing::warn!(
event = "storage_pressure",
level = level.label(),
kind,
bytes,
limit_bytes,
pct
),
PressureLevel::Critical => tracing::error!(
event = "storage_pressure",
level = level.label(),
kind,
bytes,
limit_bytes,
pct
),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn thresholds_and_highest_severity() {
for (bytes, expected) in [
(69, vec![]),
(70, vec![PressureLevel::Warn]),
(89, vec![PressureLevel::Warn]),
(90, vec![PressureLevel::Critical]),
(120, vec![PressureLevel::Critical]),
] {
assert_eq!(observe(PressureState::default(), bytes, 100, 0).1, expected);
}
assert_eq!(
observe(PressureState::default(), 0, 0, 0).1,
vec![PressureLevel::Critical]
);
assert_eq!(
observe(PressureState::default(), u64::MAX, u64::MAX, 0).1,
vec![PressureLevel::Critical]
);
}
#[test]
fn hysteresis_clears_at_five_points() {
let (warn, _) = observe(PressureState::default(), 70, 100, 0);
assert_eq!(
observe(warn, 66, 100, ALERT_INTERVAL_MS).1,
vec![PressureLevel::Warn]
);
assert!(observe(warn, 65, 100, ALERT_INTERVAL_MS).1.is_empty());
let (critical, _) = observe(PressureState::default(), 90, 100, 0);
assert_eq!(
observe(critical, 86, 100, ALERT_INTERVAL_MS).1,
vec![PressureLevel::Critical]
);
assert_eq!(
observe(critical, 85, 100, ALERT_INTERVAL_MS).1,
vec![PressureLevel::Warn]
);
}
#[test]
fn rate_limit_survives_clear_reentry_and_clock_rollback() {
let (state, _) = observe(PressureState::default(), 70, 100, 100);
let (cleared, _) = observe(state, 0, 100, 101);
for at in [0, 102, ALERT_INTERVAL_MS + 99] {
assert!(observe(cleared, 70, 100, at).1.is_empty());
}
assert_eq!(
observe(cleared, 70, 100, ALERT_INTERVAL_MS + 100).1,
vec![PressureLevel::Warn]
);
assert_eq!(
observe(state, 90, 100, 101).1,
vec![PressureLevel::Critical]
);
}
}