use std::path::{Path, PathBuf};
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct OomCounts {
pub oom: u64,
pub oom_kill: u64,
}
pub fn read_oom_counts(cgroup: &Path) -> Option<OomCounts> {
let text = std::fs::read_to_string(cgroup.join("memory.events")).ok()?;
let count = |name: &str| -> u64 {
text.lines()
.filter_map(|l| l.strip_prefix(name))
.filter_map(|n| n.trim().parse::<u64>().ok())
.next()
.unwrap_or(0)
};
Some(OomCounts {
oom: count("oom "),
oom_kill: count("oom_kill "),
})
}
pub fn classify_oom(cgroup: &Path, before: OomCounts) -> bool {
let Some(now) = read_oom_counts(cgroup) else {
return false;
};
now.oom_kill > before.oom_kill
}
#[cfg(target_os = "linux")]
pub fn own_cgroup() -> Option<PathBuf> {
let text = std::fs::read_to_string("/proc/self/cgroup").ok()?;
let rel = text.lines().find_map(|l| l.strip_prefix("0::"))?.trim();
let dir = PathBuf::from("/sys/fs/cgroup").join(rel.trim_start_matches('/'));
dir.exists().then_some(dir)
}
#[cfg(not(target_os = "linux"))]
pub fn own_cgroup() -> Option<PathBuf> {
None
}
pub struct OomWatch {
cgroup: Option<PathBuf>,
before: OomCounts,
}
impl OomWatch {
pub fn start() -> Self {
let cgroup = own_cgroup();
let before = cgroup
.as_deref()
.and_then(read_oom_counts)
.unwrap_or_default();
Self { cgroup, before }
}
pub fn record(&self, job_dir: &Path) {
let Some(cgroup) = self.cgroup.as_deref() else {
return;
};
if classify_oom(cgroup, self.before) {
mark_oom(job_dir);
}
}
}
pub fn was_oom_killed(job_dir: &Path) -> bool {
oom_evidence(job_dir)
}
pub fn oom_evidence(job_dir: &Path) -> bool {
job_dir.join("oom").exists()
}
pub fn mark_oom(job_dir: &Path) {
std::fs::write(job_dir.join("oom"), b"1").ok();
}
pub fn clear_oom(job_dir: &Path) {
std::fs::remove_file(job_dir.join("oom")).ok();
}
pub fn mark_user_kill(job_dir: &Path) {
std::fs::write(job_dir.join("killed-by-user"), b"1").ok();
}
pub fn was_user_killed(job_dir: &Path) -> bool {
job_dir.join("killed-by-user").exists()
}
pub fn clear_user_kill(job_dir: &Path) {
std::fs::remove_file(job_dir.join("killed-by-user")).ok();
}
pub fn oom_evidence_is_available() -> bool {
own_cgroup().is_some()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_out_of_memory_record_is_read_back_whatever_an_earlier_qex_wrote() {
let dir = std::env::temp_dir().join(format!("qex-oom-{}", std::process::id()));
std::fs::remove_dir_all(&dir).ok();
std::fs::create_dir_all(&dir).unwrap();
assert!(
!was_oom_killed(&dir),
"a new job has no out-of-memory record"
);
mark_oom(&dir);
assert!(was_oom_killed(&dir));
clear_oom(&dir);
assert!(!was_oom_killed(&dir));
for text in ["1", "job", "machine", "session"] {
std::fs::write(dir.join("oom"), text.as_bytes()).unwrap();
assert!(
was_oom_killed(&dir),
"a record of an earlier qex that holds `{text}` names a kill for memory"
);
}
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn the_mark_of_a_kill_by_a_command_can_be_cleared() {
let dir = std::env::temp_dir().join(format!("qex-userkill-{}", std::process::id()));
std::fs::remove_dir_all(&dir).ok();
std::fs::create_dir_all(&dir).unwrap();
assert!(!was_user_killed(&dir));
mark_user_kill(&dir);
assert!(was_user_killed(&dir));
clear_user_kill(&dir);
assert!(!was_user_killed(&dir));
std::fs::remove_dir_all(&dir).ok();
}
fn a_cgroup_with_events(name: &str, events: &str) -> std::path::PathBuf {
let dir = std::env::temp_dir().join(format!("qex-events-{}-{name}", std::process::id()));
std::fs::remove_dir_all(&dir).ok();
std::fs::create_dir_all(&dir).unwrap();
std::fs::write(dir.join("memory.events"), events.as_bytes()).unwrap();
dir
}
#[test]
fn a_new_kill_names_a_kill_for_memory() {
let dir = a_cgroup_with_events("kill", "low 0\nhigh 0\nmax 3\noom 1\noom_kill 1\n");
assert!(classify_oom(&dir, OomCounts::default()));
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn a_cgroup_with_no_kill_gives_no_answer() {
let dir = a_cgroup_with_events("nokill", "low 0\nhigh 0\nmax 0\noom 0\noom_kill 0\n");
assert!(!classify_oom(&dir, OomCounts::default()));
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn a_limit_that_stopped_no_process_gives_no_answer() {
let dir = a_cgroup_with_events("noproc", "low 0\nhigh 0\nmax 5\noom 2\noom_kill 0\n");
assert!(!classify_oom(&dir, OomCounts::default()));
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn a_cgroup_with_no_events_file_gives_no_answer() {
let dir = std::env::temp_dir().join(format!("qex-events-{}-none", std::process::id()));
std::fs::remove_dir_all(&dir).ok();
std::fs::create_dir_all(&dir).unwrap();
assert!(!classify_oom(&dir, OomCounts::default()));
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn a_count_that_qex_cannot_read_is_zero() {
let dir = a_cgroup_with_events("unread", "oom what\noom_kill later\n");
assert_eq!(
read_oom_counts(&dir),
Some(OomCounts {
oom: 0,
oom_kill: 0
}),
"a count that qex cannot read must not become evidence"
);
assert!(!classify_oom(&dir, OomCounts::default()));
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn a_count_from_before_the_attempt_gives_no_answer() {
let dir = a_cgroup_with_events("before", "low 0\nhigh 0\nmax 0\noom 0\noom_kill 4\n");
assert!(
!classify_oom(
&dir,
OomCounts {
oom: 0,
oom_kill: 4
}
),
"a count that this attempt did not raise says nothing about it"
);
assert!(classify_oom(
&dir,
OomCounts {
oom: 0,
oom_kill: 3
}
));
std::fs::remove_dir_all(&dir).ok();
}
#[test]
fn the_answer_of_an_attempt_reaches_the_record() {
let job = std::env::temp_dir().join(format!("qex-record-{}", std::process::id()));
std::fs::remove_dir_all(&job).ok();
std::fs::create_dir_all(&job).unwrap();
let cgroup = a_cgroup_with_events("record", "oom 0\noom_kill 0\n");
let before = read_oom_counts(&cgroup).unwrap();
std::fs::write(cgroup.join("memory.events"), b"oom 0\noom_kill 1\n").unwrap();
if classify_oom(&cgroup, before) {
mark_oom(&job);
}
assert!(
job.join("oom").exists(),
"the supervisor must write the record, and not leave the answer in the cgroup"
);
assert!(was_oom_killed(&job));
std::fs::remove_dir_all(&job).ok();
std::fs::remove_dir_all(&cgroup).ok();
}
}