mirror of
https://github.com/alphaXiv/OpenResearch.git
synced 2026-10-02 01:34:34 +08:00
OR-347: Kill the whole Codex app-server process group on teardown (#484)
* Kill the whole Codex app-server process group on teardown (OR-347) npm installs run `node codex.js app-server`, which spawns the native app-server in the same process group. orx only SIGKILLed the wrapper, so the native server survived holding the thread writer lock, and a replacement server's thread/resume failed with "already has an active writer". Teardown now signals the whole group, once, with a bounded reap. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Detect the killed child by stdout EOF in the teardown test A pid probe succeeds on an unreaped zombie, so the test could fail when the new parent reaps slowly; EOF proves every writer in the group exited. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
1483698805
commit
8f2fc4f514
+108
-15
@@ -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<Child>,
|
||||
/// 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<ChildStdin>,
|
||||
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<CodexClient>, 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
|
||||
|
||||
Reference in New Issue
Block a user