mirror of
https://github.com/openai/codex.git
synced 2026-08-23 13:09:46 +00:00
Preserve child waiters during process termination (#37498)
## Why Aborting the child waiter during termination can leave an exited PTY child unreaped and prevent the session from recording its exit status. ## What changed - Detach the child waiter instead of aborting it when terminating or dropping a `ProcessHandle`, while continuing to abort the I/O helper tasks. - Cover explicit termination and handle drops for pipe processes, plus queued PTY waiters and PTY process-group termination on Unix. GitOrigin-RevId: 313787f4edec80b30cf5a78f0cc6733f1520d9a4
This commit is contained in:
@@ -242,7 +242,7 @@ impl ProcessHandle {
|
||||
result
|
||||
}
|
||||
|
||||
/// Attempts to kill the child and abort helper tasks.
|
||||
/// Attempts to kill the child and abort I/O helper tasks.
|
||||
pub fn terminate(&self) {
|
||||
self.request_terminate();
|
||||
|
||||
@@ -264,7 +264,8 @@ impl ProcessHandle {
|
||||
if let Ok(mut h) = self.wait_handle.lock()
|
||||
&& let Some(handle) = h.take()
|
||||
{
|
||||
handle.abort();
|
||||
// Even a queued PTY waiter must run to reap the terminated child.
|
||||
drop(handle);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -827,6 +827,140 @@ async fn pipe_terminate_aborts_detached_readers() -> anyhow::Result<()> {
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn pipe_terminate_reaps_child() -> anyhow::Result<()> {
|
||||
let env_map: HashMap<String, String> = std::env::vars().collect();
|
||||
let command = if cfg!(windows) {
|
||||
"ping -n 60 127.0.0.1 > NUL"
|
||||
} else {
|
||||
"sleep 60"
|
||||
};
|
||||
let (program, args) = shell_command(command);
|
||||
let spawned = spawn_pipe_process(&program, &args, Path::new("."), &env_map, &None, &[]).await?;
|
||||
let (session, _output_rx, exit_rx) = combine_spawned_output(spawned);
|
||||
|
||||
session.terminate();
|
||||
|
||||
let exit_code = tokio::time::timeout(tokio::time::Duration::from_secs(5), exit_rx)
|
||||
.await
|
||||
.map_err(|_| anyhow::anyhow!("timed out waiting for terminated child to be reaped"))?
|
||||
.map_err(|_| anyhow::anyhow!("child waiter was aborted before reaping"))?;
|
||||
assert_eq!(session.exit_code(), Some(exit_code));
|
||||
assert!(session.has_exited());
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn pipe_drop_reaps_child() -> anyhow::Result<()> {
|
||||
let env_map: HashMap<String, String> = std::env::vars().collect();
|
||||
let command = if cfg!(windows) {
|
||||
"ping -n 60 127.0.0.1 > NUL"
|
||||
} else {
|
||||
"sleep 60"
|
||||
};
|
||||
let (program, args) = shell_command(command);
|
||||
let spawned = spawn_pipe_process(&program, &args, Path::new("."), &env_map, &None, &[]).await?;
|
||||
let (session, _output_rx, exit_rx) = combine_spawned_output(spawned);
|
||||
|
||||
drop(session);
|
||||
|
||||
tokio::time::timeout(tokio::time::Duration::from_secs(5), exit_rx)
|
||||
.await
|
||||
.map_err(|_| anyhow::anyhow!("timed out waiting for dropped child to be reaped"))?
|
||||
.map_err(|_| anyhow::anyhow!("child waiter was aborted before reaping"))?;
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[test]
|
||||
fn pty_terminate_reaps_child_when_waiter_is_queued() -> anyhow::Result<()> {
|
||||
let runtime = tokio::runtime::Builder::new_current_thread()
|
||||
.enable_all()
|
||||
.max_blocking_threads(1)
|
||||
.build()?;
|
||||
let (blocker_started_tx, blocker_started_rx) = std::sync::mpsc::channel();
|
||||
let (release_blocker_tx, release_blocker_rx) = std::sync::mpsc::channel();
|
||||
// Keep the PTY's blocking child waiter queued until after termination.
|
||||
runtime.spawn_blocking(move || {
|
||||
let _ = blocker_started_tx.send(());
|
||||
let _ = release_blocker_rx.recv_timeout(std::time::Duration::from_secs(5));
|
||||
});
|
||||
blocker_started_rx.recv_timeout(std::time::Duration::from_secs(2))?;
|
||||
|
||||
let pid_file = std::env::temp_dir().join(format!(
|
||||
"codex-pty-reap-{}-{}",
|
||||
std::process::id(),
|
||||
std::time::SystemTime::now()
|
||||
.duration_since(std::time::UNIX_EPOCH)?
|
||||
.as_nanos()
|
||||
));
|
||||
runtime.block_on(async {
|
||||
let mut env_map: HashMap<String, String> = std::env::vars().collect();
|
||||
env_map.insert(
|
||||
"CODEX_PTY_TEST_PID_FILE".to_string(),
|
||||
pid_file.display().to_string(),
|
||||
);
|
||||
let (program, args) =
|
||||
shell_command("printf '%s' \"$$\" > \"$CODEX_PTY_TEST_PID_FILE\"; sleep 60");
|
||||
let spawned = spawn_pty_process(
|
||||
&program,
|
||||
&args,
|
||||
Path::new("."),
|
||||
&env_map,
|
||||
&None,
|
||||
TerminalSize::default(),
|
||||
&[],
|
||||
)
|
||||
.await?;
|
||||
let (session, _output_rx, exit_rx) = combine_spawned_output(spawned);
|
||||
|
||||
let child_pid = tokio::time::timeout(std::time::Duration::from_secs(2), async {
|
||||
loop {
|
||||
if let Ok(pid) = std::fs::read_to_string(&pid_file)
|
||||
&& let Ok(pid) = pid.parse::<libc::pid_t>()
|
||||
{
|
||||
return pid;
|
||||
}
|
||||
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.map_err(|_| anyhow::anyhow!("timed out waiting for the PTY child PID"))?;
|
||||
std::fs::remove_file(&pid_file)?;
|
||||
|
||||
session.terminate();
|
||||
drop(session);
|
||||
release_blocker_tx.send(())?;
|
||||
let _ = tokio::time::timeout(std::time::Duration::from_secs(2), exit_rx)
|
||||
.await
|
||||
.map_err(|_| anyhow::anyhow!("timed out waiting for the PTY child waiter"))?;
|
||||
|
||||
// A returned PID proves an exited child was still an unreaped zombie.
|
||||
let wait_result = tokio::time::timeout(std::time::Duration::from_secs(2), async {
|
||||
loop {
|
||||
let result =
|
||||
unsafe { libc::waitpid(child_pid, std::ptr::null_mut(), libc::WNOHANG) };
|
||||
if result != 0 {
|
||||
return result;
|
||||
}
|
||||
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
|
||||
}
|
||||
})
|
||||
.await
|
||||
.map_err(|_| anyhow::anyhow!("timed out waiting for the PTY child to exit"))?;
|
||||
let wait_error = std::io::Error::last_os_error();
|
||||
assert_eq!(
|
||||
wait_result, -1,
|
||||
"PTY child {child_pid} remained an unreaped zombie"
|
||||
);
|
||||
assert_eq!(wait_error.raw_os_error(), Some(libc::ECHILD));
|
||||
|
||||
Ok(())
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
|
||||
async fn pty_terminate_kills_background_children_in_same_process_group() -> anyhow::Result<()> {
|
||||
@@ -844,7 +978,7 @@ async fn pty_terminate_kills_background_children_in_same_process_group() -> anyh
|
||||
&[],
|
||||
)
|
||||
.await?;
|
||||
let (session, mut output_rx, _exit_rx) = combine_spawned_output(spawned);
|
||||
let (session, mut output_rx, exit_rx) = combine_spawned_output(spawned);
|
||||
|
||||
let bg_pid = match wait_for_marker_pid(&mut output_rx, marker, /*timeout_ms*/ 2_000).await {
|
||||
Ok(pid) => pid,
|
||||
@@ -860,6 +994,12 @@ async fn pty_terminate_kills_background_children_in_same_process_group() -> anyh
|
||||
|
||||
session.terminate();
|
||||
|
||||
tokio::time::timeout(tokio::time::Duration::from_secs(3), exit_rx)
|
||||
.await
|
||||
.map_err(|_| anyhow::anyhow!("timed out waiting for terminated PTY child to be reaped"))?
|
||||
.map_err(|_| anyhow::anyhow!("PTY child waiter was aborted before reaping"))?;
|
||||
assert!(session.has_exited());
|
||||
|
||||
let exited = wait_for_process_exit(bg_pid, /*timeout_ms*/ 3_000).await?;
|
||||
if !exited {
|
||||
let _ = unsafe { libc::kill(bg_pid, libc::SIGKILL) };
|
||||
|
||||
Reference in New Issue
Block a user