| 1 | //! Process-tree containment for a spawned child and everything it starts. |
| 2 | //! |
| 3 | //! Moved out of `hooks/executor.rs` (where it was `HookProcessTree` / |
| 4 | //! `WindowsHookJob`) so the hook runner and the extension host share one |
| 5 | //! implementation. Killing only the immediate child can leave the real work |
| 6 | //! alive (a hook's shell runtime, a host plugin's `child_process`), so: |
| 7 | //! |
| 8 | //! * **Unix:** the child must be spawned as the leader of its own process |
| 9 | //! group (`process_group(0)`); the tree is that group, and it is SIGKILLed |
| 10 | //! as a whole. |
| 11 | //! * **Windows:** the child is assigned to a Job Object configured with |
| 12 | //! `JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE`, so closing the job kills every |
| 13 | //! process in it. Hooks are created suspended and resumed only after the |
| 14 | //! assignment. Tokio children are assigned after spawning and can start |
| 15 | //! descendants before assignment. This bounds ordinary process lifetimes; |
| 16 | //! it is not a security sandbox. |
| 17 | //! |
| 18 | //! Dropping the guard kills the tree. That is deliberate: the guard's lifetime |
| 19 | //! is the tree's lifetime, unless [`ProcessTree::release`] explicitly lets the |
| 20 | //! tree outlive it. |
| 21 | |
| 22 | #[cfg(windows)] |
| 23 | use windows::Win32::Foundation::{CloseHandle, HANDLE}; |
| 24 | #[cfg(windows)] |
| 25 | use windows::Win32::System::JobObjects::{ |
| 26 | AssignProcessToJobObject, CreateJobObjectW, JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE, |
| 27 | JOB_OBJECT_LIMIT_PROCESS_MEMORY, JOBOBJECT_EXTENDED_LIMIT_INFORMATION, |
| 28 | JobObjectExtendedLimitInformation, SetInformationJobObject, TerminateJobObject, |
| 29 | }; |
| 30 | #[cfg(windows)] |
| 31 | use windows::core::PCWSTR; |
| 32 | |
| 33 | #[cfg(windows)] |
| 34 | use windows::Win32::System::Diagnostics::ToolHelp::{ |
| 35 | CreateToolhelp32Snapshot, TH32CS_SNAPTHREAD, THREADENTRY32, Thread32First, Thread32Next, |
| 36 | }; |
| 37 | #[cfg(windows)] |
| 38 | use windows::Win32::System::Threading::{OpenThread, ResumeThread, THREAD_SUSPEND_RESUME}; |
| 39 | |
| 40 | /// Owns the process tree rooted at one spawned child. |
| 41 | pub(crate) struct ProcessTree { |
| 42 | #[cfg(unix)] |
| 43 | pgid: libc::pid_t, |
| 44 | #[cfg(windows)] |
| 45 | job: WindowsJob, |
| 46 | } |
| 47 | |
| 48 | // SAFETY (Windows): the job handle is an owned kernel handle; it is only |
| 49 | // used through `&self` calls that the OS serializes, and closed once in Drop. |
| 50 | #[cfg(windows)] |
| 51 | unsafe impl Send for ProcessTree {} |
| 52 | #[cfg(windows)] |
| 53 | unsafe impl Sync for ProcessTree {} |
| 54 | |
| 55 | impl ProcessTree { |
| 56 | /// Contain a `std::process::Child` (spawned with `process_group(0)` on Unix). |
| 57 | pub(crate) fn attach(child: &std::process::Child) -> std::io::Result<Self> { |
| 58 | #[cfg(windows)] |
| 59 | { |
| 60 | use std::os::windows::io::AsRawHandle; |
| 61 | Self::attach_parts(child.id(), child.as_raw_handle()) |
| 62 | } |
| 63 | #[cfg(not(windows))] |
| 64 | { |
| 65 | Self::attach_parts(child.id()) |
| 66 | } |
| 67 | } |
| 68 | |
| 69 | /// Contain a `tokio::process::Child` (spawned with `process_group(0)` on |
| 70 | /// Unix). Fails if the child has already been reaped. |
| 71 | pub(crate) fn attach_tokio(child: &tokio::process::Child) -> std::io::Result<Self> { |
| 72 | let pid = child |
| 73 | .id() |
| 74 | .ok_or_else(|| std::io::Error::other("child already exited"))?; |
| 75 | #[cfg(windows)] |
| 76 | { |
| 77 | let handle = child |
| 78 | .raw_handle() |
| 79 | .ok_or_else(|| std::io::Error::other("child already exited"))?; |
| 80 | Self::attach_parts(pid, handle) |
| 81 | } |
| 82 | #[cfg(not(windows))] |
| 83 | { |
| 84 | Self::attach_parts(pid) |
| 85 | } |
| 86 | } |
| 87 | |
| 88 | /// Typed owned handle from a suspended Core-selected process. The caller |
| 89 | /// attaches/caps the Job before any process thread is resumed. |
| 90 | #[cfg(windows)] |
| 91 | pub(crate) fn attach_windows_handle( |
| 92 | pid: u32, |
| 93 | process: &std::os::windows::io::OwnedHandle, |
| 94 | ) -> std::io::Result<Self> { |
| 95 | use std::os::windows::io::AsRawHandle; |
| 96 | Self::attach_parts(pid, process.as_raw_handle()) |
| 97 | } |
| 98 | |
| 99 | #[cfg(windows)] |
| 100 | fn attach_parts(_pid: u32, handle: std::os::windows::io::RawHandle) -> std::io::Result<Self> { |
| 101 | Ok(Self { |
| 102 | job: WindowsJob::attach(handle)?, |
| 103 | }) |
| 104 | } |
| 105 | |
| 106 | #[cfg(not(windows))] |
| 107 | #[allow(clippy::unnecessary_wraps)] |
| 108 | fn attach_parts(pid: u32) -> std::io::Result<Self> { |
| 109 | #[cfg(unix)] |
| 110 | { |
| 111 | Ok(Self { |
| 112 | pgid: pid as libc::pid_t, |
| 113 | }) |
| 114 | } |
| 115 | #[cfg(not(unix))] |
| 116 | { |
| 117 | let _ = pid; |
| 118 | Ok(Self {}) |
| 119 | } |
| 120 | } |
| 121 | |
| 122 | /// Kill every process in the tree. A tree that is already gone is not an |
| 123 | /// error. On failure the caller should fall back to killing the child. |
| 124 | pub(crate) fn kill(&self) -> std::io::Result<()> { |
| 125 | #[cfg(unix)] |
| 126 | { |
| 127 | // SAFETY: kill(2) dereferences no pointers; a negative pid names |
| 128 | // the process group. |
| 129 | let result = unsafe { libc::kill(-self.pgid, libc::SIGKILL) }; |
| 130 | if result != 0 { |
| 131 | let error = std::io::Error::last_os_error(); |
| 132 | if error.raw_os_error() != Some(libc::ESRCH) { |
| 133 | return Err(error); |
| 134 | } |
| 135 | } |
| 136 | Ok(()) |
| 137 | } |
| 138 | #[cfg(windows)] |
| 139 | { |
| 140 | self.job.terminate() |
| 141 | } |
| 142 | #[cfg(not(any(unix, windows)))] |
| 143 | { |
| 144 | Err(std::io::Error::other( |
| 145 | "process-tree containment is not supported on this platform", |
| 146 | )) |
| 147 | } |
| 148 | } |
| 149 | } |
| 150 | |
| 151 | impl ProcessTree { |
| 152 | /// Windows: cap the committed memory of every process in the job at |
| 153 | /// `bytes` each (`JOB_OBJECT_LIMIT_PROCESS_MEMORY`); an allocation past |
| 154 | /// it fails. Kill-on-close stays set. The extension host's memory cap. |
| 155 | #[cfg(windows)] |
| 156 | pub(crate) fn limit_process_memory(&self, bytes: u64) -> std::io::Result<()> { |
| 157 | self.job.limit_process_memory(bytes) |
| 158 | } |
| 159 | |
| 160 | /// Give up containment without killing anything: the tree outlives the |
| 161 | /// guard. For a command that exited on its own and may have deliberately |
| 162 | /// left something running. |
| 163 | pub(crate) fn release(self) { |
| 164 | #[cfg(windows)] |
| 165 | { |
| 166 | // Closing the handle kills the job unless the limit is cleared |
| 167 | // first. If clearing fails, the close still kills: the safe side. |
| 168 | let _ = self.job.clear_kill_on_close(); |
| 169 | } |
| 170 | #[cfg(not(windows))] |
| 171 | std::mem::forget(self); |
| 172 | } |
| 173 | } |
| 174 | |
| 175 | impl Drop for ProcessTree { |
| 176 | fn drop(&mut self) { |
| 177 | #[cfg(unix)] |
| 178 | // SAFETY: kill(2) dereferences no pointers. |
| 179 | unsafe { |
| 180 | // The leader may have exited while a descendant still holds an |
| 181 | // inherited pipe. Reaping the group keeps lifetimes bounded. |
| 182 | let _ = libc::kill(-self.pgid, libc::SIGKILL); |
| 183 | } |
| 184 | // On Windows, dropping `WindowsJob` closes a Job Object configured |
| 185 | // with JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE. |
| 186 | } |
| 187 | } |
| 188 | |
| 189 | /// Start a synchronous child inside the same process group / Job used by hooks. |
| 190 | /// Windows starts suspended and resumes only after Job assignment. The caller |
| 191 | /// still owns environment policy and bounded waits; this is lifetime containment. |
| 192 | pub(crate) fn spawn_contained_std( |
| 193 | command: &mut std::process::Command, |
| 194 | ) -> std::io::Result<(std::process::Child, ProcessTree)> { |
| 195 | use wait_timeout::ChildExt; |
| 196 | #[cfg(unix)] |
| 197 | { |
| 198 | use std::os::unix::process::CommandExt; |
| 199 | command.process_group(0); |
| 200 | } |
| 201 | #[cfg(windows)] |
| 202 | { |
| 203 | use std::os::windows::process::CommandExt; |
| 204 | command.creation_flags(0x0000_0004 | 0x0800_0000); // suspended, no window |
| 205 | } |
| 206 | let mut child = command.spawn()?; |
| 207 | let tree = match ProcessTree::attach(&child) { |
| 208 | Ok(tree) => tree, |
| 209 | Err(error) => { |
| 210 | let _ = child.kill(); |
| 211 | let _ = child.wait_timeout(std::time::Duration::from_millis(250)); |
| 212 | return Err(error); |
| 213 | } |
| 214 | }; |
| 215 | #[cfg(windows)] |
| 216 | if let Err(error) = resume_windows_process(&child) { |
| 217 | drop(tree); |
| 218 | let _ = child.kill(); |
| 219 | let _ = child.wait_timeout(std::time::Duration::from_millis(250)); |
| 220 | return Err(error); |
| 221 | } |
| 222 | Ok((child, tree)) |
| 223 | } |
| 224 | |
| 225 | #[cfg(windows)] |
| 226 | fn resume_windows_process(child: &std::process::Child) -> std::io::Result<()> { |
| 227 | let snapshot = |
| 228 | // SAFETY: returned handle is owned here; closed before return. |
| 229 | unsafe { CreateToolhelp32Snapshot(TH32CS_SNAPTHREAD, 0).map_err(windows_io_error)? }; |
| 230 | let result = (|| { |
| 231 | let mut entry = THREADENTRY32 { |
| 232 | dwSize: std::mem::size_of::<THREADENTRY32>() as u32, |
| 233 | ..Default::default() |
| 234 | }; |
| 235 | // SAFETY: `entry` is live with dwSize initialized above. |
| 236 | let mut next = unsafe { Thread32First(snapshot, &mut entry) }; |
| 237 | let mut resumed = 0usize; |
| 238 | while next.is_ok() { |
| 239 | if entry.th32OwnerProcessID == child.id() { |
| 240 | // SAFETY: returned handle is owned here; closed below. |
| 241 | let thread = unsafe { |
| 242 | OpenThread(THREAD_SUSPEND_RESUME, false, entry.th32ThreadID) |
| 243 | .map_err(windows_io_error)? |
| 244 | }; |
| 245 | // SAFETY: `thread` is a live owned handle. |
| 246 | let resume_result = unsafe { ResumeThread(thread) }; |
| 247 | // SAFETY: `thread` is owned here and not used after. |
| 248 | let close_result = unsafe { CloseHandle(thread).map_err(windows_io_error) }; |
| 249 | if resume_result == u32::MAX { |
| 250 | return Err(std::io::Error::last_os_error()); |
| 251 | } |
| 252 | close_result?; |
| 253 | resumed += 1; |
| 254 | } |
| 255 | // SAFETY: `entry` is live with dwSize initialized above. |
| 256 | next = unsafe { Thread32Next(snapshot, &mut entry) }; |
| 257 | } |
| 258 | if resumed == 0 { |
| 259 | return Err(std::io::Error::other( |
| 260 | "suspended process had no resumable thread", |
| 261 | )); |
| 262 | } |
| 263 | Ok(()) |
| 264 | })(); |
| 265 | // SAFETY: `snapshot` is owned here and not used after. |
| 266 | let close_result = unsafe { CloseHandle(snapshot).map_err(windows_io_error) }; |
| 267 | result?; |
| 268 | close_result |
| 269 | } |
| 270 | |
| 271 | /// `Command::output()` for a child whose whole process tree dies with the |
| 272 | /// returned future. Dropping it — a caller's `timeout` elapsing, or a |
| 273 | /// cancelled tool call — kills the child and everything it started, where a |
| 274 | /// bare `output()` left them running. Stdin is closed: callers run |
| 275 | /// non-interactive work that must never wait on input. |
| 276 | pub(crate) async fn contained_output( |
| 277 | cmd: &mut tokio::process::Command, |
| 278 | ) -> std::io::Result<std::process::Output> { |
| 279 | contained_output_until(cmd, std::future::pending()) |
| 280 | .await |
| 281 | .map(|run| run.output) |
| 282 | } |
| 283 | |
| 284 | /// What [`contained_output_until`] collected. |
| 285 | pub(crate) struct ContainedOutput { |
| 286 | pub(crate) output: std::process::Output, |
| 287 | /// `stop` fired before the command exited: the tree was killed, and |
| 288 | /// `output` holds everything it wrote until then. |
| 289 | pub(crate) stopped: bool, |
| 290 | } |
| 291 | |
| 292 | /// [`contained_output`], plus a `stop` future (a deadline, a cancel token). |
| 293 | /// When `stop` fires first the tree is killed and the output written so far is |
| 294 | /// still returned, so a hung command leaves evidence of where it hung. |
| 295 | /// |
| 296 | /// A command that exits on its own is released, not killed: whatever it |
| 297 | /// deliberately left running with its output redirected (`nohup server >log |
| 298 | /// &`) keeps running, as it did with a bare `output()`. Only a dropped future |
| 299 | /// or `stop` kills the tree. Until then the group is also listed for |
| 300 | /// [`kill_contained_trees_for_exit`], because a terminating signal ends |
| 301 | /// Codewhale through `process::exit`, where no destructor runs, and a child in |
| 302 | /// its own process group no longer receives the terminal's Ctrl+C itself. |
| 303 | pub(crate) async fn contained_output_until( |
| 304 | cmd: &mut tokio::process::Command, |
| 305 | stop: impl std::future::Future<Output = ()>, |
| 306 | ) -> std::io::Result<ContainedOutput> { |
| 307 | contained_run(cmd, None, stop, ProcessTree::attach_tokio).await |
| 308 | } |
| 309 | |
| 310 | /// [`contained_output`] for a command that reads a request on stdin (a plugin |
| 311 | /// tool's JSON input): `input` is written while the output pipes drain, then |
| 312 | /// stdin is closed. |
| 313 | pub(crate) async fn contained_output_with_input( |
| 314 | cmd: &mut tokio::process::Command, |
| 315 | input: Vec<u8>, |
| 316 | ) -> std::io::Result<std::process::Output> { |
| 317 | contained_run( |
| 318 | cmd, |
| 319 | Some(input), |
| 320 | std::future::pending(), |
| 321 | ProcessTree::attach_tokio, |
| 322 | ) |
| 323 | .await |
| 324 | .map(|run| run.output) |
| 325 | } |
| 326 | |
| 327 | async fn contained_run( |
| 328 | cmd: &mut tokio::process::Command, |
| 329 | input: Option<Vec<u8>>, |
| 330 | stop: impl std::future::Future<Output = ()>, |
| 331 | attach: impl FnOnce(&tokio::process::Child) -> std::io::Result<ProcessTree>, |
| 332 | ) -> std::io::Result<ContainedOutput> { |
| 333 | contained_run_with_limits(cmd, input, stop, attach, None).await |
| 334 | } |
| 335 | |
| 336 | /// Strict capture bounds for the admitted script runner. Overflow refuses the result; |
| 337 | /// it never presents a truncated ToolResult as success. Same containment driver. |
| 338 | pub(crate) async fn contained_output_with_input_bounded( |
| 339 | cmd: &mut tokio::process::Command, |
| 340 | input: Vec<u8>, |
| 341 | stdout_limit: usize, |
| 342 | stderr_limit: usize, |
| 343 | stop: impl std::future::Future<Output = ()>, |
| 344 | ) -> std::io::Result<ContainedOutput> { |
| 345 | contained_run_with_limits( |
| 346 | cmd, |
| 347 | Some(input), |
| 348 | stop, |
| 349 | ProcessTree::attach_tokio, |
| 350 | Some((stdout_limit, stderr_limit)), |
| 351 | ) |
| 352 | .await |
| 353 | } |
| 354 | async fn contained_run_with_limits( |
| 355 | cmd: &mut tokio::process::Command, |
| 356 | input: Option<Vec<u8>>, |
| 357 | stop: impl std::future::Future<Output = ()>, |
| 358 | attach: impl FnOnce(&tokio::process::Child) -> std::io::Result<ProcessTree>, |
| 359 | limits: Option<(usize, usize)>, |
| 360 | ) -> std::io::Result<ContainedOutput> { |
| 361 | use std::process::Stdio; |
| 362 | cmd.stdin(if input.is_some() { |
| 363 | Stdio::piped() |
| 364 | } else { |
| 365 | Stdio::null() |
| 366 | }) |
| 367 | .stdout(Stdio::piped()) |
| 368 | .stderr(Stdio::piped()) |
| 369 | .kill_on_drop(true); |
| 370 | #[cfg(unix)] |
| 371 | cmd.process_group(0); |
| 372 | let mut child = cmd.spawn()?; |
| 373 | // Never report a contained run when attachment failed. The child's |
| 374 | // kill-on-drop guard still cleans up the direct child on this error path. |
| 375 | let tree = attach(&child)?; |
| 376 | // Declared after `tree`, so a dropped future unlists the group first and |
| 377 | // the tree guard then kills it. |
| 378 | #[cfg(unix)] |
| 379 | let _listed = child.id().map(ContainedExitListing::new); |
| 380 | let mut stdout_pipe = child.stdout.take(); |
| 381 | let mut stderr_pipe = child.stderr.take(); |
| 382 | let stdin_pipe = child.stdin.take(); |
| 383 | let mut stdout = Vec::new(); |
| 384 | let mut stderr = Vec::new(); |
| 385 | let exited = { |
| 386 | // Written alongside the drains: a child that answers before it has |
| 387 | // read all of its input would otherwise deadlock on a full pipe. |
| 388 | let feed = async move { |
| 389 | use tokio::io::AsyncWriteExt; |
| 390 | if let (Some(mut pipe), Some(input)) = (stdin_pipe, input) { |
| 391 | // A child that exits without reading its input closes the |
| 392 | // pipe; that is its answer, not a run failure. |
| 393 | if pipe.write_all(&input).await.is_ok() { |
| 394 | let _ = pipe.shutdown().await; |
| 395 | } |
| 396 | } |
| 397 | }; |
| 398 | let run = async { |
| 399 | let out = drain_pipe_with_limit( |
| 400 | stdout_pipe.as_mut(), |
| 401 | &mut stdout, |
| 402 | limits.map(|limits| limits.0), |
| 403 | ); |
| 404 | let err = drain_pipe_with_limit( |
| 405 | stderr_pipe.as_mut(), |
| 406 | &mut stderr, |
| 407 | limits.map(|limits| limits.1), |
| 408 | ); |
| 409 | let wait = async { |
| 410 | let status = child.wait().await?; |
| 411 | if limits.is_some() { |
| 412 | let _ = tree.kill(); |
| 413 | } |
| 414 | Ok::<_, std::io::Error>(status) |
| 415 | }; |
| 416 | if limits.is_some() { |
| 417 | let (_, _, status, ()) = tokio::try_join!(out, err, wait, async { |
| 418 | feed.await; |
| 419 | Ok::<(), std::io::Error>(()) |
| 420 | })?; |
| 421 | Ok(status) |
| 422 | } else { |
| 423 | let (out, err, status, ()) = tokio::join!(out, err, wait, feed); |
| 424 | out?; |
| 425 | err?; |
| 426 | status |
| 427 | } |
| 428 | }; |
| 429 | tokio::select! { |
| 430 | status = run => match status { |
| 431 | Ok(status) => Some(status), |
| 432 | Err(error) => { |
| 433 | let _ = tree.kill(); let _ = child.start_kill(); |
| 434 | let _ = tokio::time::timeout(std::time::Duration::from_secs(2), child.wait()).await; |
| 435 | return Err(error); |
| 436 | } |
| 437 | }, |
| 438 | () = stop => None, |
| 439 | } |
| 440 | }; |
| 441 | let (status, stopped) = match exited { |
| 442 | Some(status) => { |
| 443 | if limits.is_some() { |
| 444 | drop(tree); |
| 445 | } else { |
| 446 | tree.release(); |
| 447 | } |
| 448 | (status, false) |
| 449 | } |
| 450 | None => { |
| 451 | drop(tree); |
| 452 | let _ = child.start_kill(); |
| 453 | let status = if limits.is_some() { |
| 454 | tokio::time::timeout(std::time::Duration::from_secs(2), child.wait()) |
| 455 | .await |
| 456 | .map_err(|_| { |
| 457 | std::io::Error::other("execution child did not reap after cancellation") |
| 458 | })?? |
| 459 | } else { |
| 460 | child.wait().await? |
| 461 | }; |
| 462 | // Collect what is still buffered in the pipes. Bounded: a process |
| 463 | // that escaped the tree may still hold one open. |
| 464 | let _ = tokio::time::timeout(std::time::Duration::from_secs(1), async { |
| 465 | let _ = tokio::join!( |
| 466 | drain_pipe_with_limit( |
| 467 | stdout_pipe.as_mut(), |
| 468 | &mut stdout, |
| 469 | limits.map(|limits| limits.0) |
| 470 | ), |
| 471 | drain_pipe_with_limit( |
| 472 | stderr_pipe.as_mut(), |
| 473 | &mut stderr, |
| 474 | limits.map(|limits| limits.1) |
| 475 | ), |
| 476 | ); |
| 477 | }) |
| 478 | .await; |
| 479 | (status, true) |
| 480 | } |
| 481 | }; |
| 482 | Ok(ContainedOutput { |
| 483 | output: std::process::Output { |
| 484 | status, |
| 485 | stdout, |
| 486 | stderr, |
| 487 | }, |
| 488 | stopped, |
| 489 | }) |
| 490 | } |
| 491 | |
| 492 | /// Read each chunk under the selected capture bound; overflow is an error. |
| 493 | async fn drain_pipe_with_limit<R: tokio::io::AsyncRead + Unpin>( |
| 494 | pipe: Option<&mut R>, |
| 495 | into: &mut Vec<u8>, |
| 496 | limit: Option<usize>, |
| 497 | ) -> std::io::Result<()> { |
| 498 | use tokio::io::AsyncReadExt; |
| 499 | let Some(pipe) = pipe else { |
| 500 | return Ok(()); |
| 501 | }; |
| 502 | let mut chunk = [0_u8; 8 * 1024]; |
| 503 | loop { |
| 504 | let read = pipe.read(&mut chunk).await?; |
| 505 | if read == 0 { |
| 506 | return Ok(()); |
| 507 | } |
| 508 | if limit.is_some_and(|limit| into.len().saturating_add(read) > limit) { |
| 509 | return Err(std::io::Error::other( |
| 510 | "execution output exceeded its capture limit", |
| 511 | )); |
| 512 | } |
| 513 | into.extend_from_slice(&chunk[..read]); |
| 514 | } |
| 515 | } |
| 516 | |
| 517 | /// Process groups of [`contained_output_until`] runs still in flight. |
| 518 | #[cfg(unix)] |
| 519 | static CONTAINED_EXIT_GROUPS: std::sync::OnceLock< |
| 520 | std::sync::Mutex<std::collections::HashSet<u32>>, |
| 521 | > = std::sync::OnceLock::new(); |
| 522 | |
| 523 | #[cfg(unix)] |
| 524 | fn contained_exit_groups() -> std::sync::MutexGuard<'static, std::collections::HashSet<u32>> { |
| 525 | CONTAINED_EXIT_GROUPS |
| 526 | .get_or_init(Default::default) |
| 527 | .lock() |
| 528 | .unwrap_or_else(std::sync::PoisonError::into_inner) |
| 529 | } |
| 530 | |
| 531 | /// Lists one in-flight process group until dropped. |
| 532 | #[cfg(unix)] |
| 533 | struct ContainedExitListing(u32); |
| 534 | |
| 535 | #[cfg(unix)] |
| 536 | impl ContainedExitListing { |
| 537 | fn new(process_group_id: u32) -> Self { |
| 538 | contained_exit_groups().insert(process_group_id); |
| 539 | Self(process_group_id) |
| 540 | } |
| 541 | } |
| 542 | |
| 543 | #[cfg(unix)] |
| 544 | impl Drop for ContainedExitListing { |
| 545 | fn drop(&mut self) { |
| 546 | contained_exit_groups().remove(&self.0); |
| 547 | } |
| 548 | } |
| 549 | |
| 550 | /// Kill every in-flight contained tree. The process-wide signal path calls |
| 551 | /// this immediately before `process::exit`, where Rust destructors cannot run. |
| 552 | #[cfg(unix)] |
| 553 | pub(crate) fn kill_contained_trees_for_exit() { |
| 554 | let groups = contained_exit_groups().drain().collect::<Vec<_>>(); |
| 555 | for process_group_id in groups { |
| 556 | if let Ok(process_group_id) = libc::pid_t::try_from(process_group_id) { |
| 557 | // SAFETY: the id is a child spawned with `process_group(0)` whose |
| 558 | // run is still in flight. A negative pid names that group, never |
| 559 | // Codewhale's own. |
| 560 | unsafe { |
| 561 | libc::kill(-process_group_id, libc::SIGKILL); |
| 562 | } |
| 563 | } |
| 564 | } |
| 565 | } |
| 566 | |
| 567 | #[cfg(windows)] |
| 568 | struct WindowsJob { |
| 569 | handle: HANDLE, |
| 570 | } |
| 571 | |
| 572 | #[cfg(windows)] |
| 573 | impl WindowsJob { |
| 574 | fn attach(child: std::os::windows::io::RawHandle) -> std::io::Result<Self> { |
| 575 | // SAFETY: returned handle is owned by the new wrapper. |
| 576 | let handle = unsafe { CreateJobObjectW(None, PCWSTR::null()).map_err(windows_io_error)? }; |
| 577 | let job = Self { handle }; |
| 578 | let mut limits = JOBOBJECT_EXTENDED_LIMIT_INFORMATION::default(); |
| 579 | limits.BasicLimitInformation.LimitFlags = JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE; |
| 580 | |
| 581 | // SAFETY: `limits` is live with matching size; both handles are live. |
| 582 | unsafe { |
| 583 | SetInformationJobObject( |
| 584 | job.handle, |
| 585 | JobObjectExtendedLimitInformation, |
| 586 | &limits as *const _ as *const core::ffi::c_void, |
| 587 | std::mem::size_of::<JOBOBJECT_EXTENDED_LIMIT_INFORMATION>() as u32, |
| 588 | ) |
| 589 | .map_err(windows_io_error)?; |
| 590 | AssignProcessToJobObject(job.handle, HANDLE(child)).map_err(windows_io_error)?; |
| 591 | } |
| 592 | Ok(job) |
| 593 | } |
| 594 | |
| 595 | fn limit_process_memory(&self, bytes: u64) -> std::io::Result<()> { |
| 596 | let mut limits = JOBOBJECT_EXTENDED_LIMIT_INFORMATION::default(); |
| 597 | limits.BasicLimitInformation.LimitFlags = |
| 598 | JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE | JOB_OBJECT_LIMIT_PROCESS_MEMORY; |
| 599 | limits.ProcessMemoryLimit = usize::try_from(bytes).unwrap_or(usize::MAX); |
| 600 | // SAFETY: `limits` is live with matching size; the handle is live. |
| 601 | unsafe { |
| 602 | SetInformationJobObject( |
| 603 | self.handle, |
| 604 | JobObjectExtendedLimitInformation, |
| 605 | &limits as *const _ as *const core::ffi::c_void, |
| 606 | std::mem::size_of::<JOBOBJECT_EXTENDED_LIMIT_INFORMATION>() as u32, |
| 607 | ) |
| 608 | .map_err(windows_io_error) |
| 609 | } |
| 610 | } |
| 611 | |
| 612 | fn clear_kill_on_close(&self) -> std::io::Result<()> { |
| 613 | let limits = JOBOBJECT_EXTENDED_LIMIT_INFORMATION::default(); |
| 614 | // SAFETY: `limits` is live with matching size; the handle is live. |
| 615 | unsafe { |
| 616 | SetInformationJobObject( |
| 617 | self.handle, |
| 618 | JobObjectExtendedLimitInformation, |
| 619 | &limits as *const _ as *const core::ffi::c_void, |
| 620 | std::mem::size_of::<JOBOBJECT_EXTENDED_LIMIT_INFORMATION>() as u32, |
| 621 | ) |
| 622 | .map_err(windows_io_error) |
| 623 | } |
| 624 | } |
| 625 | |
| 626 | fn terminate(&self) -> std::io::Result<()> { |
| 627 | // SAFETY: `self.handle` is a live owned job handle. |
| 628 | unsafe { TerminateJobObject(self.handle, 1).map_err(windows_io_error) } |
| 629 | } |
| 630 | } |
| 631 | |
| 632 | #[cfg(windows)] |
| 633 | impl Drop for WindowsJob { |
| 634 | fn drop(&mut self) { |
| 635 | // SAFETY: `self.handle` is owned here; Drop runs once. |
| 636 | unsafe { |
| 637 | let _ = CloseHandle(self.handle); |
| 638 | } |
| 639 | } |
| 640 | } |
| 641 | |
| 642 | #[cfg(windows)] |
| 643 | pub(crate) fn windows_io_error(error: windows::core::Error) -> std::io::Error { |
| 644 | std::io::Error::other(error) |
| 645 | } |
| 646 | |
| 647 | /// Test support: wait for `pid` to stop existing. `true` once it is gone. |
| 648 | /// A zombie counts as gone: it was killed and only awaits a reaper, which a |
| 649 | /// container whose PID 1 does not reap orphans may never provide. |
| 650 | #[cfg(all(test, unix))] |
| 651 | pub(crate) fn wait_for_pid_exit(pid: libc::pid_t, within: std::time::Duration) -> bool { |
| 652 | let deadline = std::time::Instant::now() + within; |
| 653 | loop { |
| 654 | // SAFETY: signal 0 only checks that the process exists. |
| 655 | if unsafe { libc::kill(pid, 0) } != 0 || is_zombie(pid) { |
| 656 | return true; |
| 657 | } |
| 658 | if std::time::Instant::now() >= deadline { |
| 659 | // Do not leak the fixture when the assertion is about to fail. |
| 660 | // SAFETY: kill(2) dereferences no pointers. |
| 661 | unsafe { libc::kill(pid, libc::SIGKILL) }; |
| 662 | return false; |
| 663 | } |
| 664 | std::thread::sleep(std::time::Duration::from_millis(50)); |
| 665 | } |
| 666 | } |
| 667 | |
| 668 | /// Linux only (`/proc`); elsewhere launchd/init reaps orphans promptly. |
| 669 | #[cfg(all(test, unix))] |
| 670 | fn is_zombie(pid: libc::pid_t) -> bool { |
| 671 | std::fs::read_to_string(format!("/proc/{pid}/stat")) |
| 672 | .ok() |
| 673 | .and_then(|stat| { |
| 674 | stat.rsplit_once(')') |
| 675 | .map(|(_, rest)| rest.trim_start().starts_with('Z')) |
| 676 | }) |
| 677 | .unwrap_or(false) |
| 678 | } |
| 679 | |
| 680 | /// Test support: read the pid a fixture wrote to `path`, waiting for it. |
| 681 | #[cfg(all(test, unix))] |
| 682 | pub(crate) fn read_pid_file(path: &std::path::Path, within: std::time::Duration) -> libc::pid_t { |
| 683 | let deadline = std::time::Instant::now() + within; |
| 684 | loop { |
| 685 | if let Some(pid) = parse_pid_file(path) { |
| 686 | return pid; |
| 687 | } |
| 688 | assert!( |
| 689 | std::time::Instant::now() < deadline, |
| 690 | "fixture never wrote {}", |
| 691 | path.display() |
| 692 | ); |
| 693 | std::thread::sleep(std::time::Duration::from_millis(20)); |
| 694 | } |
| 695 | } |
| 696 | |
| 697 | #[cfg(all(test, unix))] |
| 698 | fn parse_pid_file(path: &std::path::Path) -> Option<libc::pid_t> { |
| 699 | std::fs::read_to_string(path) |
| 700 | .ok() |
| 701 | .and_then(|text| text.trim().parse().ok()) |
| 702 | } |
| 703 | |
| 704 | /// Test support: drive `run` until the fixture has written its pid to `path`, |
| 705 | /// then drop `run`. No wall-clock guess about how long the fixture needs to |
| 706 | /// start; `run` finishing first is a fixture bug. |
| 707 | #[cfg(all(test, unix))] |
| 708 | pub(crate) async fn drop_once_pid_written<F: std::future::Future>( |
| 709 | run: F, |
| 710 | path: &std::path::Path, |
| 711 | ) -> libc::pid_t { |
| 712 | let written = async { |
| 713 | loop { |
| 714 | if let Some(pid) = parse_pid_file(path) { |
| 715 | return pid; |
| 716 | } |
| 717 | tokio::time::sleep(std::time::Duration::from_millis(20)).await; |
| 718 | } |
| 719 | }; |
| 720 | tokio::time::timeout(std::time::Duration::from_secs(60), async { |
| 721 | tokio::select! { |
| 722 | _ = run => panic!("fixture exited before writing {}", path.display()), |
| 723 | pid = written => pid, |
| 724 | } |
| 725 | }) |
| 726 | .await |
| 727 | .unwrap_or_else(|_| panic!("fixture never wrote {}", path.display())) |
| 728 | } |
| 729 | |
| 730 | #[cfg(all(test, unix))] |
| 731 | mod tests { |
| 732 | use super::*; |
| 733 | use std::time::Duration; |
| 734 | |
| 735 | #[tokio::test] |
| 736 | async fn failed_containment_attachment_is_reported_and_kills_the_child() { |
| 737 | let pid = std::cell::Cell::new(None); |
| 738 | let mut cmd = tokio::process::Command::new("/bin/sleep"); |
| 739 | cmd.arg("300"); |
| 740 | let result = tokio::time::timeout( |
| 741 | Duration::from_secs(1), |
| 742 | contained_run(&mut cmd, None, std::future::pending(), |child| { |
| 743 | pid.set(child.id()); |
| 744 | Err(std::io::Error::new( |
| 745 | std::io::ErrorKind::PermissionDenied, |
| 746 | "injected containment attachment failure", |
| 747 | )) |
| 748 | }), |
| 749 | ) |
| 750 | .await; |
| 751 | let pid = libc::pid_t::try_from(pid.get().expect("child spawned")).expect("pid"); |
| 752 | // Let Tokio reap while checking cleanup; the helper kills the fixture |
| 753 | // itself if the child ever outlives this failed run. |
| 754 | assert!( |
| 755 | tokio::task::spawn_blocking(move || wait_for_pid_exit(pid, Duration::from_secs(5))) |
| 756 | .await |
| 757 | .expect("cleanup probe"), |
| 758 | "the direct child outlived failed containment attachment" |
| 759 | ); |
| 760 | let Err(error) = result |
| 761 | .expect("attachment failure must return without running the command to completion") |
| 762 | else { |
| 763 | panic!("a containment failure must not be accepted as a successful run"); |
| 764 | }; |
| 765 | assert_eq!(error.kind(), std::io::ErrorKind::PermissionDenied); |
| 766 | assert_eq!(error.to_string(), "injected containment attachment failure"); |
| 767 | } |
| 768 | |
| 769 | #[tokio::test] |
| 770 | async fn dropping_contained_output_kills_the_whole_tree() { |
| 771 | let tmp = tempfile::tempdir().expect("tempdir"); |
| 772 | let pid_file = tmp.path().join("grandchild.pid"); |
| 773 | let mut cmd = tokio::process::Command::new("/bin/sh"); |
| 774 | cmd.arg("-c") |
| 775 | .arg("sleep 300 & echo $! > grandchild.pid; wait") |
| 776 | .current_dir(tmp.path()); |
| 777 | let grandchild = drop_once_pid_written(contained_output(&mut cmd), &pid_file).await; |
| 778 | assert!( |
| 779 | wait_for_pid_exit(grandchild, Duration::from_secs(5)), |
| 780 | "a process started by the dropped command is still running" |
| 781 | ); |
| 782 | } |
| 783 | |
| 784 | #[tokio::test] |
| 785 | async fn contained_output_captures_output_and_closes_stdin() { |
| 786 | let mut cmd = tokio::process::Command::new("/bin/sh"); |
| 787 | cmd.arg("-c").arg("cat; echo done"); |
| 788 | let output = tokio::time::timeout(Duration::from_secs(10), contained_output(&mut cmd)) |
| 789 | .await |
| 790 | .expect("stdin must be closed, not inherited") |
| 791 | .expect("run"); |
| 792 | assert!(output.status.success()); |
| 793 | assert_eq!(String::from_utf8_lossy(&output.stdout), "done\n"); |
| 794 | } |
| 795 | |
| 796 | /// The input reaches stdin, and stdin is closed after it. |
| 797 | #[tokio::test] |
| 798 | async fn contained_output_with_input_feeds_then_closes_stdin() { |
| 799 | let mut cmd = tokio::process::Command::new("/bin/sh"); |
| 800 | cmd.arg("-c").arg("cat; echo done"); |
| 801 | let output = tokio::time::timeout( |
| 802 | Duration::from_secs(10), |
| 803 | contained_output_with_input(&mut cmd, b"request\n".to_vec()), |
| 804 | ) |
| 805 | .await |
| 806 | .expect("stdin must be closed after the input") |
| 807 | .expect("run"); |
| 808 | assert!(output.status.success()); |
| 809 | assert_eq!(String::from_utf8_lossy(&output.stdout), "request\ndone\n"); |
| 810 | } |
| 811 | |
| 812 | /// A command that exits on its own is not reaped along with what it |
| 813 | /// deliberately left running, as with a bare `output()`. |
| 814 | #[tokio::test] |
| 815 | async fn contained_output_leaves_a_detached_daemon_of_a_clean_exit_running() { |
| 816 | let tmp = tempfile::tempdir().expect("tempdir"); |
| 817 | let mut cmd = tokio::process::Command::new("/bin/sh"); |
| 818 | cmd.arg("-c") |
| 819 | .arg("sleep 300 >/dev/null 2>&1 & echo $! > daemon.pid") |
| 820 | .current_dir(tmp.path()); |
| 821 | let output = contained_output(&mut cmd).await.expect("run"); |
| 822 | assert!(output.status.success()); |
| 823 | let daemon = read_pid_file(&tmp.path().join("daemon.pid"), Duration::from_secs(5)); |
| 824 | // `wait_for_pid_exit` kills the fixture when it times out. |
| 825 | assert!( |
| 826 | !wait_for_pid_exit(daemon, Duration::from_millis(500)), |
| 827 | "the daemon a clean command left behind was killed" |
| 828 | ); |
| 829 | } |
| 830 | |
| 831 | /// `stop` kills the tree but keeps what the command already wrote. |
| 832 | #[tokio::test] |
| 833 | async fn stopped_contained_output_keeps_partial_output_and_kills_the_tree() { |
| 834 | let tmp = tempfile::tempdir().expect("tempdir"); |
| 835 | let mut cmd = tokio::process::Command::new("/bin/sh"); |
| 836 | cmd.arg("-c") |
| 837 | .arg("echo before-stop; sleep 300 & echo $! > grandchild.pid; wait") |
| 838 | .current_dir(tmp.path()); |
| 839 | let pid_file = tmp.path().join("grandchild.pid"); |
| 840 | let stop = async { |
| 841 | while parse_pid_file(&pid_file).is_none() { |
| 842 | tokio::time::sleep(Duration::from_millis(20)).await; |
| 843 | } |
| 844 | }; |
| 845 | let run = tokio::time::timeout( |
| 846 | Duration::from_secs(60), |
| 847 | contained_output_until(&mut cmd, stop), |
| 848 | ) |
| 849 | .await |
| 850 | .expect("stop must end the run") |
| 851 | .expect("run"); |
| 852 | assert!(run.stopped); |
| 853 | assert_eq!(String::from_utf8_lossy(&run.output.stdout), "before-stop\n"); |
| 854 | let grandchild = read_pid_file(&pid_file, Duration::from_secs(5)); |
| 855 | assert!( |
| 856 | wait_for_pid_exit(grandchild, Duration::from_secs(5)), |
| 857 | "a process started by the stopped command is still running" |
| 858 | ); |
| 859 | } |
| 860 | |
| 861 | /// An in-flight run is listed for the signal-exit kill, and unlisted once |
| 862 | /// the future is gone. |
| 863 | #[tokio::test] |
| 864 | async fn in_flight_contained_run_is_listed_for_signal_exit() { |
| 865 | let tmp = tempfile::tempdir().expect("tempdir"); |
| 866 | let mut cmd = tokio::process::Command::new("/bin/sh"); |
| 867 | cmd.arg("-c") |
| 868 | .arg("echo $$ > leader.pid; sleep 300") |
| 869 | .current_dir(tmp.path()); |
| 870 | let pid_file = tmp.path().join("leader.pid"); |
| 871 | let mut run = Box::pin(contained_output(&mut cmd)); |
| 872 | let leader = drop_once_pid_written( |
| 873 | async { |
| 874 | run.as_mut().await.ok(); |
| 875 | }, |
| 876 | &pid_file, |
| 877 | ) |
| 878 | .await; |
| 879 | let leader_group = u32::try_from(leader).expect("pid"); |
| 880 | assert!(contained_exit_groups().contains(&leader_group)); |
| 881 | drop(run); |
| 882 | assert!(!contained_exit_groups().contains(&leader_group)); |
| 883 | // Let this current-thread runtime poll Tokio's orphan reaper while the |
| 884 | // cleanup probe waits; kill(pid, 0) still sees an unreaped macOS child. |
| 885 | assert!( |
| 886 | tokio::task::spawn_blocking(move || wait_for_pid_exit(leader, Duration::from_secs(5))) |
| 887 | .await |
| 888 | .expect("cleanup probe") |
| 889 | ); |
| 890 | } |
| 891 | } |
| 892 | |
| 893 | #[cfg(all(test, unix))] |
| 894 | mod admitted_execution_tests { |
| 895 | use super::*; |
| 896 | #[tokio::test] |
| 897 | async fn bounded_script_driver_refuses_overflow_without_deadlock() { |
| 898 | let mut command = tokio::process::Command::new("sh"); |
| 899 | command.args(["-c", "while :; do printf 1234567890; done"]); |
| 900 | let result = tokio::time::timeout( |
| 901 | std::time::Duration::from_secs(3), |
| 902 | contained_output_with_input_bounded( |
| 903 | &mut command, |
| 904 | Vec::new(), |
| 905 | 64, |
| 906 | 64, |
| 907 | std::future::pending(), |
| 908 | ), |
| 909 | ) |
| 910 | .await |
| 911 | .expect("overflow must not wait for execution deadline"); |
| 912 | assert!(result.is_err()); |
| 913 | } |
| 914 | #[tokio::test] |
| 915 | async fn bounded_script_driver_feeds_and_drains_concurrently() { |
| 916 | let mut command = tokio::process::Command::new("sh"); |
| 917 | command.args(["-c", "printf answer; cat >/dev/null"]); |
| 918 | let result = tokio::time::timeout( |
| 919 | std::time::Duration::from_secs(3), |
| 920 | contained_output_with_input_bounded( |
| 921 | &mut command, |
| 922 | vec![b'x'; 256 * 1024], |
| 923 | 64, |
| 924 | 64, |
| 925 | std::future::pending(), |
| 926 | ), |
| 927 | ) |
| 928 | .await |
| 929 | .unwrap() |
| 930 | .unwrap(); |
| 931 | assert!(!result.stopped); |
| 932 | assert!(result.output.status.success()); |
| 933 | assert_eq!(result.output.stdout, b"answer"); |
| 934 | } |
| 935 | } |
| 936 |