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
14 changes: 13 additions & 1 deletion plugins/codex/scripts/lib/app-server.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -287,11 +287,23 @@ class BrokerCodexAppServerClient extends AppServerClientBase {
const target = parseBrokerEndpoint(this.endpoint);
this.socket = net.createConnection({ path: target.path });
this.socket.setEncoding("utf8");
this.socket.on("connect", resolve);
const timeoutMs = this.options.connectTimeoutMs ?? 2000;
const timer = setTimeout(() => {
const error = Object.assign(new Error("Timed out connecting to the Codex app-server broker."), {
code: "ETIMEDOUT"
});
this.socket.destroy();
reject(error);
}, timeoutMs);
this.socket.on("connect", () => {
clearTimeout(timer);
resolve();
});
this.socket.on("data", (chunk) => {
this.handleChunk(chunk);
});
this.socket.on("error", (error) => {
clearTimeout(timer);
if (!this.exitResolved) {
reject(error);
}
Expand Down
43 changes: 37 additions & 6 deletions plugins/codex/scripts/lib/broker-lifecycle.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -24,18 +24,49 @@ function connectToEndpoint(endpoint) {
export async function waitForBrokerEndpoint(endpoint, timeoutMs = 2000) {
const start = Date.now();
while (Date.now() - start < timeoutMs) {
const remainingMs = timeoutMs - (Date.now() - start);
if (remainingMs <= 0) {
break;
}
const ready = await new Promise((resolve) => {
let settled = false;
const finish = (value) => {
if (settled) {
return;
}
settled = true;
resolve(value);
};

const socket = connectToEndpoint(endpoint);
socket.on("connect", () => {
socket.end();
resolve(true);
});
socket.on("error", () => resolve(false));
const attemptTimeoutMs = Math.max(1, Math.min(100, remainingMs));
const timer = setTimeout(() => {
socket.destroy();
finish(false);
}, attemptTimeoutMs);

const onDone = (value) => {
clearTimeout(timer);
if (value) {
socket.end();
} else {
socket.destroy();
}
finish(value);
};

socket.setTimeout(attemptTimeoutMs, () => onDone(false));
socket.on("connect", () => onDone(true));
socket.on("error", () => onDone(false));
});
if (ready) {
return true;
}
await new Promise((resolve) => setTimeout(resolve, 50));
const waitMs = Math.min(50, Math.max(0, timeoutMs - (Date.now() - start)));
if (waitMs <= 0) {
break;
}
await new Promise((resolve) => setTimeout(resolve, waitMs));
}
return false;
}
Expand Down
3 changes: 2 additions & 1 deletion plugins/codex/scripts/lib/codex.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -621,7 +621,8 @@ async function withAppServer(cwd, fn) {
const brokerRequested = client?.transport === "broker" || Boolean(process.env[BROKER_ENDPOINT_ENV]);
const shouldRetryDirect =
(client?.transport === "broker" && error?.rpcCode === BROKER_BUSY_RPC_CODE) ||
(brokerRequested && (error?.code === "ENOENT" || error?.code === "ECONNREFUSED"));
(brokerRequested &&
(error?.code === "ENOENT" || error?.code === "ECONNREFUSED" || error?.code === "ETIMEDOUT"));

if (client) {
await client.close().catch(() => {});
Expand Down
31 changes: 31 additions & 0 deletions tests/broker-lifecycle.test.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
import net from "node:net";
import test from "node:test";
import assert from "node:assert/strict";

import { waitForBrokerEndpoint } from "../plugins/codex/scripts/lib/broker-lifecycle.mjs";

function hangingConnection() {
return {
setTimeout() {},
on() {
return this;
},
removeAllListeners() {},
end() {},
destroy() {}
};
}

test("waitForBrokerEndpoint returns false when connect hangs past the timeout", { timeout: 2000 }, async (t) => {
const originalCreateConnection = net.createConnection;
net.createConnection = hangingConnection;
t.after(() => {
net.createConnection = originalCreateConnection;
});

const started = Date.now();
const ready = await waitForBrokerEndpoint("unix:/tmp/codex-hung-broker.sock", 150);

assert.equal(ready, false);
assert.equal(Date.now() - started < 1000, true);
});