diff --git a/README.md b/README.md index e0ea536dc..ff5e22c7f 100644 --- a/README.md +++ b/README.md @@ -467,6 +467,14 @@ Beyond the imports, this fork carries fixes for defects the imports themselves s every MCP server under it — running. `taskkill /T /F` walks the tree there instead - a state write that Windows briefly refuses — a scanner or indexer holding the file open, which surfaces as `EPERM`/`EBUSY` on the replacing rename — is retried instead of losing the record +- a live broker that is *refusing* connections — a full accept backlog, or a Windows named pipe + with no free instance — is given a longer window to start accepting again and reused if it does, + instead of being replaced by a second broker with its own app-server and MCP servers. A broker + that is simply gone still costs nothing: only the probe outcomes that mean "listening, but not + right now" buy that wait ([#768](https://github.com/openai/codex-plugin-cc/pull/768) raises the + duplicate-broker problem upstream; its own two fixes — never tearing down a live broker, and + serializing the check-then-create window — were already here, and its blanket 3s probe is not + taken, since it would charge every gone broker for the rare busy one) - the app-server typecheck (`npm run build`) passes Each of those came out of an adversarial review of the merges, re-run after every round of fixes; diff --git a/plugins/codex/scripts/lib/broker-lifecycle.mjs b/plugins/codex/scripts/lib/broker-lifecycle.mjs index 3d7912afb..89c8fa419 100644 --- a/plugins/codex/scripts/lib/broker-lifecycle.mjs +++ b/plugins/codex/scripts/lib/broker-lifecycle.mjs @@ -18,7 +18,6 @@ import { import { withLock } from "./locking.mjs"; import { getProcessIdentity, - isPidAlive, isProcessRunning, isProcessTreeRunning, isValidPid, @@ -75,19 +74,32 @@ function probeEndpoint(endpoint, timeoutMs) { }); } -async function waitForBrokerEndpoint(endpoint, timeoutMs = 2000) { +/** + * Retry probeEndpoint() until the deadline, reporting how the last attempt ended. + * + * The outcome matters to the caller that has to decide what a silent endpoint means: "nothing is + * there" (ECONNREFUSED, ENOENT) and "something is there but would not take this connection right + * now" (a timeout, EAGAIN from a full accept backlog, EBUSY from a Windows named pipe with no free + * instance) are the same boolean and opposite situations. + */ +async function probeBrokerEndpoint(endpoint, timeoutMs) { const deadline = Date.now() + timeoutMs; + let outcome = "timeout"; while (Date.now() < deadline) { - const probe = await probeEndpoint(endpoint, Math.min(150, deadline - Date.now())); - if (probe === "connect") { - return true; + outcome = await probeEndpoint(endpoint, Math.min(150, deadline - Date.now())); + if (outcome === "connect") { + return { ready: true, outcome }; } - if (probe === "invalid") { - return false; + if (outcome === "invalid") { + return { ready: false, outcome }; } await new Promise((resolve) => setTimeout(resolve, 50)); } - return false; + return { ready: false, outcome }; +} + +async function waitForBrokerEndpoint(endpoint, timeoutMs = 2000) { + return (await probeBrokerEndpoint(endpoint, timeoutMs)).ready; } export async function sendBrokerShutdown(endpoint, options = {}) { @@ -654,25 +666,61 @@ function withBrokerLock(cwd, options, action) { ); } +// How long a broker that is listening but refusing connections gets to start accepting again. +// +// Only the outcomes in BUSY_ENDPOINT_OUTCOMES buy this wait. A connect() to a Unix socket is +// completed by the kernel as soon as the peer is listening — the server does not have to call +// accept() — so a broker that is merely mid-turn answers the fast probe in microseconds and never +// reaches here. What does reach here is an endpoint that is present and refusing: a full accept +// backlog (EAGAIN), or a Windows named pipe with no free instance (EBUSY, ERROR_PIPE_BUSY), where +// spawning a replacement would start a second app-server, and every MCP server under it, for a +// broker that is about to be available again. +const BUSY_BROKER_PROBE_TIMEOUT_MS = 1500; + +// "Listening, but not taking this connection right now." Everything else a probe can report — +// ECONNREFUSED, ENOENT, an unparsable endpoint — means nothing is there to wait for, and waiting +// on it would make the common case (a broker that is really gone) pay for the rare one. +const BUSY_ENDPOINT_OUTCOMES = new Set(["timeout", "EAGAIN", "EBUSY"]); + async function ensureBrokerSessionLocked(cwd, options = {}) { const shutdownOptions = { ...options, killProcess: options.killProcess ?? terminateProcessTree }; + const probe = options.probeBrokerEndpoint ?? probeBrokerEndpoint; const existing = loadBrokerSession(cwd); - if (existing?.endpoint && (await waitForBrokerEndpoint(existing.endpoint, 150))) { + const fastProbe = existing?.endpoint + ? await probe(existing.endpoint, 150) + : { ready: false, outcome: "invalid" }; + if (fastProbe.ready) { + return existing; + } + + // Pinned to the identity recorded at spawn time, like the shutdown path: a reused pid must read + // as gone rather than keep an abandoned record alive. + const existingIsAlive = existing + ? isProcessRunning(existing.pid, { identity: existing.processIdentity ?? undefined }) + : false; + + // Refusing, not gone: give it the longer window. Answering there means it is the workspace's + // broker and back in service, so it is reused as it stands. + if ( + existingIsAlive && + BUSY_ENDPOINT_OUTCOMES.has(fastProbe.outcome) && + (await probe(existing.endpoint, options.busyProbeTimeoutMs ?? BUSY_BROKER_PROBE_TIMEOUT_MS)).ready + ) { return existing; } - // Only reclaim a broker we can prove is gone. The probe above waits 150ms, which a live but busy - // broker can miss, and forcing a shutdown here on that alone would kill (or orphan) a broker that - // another, unrelated session is still using — the ownership-verified teardown below can force a - // kill once `options.killProcess` is set, so this gate keeps that path from firing on a broker - // that only failed to answer within 150ms. + // Only reclaim a broker we can prove is gone. The probes above prove nothing about the process — + // a wedged broker stays silent through both — and forcing a shutdown on that alone would kill (or + // orphan) a broker that another, unrelated session is still using; the ownership-verified teardown + // below can force a kill once `options.killProcess` is set, so this gate keeps that path from + // firing on a broker that merely failed to answer. // - // A live one is left exactly as it is. Once the replacement below takes over it has no clients, - // so its own idle shutdown reclaims both the process and its files. - if (existing && !isPidAlive(existing.pid)) { + // A live one that stayed silent is left exactly as it is, and replaced. Once the replacement below + // takes over it has no clients, so its own idle shutdown reclaims both the process and its files. + if (existing && !existingIsAlive) { await shutdownBrokerSessionLocked(cwd, shutdownOptions); } diff --git a/plugins/codex/scripts/lib/lifecycle-limits.mjs b/plugins/codex/scripts/lib/lifecycle-limits.mjs index 65fd9cd8e..f28ac42d5 100644 --- a/plugins/codex/scripts/lib/lifecycle-limits.mjs +++ b/plugins/codex/scripts/lib/lifecycle-limits.mjs @@ -3,6 +3,7 @@ import process from "node:process"; const DEFAULT_BROKER_IDLE_SHUTDOWN_MS = 10 * 60 * 1000; const DEFAULT_BROKER_STARTUP_TIMEOUT_MS = 5 * 60 * 1000; const DEFAULT_WORKER_TTL_MS = 24 * 60 * 60 * 1000; +const DEFAULT_TURN_INTERRUPT_BUDGET_MS = 2200; /** `setTimeout` truncates anything larger to a 32-bit int, firing almost immediately instead. */ const MAX_TIMEOUT_MS = 2 ** 31 - 1; @@ -67,6 +68,18 @@ export function brokerStartupTimeoutMs(env = process.env) { return readDurationMs(env.CODEX_BROKER_STARTUP_TIMEOUT_MS, DEFAULT_BROKER_STARTUP_TIMEOUT_MS); } +/** + * How long a SessionEnd hook run may spend interrupting turns, in total. + * + * The reaper's dead-turn interrupts and the session's own share this one budget so they cannot + * stack past the hook timeout, and each job gets a slice of what is left. The default suits a + * handful of jobs on a responsive machine; a workspace with many active jobs, or a loaded machine + * where each attempt takes longer, may need more. `0` disables turn interrupts entirely. + */ +export function turnInterruptBudgetMs(env = process.env) { + return readDurationMs(env.CODEX_TURN_INTERRUPT_BUDGET_MS, DEFAULT_TURN_INTERRUPT_BUDGET_MS); +} + /** * Ceiling on a detached background worker's wall-clock lifetime. * diff --git a/plugins/codex/scripts/session-lifecycle-hook.mjs b/plugins/codex/scripts/session-lifecycle-hook.mjs index 28e0b16d6..c491494d4 100644 --- a/plugins/codex/scripts/session-lifecycle-hook.mjs +++ b/plugins/codex/scripts/session-lifecycle-hook.mjs @@ -5,7 +5,7 @@ import process from "node:process"; import { forceKillProcessTree, isPidAlive, terminateProcessTree } from "./lib/process.mjs"; import { reconcileJobLiveness } from "./lib/job-control.mjs"; -import { brokerIdleShutdownMs } from "./lib/lifecycle-limits.mjs"; +import { brokerIdleShutdownMs, turnInterruptBudgetMs } from "./lib/lifecycle-limits.mjs"; import { BROKER_ENDPOINT_ENV } from "./lib/app-server.mjs"; import { LOG_FILE_ENV, @@ -28,7 +28,7 @@ const PLUGIN_DATA_ENV = "CLAUDE_PLUGIN_DATA"; // (identity waits are capped per job and reserve room for the interrupt // RPCs themselves), then a short exit grace for killed workers, then the // shutdown exchange with its one retry. -const TURN_INTERRUPT_BUDGET_MS = 2200; +const TURN_INTERRUPT_BUDGET_MS = turnInterruptBudgetMs(); const TURN_IDENTITY_WAIT_MS = 500; const TURN_INTERRUPT_RESERVE_MS = 1000; diff --git a/tests/broker-lifecycle.test.mjs b/tests/broker-lifecycle.test.mjs index 18b95c4ec..adb751977 100644 --- a/tests/broker-lifecycle.test.mjs +++ b/tests/broker-lifecycle.test.mjs @@ -214,7 +214,9 @@ test("shutdown request always uses a finite deadline", { skip: process.platform const sessionDir = makeTempDir("cxc-unresponsive-"); const socketPath = path.join(sessionDir, "broker.sock"); const sockets = new Set(); + let accepted = false; const server = net.createServer((socket) => { + accepted = true; sockets.add(socket); socket.on("close", () => sockets.delete(socket)); socket.on("data", () => { @@ -229,24 +231,37 @@ test("shutdown request always uses a finite deadline", { skip: process.platform await new Promise((resolve) => server.close(resolve)); }); - for (const timeoutMs of [0, 40]) { - const startedAt = Date.now(); - const response = await sendBrokerShutdown(`unix:${socketPath}`, { - instanceToken: "instance-token-1234567890", - timeoutMs - }); - // A broker that accepts the connection and then says nothing is ambiguous: - // it may be alive and busy with its refusal lost on the wire. The outcome - // has to report that — not delivered, not refused, and deliberately not - // unreachable — so teardown never reaps a broker that may still be serving - // someone. - assert.equal(response.delivered, false); - assert.equal(response.refused, false); - assert.equal(response.unreachable, false); - assert.equal(response.result, null); - assert.equal(response.error, null); - assert.ok(Date.now() - startedAt < 500, "shutdown request exceeded its deadline"); - } + // `timeoutMs: 0` is the case that has to settle at all: socket.setTimeout(0) disables the timer + // outright, so without the clamp inside sendBrokerShutdown() this call would wait on a silent + // peer forever. Racing it against a watchdog states exactly that, and fails in seconds instead + // of hanging the suite. The watchdog is deliberately far longer than the work: what is under + // test is finite versus infinite, not fast versus slow. + let watchdog; + const settled = await Promise.race([ + sendBrokerShutdown(`unix:${socketPath}`, { instanceToken: "instance-token-1234567890", timeoutMs: 0 }), + new Promise((resolve) => { + watchdog = setTimeout(() => resolve("never settled"), 30000); + }) + ]); + clearTimeout(watchdog); + assert.notEqual(settled, "never settled", "a zero timeout must not become no timeout"); + + // A broker that accepts the connection and then says nothing is ambiguous: it may be alive and + // busy with its refusal lost on the wire. The outcome has to report that — not delivered, not + // refused, and deliberately not unreachable — so teardown never reaps a broker that may still be + // serving someone. The budget here is generous on purpose: with a 1ms deadline the connect + // itself loses the race on a loaded machine, and `unreachable` then reports the truth about a + // scenario the test never managed to set up. + const response = await sendBrokerShutdown(`unix:${socketPath}`, { + instanceToken: "instance-token-1234567890", + timeoutMs: 1000 + }); + assert.equal(accepted, true, "the server never accepted the connection this case is about"); + assert.equal(response.delivered, false); + assert.equal(response.refused, false); + assert.equal(response.unreachable, false); + assert.equal(response.result, null); + assert.equal(response.error, null); assert.equal(fs.existsSync(socketPath), true); }); @@ -1297,3 +1312,162 @@ test("clearBrokerSession deletes the same record loadBrokerSession() returned, n assert.equal(loadBrokerSession(workspace), null); }); }); + +// A broker script that records the fact it was launched. Spawning it at all means the +// workspace's existing broker was written off and replaced. It is written outside any broker +// session dir so production teardown can still rmdir that dir. +function writeSpawnMarkerBroker() { + // Not a cxc-* name: that prefix is what resolveOwnedSessionDir() reads as a broker session dir. + const dir = makeTempDir("marker-broker-"); + const scriptPath = path.join(dir, "marker-broker.mjs"); + const markerPath = path.join(dir, "spawned.marker"); + fs.writeFileSync( + scriptPath, + [ + 'import fs from "node:fs";', + `fs.appendFileSync(${JSON.stringify(markerPath)}, "spawned\\n");`, + "setTimeout(() => {}, 60000);" + ].join("\n") + ); + return { scriptPath, spawned: () => fs.existsSync(markerPath) }; +} + +async function waitUntil(predicate, timeoutMs = 10000) { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (predicate()) { + return true; + } + await new Promise((resolve) => setTimeout(resolve, 25)); + } + return predicate(); +} + +// The probe outcome decides everything here, and the outcomes worth deciding on cannot be produced +// on demand from a test: a Unix connect() succeeds the moment the peer listens, so "listening but +// refusing" needs a saturated accept backlog, and EBUSY needs a Windows named pipe. Injecting the +// probe keeps these tests about the decision, which is what changed. +function scriptedProbe(outcomes) { + const calls = []; + const probe = async (endpoint, timeoutMs) => { + calls.push({ endpoint, timeoutMs }); + const next = outcomes[calls.length - 1] ?? { ready: false, outcome: "ECONNREFUSED" }; + return next; + }; + return { probe, calls }; +} + +function recordedSession(sessionDir, { pid, endpoint }) { + return { + endpoint: endpoint ?? `unix:${path.join(sessionDir, "broker.sock")}`, + pid, + pidFile: null, + logFile: null, + sessionDir, + instanceToken: "probe-outcome-token" + }; +} + +test("a live broker that is refusing connections is given the longer window and reused", async () => { + const workspace = makeTempDir(); + const sessionDir = makeTempDir("cxc-refusing-"); + const session = recordedSession(sessionDir, { pid: process.pid }); + saveBrokerSession(workspace, session); + + // EAGAIN is what a full accept backlog reports: the broker is there and will take the next + // connection, so replacing it would start a second app-server for nothing. + const { probe, calls } = scriptedProbe([ + { ready: false, outcome: "EAGAIN" }, + { ready: true, outcome: "connect" } + ]); + const marker = writeSpawnMarkerBroker(); + const result = await ensureBrokerSession(workspace, { + scriptPath: marker.scriptPath, + probeBrokerEndpoint: probe + }); + + assert.equal(result?.endpoint, session.endpoint, "the refusing broker must be handed back"); + assert.equal(marker.spawned(), false, "it must not be replaced by a second broker"); + assert.deepEqual(loadBrokerSession(workspace), session, "its persisted session must be untouched"); + assert.equal(calls.length, 2, "the longer window must be entered exactly once"); + assert.deepEqual( + calls.map((call) => call.timeoutMs), + [150, 1500], + "the fast probe's budget, then the shipped default for the longer one" + ); +}); + +test("a broker that is simply gone never pays for the longer window", async () => { + const workspace = makeTempDir(); + const sessionDir = makeTempDir("cxc-gone-"); + // Alive pid, nothing listening: the broker's own graceful shutdown closes its listener seconds + // before the process exits, and a record can outlive its socket. Waiting on that would charge + // the common case for the rare one. + const session = recordedSession(sessionDir, { pid: process.pid }); + saveBrokerSession(workspace, session); + + const { probe, calls } = scriptedProbe([{ ready: false, outcome: "ECONNREFUSED" }]); + const marker = writeSpawnMarkerBroker(); + await ensureBrokerSession(workspace, { + scriptPath: marker.scriptPath, + probeBrokerEndpoint: probe, + timeoutMs: 300 + }).catch(() => {}); + + assert.equal(calls.length, 1, "ECONNREFUSED must not buy a second probe"); + assert.equal(await waitUntil(() => marker.spawned()), true, "the gone broker must be replaced"); +}); + +test("a broker that keeps refusing through the longer window is replaced", async () => { + const workspace = makeTempDir(); + const sessionDir = makeTempDir("cxc-wedged-"); + const session = recordedSession(sessionDir, { pid: process.pid }); + saveBrokerSession(workspace, session); + + // Wedged rather than momentarily busy. Handing this one back would block every later command + // behind it, so the replacement must still run once the longer window has had its say. + const { probe, calls } = scriptedProbe([ + { ready: false, outcome: "timeout" }, + { ready: false, outcome: "timeout" } + ]); + const marker = writeSpawnMarkerBroker(); + await ensureBrokerSession(workspace, { + scriptPath: marker.scriptPath, + probeBrokerEndpoint: probe, + busyProbeTimeoutMs: 300, + timeoutMs: 300 + }).catch(() => {}); + + assert.equal(calls.length, 2, "the longer window must have been tried"); + assert.equal(await waitUntil(() => marker.spawned()), true, "a broker that stays silent must be replaced"); +}); + +test("a dead broker's record is reclaimed rather than waited on", { skip: process.platform === "win32" }, async () => { + const workspace = makeTempDir(); + const sessionDir = makeTempDir("cxc-dead-"); + const socketPath = path.join(sessionDir, "broker.sock"); + + // A real dead broker: the socket file it left behind is still on disk, its pid is not. + const holder = await spawnSocketHolder(socketPath); + const deadPid = holder.pid; + holder.kill("SIGKILL"); + await new Promise((resolve) => holder.once("exit", resolve)); + + const session = recordedSession(sessionDir, { pid: deadPid }); + saveBrokerSession(workspace, session); + + // A timeout is the one outcome that would buy the longer window from a live broker, so this is + // where the liveness gate has to do the work: the pid is gone, and waiting on a dead broker's + // endpoint delays the replacement for nothing. + const { probe, calls } = scriptedProbe([{ ready: false, outcome: "timeout" }]); + const marker = writeSpawnMarkerBroker(); + await ensureBrokerSession(workspace, { + scriptPath: marker.scriptPath, + probeBrokerEndpoint: probe, + timeoutMs: 300, + killProcess: () => {} + }).catch(() => {}); + + assert.equal(calls.length, 1, "a dead broker must not buy a second probe"); + assert.equal(await waitUntil(() => marker.spawned()), true, "it must be replaced"); +}); diff --git a/tests/lifecycle-limits.test.mjs b/tests/lifecycle-limits.test.mjs index 81daea9a7..de4f7c3e9 100644 --- a/tests/lifecycle-limits.test.mjs +++ b/tests/lifecycle-limits.test.mjs @@ -6,6 +6,7 @@ import { brokerIdleShutdownMs, brokerStartupTimeoutMs, disarmTimeout, + turnInterruptBudgetMs, workerTtlMs } from "../plugins/codex/scripts/lib/lifecycle-limits.mjs"; import { isPidAlive } from "../plugins/codex/scripts/lib/process.mjs"; @@ -53,6 +54,22 @@ test("worker ttl falls back to the default on unusable input", () => { } }); +test("turn interrupt budget defaults to 2200ms", () => { + assert.equal(turnInterruptBudgetMs({}), 2200); + assert.equal(turnInterruptBudgetMs({ CODEX_TURN_INTERRUPT_BUDGET_MS: "" }), 2200); +}); + +test("turn interrupt budget honours an override and can be disabled", () => { + assert.equal(turnInterruptBudgetMs({ CODEX_TURN_INTERRUPT_BUDGET_MS: "8000" }), 8000); + assert.equal(turnInterruptBudgetMs({ CODEX_TURN_INTERRUPT_BUDGET_MS: "0" }), 0); +}); + +test("turn interrupt budget falls back to the default on unusable input", () => { + for (const raw of ["soon", "-5", "NaN", " "]) { + assert.equal(turnInterruptBudgetMs({ CODEX_TURN_INTERRUPT_BUDGET_MS: raw }), 2200); + } +}); + test("a disabled limit never fires", async () => { // setTimeout(fn, 0) means "next tick", but 0 is our documented way to disable a limit. Passing // it straight through would turn "no ceiling" into "terminate immediately". diff --git a/tests/runtime.test.mjs b/tests/runtime.test.mjs index 8c9dff126..1c0e34276 100644 --- a/tests/runtime.test.mjs +++ b/tests/runtime.test.mjs @@ -3991,7 +3991,12 @@ test("a hung turn interrupt cannot starve later jobs' interrupt attempts", async cwd: repo, env: { ...process.env, - CODEX_COMPANION_APP_SERVER_ENDPOINT: `unix:${socketPath}` + CODEX_COMPANION_APP_SERVER_ENDPOINT: `unix:${socketPath}`, + // What is under test is that the budget is *shared*, not how large it is: the hung + // interrupt consumes its whole slice, so with the 2200ms default a loaded machine can + // spend the remainder before the second job is ever attempted, and the test reports a + // starvation that did not happen. The slicing is identical at any budget. + CODEX_TURN_INTERRUPT_BUDGET_MS: "8000" }, input: JSON.stringify({ hook_event_name: "SessionEnd", @@ -4139,7 +4144,12 @@ test("a hung dead-worker reap interrupt cannot starve the session's own interrup cwd: repo, env: { ...process.env, - CODEX_COMPANION_APP_SERVER_ENDPOINT: `unix:${socketPath}` + CODEX_COMPANION_APP_SERVER_ENDPOINT: `unix:${socketPath}`, + // What is under test is that the budget is *shared*, not how large it is: the hung + // interrupt consumes its whole slice, so with the 2200ms default a loaded machine can + // spend the remainder before the second job is ever attempted, and the test reports a + // starvation that did not happen. The slicing is identical at any budget. + CODEX_TURN_INTERRUPT_BUDGET_MS: "8000" }, input: JSON.stringify({ hook_event_name: "SessionEnd",