Skip to content
Open
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
110 changes: 92 additions & 18 deletions plugins/codex/scripts/app-server-broker.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@ import { BROKER_BUSY_RPC_CODE, CodexAppServerClient } from "./lib/app-server.mjs
import { parseBrokerEndpoint } from "./lib/broker-endpoint.mjs";

const STREAMING_METHODS = new Set(["turn/start", "review/start", "thread/compact/start"]);
const BROKER_IDLE_TIMEOUT_ENV = "CODEX_COMPANION_BROKER_IDLE_TIMEOUT_MS";
const DEFAULT_BROKER_IDLE_TIMEOUT_MS = 5 * 60 * 1000;

function buildStreamThreadIds(method, params, result) {
const threadIds = new Set();
Expand Down Expand Up @@ -45,6 +47,14 @@ function writePidFile(pidFile) {
fs.writeFileSync(pidFile, `${process.pid}\n`, "utf8");
}

function resolveIdleTimeoutMs(value) {
const parsed = Number(value);
if (!Number.isFinite(parsed) || parsed <= 0) {
return DEFAULT_BROKER_IDLE_TIMEOUT_MS;
}
return Math.max(1, Math.floor(parsed));
}

async function main() {
const [subcommand, ...argv] = process.argv.slice(2);
if (subcommand !== "serve") {
Expand All @@ -70,6 +80,28 @@ async function main() {
let activeStreamSocket = null;
let activeStreamThreadIds = null;
const sockets = new Set();
const idleTimeoutMs = resolveIdleTimeoutMs(process.env[BROKER_IDLE_TIMEOUT_ENV]);
let idleTimer = null;
let shutdownPromise = null;

function cancelIdleShutdown() {
if (idleTimer) {
clearTimeout(idleTimer);
idleTimer = null;
}
}

function isFullyIdle() {
return sockets.size === 0 && !activeRequestSocket && !activeStreamSocket;
}

function isBusyForShutdown(requestingSocket) {
return Boolean(
activeRequestSocket ||
activeStreamSocket ||
[...sockets].some((socket) => socket !== requestingSocket)
);
}

function clearSocketOwnership(socket) {
if (activeRequestSocket === socket) {
Expand Down Expand Up @@ -99,23 +131,57 @@ async function main() {
}
}

async function shutdown(server) {
for (const socket of sockets) {
socket.end();
}
await appClient.close().catch(() => {});
await new Promise((resolve) => server.close(resolve));
if (listenTarget.kind === "unix" && fs.existsSync(listenTarget.path)) {
fs.unlinkSync(listenTarget.path);
function shutdown(server) {
if (!shutdownPromise) {
shutdownPromise = (async () => {
cancelIdleShutdown();
const serverClosed = new Promise((resolve) => server.close(resolve));
if (listenTarget.kind === "unix") {
fs.rmSync(listenTarget.path, { force: true });
}
for (const socket of sockets) {
socket.end();
}
await appClient.close().catch(() => {});
await serverClosed;
if (pidFile) {
fs.rmSync(pidFile, { force: true });
}
})();
}
if (pidFile && fs.existsSync(pidFile)) {
fs.unlinkSync(pidFile);
return shutdownPromise;
}

function scheduleIdleShutdown(server) {
cancelIdleShutdown();
if (!isFullyIdle() || shutdownPromise) {
return;
}

idleTimer = setTimeout(() => {
idleTimer = null;
if (!isFullyIdle() || shutdownPromise) {
return;
}
shutdown(server).then(
() => process.exit(0),
(error) => {
process.stderr.write(`${error instanceof Error ? error.message : String(error)}\n`);
process.exit(1);
}
);
}, idleTimeoutMs);
idleTimer.unref?.();
}

appClient.setNotificationHandler(routeNotification);

const server = net.createServer((socket) => {
if (shutdownPromise) {
socket.destroy();
return;
}
cancelIdleShutdown();
sockets.add(socket);
socket.setEncoding("utf8");
let buffer = "";
Expand Down Expand Up @@ -158,6 +224,13 @@ async function main() {
}

if (message.id !== undefined && message.method === "broker/shutdown") {
if (isBusyForShutdown(socket)) {
send(socket, {
id: message.id,
error: buildJsonRpcError(BROKER_BUSY_RPC_CODE, "Shared Codex broker is busy.")
});
continue;
}
send(socket, { id: message.id, result: {} });
await shutdown(server);
process.exit(0);
Expand Down Expand Up @@ -222,15 +295,15 @@ async function main() {
}
});

socket.on("close", () => {
sockets.delete(socket);
const releaseSocket = () => {
const removed = sockets.delete(socket);
clearSocketOwnership(socket);
});

socket.on("error", () => {
sockets.delete(socket);
clearSocketOwnership(socket);
});
if (removed) {
scheduleIdleShutdown(server);
}
};
socket.on("close", releaseSocket);
socket.on("error", releaseSocket);
});

process.on("SIGTERM", async () => {
Expand All @@ -243,6 +316,7 @@ async function main() {
process.exit(0);
});

server.on("listening", () => scheduleIdleShutdown(server));
server.listen(listenTarget.path);
}

Expand Down
39 changes: 34 additions & 5 deletions plugins/codex/scripts/lib/app-server.mjs
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
/**
* @typedef {Error & { data?: unknown, rpcCode?: number }} ProtocolError
* @typedef {Error & { code?: string, codexTransport?: "broker", data?: unknown, rpcCode?: number }} ProtocolError
* @typedef {import("./app-server-protocol").AppServerMethod} AppServerMethod
* @typedef {import("./app-server-protocol").AppServerNotification} AppServerNotification
* @typedef {import("./app-server-protocol").AppServerNotificationHandler} AppServerNotificationHandler
Expand All @@ -13,7 +13,7 @@ import process from "node:process";
import { spawn } from "node:child_process";
import readline from "node:readline";
import { parseBrokerEndpoint } from "./broker-endpoint.mjs";
import { ensureBrokerSession, loadBrokerSession } from "./broker-lifecycle.mjs";
import { ensureBrokerSession, isBrokerSessionReady, loadBrokerSession } from "./broker-lifecycle.mjs";
import { terminateProcessTree } from "./process.mjs";

const PLUGIN_MANIFEST_URL = new URL("../../.claude-plugin/plugin.json", import.meta.url);
Expand Down Expand Up @@ -54,6 +54,22 @@ function createProtocolError(message, data) {
return error;
}

function createBrokerConnectionClosedError() {
const error = /** @type {ProtocolError} */ (new Error("codex app-server connection closed."));
error.code = "ECONNRESET";
return error;
}

function markBrokerTransportError(error) {
const marked = /** @type {ProtocolError} */ (error instanceof Error ? error : new Error(String(error)));
marked.codexTransport = "broker";
return marked;
}

export function isBrokerTransportError(error) {
return error instanceof Error && /** @type {ProtocolError} */ (error).codexTransport === "broker";
}

class AppServerClientBase {
constructor(cwd, options = {}) {
this.cwd = cwd;
Expand Down Expand Up @@ -175,6 +191,11 @@ class AppServerClientBase {
this.resolveExit(undefined);
}

async waitForExit() {
await this.exitPromise;
throw this.exitError ?? new Error("codex app-server connection closed.");
}

sendMessage(_message) {
throw new Error("sendMessage must be implemented by subclasses.");
}
Expand Down Expand Up @@ -298,7 +319,7 @@ class BrokerCodexAppServerClient extends AppServerClientBase {
this.handleExit(error);
});
this.socket.on("close", () => {
this.handleExit(this.exitError);
this.handleExit(this.exitError ?? createBrokerConnectionClosedError());
});
});

Expand Down Expand Up @@ -338,7 +359,8 @@ export class CodexAppServerClient {
if (!options.disableBroker) {
brokerEndpoint = options.brokerEndpoint ?? options.env?.[BROKER_ENDPOINT_ENV] ?? process.env[BROKER_ENDPOINT_ENV] ?? null;
if (!brokerEndpoint && options.reuseExistingBroker) {
brokerEndpoint = loadBrokerSession(cwd)?.endpoint ?? null;
const brokerSession = loadBrokerSession(cwd);
brokerEndpoint = (await isBrokerSessionReady(brokerSession)) ? brokerSession.endpoint : null;
}
if (!brokerEndpoint && !options.reuseExistingBroker) {
const brokerSession = await ensureBrokerSession(cwd, { env: options.env });
Expand All @@ -348,7 +370,14 @@ export class CodexAppServerClient {
const client = brokerEndpoint
? new BrokerCodexAppServerClient(cwd, { ...options, brokerEndpoint })
: new SpawnedCodexAppServerClient(cwd, options);
await client.initialize();
try {
await client.initialize();
} catch (error) {
if (client.transport === "broker") {
throw markBrokerTransportError(error);
}
throw error;
}
return client;
}
}
Loading