Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 45 additions & 2 deletions codex-rs/utils/pty/src/child.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
//! Local child ownership and platform selection, independent of process transports.
//!
//! Every child exposes Tokio stdio handles. Native children retain their PID until
//! Every child exposes Tokio stdio handles. Reap-only children retain ownership until
//! reaped, including when a wait is cancelled or the async runtime shuts down.

use std::io;
Expand All @@ -15,6 +15,10 @@ use tokio::process::ChildStdout;
#[path = "macos_child.rs"]
pub(super) mod macos;

#[cfg(unix)]
#[path = "child_reaper.rs"]
pub(super) mod reaper;

/// A local subprocess with owned stdio and cancellation-safe exit handling.
pub struct Child {
pub(super) inner: ChildKind,
Expand All @@ -25,6 +29,11 @@ pub struct Child {

pub(super) enum ChildKind {
Tokio(tokio::process::Child),
#[cfg(unix)]
TokioReapOnly {
child: Option<tokio::process::Child>,
reaper: std::sync::mpsc::Sender<reaper::ChildToReap>,
},
#[cfg(target_os = "macos")]
Native(macos::NativeChild),
}
Expand All @@ -33,6 +42,10 @@ impl Child {
pub fn id(&self) -> Option<u32> {
match &self.inner {
ChildKind::Tokio(child) => child.id(),
#[cfg(unix)]
ChildKind::TokioReapOnly { child, .. } => {
child.as_ref().and_then(tokio::process::Child::id)
}
#[cfg(target_os = "macos")]
ChildKind::Native(child) => child.id(),
}
Expand All @@ -43,12 +56,20 @@ impl Child {
self.stdin.take();
match &mut self.inner {
ChildKind::Tokio(child) => child.wait().await,
#[cfg(unix)]
ChildKind::TokioReapOnly { child, .. } => {
child
.as_mut()
.ok_or_else(|| io::Error::other("child already transferred to reaper"))?
.wait()
.await
}
#[cfg(target_os = "macos")]
ChildKind::Native(child) => child.wait().await,
}
}

/// Drain both output pipes while waiting, retaining kill-on-drop on cancellation.
/// Drain both output pipes while waiting, retaining the configured drop policy on cancellation.
pub async fn wait_with_output(mut self) -> io::Result<std::process::Output> {
let mut stdout = self.stdout.take();
let mut stderr = self.stderr.take();
Expand Down Expand Up @@ -85,6 +106,14 @@ impl Child {
self.stdin.take();
match &mut self.inner {
ChildKind::Tokio(child) => child.kill().await,
#[cfg(unix)]
ChildKind::TokioReapOnly { child, .. } => {
child
.as_mut()
.ok_or_else(|| io::Error::other("child already transferred to reaper"))?
.kill()
.await
}
#[cfg(target_os = "macos")]
ChildKind::Native(child) => child.kill().await,
}
Expand All @@ -94,3 +123,17 @@ impl Child {
#[cfg(all(test, unix))]
#[path = "child_tests.rs"]
mod tests;

impl Drop for ChildKind {
fn drop(&mut self) {
#[cfg(unix)]
if let Self::TokioReapOnly { child, reaper } = self
&& let Some(child) = child.take()
&& child.id().is_some()
{
// Transfer the handle, not only its PID: Tokio must not enqueue it
// into a stopped runtime's orphan queue before our worker reaps it.
let _ = reaper.send(reaper::ChildToReap::Tokio(child));
}
}
}
51 changes: 49 additions & 2 deletions codex-rs/utils/pty/src/child_command.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@
//!
//! The wrapped Tokio command is private: callers cannot install callbacks or
//! change settings that the native backend cannot inspect. Children receive only
//! explicitly supplied environment variables and kill-on-drop. Stdio, descriptor
//! explicitly supplied environment variables and default to kill-on-drop. Stdio, descriptor
//! inheritance, and compatibility fallbacks are configured independently.

use std::ffi::OsStr;
Expand Down Expand Up @@ -36,6 +36,19 @@ pub enum SpawnFallback {
ReturnError,
}

/// Cleanup performed when a child handle is dropped before its exit is collected.
#[derive(Clone, Copy)]
pub(crate) enum ChildDropPolicy {
/// Kill the direct child, then reap it. This is the default.
KillAndReap,
/// Reap the child when it exits, leaving termination to the caller.
#[allow(
dead_code,
reason = "Used by the pipe adapter in the next stacked change."
)]
ReapOnly,
}

/// An explicit child stdin, including a socket used for bidirectional fd transfer.
pub enum ChildStdin {
Piped,
Expand All @@ -50,6 +63,7 @@ pub struct Command {
pub(crate) descriptor_policy: DescriptorPolicy,
pub(crate) fallback: SpawnFallback,
pub(crate) stdin: ChildStdin,
pub(crate) drop_policy: ChildDropPolicy,
#[cfg(unix)]
pub(crate) arg0: Option<OsString>,
}
Expand All @@ -69,6 +83,7 @@ impl Command {
descriptor_policy: DescriptorPolicy::Inherit,
fallback: SpawnFallback::Compatible,
stdin: ChildStdin::Piped,
drop_policy: ChildDropPolicy::KillAndReap,
#[cfg(unix)]
arg0: None,
}
Expand Down Expand Up @@ -108,6 +123,20 @@ impl Command {
self
}

/// Choose who owns termination when the child handle is dropped.
#[allow(
dead_code,
reason = "Used by the pipe adapter in the next stacked change."
)]
pub(crate) fn drop_policy(&mut self, policy: ChildDropPolicy) -> &mut Self {
self.drop_policy = policy;
self.inner.kill_on_drop(match policy {
ChildDropPolicy::KillAndReap => true,
ChildDropPolicy::ReapOnly => false,
});
self
}

pub fn stdin(&mut self, stdin: ChildStdin) -> &mut Self {
self.stdin = stdin;
self
Expand Down Expand Up @@ -170,12 +199,30 @@ impl Command {
});
}
}
#[cfg(unix)]
let reaper = match self.drop_policy {
ChildDropPolicy::KillAndReap => None,
ChildDropPolicy::ReapOnly => Some(crate::child::reaper::sender()?),
};
let mut child = self.inner.spawn()?;
Ok(Child {
stdin: child.stdin.take(),
stdout: child.stdout.take(),
stderr: child.stderr.take(),
inner: ChildKind::Tokio(child),
inner: {
#[cfg(unix)]
{
match reaper {
Some(reaper) => ChildKind::TokioReapOnly {
child: Some(child),
reaper,
},
None => ChildKind::Tokio(child),
}
}
#[cfg(not(unix))]
ChildKind::Tokio(child)
},
})
}
}
Expand Down
68 changes: 68 additions & 0 deletions codex-rs/utils/pty/src/child_reaper.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
//! Runtime-independent cleanup for children whose termination is externally owned.
//!
//! Initialize the shared worker before launch so dropping a live child never
//! creates a thread or waits for its exit. Poll only transferred PIDs, retaining
//! ownership until reaped, so one live child cannot delay cleanup of another.

use std::io;
#[cfg(target_os = "macos")]
use std::ptr;
use std::sync::Mutex;
use std::sync::PoisonError;
use std::sync::mpsc;
use std::time::Duration;

/// Exclusive ownership transferred by a child handle after its caller drops it.
pub(crate) enum ChildToReap {
#[cfg(target_os = "macos")]
Native(libc::pid_t),
Tokio(tokio::process::Child),
}

static REAPER: Mutex<Option<mpsc::Sender<ChildToReap>>> = Mutex::new(None);
const REAP_BATCH_SIZE: usize = 64;

/// Reserve the shared reaper, returning thread creation failures before launch.
pub(crate) fn sender() -> io::Result<mpsc::Sender<ChildToReap>> {
let mut reaper = REAPER.lock().unwrap_or_else(PoisonError::into_inner);
if let Some(sender) = &*reaper {
return Ok(sender.clone());
}
let (sender, receiver) = mpsc::channel();
std::thread::Builder::new()
.name("codex-child-reaper".into())
.spawn(move || {
let mut pending = Vec::new();
loop {
if pending.is_empty() {
let Ok(pid) = receiver.recv() else { return };
pending.push(pid);
} else {
match receiver.recv_timeout(Duration::from_millis(10)) {
Ok(pid) => pending.push(pid),
Err(mpsc::RecvTimeoutError::Timeout) => {}
Err(mpsc::RecvTimeoutError::Disconnected) => return,
}
}
// Amortize scans across bursts without starving reaping when drops continue.
pending.extend(receiver.try_iter().take(REAP_BATCH_SIZE - 1));
pending.retain_mut(|child| match child {
#[cfg(target_os = "macos")]
ChildToReap::Native(pid) => {
// SAFETY: Drop transferred exclusive ownership of this PID.
let result = unsafe { libc::waitpid(*pid, ptr::null_mut(), libc::WNOHANG) };
result == 0
|| (result == -1
&& io::Error::last_os_error().kind() == io::ErrorKind::Interrupted)
}
ChildToReap::Tokio(child) => match child.try_wait() {
Ok(None) => true,
Ok(Some(_)) => false,
Err(error) => error.kind() == io::ErrorKind::Interrupted,
},
});
}
})?;
*reaper = Some(sender.clone());
Ok(sender)
}
97 changes: 96 additions & 1 deletion codex-rs/utils/pty/src/child_tests.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
//! Regression coverage for output-pipe lifetimes while waiting for a child.
//! Regression coverage for output-pipe lifetimes and runtime-independent reaping.

use std::future::Future;
use std::future::poll_fn;
Expand Down Expand Up @@ -66,3 +66,98 @@ async fn wait_with_output_keeps_eof_pipes_open_until_exit() -> anyhow::Result<()
);
Ok(())
}

#[test]
fn non_killing_drop_reaps_after_runtime_shutdown() -> anyhow::Result<()> {
use crate::child_command::ChildDropPolicy;
use std::time::Duration;
if std::env::var_os("CODEX_TEST_ISOLATED_REAPER").is_none() {
// Other Tokio tests must not drain this process's global orphan queue.
let output = std::process::Command::new(std::env::current_exe()?)
.args([
"--exact",
"child::tests::non_killing_drop_reaps_after_runtime_shutdown",
"--nocapture",
])
.env("CODEX_TEST_ISOLATED_REAPER", "1")
.output()?;
assert!(
output.status.success(),
"isolated reaper test failed:\n{}\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr),
);
return Ok(());
}
for program in ["cat", "/bin/cat"] {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
let mut child = runtime.block_on(async {
let mut command = crate::Command::new(program);
command.drop_policy(ChildDropPolicy::ReapOnly);
command.spawn()
})?;
let stdin = child.stdin.take();
let pid = child.id().expect("live PID");
let next = runtime.block_on(async {
let mut command = crate::Command::new(program);
command.drop_policy(ChildDropPolicy::ReapOnly);
command.spawn()
})?;
let next_pid = next.id().expect("live PID");
drop(runtime);
drop(child);
drop(next);
// The shared reaper must collect the second child while the first remains alive.
wait_until_reaped(next_pid)?;
// Keep stdin open: cat must remain alive even after its owner and runtime go away.
std::thread::sleep(Duration::from_millis(50));
let mut info = std::mem::MaybeUninit::<libc::siginfo_t>::zeroed();
// SAFETY: waitid only observes our child without reaping it, with writable storage.
assert_eq!(
unsafe {
libc::waitid(
libc::P_PID,
pid,
info.as_mut_ptr(),
libc::WEXITED | libc::WNOHANG | libc::WNOWAIT,
)
},
0
);
// SAFETY: waitid succeeded and initialized the zeroed signal information.
assert_eq!(unsafe { info.assume_init().si_pid() }, 0);
drop(stdin);
wait_until_reaped(pid)?;
}
Ok(())
}

fn wait_until_reaped(pid: u32) -> anyhow::Result<()> {
use std::io;
use std::time::Duration;
let mut info = std::mem::MaybeUninit::<libc::siginfo_t>::zeroed();
let deadline = std::time::Instant::now() + Duration::from_secs(5);
loop {
// SAFETY: WNOWAIT prevents the test from stealing the reaper's child.
if unsafe {
libc::waitid(
libc::P_PID,
pid,
info.as_mut_ptr(),
libc::WEXITED | libc::WNOHANG | libc::WNOWAIT,
)
} == -1
{
assert_eq!(
io::Error::last_os_error().raw_os_error(),
Some(libc::ECHILD)
);
break;
}
anyhow::ensure!(std::time::Instant::now() < deadline, "child was not reaped");
std::thread::sleep(Duration::from_millis(10));
}
Ok(())
}
Loading
Loading