use std::path::PathBuf;
use std::sync::atomic::{AtomicU8, AtomicU64, Ordering};
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub enum Level {
Ok = 0,
Warn = 1,
Shed = 2,
}
impl Level {
pub fn as_str(self) -> &'static str {
match self {
Level::Ok => "ok",
Level::Warn => "warn",
Level::Shed => "shed",
}
}
fn from_u8(v: u8) -> Level {
match v {
2 => Level::Shed,
1 => Level::Warn,
_ => Level::Ok,
}
}
}
pub struct Pressure {
disk_path: Option<PathBuf>,
shed_below: u64,
level: AtomicU8,
cause: AtomicU8,
last_check_ms: AtomicU64,
pub disk_free: AtomicU64,
}
const RECHECK_MS: u64 = 2_000;
impl Pressure {
pub fn refusal(&self, low_priority: bool) -> Option<String> {
match self.level() {
Level::Shed => Some(format!(
"{} pressure (shedding new work; in-flight work drains)",
self.cause()
)),
Level::Warn if low_priority => Some(format!(
"{} pressure (low-priority work sheds at warn)",
self.cause()
)),
_ => None,
}
}
}
impl Pressure {
pub fn new(disk_path: Option<PathBuf>, shed_below: u64) -> Pressure {
Pressure {
disk_path,
shed_below,
level: AtomicU8::new(0),
cause: AtomicU8::new(0),
last_check_ms: AtomicU64::new(0),
disk_free: AtomicU64::new(u64::MAX),
}
}
pub fn level(&self) -> Level {
let now = crate::state::now_ms();
let last = self.last_check_ms.load(Ordering::Relaxed);
if now.saturating_sub(last) >= RECHECK_MS
&& self
.last_check_ms
.compare_exchange(last, now, Ordering::Relaxed, Ordering::Relaxed)
.is_ok()
{
let (level, cause) = self.assess();
self.level.store(level as u8, Ordering::Relaxed);
self.cause.store(cause, Ordering::Relaxed);
}
Level::from_u8(self.level.load(Ordering::Relaxed))
}
pub fn shedding(&self) -> bool {
self.level() == Level::Shed
}
pub fn cause(&self) -> &'static str {
match self.cause.load(Ordering::Relaxed) {
1 => "disk",
2 => "memory",
_ => "none",
}
}
fn assess(&self) -> (Level, u8) {
if let Some(p) = &self.disk_path
&& self.shed_below > 0
&& let Some(free) = free_bytes(p)
{
self.disk_free.store(free, Ordering::Relaxed);
if free < self.shed_below {
return (Level::Shed, 1);
}
if free < self.shed_below.saturating_mul(2) {
return (Level::Warn, 1);
}
}
if crate::supervisor::cgroup::under_memory_pressure() {
return (Level::Shed, 2);
}
(Level::Ok, 0)
}
}
pub fn free_bytes(path: &std::path::Path) -> Option<u64> {
use std::os::unix::ffi::OsStrExt;
let c = std::ffi::CString::new(path.as_os_str().as_bytes()).ok()?;
let mut sv: libc::statvfs = unsafe { std::mem::zeroed() };
if unsafe { libc::statvfs(c.as_ptr(), &mut sv) } != 0 {
return None;
}
Some((sv.f_bavail as u64).saturating_mul(sv.f_frsize as u64))
}
pub fn parse_bytes(s: &str) -> Result<u64, String> {
let t = s.trim();
if t.is_empty() {
return Err("empty size".into());
}
let lower = t.to_ascii_lowercase();
let (num, mult) = if let Some(n) = lower
.strip_suffix("gib")
.or_else(|| lower.strip_suffix("gb"))
.or_else(|| lower.strip_suffix('g'))
{
(n, 1u64 << 30)
} else if let Some(n) = lower
.strip_suffix("mib")
.or_else(|| lower.strip_suffix("mb"))
.or_else(|| lower.strip_suffix('m'))
{
(n, 1u64 << 20)
} else if let Some(n) = lower
.strip_suffix("kib")
.or_else(|| lower.strip_suffix("kb"))
.or_else(|| lower.strip_suffix('k'))
{
(n, 1u64 << 10)
} else {
(lower.as_str(), 1u64)
};
let v: f64 = num
.trim()
.parse()
.map_err(|_| format!("invalid size {s:?} (want e.g. 256MB, 1.5GiB, or bytes)"))?;
if !(v.is_finite() && v >= 0.0) {
return Err(format!("invalid size {s:?}"));
}
Ok((v * mult as f64) as u64)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn refusal_gives_low_priority_its_teeth_one_level_early() {
let free = free_bytes(std::path::Path::new("/")).expect("statvfs /");
let warn_band = Pressure::new(Some("/".into()), free * 2 / 3);
assert_eq!(warn_band.level(), Level::Warn);
assert!(
warn_band.refusal(false).is_none(),
"normal work admits at warn"
);
let msg = warn_band.refusal(true).expect("low sheds at warn");
assert!(msg.contains("low-priority"), "{msg}");
let shedding = Pressure::new(Some("/".into()), u64::MAX);
assert_eq!(shedding.level(), Level::Shed);
assert!(shedding.refusal(false).is_some());
assert!(shedding.refusal(true).is_some());
let none = Pressure::new(None, 0);
assert_eq!(none.level(), Level::Ok);
assert!(none.refusal(true).is_none());
}
#[test]
fn sizes_parse_and_the_root_filesystem_reports_headroom() {
assert_eq!(parse_bytes("256MB").unwrap(), 256 << 20);
assert_eq!(parse_bytes("256MiB").unwrap(), 256 << 20);
assert_eq!(
parse_bytes("1.5GiB").unwrap(),
(1.5 * (1u64 << 30) as f64) as u64
);
assert_eq!(parse_bytes("1024").unwrap(), 1024);
assert_eq!(parse_bytes("0").unwrap(), 0);
assert!(parse_bytes("lots").is_err());
assert!(free_bytes(std::path::Path::new("/")).unwrap() > 0);
}
#[test]
fn levels_change_at_the_declared_thresholds() {
let p = Pressure::new(None, 256 << 20);
assert_eq!(p.assess().0, Level::Ok);
}
}