diff --git a/src/local/codex.rs b/src/local/codex.rs index ffd6995f..769c1cdd 100644 --- a/src/local/codex.rs +++ b/src/local/codex.rs @@ -15,7 +15,7 @@ use std::collections::HashMap; use std::process::Stdio; -use std::sync::atomic::{AtomicI64, Ordering}; +use std::sync::atomic::{AtomicBool, AtomicI64, Ordering}; use std::sync::Arc; use std::time::Duration; @@ -38,6 +38,8 @@ const HISTORY_READ_TIMEOUT: Duration = Duration::from_secs(2); const INTERRUPT_TIMEOUT: Duration = Duration::from_secs(5); /// A healthy app-server answers `initialize` immediately. const HANDSHAKE_TIMEOUT: Duration = Duration::from_secs(10); +/// `ensure` reaps under the host locks — a wedged child must not stall every bring-up. +const TERMINATION_TIMEOUT: Duration = Duration::from_secs(5); /// One inbound line, classified. JSON-RPC over one stream: a message with both /// `id` and `method` is a server→client *request* (approvals — must be @@ -170,6 +172,12 @@ pub enum TurnEvent { /// A live JSON-RPC connection to one session's `codex app-server` child. pub struct CodexClient { child: Mutex, + /// npm installs run a node wrapper whose native child survives a SIGKILL + /// to the wrapper alone, so teardown must signal the whole group. + #[cfg(unix)] + process_group_id: u32, + /// Once set, the group id may be reused — never signal it again. + terminated: AtomicBool, stdin: Mutex, next_id: AtomicI64, /// Our outstanding requests. Sync mutex: touched from the reader task and @@ -324,6 +332,30 @@ impl CodexClient { } } + /// Kill the child's whole process group and reap it; an orphaned native + /// app-server would keep the thread's writer lock. + async fn terminate(&self) { + // Locked first so a racing caller returns only after the reap. + let mut child = self.child.lock().await; + if self.terminated.swap(true, Ordering::AcqRel) { + return; + } + #[cfg(unix)] + if let Ok(process_group_id) = i32::try_from(self.process_group_id) { + // SAFETY: spawn assigns the child its own process group before this id is captured. + unsafe { + libc::kill(-process_group_id, libc::SIGKILL); + } + } + let _ = child.start_kill(); + if tokio::time::timeout(TERMINATION_TIMEOUT, child.wait()) + .await + .is_err() + { + eprintln!("orx up: timed out reaping codex app-server"); + } + } + async fn write_line(&self, msg: &Value) -> Result<()> { let mut stdin = self.stdin.lock().await; stdin @@ -548,7 +580,7 @@ async fn read_loop(client: Arc, stdout: tokio::process::ChildStdout // would eat a request timeout per turn until restart), then drop pending // senders so callers classify this as an ambiguous transport close rather // than a structured server rejection that would be safe to replay. - let _ = client.child.lock().await.kill().await; + client.terminate().await; client.pending.lock().unwrap().clear(); client.unanswered.lock().unwrap().clear(); if let Some(tx) = client.turn.lock().unwrap().as_ref() { @@ -595,13 +627,17 @@ async fn spawn_client( cmd.env("ORX_DATA_DIR", &dir); } // Own process group: a terminal SIGINT reaches orx up alone, which then - // tears the child down deliberately (kill_on_drop / shutdown()). + // tears the child's group down deliberately (terminate()). #[cfg(unix)] cmd.process_group(0); let mut child = cmd .spawn() .map_err(|e| anyhow!("Could not spawn {} app-server: {}", bin.display(), e))?; + #[cfg(unix)] + let process_group_id = child + .id() + .ok_or_else(|| anyhow!("codex app-server: spawned child has no process id"))?; let stdout = child .stdout .take() @@ -613,6 +649,9 @@ async fn spawn_client( let client = Arc::new(CodexClient { child: Mutex::new(child), + #[cfg(unix)] + process_group_id, + terminated: AtomicBool::new(false), stdin: Mutex::new(stdin), next_id: AtomicI64::new(1), pending: std::sync::Mutex::new(HashMap::new()), @@ -720,7 +759,7 @@ impl CodexHost { } } if let Some(stale) = guard.remove(session_id) { - let _ = stale.child.lock().await.kill().await; + stale.terminate().await; } } let host = self.clone(); @@ -738,17 +777,17 @@ impl CodexHost { { let existing = existing.clone(); drop(guard); - let _ = client.child.lock().await.kill().await; + client.terminate().await; return Ok(existing); } } if let Some(stale) = guard.remove(&session) { - let _ = stale.child.lock().await.kill().await; + stale.terminate().await; } guard.insert(session.clone(), client.clone()); } if let Err(e) = handshake(&client).await { - let _ = client.child.lock().await.kill().await; + client.terminate().await; let mut guard = host.inner.lock().await; if guard.get(&session).is_some_and(|c| Arc::ptr_eq(c, &client)) { guard.remove(&session); @@ -766,11 +805,13 @@ impl CodexHost { let mut guard = self.inner.lock().await; let client = guard.get(session_id)?; if matches!(client.child.lock().await.try_wait(), Ok(None)) { - Some(client.clone()) - } else { - guard.remove(session_id); - None + return Some(client.clone()); } + let dead = guard.remove(session_id)?; + drop(guard); + // The wrapper can die alone while its native child lives on. + dead.terminate().await; + None } /// Interrupt and harvest the in-flight turn while its child is reachable, @@ -796,7 +837,7 @@ impl CodexHost { } }; if let Some(retired) = retired { - let _ = retired.child.lock().await.kill().await; + retired.terminate().await; } items } @@ -804,14 +845,14 @@ impl CodexHost { /// Kill and reap one session's child (on session delete). pub async fn kill_session(&self, session_id: &str) { if let Some(client) = self.inner.lock().await.remove(session_id) { - let _ = client.child.lock().await.kill().await; + client.terminate().await; } } - /// Kill and reap every child (also happens via kill_on_drop on exit). + /// Kill and reap every child's process group. pub async fn shutdown(&self) { for (_, client) in self.inner.lock().await.drain() { - let _ = client.child.lock().await.kill().await; + client.terminate().await; } } } @@ -820,6 +861,58 @@ impl CodexHost { mod tests { use super::*; + #[cfg(unix)] + #[tokio::test] + async fn terminate_kills_the_native_child_behind_a_wrapper() { + // The npm `codex.js` shape: a wrapper whose child outlives its SIGKILL. + let mut command = Command::new("sh"); + command + .args(["-c", "sleep 30 & echo $!; wait"]) + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::null()) + .process_group(0); + let mut child = command.spawn().unwrap(); + let process_group_id = child.id().unwrap(); + let stdin = child.stdin.take().unwrap(); + let mut stdout = BufReader::new(child.stdout.take().unwrap()); + let mut native_pid = String::new(); + stdout.read_line(&mut native_pid).await.unwrap(); + let client = CodexClient { + child: Mutex::new(child), + process_group_id, + terminated: AtomicBool::new(false), + stdin: Mutex::new(stdin), + next_id: AtomicI64::new(1), + pending: std::sync::Mutex::new(HashMap::new()), + turn: std::sync::Mutex::new(None), + unanswered: std::sync::Mutex::new(HashMap::new()), + active_turn: std::sync::Mutex::new(None), + resumed_thread: std::sync::Mutex::new(None), + thread_model: std::sync::Mutex::new(None), + last_collab_mode: std::sync::Mutex::new(None), + native_store: NativeStore::Isolated, + }; + + client.terminate().await; + // A second call must not signal the (possibly reused) group again. + client.terminate().await; + assert!(client.terminated.load(Ordering::Acquire)); + + // EOF, not a pid probe: a killed child can linger as an unreaped zombie. + let mut rest = Vec::new(); + let eof = tokio::time::timeout( + Duration::from_secs(1), + tokio::io::AsyncReadExt::read_to_end(&mut stdout, &mut rest), + ) + .await; + assert!( + eof.is_ok(), + "native child {} survived terminate", + native_pid.trim() + ); + } + #[test] fn classify_discriminates_the_three_wire_shapes() { // Server→client request: id + method. Id kept as raw Value. (Fixture