use std::io::Write;
use std::os::unix::process::CommandExt;
use std::process::Stdio;
use std::sync::atomic::{AtomicBool, AtomicU32, Ordering};
use std::sync::{Mutex, OnceLock};
const REAPER_SCRIPT: &str = r#"
set -f
ids=""
while IFS= read -r id; do
case "$id" in
"") ;;
*) ids="${ids}${id}
" ;;
esac
done
printf '%s' "$ids" | while IFS= read -r id; do
container rm --force "$id" >/dev/null 2>&1
done
"#;
static REAPER: OnceLock<Mutex<Option<std::process::ChildStdin>>> = OnceLock::new();
static SPAWN_FAILURES: AtomicU32 = AtomicU32::new(0);
static EXHAUSTED_WARNED: AtomicBool = AtomicBool::new(false);
const MAX_SPAWN_FAILURES: u32 = 3;
fn is_valid_container_id(id: &str) -> bool {
crate::core::util::is_valid_container_id(id)
}
pub(crate) fn register(id: &str) {
if !is_valid_container_id(id) {
tracing::warn!("watchdog refused invalid container id: {id:?}");
return;
}
let mut guard = REAPER
.get_or_init(|| Mutex::new(None))
.lock()
.expect("reaper mutex must not be poisoned while registering a container");
if guard.is_none() {
if SPAWN_FAILURES.load(Ordering::Relaxed) >= MAX_SPAWN_FAILURES {
warn_exhausted_once();
return;
}
*guard = try_spawn_reaper();
if guard.is_none() {
if SPAWN_FAILURES.load(Ordering::Relaxed) >= MAX_SPAWN_FAILURES {
warn_exhausted_once();
}
return;
}
}
if write_id(
guard.as_mut().expect("reaper stdin present after spawn"),
id,
) {
return;
}
tracing::warn!("watchdog reaper pipe closed; container {id} was not registered");
*guard = None;
if SPAWN_FAILURES.load(Ordering::Relaxed) >= MAX_SPAWN_FAILURES {
warn_exhausted_once();
return;
}
*guard = try_spawn_reaper();
if let Some(w) = guard.as_mut() {
if !write_id(w, id) {
tracing::warn!("watchdog reaper pipe closed; container {id} was not registered");
*guard = None;
}
} else if SPAWN_FAILURES.load(Ordering::Relaxed) >= MAX_SPAWN_FAILURES {
warn_exhausted_once();
}
}
fn write_id(w: &mut std::process::ChildStdin, id: &str) -> bool {
writeln!(w, "{id}").and_then(|()| w.flush()).is_ok()
}
fn warn_exhausted_once() {
if !EXHAUSTED_WARNED.swap(true, Ordering::Relaxed) {
tracing::warn!(
"watchdog reaper spawn exhausted after {MAX_SPAWN_FAILURES} failures; further registration is disabled"
);
}
}
fn try_spawn_reaper() -> Option<std::process::ChildStdin> {
let child = std::process::Command::new("/bin/sh")
.arg("-c")
.arg(REAPER_SCRIPT)
.stdin(Stdio::piped())
.stdout(Stdio::null())
.stderr(Stdio::null())
.process_group(0)
.spawn();
match child {
Ok(mut c) => match c.stdin.take() {
Some(stdin) => Some(stdin),
None => {
SPAWN_FAILURES.fetch_add(1, Ordering::Relaxed);
tracing::warn!("failed to spawn watchdog reaper: stdin pipe was not created");
None
}
},
Err(err) => {
SPAWN_FAILURES.fetch_add(1, Ordering::Relaxed);
tracing::warn!("failed to spawn watchdog reaper: {err}");
None
}
}
}
#[cfg(test)]
mod tests {
use super::is_valid_container_id;
#[test]
fn accepts_valid_container_ids() {
assert!(
is_valid_container_id("c-1-2-3"),
"デフォルト形式の ID は適合するべき"
);
assert!(
is_valid_container_id("my_container.test-01"),
"許可文字のみの名前は適合するべき"
);
assert!(
is_valid_container_id("A1"),
"英数字 2 文字以上の名前は適合するべき"
);
}
#[test]
fn rejects_invalid_container_ids() {
assert!(!is_valid_container_id(""), "空文字は拒否するべき");
assert!(
!is_valid_container_id("A"),
"1 文字の ID は拒否するべき (nameValid は実質 2 文字以上)"
);
assert!(
!is_valid_container_id("has space"),
"空白を含む ID は拒否するべき"
);
assert!(
!is_valid_container_id("glob*"),
"glob 文字を含む ID は拒否するべき"
);
assert!(
!is_valid_container_id("line\nbreak"),
"改行を含む ID は拒否するべき"
);
assert!(
!is_valid_container_id("q?"),
"クエスチョンマークを含む ID は拒否するべき"
);
assert!(
!is_valid_container_id("a/b"),
"スラッシュを含む ID は拒否するべき"
);
}
}