mod ledger;
mod liveness;
mod mode;
mod sweep;
mod watchdog;
pub(crate) use ledger::{Ledger, RunRecord};
pub(crate) use mode::ReaperMode;
use std::path::{Path, PathBuf};
use std::sync::{Arc, OnceLock};
use crate::backend::SandboxBackend;
use crate::run_id::RunId;
pub(crate) fn record_create(
ledger: &Ledger,
backend_name: &str,
msb_path: Option<String>,
pid: u32,
started_iso: String,
name: &str,
keep_alive: bool,
) -> bool {
let wrote_record = if ledger.read_record().is_none() {
let record = RunRecord {
pid,
started_iso,
backend: backend_name.to_string(),
msb_path,
};
let _ = ledger.write_record(&record);
true
} else {
false
};
if !keep_alive {
let _ = ledger.append_sandbox(name);
}
wrote_record
}
pub(crate) fn record_stop(ledger: &Ledger, name: &str, keep_alive: bool) {
if keep_alive {
return;
}
ledger.remove_sandbox_and_prune_if_empty(name);
}
struct State {
mode: ReaperMode,
cache_dir: PathBuf,
ledger: Ledger,
watchdog_spawned: OnceLock<()>,
}
static STATE: OnceLock<State> = OnceLock::new();
fn state() -> &'static State {
STATE.get_or_init(|| {
let mode = ReaperMode::from_env();
let cache_dir = crate::cache_dir::dir();
let ledger = Ledger::new(&cache_dir, RunId::value());
State {
mode,
cache_dir,
ledger,
watchdog_spawned: OnceLock::new(),
}
})
}
fn ledger_for(cache_dir_override: Option<&Path>) -> Ledger {
match cache_dir_override {
Some(dir) => Ledger::new(dir, RunId::value()),
None => state().ledger.clone(),
}
}
pub(crate) fn sweep(backend: &Arc<dyn SandboxBackend>) {
let st = state();
if !st.mode.sweep_enabled() {
return;
}
sweep::run(backend, &st.cache_dir, RunId::value());
}
pub(crate) fn before_create(
backend: &Arc<dyn SandboxBackend>,
name: &str,
keep_alive: bool,
cache_dir_override: Option<&Path>,
) {
let st = state();
if !st.mode.ledger_enabled() {
return;
}
let ledger = ledger_for(cache_dir_override);
let pid = std::process::id();
let started_iso = liveness::process_started_iso(pid).unwrap_or_default();
let _ = record_create(
&ledger,
backend.name(),
backend
.backend_binary_path()
.map(|p| p.display().to_string()),
pid,
started_iso,
name,
keep_alive,
);
if st.mode.watchdog_enabled() {
st.watchdog_spawned
.get_or_init(|| watchdog::spawn(&st.cache_dir, backend, &st.ledger));
}
}
pub(crate) fn after_stop(name: &str, keep_alive: bool, cache_dir_override: Option<&Path>) {
let st = state();
if !st.mode.ledger_enabled() {
return;
}
record_stop(&ledger_for(cache_dir_override), name, keep_alive);
}
pub(crate) fn before_ensure_network(id: &str, cache_dir_override: Option<&Path>) {
let st = state();
if !st.mode.ledger_enabled() {
return;
}
let _ = ledger_for(cache_dir_override).append_network(id);
}
pub(crate) fn after_remove_network(id: &str, cache_dir_override: Option<&Path>) {
let st = state();
if !st.mode.ledger_enabled() {
return;
}
ledger_for(cache_dir_override).remove_network_and_prune_if_empty(id);
}
#[cfg(test)]
mod tests {
use super::*;
use std::path::PathBuf;
fn temp_cache_dir(label: &str) -> PathBuf {
let dir = std::env::temp_dir().join(format!(
"rz-reaper-facade-{label}-{}-{}",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
std::fs::create_dir_all(&dir).unwrap();
dir
}
#[test]
fn record_create_writes_the_record_exactly_once_and_appends_every_call() {
let cache = temp_cache_dir("record-once");
let ledger = Ledger::new(&cache, "run1");
let wrote_first = record_create(
&ledger,
"docker",
None,
4242,
"2025-01-01T00:00:00Z".to_string(),
"rz-run1-0",
false,
);
assert!(wrote_first, "the first call must write the record");
let record = ledger.read_record().expect("record must exist");
assert_eq!(record.pid, 4242);
assert_eq!(record.backend, "docker");
let wrote_second = record_create(
&ledger,
"docker",
None,
4242,
"2025-01-01T00:00:00Z".to_string(),
"rz-run1-1",
false,
);
assert!(!wrote_second, "a later call must not rewrite the record");
assert_eq!(
ledger.sandbox_names(),
vec!["rz-run1-0".to_string(), "rz-run1-1".to_string()],
"every call appends its own sandbox name, record-writing or not"
);
}
#[test]
fn record_create_skips_the_ledger_append_for_a_keep_alive_spec() {
let cache = temp_cache_dir("keep-alive-skip");
let ledger = Ledger::new(&cache, "run1");
record_create(
&ledger,
"docker",
None,
1,
"2025-01-01T00:00:00Z".to_string(),
"rz-reuse-0",
true,
);
assert!(
ledger.sandbox_names().is_empty(),
"a keep_alive sandbox must never be listed in the reaping ledger"
);
}
#[test]
fn record_stop_removes_the_name_and_clean_shutdown_deletes_files_once_empty() {
let cache = temp_cache_dir("stop-clean");
let ledger = Ledger::new(&cache, "run1");
record_create(
&ledger,
"docker",
None,
1,
"2025-01-01T00:00:00Z".to_string(),
"rz-run1-0",
false,
);
assert!(ledger.record_path().exists());
record_stop(&ledger, "rz-run1-0", false);
assert!(
!ledger.record_path().exists(),
"the last sandbox stopping with no networks must trigger clean-shutdown deletion"
);
assert!(!ledger.sandboxes_path().exists());
}
#[test]
fn record_stop_on_a_keep_alive_name_never_touches_the_ledger() {
let cache = temp_cache_dir("stop-keep-alive");
let ledger = Ledger::new(&cache, "run1");
record_create(
&ledger,
"docker",
None,
1,
"2025-01-01T00:00:00Z".to_string(),
"rz-run1-0",
false,
);
record_stop(&ledger, "rz-run1-0", true); assert_eq!(
ledger.sandbox_names(),
vec!["rz-run1-0".to_string()],
"a keep_alive stop must not remove a name it never listed"
);
}
#[test]
fn record_stop_leaves_files_in_place_while_a_network_is_still_listed() {
let cache = temp_cache_dir("stop-network-remains");
let ledger = Ledger::new(&cache, "run1");
record_create(
&ledger,
"docker",
None,
1,
"2025-01-01T00:00:00Z".to_string(),
"rz-run1-0",
false,
);
ledger.append_network("net-1").unwrap();
record_stop(&ledger, "rz-run1-0", false);
assert!(
ledger.record_path().exists(),
"files survive while a network is still listed, even with no sandboxes left"
);
}
struct NoopBackend;
#[async_trait::async_trait]
impl crate::backend::SandboxBackend for NoopBackend {
fn name(&self) -> &str {
"noop"
}
fn supports_native_networks(&self) -> bool {
false
}
async fn create(
&self,
_spec: crate::model::ContainerSpec,
) -> crate::error::Result<Box<dyn crate::backend::SandboxHandle>> {
unimplemented!()
}
async fn start(
&self,
_handle: &dyn crate::backend::SandboxHandle,
) -> crate::error::Result<()> {
unimplemented!()
}
async fn stop(
&self,
_handle: &dyn crate::backend::SandboxHandle,
) -> crate::error::Result<()> {
unimplemented!()
}
async fn remove(
&self,
_handle: &dyn crate::backend::SandboxHandle,
) -> crate::error::Result<()> {
unimplemented!()
}
async fn exec(
&self,
_handle: &dyn crate::backend::SandboxHandle,
_cmd: &[String],
) -> crate::error::Result<crate::model::ExecResult> {
unimplemented!()
}
async fn logs(
&self,
_handle: &dyn crate::backend::SandboxHandle,
) -> crate::error::Result<String> {
unimplemented!()
}
async fn follow_logs(
&self,
_handle: &dyn crate::backend::SandboxHandle,
_consumer: Box<dyn Fn(String) + Send + Sync>,
) -> crate::error::Result<crate::backend::FollowHandle> {
unimplemented!()
}
async fn ensure_network(&self, _network_id: &str) -> crate::error::Result<()> {
Ok(())
}
async fn remove_network(&self, _network_id: &str) -> crate::error::Result<()> {
Ok(())
}
fn cleanup_sync(&self, _container_id: &str) {}
fn remove_by_name(&self, _name: &str) {}
fn watchdog_kill_command(&self) -> Vec<String> {
vec!["true".to_string()]
}
}
#[test]
fn concurrent_before_create_callers_spawn_the_watchdog_at_most_once_ever() {
if !state().mode.watchdog_enabled() {
eprintln!(
"skipping: RIGHTSIZE_REAPER does not have the watchdog enabled in this \
environment"
);
return;
}
let backend: Arc<dyn crate::backend::SandboxBackend> = Arc::new(NoopBackend);
const THREADS: usize = 16;
let barrier = Arc::new(std::sync::Barrier::new(THREADS));
let handles: Vec<_> = (0..THREADS)
.map(|i| {
let backend = backend.clone();
let barrier = barrier.clone();
std::thread::spawn(move || {
barrier.wait();
before_create(&backend, &format!("rz-watchdog-race-test-{i}"), false, None);
})
})
.collect();
for h in handles {
h.join().expect("before_create must not panic");
}
let spawned = watchdog::SPAWN_CALL_COUNT.load(std::sync::atomic::Ordering::SeqCst);
assert!(
spawned <= 1,
"the watchdog process must be forked AT MOST ONCE across this process's \
entire lifetime, even with {THREADS} threads racing before_create \
concurrently; observed {spawned} real spawns — a second spawn means its \
losing Child's stdin was closed immediately, which would have fired that \
watchdog's reap logic against every sandbox already in the shared ledger"
);
}
#[test]
fn forced_concurrent_interleaving_shows_the_old_gate_red_and_the_new_gate_green() {
let cache = temp_cache_dir("toctou-forced-interleaving");
let ledger = Arc::new(Ledger::new(&cache, "toctou-run"));
const THREADS: usize = 16;
let read_barrier = Arc::new(std::sync::Barrier::new(THREADS));
let saw_no_record = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let handles: Vec<_> = (0..THREADS)
.map(|i| {
let ledger = ledger.clone();
let read_barrier = read_barrier.clone();
let saw_no_record = saw_no_record.clone();
std::thread::spawn(move || {
let existing = ledger.read_record();
read_barrier.wait();
if existing.is_none() {
saw_no_record.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
let record = RunRecord {
pid: i as u32,
started_iso: String::new(),
backend: "docker".to_string(),
msb_path: None,
};
let _ = ledger.write_record(&record);
}
})
})
.collect();
for h in handles {
h.join().unwrap();
}
let wrote_record_true_count = saw_no_record.load(std::sync::atomic::Ordering::SeqCst);
assert_eq!(
wrote_record_true_count, THREADS,
"forcing every thread through Ledger::read_record before any of them calls \
write_record must make every single one observe \"no record yet\" — this \
is the exact TOCTOU window before_create's watchdog gate used to depend \
on staying uncontended"
);
assert!(
wrote_record_true_count > 1,
"sanity: the old gate's failure mode requires more than one thread to see \
wrote_record = true"
);
let gate: OnceLock<()> = OnceLock::new();
let spawn_count = Arc::new(std::sync::atomic::AtomicUsize::new(0));
for _ in 0..wrote_record_true_count {
let spawn_count = spawn_count.clone();
gate.get_or_init(|| {
spawn_count.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
});
}
assert_eq!(
spawn_count.load(std::sync::atomic::Ordering::SeqCst),
1,
"the OnceLock-gated spawn must run exactly once even when \
{wrote_record_true_count} callers all believed they were first"
);
}
}