use std::path::{Path, PathBuf};
use std::time::{Duration, SystemTime};
const PORTABLE_PREFIX: &str = "kglite_portable_";
const MEMORY_LIMIT_PREFIX: &str = "kglite_spill_";
const SPILL_PREFIXES: [&str; 2] = [PORTABLE_PREFIX, MEMORY_LIMIT_PREFIX];
const ORPHAN_MIN_AGE: Duration = Duration::from_secs(3600);
fn spill_root() -> PathBuf {
match std::env::var_os("KGLITE_TMPDIR") {
Some(dir) if !dir.is_empty() => PathBuf::from(dir),
_ => std::env::temp_dir(),
}
}
static NEXT_SPILL_DIR: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
fn spill_dir_name(prefix: &str, pid: u32, nanos: u128, seq: u64) -> String {
format!("{prefix}{pid}_{nanos:x}{seq:016x}")
}
fn mint_spill_dir(prefix: &str) -> PathBuf {
sweep_once();
spill_root().join(spill_dir_name(
prefix,
std::process::id(),
SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.unwrap_or_default()
.as_nanos(),
NEXT_SPILL_DIR.fetch_add(1, std::sync::atomic::Ordering::Relaxed),
))
}
pub(super) fn portable_temp_dir() -> PathBuf {
mint_spill_dir(PORTABLE_PREFIX)
}
pub(crate) fn memory_limit_temp_dir() -> PathBuf {
mint_spill_dir(MEMORY_LIMIT_PREFIX)
}
fn parse_spill_pid(name: &str) -> Option<u32> {
let rest = SPILL_PREFIXES
.iter()
.find_map(|prefix| name.strip_prefix(prefix))?;
let (pid, suffix) = rest.split_once('_')?;
if pid.is_empty() || !pid.bytes().all(|b| b.is_ascii_digit()) {
return None;
}
if suffix.is_empty() || !suffix.bytes().all(|b| b.is_ascii_hexdigit()) {
return None;
}
match pid.parse::<u32>() {
Ok(0) | Err(_) => None,
Ok(pid) => Some(pid),
}
}
#[cfg(unix)]
fn pid_is_alive(pid: u32) -> bool {
let Ok(pid) = i32::try_from(pid) else {
return true;
};
if unsafe { libc::kill(pid, 0) } == 0 {
return true;
}
std::io::Error::last_os_error().raw_os_error() != Some(libc::ESRCH)
}
#[cfg(unix)]
#[derive(Debug, Default, PartialEq, Eq)]
struct SweepOutcome {
removed: usize,
failed: usize,
}
#[cfg(unix)]
fn is_reclaimable(path: &Path, name: &str, self_pid: u32, now: SystemTime) -> bool {
let Some(pid) = parse_spill_pid(name) else {
return false;
};
if pid == self_pid {
return false;
}
let Ok(meta) = std::fs::symlink_metadata(path) else {
return false;
};
if !meta.file_type().is_dir() {
return false;
}
let Ok(age) = meta.modified().and_then(|m| {
now.duration_since(m)
.map_err(|_| std::io::Error::other("mtime is in the future"))
}) else {
return false;
};
if age < ORPHAN_MIN_AGE {
return false;
}
!pid_is_alive(pid)
}
#[cfg(unix)]
fn sweep_orphans(root: &Path, self_pid: u32, now: SystemTime) -> SweepOutcome {
let mut outcome = SweepOutcome::default();
let Ok(entries) = std::fs::read_dir(root) else {
return outcome;
};
for entry in entries.flatten() {
let path = entry.path();
let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
continue;
};
if !is_reclaimable(&path, name, self_pid, now) {
continue;
}
match std::fs::remove_dir_all(&path) {
Ok(()) => outcome.removed += 1,
Err(_) => outcome.failed += 1,
}
}
outcome
}
#[cfg(unix)]
fn sweep_once() {
static SWEEP: std::sync::Once = std::sync::Once::new();
SWEEP.call_once(|| {
let _ = sweep_orphans(&spill_root(), std::process::id(), SystemTime::now());
});
}
#[cfg(not(unix))]
fn sweep_once() {}
#[cfg(all(test, unix))]
mod tests {
use super::*;
use std::fs;
fn mint(root: &Path, name: &str, age: Duration) -> PathBuf {
let dir = root.join(name);
fs::create_dir_all(dir.join("type_0")).unwrap();
fs::write(dir.join("type_0").join("col.bin"), b"payload").unwrap();
fs::File::open(&dir)
.unwrap()
.set_modified(SystemTime::now() - age)
.unwrap();
dir
}
fn dead_pid() -> u32 {
let mut child = std::process::Command::new("/bin/sh")
.args(["-c", "exit 0"])
.spawn()
.unwrap();
let pid = child.id();
child.wait().unwrap();
pid
}
#[test]
fn two_spills_reading_one_clock_tick_get_different_directories() {
let pid = std::process::id();
let tick = 0x1234_5678_9abc_u128;
for prefix in SPILL_PREFIXES {
let first = spill_dir_name(prefix, pid, tick, 0);
let second = spill_dir_name(prefix, pid, tick, 1);
assert_ne!(
first, second,
"two spills in one clock tick shared a directory ({prefix})"
);
assert_eq!(parse_spill_pid(&first), Some(pid));
assert_eq!(parse_spill_pid(&second), Some(pid));
}
assert_ne!(
spill_dir_name(PORTABLE_PREFIX, pid, tick, 0),
spill_dir_name(MEMORY_LIMIT_PREFIX, pid, tick, 0)
);
}
#[test]
fn parses_only_the_exact_spill_name() {
assert_eq!(parse_spill_pid("kglite_portable_1234_1a2b3c"), Some(1234));
assert_eq!(parse_spill_pid("kglite_spill_1234_1a2b3c"), Some(1234));
for prefix in SPILL_PREFIXES {
for suffix in [
"1234",
"1234_",
"_1a2b",
"1234_1a2b_3c",
"12x4_1a2b",
"+12_1a2b",
"1234_zzz",
"1234_+1a",
"0_1a2b",
"99999999999999999999_1a2b",
"notapid_x",
] {
let bad = format!("{prefix}{suffix}");
assert_eq!(parse_spill_pid(&bad), None, "{bad} parsed as a spill dir");
}
}
for bad in [
"kglite-portable-1234-1a2b",
"kglite-spill-1234-1a2b",
"kglite_spilled_1234_1a2b",
"someone-elses-data",
".DS_Store",
] {
assert_eq!(parse_spill_pid(bad), None, "{bad} parsed as a spill dir");
}
}
#[test]
fn sweeps_old_dead_pid_dirs_only() {
for prefix in SPILL_PREFIXES {
let root = tempfile::tempdir().unwrap();
let dead = dead_pid();
let old_orphan = mint(
root.path(),
&format!("{prefix}{dead}_aa"),
Duration::from_secs(7200),
);
let young_orphan = mint(
root.path(),
&format!("{prefix}{dead}_bb"),
Duration::from_secs(60),
);
let live = mint(
root.path(),
&format!("{prefix}{}_cc", std::process::id()),
Duration::from_secs(7200),
);
let junk = mint(root.path(), "not-a-spill-dir", Duration::from_secs(7200));
let malformed = mint(
root.path(),
&format!("{prefix}{dead}"),
Duration::from_secs(7200),
);
let outcome = sweep_orphans(root.path(), std::process::id(), SystemTime::now());
assert_eq!(
outcome,
SweepOutcome {
removed: 1,
failed: 0
},
"{prefix}"
);
assert!(!old_orphan.exists(), "{prefix}");
assert!(young_orphan.exists(), "{prefix}");
assert!(live.exists(), "{prefix}");
assert!(junk.exists(), "{prefix}");
assert!(malformed.exists(), "{prefix}");
}
}
#[test]
fn one_sweep_reclaims_both_producers_orphans() {
let root = tempfile::tempdir().unwrap();
let dead = dead_pid();
let old = Duration::from_secs(7200);
let from_a_load = mint(root.path(), &format!("{PORTABLE_PREFIX}{dead}_aa"), old);
let from_a_limit = mint(root.path(), &format!("{MEMORY_LIMIT_PREFIX}{dead}_bb"), old);
let outcome = sweep_orphans(root.path(), std::process::id(), SystemTime::now());
assert_eq!(
outcome,
SweepOutcome {
removed: 2,
failed: 0
}
);
assert!(!from_a_load.exists());
assert!(!from_a_limit.exists());
}
#[test]
fn spares_a_symlink_wearing_a_spill_name() {
let root = tempfile::tempdir().unwrap();
let target = tempfile::tempdir().unwrap();
fs::write(target.path().join("precious.txt"), b"keep me").unwrap();
let link = root
.path()
.join(format!("{MEMORY_LIMIT_PREFIX}{}_aa", dead_pid()));
std::os::unix::fs::symlink(target.path(), &link).unwrap();
let outcome = sweep_orphans(root.path(), std::process::id(), SystemTime::now());
assert_eq!(outcome, SweepOutcome::default());
assert!(target.path().join("precious.txt").exists());
assert!(link.symlink_metadata().is_ok());
}
#[test]
fn a_failed_removal_does_not_stop_the_sweep() {
let root = tempfile::tempdir().unwrap();
let dead = dead_pid();
let old = Duration::from_secs(7200);
let blocked = mint(root.path(), &format!("{PORTABLE_PREFIX}{dead}_aa"), old);
let removable = mint(root.path(), &format!("{MEMORY_LIMIT_PREFIX}{dead}_bb"), old);
let mut perms = fs::metadata(&blocked).unwrap().permissions();
std::os::unix::fs::PermissionsExt::set_mode(&mut perms, 0o500);
fs::set_permissions(&blocked, perms).unwrap();
let outcome = sweep_orphans(root.path(), std::process::id(), SystemTime::now());
assert_eq!(
outcome,
SweepOutcome {
removed: 1,
failed: 1
}
);
assert!(!removable.exists(), "the sweep stopped at the failure");
let mut perms = fs::metadata(&blocked).unwrap().permissions();
std::os::unix::fs::PermissionsExt::set_mode(&mut perms, 0o700);
fs::set_permissions(&blocked, perms).unwrap();
}
#[test]
fn our_own_pid_is_never_swept_however_old_the_dir_looks() {
let root = tempfile::tempdir().unwrap();
let ours: Vec<PathBuf> = SPILL_PREFIXES
.iter()
.map(|prefix| {
mint(
root.path(),
&format!("{prefix}{}_aa", std::process::id()),
Duration::from_secs(86_400),
)
})
.collect();
let outcome = sweep_orphans(root.path(), std::process::id(), SystemTime::now());
assert_eq!(outcome, SweepOutcome::default());
assert!(ours.iter().all(|dir| dir.exists()));
}
#[test]
fn kglite_tmpdir_overrides_the_default_root() {
let root = tempfile::tempdir().unwrap();
std::env::set_var("KGLITE_TMPDIR", root.path());
assert_eq!(spill_root(), root.path());
assert!(portable_temp_dir().starts_with(root.path()));
assert!(memory_limit_temp_dir().starts_with(root.path()));
std::env::set_var("KGLITE_TMPDIR", "");
assert_eq!(spill_root(), std::env::temp_dir());
std::env::remove_var("KGLITE_TMPDIR");
assert_eq!(spill_root(), std::env::temp_dir());
}
}