import { createInterface } from "node:readline"; import { ConfigurationError, parseEndpoint, validateInvocation, } from "./config.js"; const STARTUP_TIMEOUT_MS = 5_000; function writeDiagnostic(output, message) { output.write(`codex-app-server-bridge: ${message}\n`); } export async function runBridge(options) { try { validateInvocation(options.arguments); } catch (error) { const message = error instanceof ConfigurationError ? error.message : "invalid invocation"; writeDiagnostic(options.errorOutput, message); return 1; } let endpoint; try { endpoint = parseEndpoint(options.environment.CODEX_APP_SERVER_URL); } catch (error) { const message = error instanceof ConfigurationError ? error.message : "invalid configuration"; writeDiagnostic(options.errorOutput, `configuration error: ${message}`); return 1; } return await new Promise((resolve) => { let socket; let lineReader; let completed = false; let localShutdown = false; let failure = false; let startupTimer; const removeProcessListeners = () => { process.removeListener("SIGINT", handleSignal); process.removeListener("SIGTERM", handleSignal); }; const finish = (exitCode) => { if (completed) { return; } completed = true; if (startupTimer !== undefined) { clearTimeout(startupTimer); } lineReader?.close(); options.input.pause(); removeProcessListeners(); resolve(exitCode); }; const closeSocket = (code = 1000, reason = "bridge shutdown") => { if (socket.readyState === WebSocket.OPEN) { socket.close(code, reason); } else if (socket.readyState === WebSocket.CONNECTING) { try { socket.close(); } catch { // Process completion still closes a connection that has not opened. } } }; const fail = (message, closeCode = 1011) => { if (failure || completed) { return; } failure = true; writeDiagnostic(options.errorOutput, message); closeSocket(closeCode, "bridge failure"); finish(1); }; const beginLocalShutdown = () => { if (completed || localShutdown) { return; } localShutdown = true; if (startupTimer !== undefined) { clearTimeout(startupTimer); } lineReader?.close(); options.input.pause(); closeSocket(); if (socket.readyState !== WebSocket.OPEN && socket.readyState !== WebSocket.CLOSING) { finish(0); } }; function handleSignal() { beginLocalShutdown(); } process.once("SIGINT", handleSignal); process.once("SIGTERM", handleSignal); try { socket = new WebSocket(endpoint.url); socket.binaryType = "arraybuffer"; } catch { removeProcessListeners(); writeDiagnostic(options.errorOutput, `connection error for ${endpoint.display}`); resolve(1); return; } startupTimer = setTimeout(() => { fail(`startup timeout after ${STARTUP_TIMEOUT_MS} ms for ${endpoint.display}`); }, STARTUP_TIMEOUT_MS); socket.addEventListener("open", () => { if (completed) { closeSocket(); return; } if (startupTimer !== undefined) { clearTimeout(startupTimer); startupTimer = undefined; } lineReader = createInterface({ input: options.input, crlfDelay: Infinity, terminal: false, }); lineReader.on("line", (line) => { if (line.trim().length === 0) { return; } try { socket.send(line); } catch { fail(`transport error for ${endpoint.display}`); } }); lineReader.once("close", beginLocalShutdown); }); socket.addEventListener("message", (event) => { if (typeof event.data !== "string") { fail("protocol error: binary WebSocket frames are not supported", 1003); return; } options.output.write(`${event.data}\n`); }); socket.addEventListener("error", () => { if (localShutdown) { return; } fail(`connection error for ${endpoint.display}`); }); socket.addEventListener("close", () => { if (localShutdown) { finish(0); } else if (failure) { finish(1); } else { fail(`connection closed unexpectedly for ${endpoint.display}`); } }); }); } export async function main() { process.exitCode = await runBridge({ arguments: process.argv.slice(2), environment: process.env, input: process.stdin, output: process.stdout, errorOutput: process.stderr, }); }