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
8 changes: 8 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
82 changes: 65 additions & 17 deletions plugins/codex/scripts/lib/broker-lifecycle.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@ import {
import { withLock } from "./locking.mjs";
import {
getProcessIdentity,
isPidAlive,
isProcessRunning,
isProcessTreeRunning,
isValidPid,
Expand Down Expand Up @@ -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 = {}) {
Expand Down Expand Up @@ -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);
}

Expand Down
13 changes: 13 additions & 0 deletions plugins/codex/scripts/lib/lifecycle-limits.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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.
*
Expand Down
4 changes: 2 additions & 2 deletions plugins/codex/scripts/session-lifecycle-hook.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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;

Expand Down
Loading
Loading