返回 CodeWhale
process_tree.rs
根目录 / crates / tui / src / process_tree.rs
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
936 lines RUST