|
9 | 9 | //! from incidental zombies. |
10 | 10 |
|
11 | 11 | #![cfg(target_os = "linux")] |
| 12 | +#![allow(unsafe_code)] |
12 | 13 |
|
13 | 14 | use std::collections::{HashMap, HashSet}; |
14 | 15 | use std::io; |
| 16 | +use std::sync::atomic::AtomicI32; |
15 | 17 | use std::sync::atomic::{AtomicU64, Ordering}; |
16 | 18 | use std::sync::{LazyLock, Mutex, MutexGuard}; |
17 | 19 | use std::time::Duration; |
18 | 20 |
|
19 | 21 | static MANAGED_CHILDREN: LazyLock<Mutex<HashMap<i32, u64>>> = |
20 | 22 | LazyLock::new(|| Mutex::new(HashMap::new())); |
21 | 23 | static NEXT_GENERATION: AtomicU64 = AtomicU64::new(1); |
| 24 | +static REAPER_EVENT_FD: AtomicI32 = AtomicI32::new(-1); |
| 25 | +const REAPER_RECOVERY_INTERVAL: Duration = Duration::from_secs(30); |
| 26 | +#[cfg(test)] |
| 27 | +static REAPER_SCAN_COUNT: AtomicU64 = AtomicU64::new(0); |
| 28 | + |
| 29 | +extern "C" fn notify_reaper_on_sigchld(_: libc::c_int) { |
| 30 | + let fd = REAPER_EVENT_FD.load(Ordering::Acquire); |
| 31 | + if fd >= 0 { |
| 32 | + let wakeup = 1_u64; |
| 33 | + // SAFETY: write is async-signal-safe; the process-lifetime eventfd is |
| 34 | + // nonblocking, so a full counter cannot stall the signal handler. |
| 35 | + unsafe { |
| 36 | + libc::write(fd, (&raw const wakeup).cast(), size_of::<u64>()); |
| 37 | + } |
| 38 | + } |
| 39 | +} |
22 | 40 |
|
23 | 41 | /// Identity of one registry entry. The generation prevents an old waiter from |
24 | 42 | /// removing a newer child that reused the same numeric PID after reap. |
@@ -110,22 +128,87 @@ pub fn wait_until_terminal(pid: u32) -> io::Result<()> { |
110 | 128 | /// Explicitly managed children remain owned by their normal waiters. Only |
111 | 129 | /// unregistered children adopted from the workload process tree are reaped. |
112 | 130 | pub fn start_orphan_reaper() -> io::Result<()> { |
| 131 | + // This is installed before the boundary starts its runtime threads or |
| 132 | + // workload children. The descriptor and handler live until process exit. |
| 133 | + // SAFETY: eventfd has scalar arguments and returns a new descriptor. |
| 134 | + let fd = unsafe { libc::eventfd(0, libc::EFD_CLOEXEC | libc::EFD_NONBLOCK) }; |
| 135 | + if fd < 0 { |
| 136 | + return Err(io::Error::last_os_error()); |
| 137 | + } |
| 138 | + REAPER_EVENT_FD.store(fd, Ordering::Release); |
| 139 | + let action = nix::sys::signal::SigAction::new( |
| 140 | + nix::sys::signal::SigHandler::Handler(notify_reaper_on_sigchld), |
| 141 | + nix::sys::signal::SaFlags::SA_RESTART | nix::sys::signal::SaFlags::SA_NOCLDSTOP, |
| 142 | + nix::sys::signal::SigSet::empty(), |
| 143 | + ); |
| 144 | + // SAFETY: the handler only reads a lock-free atomic and writes to the |
| 145 | + // nonblocking eventfd, both valid for the remaining process lifetime. |
| 146 | + unsafe { nix::sys::signal::sigaction(nix::sys::signal::Signal::SIGCHLD, &action) } |
| 147 | + .map_err(io::Error::other)?; |
113 | 148 | std::thread::Builder::new() |
114 | 149 | .name("openshell-orphan-reaper".to_string()) |
115 | | - .spawn(|| { |
| 150 | + .spawn(move || { |
116 | 151 | loop { |
117 | 152 | if let Err(error) = reap_unmanaged_children_once() { |
118 | 153 | tracing::debug!(%error, "orphan reaper scan failed"); |
119 | 154 | } |
120 | | - std::thread::sleep(Duration::from_millis(50)); |
| 155 | + if let Err(error) = wait_for_reaper_event(fd) { |
| 156 | + tracing::debug!(%error, "orphan reaper notification failed"); |
| 157 | + std::thread::sleep(REAPER_RECOVERY_INTERVAL); |
| 158 | + } |
121 | 159 | } |
122 | 160 | }) |
123 | 161 | .map(|_| ()) |
124 | 162 | } |
125 | 163 |
|
| 164 | +fn wait_for_reaper_event(fd: libc::c_int) -> io::Result<()> { |
| 165 | + let mut notification = libc::pollfd { |
| 166 | + fd, |
| 167 | + events: libc::POLLIN, |
| 168 | + revents: 0, |
| 169 | + }; |
| 170 | + // SAFETY: poll receives one valid pollfd for the process-lifetime eventfd. |
| 171 | + let ready = unsafe { |
| 172 | + libc::poll( |
| 173 | + &raw mut notification, |
| 174 | + 1, |
| 175 | + i32::try_from(REAPER_RECOVERY_INTERVAL.as_millis()).unwrap(), |
| 176 | + ) |
| 177 | + }; |
| 178 | + if ready < 0 { |
| 179 | + let error = io::Error::last_os_error(); |
| 180 | + if error.kind() == io::ErrorKind::Interrupted { |
| 181 | + return Ok(()); |
| 182 | + } |
| 183 | + return Err(error); |
| 184 | + } |
| 185 | + if ready == 0 { |
| 186 | + return Ok(()); |
| 187 | + } |
| 188 | + if notification.revents & libc::POLLIN == 0 { |
| 189 | + return Err(io::Error::other("orphan reaper eventfd is not readable")); |
| 190 | + } |
| 191 | + let mut count = 0_u64; |
| 192 | + // Drain before scanning: a SIGCHLD arriving during the next scan remains |
| 193 | + // queued and triggers another pass. A single read drains coalesced writes. |
| 194 | + // SAFETY: read writes one u64 into valid storage; the fd is nonblocking. |
| 195 | + let read = unsafe { libc::read(fd, (&raw mut count).cast(), size_of::<u64>()) }; |
| 196 | + if read == isize::try_from(size_of::<u64>()).expect("u64 size fits isize") |
| 197 | + || (read < 0 && io::Error::last_os_error().kind() == io::ErrorKind::WouldBlock) |
| 198 | + { |
| 199 | + Ok(()) |
| 200 | + } else if read < 0 { |
| 201 | + Err(io::Error::last_os_error()) |
| 202 | + } else { |
| 203 | + Err(io::Error::other("short orphan reaper eventfd read")) |
| 204 | + } |
| 205 | +} |
| 206 | + |
126 | 207 | fn reap_unmanaged_children_once() -> io::Result<usize> { |
127 | 208 | use nix::sys::wait::{WaitPidFlag, WaitStatus, waitpid}; |
128 | 209 |
|
| 210 | + #[cfg(test)] |
| 211 | + REAPER_SCAN_COUNT.fetch_add(1, Ordering::Relaxed); |
129 | 212 | let children = direct_child_pids()?; |
130 | 213 | let registry = lock(); |
131 | 214 | let mut reaped = 0; |
@@ -170,7 +253,89 @@ mod tests { |
170 | 253 | use nix::unistd::Pid; |
171 | 254 | use std::process::Command; |
172 | 255 | use std::sync::mpsc; |
173 | | - use std::time::Duration; |
| 256 | + use std::time::{Duration, Instant}; |
| 257 | + |
| 258 | + fn fork_exiting_child() -> i32 { |
| 259 | + // SAFETY: the child does not run Rust code after fork; it exits through |
| 260 | + // libc immediately, avoiding inherited test-harness locks. |
| 261 | + let pid = unsafe { libc::fork() }; |
| 262 | + if pid == 0 { |
| 263 | + // SAFETY: terminate without unwinding or touching Rust state. |
| 264 | + unsafe { libc::_exit(0) }; |
| 265 | + } |
| 266 | + assert!(pid > 0, "fork failed: {}", io::Error::last_os_error()); |
| 267 | + pid |
| 268 | + } |
| 269 | + |
| 270 | + fn wait_until(mut condition: impl FnMut() -> bool) { |
| 271 | + let deadline = Instant::now() + Duration::from_secs(5); |
| 272 | + while !condition() { |
| 273 | + assert!(Instant::now() < deadline, "condition did not become true"); |
| 274 | + std::thread::sleep(Duration::from_millis(10)); |
| 275 | + } |
| 276 | + } |
| 277 | + |
| 278 | + #[test] |
| 279 | + fn orphan_reaper_waits_while_idle_and_reaps_only_unmanaged_children() { |
| 280 | + const HELPER_ENV: &str = "OPENSHELL_ORPHAN_REAPER_TEST_HELPER"; |
| 281 | + if std::env::var_os(HELPER_ENV).is_none() { |
| 282 | + let output = Command::new(std::env::current_exe().unwrap()) |
| 283 | + .args([ |
| 284 | + "--exact", |
| 285 | + "managed_children::tests::orphan_reaper_waits_while_idle_and_reaps_only_unmanaged_children", |
| 286 | + "--nocapture", |
| 287 | + ]) |
| 288 | + .env(HELPER_ENV, "1") |
| 289 | + .output() |
| 290 | + .expect("start isolated reaper test process"); |
| 291 | + assert!( |
| 292 | + output.status.success(), |
| 293 | + "reaper test failed: {}", |
| 294 | + String::from_utf8_lossy(&output.stdout) |
| 295 | + ); |
| 296 | + return; |
| 297 | + } |
| 298 | + |
| 299 | + start_orphan_reaper().expect("start orphan reaper"); |
| 300 | + wait_until(|| REAPER_SCAN_COUNT.load(Ordering::Relaxed) > 0); |
| 301 | + let idle_scan_count = REAPER_SCAN_COUNT.load(Ordering::Relaxed); |
| 302 | + std::thread::sleep(Duration::from_millis(150)); |
| 303 | + assert_eq!(REAPER_SCAN_COUNT.load(Ordering::Relaxed), idle_scan_count); |
| 304 | + |
| 305 | + let mut registry = lock(); |
| 306 | + let managed_pid = fork_exiting_child(); |
| 307 | + let managed = registry |
| 308 | + .register(u32::try_from(managed_pid).unwrap()) |
| 309 | + .unwrap(); |
| 310 | + drop(registry); |
| 311 | + wait_until(|| { |
| 312 | + matches!( |
| 313 | + waitid( |
| 314 | + Id::Pid(Pid::from_raw(managed_pid)), |
| 315 | + WaitPidFlag::WEXITED | WaitPidFlag::WNOHANG | WaitPidFlag::WNOWAIT |
| 316 | + ), |
| 317 | + Ok(WaitStatus::Exited(..)) |
| 318 | + ) |
| 319 | + }); |
| 320 | + wait_until(|| REAPER_SCAN_COUNT.load(Ordering::Relaxed) > idle_scan_count); |
| 321 | + assert!(matches!( |
| 322 | + waitpid(Pid::from_raw(managed_pid), None), |
| 323 | + Ok(WaitStatus::Exited(..)) |
| 324 | + )); |
| 325 | + unregister(managed); |
| 326 | + |
| 327 | + let orphan_pids: Vec<i32> = (0..8).map(|_| fork_exiting_child()).collect(); |
| 328 | + wait_until(|| { |
| 329 | + let children = direct_child_pids().expect("inspect direct children"); |
| 330 | + orphan_pids.iter().all(|pid| !children.contains(pid)) |
| 331 | + }); |
| 332 | + for pid in orphan_pids { |
| 333 | + assert_eq!( |
| 334 | + waitpid(Pid::from_raw(pid), Some(WaitPidFlag::WNOHANG)), |
| 335 | + Err(nix::errno::Errno::ECHILD) |
| 336 | + ); |
| 337 | + } |
| 338 | + } |
174 | 339 |
|
175 | 340 | #[test] |
176 | 341 | fn fast_child_remains_waitable_after_orphan_reap_attempt() { |
|
0 commit comments