From e9c746b69196ae125d70351b5963379e474d9367 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 15 Sep 2026 20:56:41 +0000 Subject: [PATCH 1/2] Potato follow-up: bounded sends, first-response timer, output commit, approval ownership, strict 101 + EOF flush Co-authored-by: Zack Jackson --- .changeset/desktop-shim-followup.md | 5 + README.md | 6 +- src/core/codex-bridge.js | 52 ++++-- src/core/desktop-shim-bridge.js | 34 ++-- test/codex-bridge.test.js | 38 +++++ test/desktop-shim.test.js | 239 ++++++++++++++++++++++++++++ 6 files changed, 347 insertions(+), 27 deletions(-) create mode 100644 .changeset/desktop-shim-followup.md diff --git a/.changeset/desktop-shim-followup.md b/.changeset/desktop-shim-followup.md new file mode 100644 index 0000000..af6a84c --- /dev/null +++ b/.changeset/desktop-shim-followup.md @@ -0,0 +1,5 @@ +--- +"grok-bot-cli": patch +--- + +Harden the Desktop shim bridge per follow-up review: bound every send/lock with the monotonic session budget (no pong/cleanup hang on a non-reading peer), track the first daemon response instead of any notification for the first-RPC deadline, commit on input or output (never fall back to stock on used stdout), answer approvals only for the turn gbot started (foreign thread/Desktop turn stays silent), require HTTP/1.1 101 plus accept, and flush EOF-tail data before the WS Close. diff --git a/README.md b/README.md index 1f80b56..98589d3 100644 --- a/README.md +++ b/README.md @@ -87,7 +87,9 @@ gbot codex desktop-shim uninstall # removes wrapper/bridge/LaunchAgent; Desktop Install copies the scripts out of the package into durable `~/.codex/bin/` paths, so uninstalling the npm package never leaves Desktop pointing into a deleted checkout. The wrapper exports install-time `CODEX_HOME` and derives the socket from it at runtime (`CODEX_APP_SERVER_SOCK`), so custom `CODEX_HOME` layouts work; `gbot` itself resolves the socket the same way (`CODEX_APP_SERVER_SOCK` wins, else `CODEX_HOME`), and `status` reports the effective path and which variable won. -Fail-open runs only before any stdin byte is consumed: the wrapper preflights the daemon socket (3 s deadline) and runs the bridge as a child — never `exec` — falling through to the real standalone `codex` for every non-`app-server` spawn, every preflight failure, and every pre-session bridge failure (exit 1), always on pristine stdio. The bridge's connect + WebSocket-upgrade handshake runs under one absolute monotonic deadline (`CODEX_BRIDGE_CONNECT_TIMEOUT`, default 10 s) with the 101 status and `Sec-WebSocket-Accept` validated, and a first-message deadline (`CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT`, default 30 s) bounds connect-to-first-daemon-message covering the first RPC — preflight (3 s) + upgrade (10 s) + first RPC (30 s) stays far below a daemon-lock hang, and none of these stages reads stdin. The commit point is byte-precise: the bridge counts every stdin byte read (buffered readahead included), so even an unforwarded blank line counts as committed. A mid-session bridge failure (exit 2+) makes the wrapper exit promptly so Desktop reconnects; the fallback never runs on half-consumed stdin. The wrapper never starts the daemon on the spawn path (a wedged daemon lock must not block Desktop) — upkeep belongs to the LaunchAgent login script and install, which share one `CODEX_HOME` for the GUI domain and the daemon they start; if the socket is absent, Desktop simply runs stock Codex until the daemon is started. +Fail-open runs only before any stdin byte is consumed and no daemon payload was forwarded: the wrapper preflights the daemon socket (3 s deadline) and runs the bridge as a child — never `exec` — falling through to the real standalone `codex` for every non-`app-server` spawn, every preflight failure, and every pre-session bridge failure (exit 1), always on pristine stdio+stdout. The bridge's connect + WebSocket-upgrade handshake runs under one absolute monotonic deadline (`CODEX_BRIDGE_CONNECT_TIMEOUT`, default 10 s) requiring `HTTP/1.1 101` plus a matching `Sec-WebSocket-Accept` (a `200` fails even with a valid hash), and a first-response deadline (`CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT`, default 30 s) bounds connect-to-first-daemon-response covering the first RPC — an unrelated notification never satisfies it — preflight (3 s) + upgrade (10 s) + first RPC (30 s) stays far below a daemon-lock hang, and none of these stages reads stdin. Every send/lock shares the same monotonic session budget (remaining first-RPC budget, else the connect timeout), so a non-reading peer cannot wedge pong writes or timeout cleanup; timeouts close and exit. The commit point is byte-precise on input OR output: the bridge counts every stdin byte read (buffered readahead included, even an unforwarded blank line) and marks stdout used on any forwarded daemon payload, so the fallback only ever runs on truly pristine stdio+stdout. Buffered EOF-tail data flushes before the WS Close. A mid-session bridge failure (exit 2+) makes the wrapper exit promptly so Desktop reconnects; the fallback never runs on half-consumed stdin or already-used stdout. The wrapper never starts the daemon on the spawn path (a wedged daemon lock must not block Desktop) — upkeep belongs to the LaunchAgent login script and install, which share one `CODEX_HOME` for the GUI domain and the daemon they start; if the socket is absent, Desktop simply runs stock Codex until the daemon is started. + +Residual risks, stated honestly: a concurrent foreign approval carrying no thread/turn ids inside `gbot`'s own `turn/start` window is indistinguishable from `gbot`'s and will still be refused — don't send to threads Desktop is actively driving. A daemon that speaks framing-valid but semantically unexpected JSON-RPC (unknown methods, id-less responses) is treated as transport; pins are to app-server schema 0.154.0. The bridge trusts the local control socket; a malicious local daemon could hold the session up to the stated budgets, not past them. Permanent tradeoff, stated plainly: Desktop's app-tools MCP (`-c` overrides on its spawn line) is not applied to the already-running managed daemon, and no config/`mcpServer`/`reload` path imports Desktop's `-c` flags — Desktop app-tools stay degraded while pointed at the shared daemon. Fully quit and relaunch ChatGPT.app after install (or login) so it inherits `CODEX_CLI_PATH`. `status` reads the macOS GUI-domain value via `launchctl getenv` (what Desktop actually inherits) alongside the calling shell's value. LaunchAgent persistence is macOS-first; elsewhere install still writes the wrapper and bridge but leaves `CODEX_CLI_PATH` for you to export. `~/.codex/bin` holds scripts only — there is no extra revert note to clean up; revert is `gbot codex desktop-shim uninstall` plus this section. @@ -114,7 +116,7 @@ Sends at `hop >= GROK_BOT_MAX_HOPS` (default 4) are refused with `reason: "hop-l - `busy` / `thread-error` / `unknown-status`: see above. - `route-not-allowed` / `hop-limit` / `experimental-disabled`: refused by operator policy, the relay bound, or the experimental-API gate. - `unsupported`: the daemon does not offer the (experimental) method `--when-busy queue` needs. -- `approval-refused` (`delivery: "accepted"`): `gbot` never approves commands or file changes on your behalf. If Codex asks while `gbot` is still connected, `send` refuses the request, exits 1, and tells you the turn id. Refusal replies are armed only once `gbot`'s own `turn/start` is in flight — an approval outstanding from another client's turn (notably ChatGPT Desktop's) during resume, a busy reject, or a queue add is recorded but never answered, so `gbot` cannot reject Desktop's approval. Residual race, stated honestly: an approval from a concurrent foreign turn arriving inside `gbot`'s own `turn/start` window is indistinguishable from `gbot`'s and will be refused; don't send to threads Desktop is actively driving. `send` disconnects as soon as the turn starts, so later approval requests stay with the daemon for a Codex client to answer; for unattended sends set `approval_policy = "never"` in the daemon's `config.toml`. +- `approval-refused` (`delivery: "accepted"`): `gbot` never approves commands or file changes on your behalf. If Codex asks while `gbot` is still connected, `send` refuses the request, exits 1, and tells you the turn id. Refusal replies are armed only once `gbot`'s own `turn/start` is in flight and only for requests naming that thread/turn — an approval outstanding from another client's turn (notably ChatGPT Desktop's) during resume, a busy reject, or a queue add is recorded but never answered, and a request naming another thread or Desktop turn inside the window is foreign-skipped, so `gbot` cannot reject Desktop's approval. Residual race, stated honestly: an approval from a concurrent foreign turn arriving inside `gbot`'s own `turn/start` window with no thread/turn ids is indistinguishable from `gbot`'s and will be refused; don't send to threads Desktop is actively driving. `send` disconnects as soon as the turn starts, so later approval requests stay with the daemon for a Codex client to answer; for unattended sends set `approval_policy = "never"` in the daemon's `config.toml`. - `transport` / `bad-response` (`delivery: "unknown"`): the connection dropped or the daemon answered off-schema after the request left. ## Talking to Grok Bot from Codex diff --git a/src/core/codex-bridge.js b/src/core/codex-bridge.js index 50bdf9f..9d8be26 100644 --- a/src/core/codex-bridge.js +++ b/src/core/codex-bridge.js @@ -282,11 +282,15 @@ export function connectCodexAppServer(path, { timeoutMs = 15000 } = {}) { const client = { refused, // Server-initiated requests are only *answered* once our own turn/start - // is in flight: anything earlier belongs to another client's turn - // (notably Desktop's) and must never be rejected on their behalf. Early - // requests are still recorded; the connection closes right after, so a - // pending request simply dies with it instead of denying someone's approval. + // is in flight, and then only when they name our own thread/turn: + // anything earlier — or naming another thread or Desktop turn — belongs + // to another client and must never be rejected on their behalf. Early + // and foreign requests are still recorded; the connection closes right + // after, so a pending request simply dies with it instead of denying + // someone's approval. answerServerRequests: false, + expectedThreadId: null, + expectedTurnId: null, request(method, params) { const id = nextId++; return new Promise((res, rej) => { @@ -327,7 +331,20 @@ export function connectCodexAppServer(path, { timeoutMs = 15000 } = {}) { const onMessage = (msg) => { if (msg.id != null && msg.method) { - const answered = client.answerServerRequests; + let answered = client.answerServerRequests; + // Approval ownership: only answer requests for the turn gbot itself + // started. A request naming another thread — or naming a turn that is + // not ours (including any named Desktop turn while our turn id is + // still unknown) — is recorded foreign-silent, never rejected. + if (answered) { + const params = msg.params && typeof msg.params === "object" && !Array.isArray(msg.params) ? msg.params : null; + if (params) { + const threadId = params.threadId ?? params.thread_id; + const turnId = params.turnId ?? params.turn_id ?? params.turn?.id; + if (threadId != null && threadId !== client.expectedThreadId) answered = false; + else if (turnId != null && turnId !== client.expectedTurnId) answered = false; + } + } refused.push({ id: msg.id, method: msg.method, params: msg.params, answered }); if (!answered) return; sendJson({ @@ -893,11 +910,14 @@ async function sendToCodexThreadInner(threadId, text, { env, envelope, whenBusy // the daemon exposes no compare-and-start request. Upgrade path: thread/queue/add + thread/queue/start once stable. // Scope refusals to this turn: server requests from earlier calls belong to another context. const seenRefused = client.refused.length; - // Arm refusal replies only now that our own turn/start is in flight: an - // approval arriving in this window plausibly belongs to our turn. Anything - // recorded earlier (resume, busy reject, queue) stays unanswered so gbot - // never rejects another client's approval. Residual race: a concurrent - // foreign approval inside this window is indistinguishable from ours. + // Arm refusal replies only now that our own turn/start is in flight, scoped + // to our own thread/turn: anything recorded earlier (resume, busy reject, + // queue) — or naming another thread or Desktop turn — stays unanswered so + // gbot never rejects another client's approval. Residual race, stated in + // the README: a concurrent foreign approval with no thread/turn ids inside + // this window is indistinguishable from ours. + client.expectedThreadId = threadId; + client.expectedTurnId = null; client.answerServerRequests = true; let turn; try { @@ -924,8 +944,18 @@ async function sendToCodexThreadInner(threadId, text, { env, envelope, whenBusy { delivery: "unknown", reason: "bad-response", threadId, envelope }, ); } + client.expectedTurnId = turnId; const freshRefused = client.refused.slice(seenRefused) - .filter((r) => r.answered && (!r.params || r.params.threadId == null || r.params.threadId === threadId)); + .filter((r) => { + if (!r.answered) return false; + const params = r.params && typeof r.params === "object" && !Array.isArray(r.params) ? r.params : null; + if (!params) return true; + const refusedThreadId = params.threadId ?? params.thread_id; + const refusedTurnId = params.turnId ?? params.turn_id ?? params.turn?.id; + if (refusedThreadId != null && refusedThreadId !== threadId) return false; + if (refusedTurnId != null && refusedTurnId !== turnId) return false; + return true; + }); if (freshRefused.length) { const methods = freshRefused.map((r) => r.method).join(", "); throw new CodexSendError( diff --git a/src/core/desktop-shim-bridge.js b/src/core/desktop-shim-bridge.js index 8be5823..a27f8dd 100644 --- a/src/core/desktop-shim-bridge.js +++ b/src/core/desktop-shim-bridge.js @@ -2,21 +2,27 @@ // Codex app-server control socket. Vendored from uploads/codex-stdio-to-daemon-ws.py // (ChatGPT Desktop -> managed daemon shim) with fail-open deltas, marked below: // (1) ws_connect enforces one absolute monotonic deadline for connect + HTTP -// upgrade (CODEX_BRIDGE_CONNECT_TIMEOUT, default 10s), validates the 101 status -// and Sec-WebSocket-Accept, and retains bytes coalesced after the upgrade -// headers instead of discarding the first data frame. +// upgrade (CODEX_BRIDGE_CONNECT_TIMEOUT, default 10s), requires HTTP/1.1 101 +// plus a matching Sec-WebSocket-Accept (a 200 fails even with a valid hash), +// and retains bytes coalesced after the upgrade headers instead of discarding +// the first data frame. // (2) read_message assembles continuation frames; control frames answer inline; -// socket writes are serialized with a lock shared by both threads. -// (3) A first-message deadline (CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT, default -// 30s) bounds connect-to-first-daemon-message covering the first RPC; the -// preflight (3s) + upgrade (10s) + first RPC (30s) budget stays far below a -// daemon-lock hang. Exits are 0 served, 1 pre-stdio (wrapper falls back to -// real Codex), 2 mid-session (wrapper exits so Desktop reconnects). -// (4) Byte-precise commit: os.read into a userspace line buffer counts every -// stdin byte (buffered readline could readahead past the commit point), so the -// fallback only ever runs on a truly pristine stdio. The reader signals a done -// event so WS close wakes the select loop; no path leaves a thread hung. +// every socket send/lock is bounded by the same monotonic session budget +// (remaining first-RPC budget, else the connect timeout) so a non-reading peer +// cannot wedge pong writes or timeout cleanup; timeouts close and exit. +// (3) A first-response deadline (CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT, default +// 30s) bounds connect-to-first-daemon-response covering the first RPC; an +// unrelated notification never disables it. The preflight (3s) + upgrade (10s) +// + first RPC (30s) budget stays far below a daemon-lock hang. Exits are +// 0 served, 1 pre-stdio (wrapper falls back to real Codex on pristine +// stdio+stdout), 2 mid-session (wrapper exits so Desktop reconnects). +// (4) Byte-precise commit on input OR output: os.read into a userspace line +// buffer counts every stdin byte (even an unforwarded blank line), and any +// forwarded daemon payload marks stdout used, so the fallback only ever runs +// on truly pristine stdio+stdout. Buffered EOF-tail data flushes before the +// WS Close. The reader signals a done event so WS close wakes the select loop; +// no path leaves a thread hung. // Desktop stdio hangs with stock `codex app-server proxy`, so this custom // bridge is required. Install writes it to the CODEX_HOME bin dir. Never reads // Desktop private pipes (CODEX_APP_TOOLS_PIPE_PATH, /tmp/codex-browser-use/). -export const BRIDGE_SOURCE = "#!/usr/bin/env python3\n\"\"\"Stdio JSONL (Desktop) <-> WebSocket-over-unix (managed app-server daemon).\nReversible: remove CODEX_CLI_PATH wrapper. Does not read Desktop private pipes.\n\"\"\"\nfrom __future__ import annotations\nimport base64, hashlib, json, os, select, socket, struct, sys, threading, time\n\nSOCK = os.environ.get(\n \"CODEX_APP_SERVER_SOCK\",\n os.path.expanduser(\"~/.codex/app-server-control/app-server-control.sock\"),\n)\nLOG = os.environ.get(\"CODEX_STDIO_BRIDGE_LOG\", os.path.expanduser(\"~/.codex/bin/codex-stdio-to-daemon-ws.log\"))\nWS_GUID = \"258EAFA5-E914-47DA-95CA-C5AB0DC85B11\"\n\n# Wrapper contract: 0 served the session, 1 failed before any stdin was\n# consumed (wrapper falls back to real Codex on pristine stdio), 2 failed\n# mid-session (wrapper exits promptly so Desktop reconnects; the fallback\n# must never run on half-consumed stdin).\nEXIT_SERVED = 0\nEXIT_PRE_STDIO = 1\nEXIT_MID_SESSION = 2\n\nCONNECT_TIMEOUT_S = float(os.environ.get(\"CODEX_BRIDGE_CONNECT_TIMEOUT\", \"10\"))\n# Bound from connect to the first daemon message: a daemon that completes the\n# upgrade but never answers the first RPC must not hold Desktop past this.\nFIRST_MESSAGE_TIMEOUT_S = float(os.environ.get(\"CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT\", \"30\"))\n\ndef log(msg: str) -> None:\n try:\n with open(LOG, \"a\") as f:\n f.write(time.strftime(\"%Y-%m-%dT%H:%M:%SZ\", time.gmtime()) + \" \" + msg + \"\\n\")\n except OSError:\n pass\n\ndef ws_connect(path: str, timeout: float = CONNECT_TIMEOUT_S):\n \"\"\"Connect + HTTP upgrade under one absolute monotonic deadline.\n\n Returns (socket, trailing) where trailing holds bytes coalesced after the\n upgrade headers (TCP may deliver the 101 and the first WS frame together;\n discarding them would lose the first message). Validates the 101 status\n and Sec-WebSocket-Accept so a non-WebSocket listener fails fast.\n \"\"\"\n deadline = time.monotonic() + timeout\n s = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)\n s.settimeout(timeout)\n try:\n s.connect(path)\n key = base64.b64encode(os.urandom(16)).decode()\n expected = base64.b64encode(hashlib.sha1((key + WS_GUID).encode()).digest()).decode()\n req = (\n f\"GET /rpc HTTP/1.1\\r\\nHost: localhost\\r\\nUpgrade: websocket\\r\\n\"\n f\"Connection: Upgrade\\r\\nSec-WebSocket-Key: {key}\\r\\nSec-WebSocket-Version: 13\\r\\n\\r\\n\"\n ).encode()\n s.sendall(req)\n buf = b\"\"\n while True:\n if b\"\\r\\n\\r\\n\" in buf:\n break\n remaining = deadline - time.monotonic()\n if remaining <= 0:\n raise TimeoutError(\"WebSocket upgrade timed out\")\n if len(buf) > 65536:\n raise ConnectionError(\"WS upgrade headers too large\")\n s.settimeout(remaining)\n chunk = s.recv(4096)\n if not chunk:\n raise ConnectionError(\"daemon closed during WS upgrade\")\n buf += chunk\n head, _, trailing = buf.partition(b\"\\r\\n\\r\\n\")\n lines = head.split(b\"\\r\\n\")\n if b\"101\" not in lines[0]:\n raise ConnectionError(f\"WS upgrade failed: {head[:200]!r}\")\n accept = None\n for line in lines[1:]:\n name, sep, value = line.partition(b\":\")\n if sep and name.strip().lower() == b\"sec-websocket-accept\":\n accept = value.strip().decode(\"latin-1\")\n if accept != expected:\n raise ConnectionError(\"WS upgrade accept mismatch\")\n except Exception:\n try:\n s.close()\n except OSError:\n pass\n raise\n s.settimeout(None)\n return s, trailing\n\ndef mask_send(sock: socket.socket, payload: bytes, opcode: int = 1) -> None:\n mask = os.urandom(4)\n ln = len(payload)\n if ln < 126:\n hdr = bytes([0x80 | opcode, 0x80 | ln])\n elif ln < 65536:\n hdr = bytes([0x80 | opcode, 0x80 | 126]) + struct.pack(\"!H\", ln)\n else:\n hdr = bytes([0x80 | opcode, 0x80 | 127]) + struct.pack(\"!Q\", ln)\n sock.sendall(hdr + mask + bytes(b ^ mask[i % 4] for i, b in enumerate(payload)))\n\ndef locked_send(sock: socket.socket, lock: threading.Lock, payload: bytes, opcode: int = 1) -> None:\n with lock:\n mask_send(sock, payload, opcode=opcode)\n\ndef _buffered_reader(sock: socket.socket, initial: bytes = b\"\"):\n buf = bytearray(initial)\n def recv(n):\n while len(buf) < n:\n chunk = sock.recv(n - len(buf))\n if not chunk:\n raise ConnectionError(\"eof\")\n buf.extend(chunk)\n out = bytes(buf[:n])\n del buf[:n]\n return out\n return recv\n\ndef read_frame(recv):\n hdr = recv(2)\n fin = bool(hdr[0] & 0x80)\n opcode = hdr[0] & 0x0F\n masked = bool(hdr[1] & 0x80)\n ln = hdr[1] & 0x7F\n if ln == 126:\n ln = struct.unpack(\"!H\", recv(2))[0]\n elif ln == 127:\n ln = struct.unpack(\"!Q\", recv(8))[0]\n mask = recv(4) if masked else b\"\"\n payload = recv(ln)\n if masked:\n payload = bytes(b ^ mask[i % 4] for i, b in enumerate(payload))\n return fin, opcode, payload\n\ndef read_message(recv, send_pong):\n \"\"\"Next complete data message, assembling continuation frames.\n\n Control frames are handled inline (pong answered, close reported) so a\n fragmented message split around a ping still reassembles.\n \"\"\"\n opcode = None\n parts = []\n while True:\n fin, op, payload = read_frame(recv)\n if op == 0x8:\n return (\"close\", None, b\"\")\n if op == 0x9:\n send_pong(payload)\n continue\n if op == 0xA:\n continue\n if op in (0x1, 0x2):\n if opcode is not None:\n raise ConnectionError(\"data frame inside fragmented message\")\n if fin:\n return (\"data\", op, payload)\n opcode = op\n parts.append(payload)\n continue\n if op == 0x0:\n if opcode is None:\n raise ConnectionError(\"stray continuation frame\")\n parts.append(payload)\n if fin:\n return (\"data\", opcode, b\"\".join(parts))\n continue\n raise ConnectionError(f\"unsupported opcode {op}\")\n\ndef stdout_writer(sock: socket.socket, recv, send_pong, done: threading.Event, first_msg: threading.Event) -> None:\n out = sys.stdout.buffer\n try:\n while True:\n kind, opcode, payload = read_message(recv, send_pong)\n if kind == \"close\":\n log(\"daemon WS close\")\n return\n if opcode == 0x1:\n # Desktop StdioConnection expects newline-delimited JSON text\n if not payload.endswith(b\"\\n\"):\n payload += b\"\\n\"\n out.write(payload)\n out.flush()\n else:\n out.write(payload)\n out.flush()\n first_msg.set()\n except Exception as e:\n log(f\"stdout_writer exit: {e}\")\n finally:\n done.set()\n\ndef main() -> int:\n log(f\"start sock={SOCK}\")\n try:\n sock, pending = ws_connect(SOCK)\n except Exception as e:\n log(f\"connect failed: {e}\")\n sys.stderr.write(f\"codex-stdio-to-daemon-ws: {e}\\n\")\n return EXIT_PRE_STDIO\n log(\"connected\")\n connected_at = time.monotonic()\n send_lock = threading.Lock()\n def send_pong(payload: bytes) -> None:\n locked_send(sock, send_lock, payload, opcode=0xA)\n recv = _buffered_reader(sock, pending)\n done = threading.Event()\n first_msg = threading.Event()\n t = threading.Thread(target=stdout_writer, args=(sock, recv, send_pong, done, first_msg), daemon=True)\n t.start()\n # Byte-precise commit point: the fallback may run only while zero stdin\n # bytes have left the pipe. Reads use os.read into a userspace line buffer\n # (buffered readline could readahead past the commit point), and every\n # byte read counts — even a blank line. first_msg bounds the first RPC.\n stdin_bytes = 0\n pending_in = b\"\"\n clean_eof = False\n try:\n fd = sys.stdin.fileno()\n while not done.is_set():\n if not first_msg.is_set() and time.monotonic() - connected_at > FIRST_MESSAGE_TIMEOUT_S:\n log(\"first daemon message timed out\")\n break\n ready, _, _ = select.select([fd], [], [], 0.5)\n if not ready:\n continue\n chunk = os.read(fd, 65536)\n if not chunk:\n clean_eof = True\n break\n stdin_bytes += len(chunk)\n pending_in += chunk\n while b\"\\n\" in pending_in:\n line, _, pending_in = pending_in.partition(b\"\\n\")\n line = line.strip()\n if not line:\n continue\n locked_send(sock, send_lock, line)\n except Exception as e:\n log(f\"stdin loop exit: {e}\")\n try:\n locked_send(sock, send_lock, b\"\", opcode=0x8)\n except Exception:\n pass\n if clean_eof:\n tail = pending_in.strip()\n if tail and not done.is_set():\n try:\n locked_send(sock, send_lock, tail)\n except Exception as e:\n log(f\"eof flush failed: {e}\")\n log(\"exit\")\n return EXIT_SERVED\n # Reader died, first RPC timed out, or stdin broke: fail open only if no\n # stdin byte was ever consumed.\n rc = EXIT_MID_SESSION if stdin_bytes else EXIT_PRE_STDIO\n log(f\"exit rc={rc}\")\n return rc\n\nif __name__ == \"__main__\":\n raise SystemExit(main())\n"; +export const BRIDGE_SOURCE = "#!/usr/bin/env python3\n\"\"\"Stdio JSONL (Desktop) <-> WebSocket-over-unix (managed app-server daemon).\nReversible: remove CODEX_CLI_PATH wrapper. Does not read Desktop private pipes.\n\"\"\"\nfrom __future__ import annotations\nimport base64, hashlib, json, os, select, socket, struct, sys, threading, time\n\nSOCK = os.environ.get(\n \"CODEX_APP_SERVER_SOCK\",\n os.path.expanduser(\"~/.codex/app-server-control/app-server-control.sock\"),\n)\nLOG = os.environ.get(\"CODEX_STDIO_BRIDGE_LOG\", os.path.expanduser(\"~/.codex/bin/codex-stdio-to-daemon-ws.log\"))\nWS_GUID = \"258EAFA5-E914-47DA-95CA-C5AB0DC85B11\"\n\n# Wrapper contract: 0 served the session, 1 failed before any stdin was\n# consumed and no daemon payload was forwarded (wrapper falls back to real\n# Codex on pristine stdio), 2 failed mid-session (wrapper exits promptly so\n# Desktop reconnects; the fallback must never run on half-consumed stdin or\n# already-used stdout).\nEXIT_SERVED = 0\nEXIT_PRE_STDIO = 1\nEXIT_MID_SESSION = 2\n\nCONNECT_TIMEOUT_S = float(os.environ.get(\"CODEX_BRIDGE_CONNECT_TIMEOUT\", \"10\"))\n# Bound from connect to the first daemon response: a daemon that completes the\n# upgrade but never answers the first RPC must not hold Desktop past this.\nFIRST_MESSAGE_TIMEOUT_S = float(os.environ.get(\"CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT\", \"30\"))\n\ndef log(msg: str) -> None:\n try:\n with open(LOG, \"a\") as f:\n f.write(time.strftime(\"%Y-%m-%dT%H:%M:%SZ\", time.gmtime()) + \" \" + msg + \"\\n\")\n except OSError:\n pass\n\ndef ws_connect(path: str, timeout: float = CONNECT_TIMEOUT_S):\n \"\"\"Connect + HTTP upgrade under one absolute monotonic deadline.\n\n Returns (socket, trailing) where trailing holds bytes coalesced after the\n upgrade headers (TCP may deliver the 101 and the first WS frame together;\n discarding them would lose the first message). Requires HTTP/1.1 101 plus\n a matching Sec-WebSocket-Accept so a 200 (or any non-101) fails fast even\n when the accept hash happens to be valid.\n \"\"\"\n deadline = time.monotonic() + timeout\n s = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)\n s.settimeout(timeout)\n try:\n s.connect(path)\n key = base64.b64encode(os.urandom(16)).decode()\n expected = base64.b64encode(hashlib.sha1((key + WS_GUID).encode()).digest()).decode()\n req = (\n f\"GET /rpc HTTP/1.1\\r\\nHost: localhost\\r\\nUpgrade: websocket\\r\\n\"\n f\"Connection: Upgrade\\r\\nSec-WebSocket-Key: {key}\\r\\nSec-WebSocket-Version: 13\\r\\n\\r\\n\"\n ).encode()\n s.sendall(req)\n buf = b\"\"\n while True:\n if b\"\\r\\n\\r\\n\" in buf:\n break\n remaining = deadline - time.monotonic()\n if remaining <= 0:\n raise TimeoutError(\"WebSocket upgrade timed out\")\n if len(buf) > 65536:\n raise ConnectionError(\"WS upgrade headers too large\")\n s.settimeout(remaining)\n chunk = s.recv(4096)\n if not chunk:\n raise ConnectionError(\"daemon closed during WS upgrade\")\n buf += chunk\n head, _, trailing = buf.partition(b\"\\r\\n\\r\\n\")\n lines = head.split(b\"\\r\\n\")\n if not lines or not lines[0].startswith(b\"HTTP/1.1 101\"):\n raise ConnectionError(f\"WS upgrade failed: {head[:200]!r}\")\n accept = None\n for line in lines[1:]:\n name, sep, value = line.partition(b\":\")\n if sep and name.strip().lower() == b\"sec-websocket-accept\":\n accept = value.strip().decode(\"latin-1\")\n if accept != expected:\n raise ConnectionError(\"WS upgrade accept mismatch\")\n except Exception:\n try:\n s.close()\n except OSError:\n pass\n raise\n s.settimeout(None)\n return s, trailing\n\ndef mask_send(sock: socket.socket, payload: bytes, opcode: int = 1) -> None:\n mask = os.urandom(4)\n ln = len(payload)\n if ln < 126:\n hdr = bytes([0x80 | opcode, 0x80 | ln])\n elif ln < 65536:\n hdr = bytes([0x80 | opcode, 0x80 | 126]) + struct.pack(\"!H\", ln)\n else:\n hdr = bytes([0x80 | opcode, 0x80 | 127]) + struct.pack(\"!Q\", ln)\n sock.sendall(hdr + mask + bytes(b ^ mask[i % 4] for i, b in enumerate(payload)))\n\ndef locked_send(sock: socket.socket, lock: threading.Lock, payload: bytes, opcode: int = 1, timeout: float = CONNECT_TIMEOUT_S) -> None:\n \"\"\"Serialized send bounded by a monotonic budget.\n\n Both the lock wait and the sendall wait share `timeout` seconds. A\n non-reading peer must never wedge pong writes or timeout cleanup: on\n expiry raise TimeoutError so the caller closes and exits without hanging.\n \"\"\"\n if timeout is None:\n timeout = CONNECT_TIMEOUT_S\n if timeout <= 0:\n raise TimeoutError(\"send budget expired\")\n if not lock.acquire(timeout=timeout):\n raise TimeoutError(\"send lock timed out\")\n try:\n sock.settimeout(timeout)\n try:\n mask_send(sock, payload, opcode=opcode)\n finally:\n try:\n sock.settimeout(None)\n except OSError:\n pass\n finally:\n lock.release()\n\ndef _buffered_reader(sock: socket.socket, initial: bytes = b\"\"):\n buf = bytearray(initial)\n def recv(n):\n while len(buf) < n:\n chunk = sock.recv(n - len(buf))\n if not chunk:\n raise ConnectionError(\"eof\")\n buf.extend(chunk)\n out = bytes(buf[:n])\n del buf[:n]\n return out\n return recv\n\ndef read_frame(recv):\n hdr = recv(2)\n fin = bool(hdr[0] & 0x80)\n opcode = hdr[0] & 0x0F\n masked = bool(hdr[1] & 0x80)\n ln = hdr[1] & 0x7F\n if ln == 126:\n ln = struct.unpack(\"!H\", recv(2))[0]\n elif ln == 127:\n ln = struct.unpack(\"!Q\", recv(8))[0]\n mask = recv(4) if masked else b\"\"\n payload = recv(ln)\n if masked:\n payload = bytes(b ^ mask[i % 4] for i, b in enumerate(payload))\n return fin, opcode, payload\n\ndef read_message(recv, send_pong):\n \"\"\"Next complete data message, assembling continuation frames.\n\n Control frames are handled inline (pong answered, close reported) so a\n fragmented message split around a ping still reassembles.\n \"\"\"\n opcode = None\n parts = []\n while True:\n fin, op, payload = read_frame(recv)\n if op == 0x8:\n return (\"close\", None, b\"\")\n if op == 0x9:\n send_pong(payload)\n continue\n if op == 0xA:\n continue\n if op in (0x1, 0x2):\n if opcode is not None:\n raise ConnectionError(\"data frame inside fragmented message\")\n if fin:\n return (\"data\", op, payload)\n opcode = op\n parts.append(payload)\n continue\n if op == 0x0:\n if opcode is None:\n raise ConnectionError(\"stray continuation frame\")\n parts.append(payload)\n if fin:\n return (\"data\", opcode, b\"\".join(parts))\n continue\n raise ConnectionError(f\"unsupported opcode {op}\")\n\ndef _is_daemon_response(payload: bytes) -> bool:\n \"\"\"True only for a JSON-RPC response (initialize answer), not a notification.\n\n The first-message timer must survive unrelated daemon notifications: only\n a message with an id and a result/error (and no method) proves the first\n RPC got a response.\n \"\"\"\n try:\n text = payload.decode(\"utf-8\").strip()\n if text.endswith(\"\\n\"):\n text = text.strip()\n msg = json.loads(text)\n except Exception:\n return False\n return isinstance(msg, dict) and msg.get(\"id\") is not None and \"method\" not in msg and (\"result\" in msg or \"error\" in msg)\n\ndef stdout_writer(sock: socket.socket, recv, send_pong, done: threading.Event, first_msg: threading.Event, output_forwarded: threading.Event) -> None:\n out = sys.stdout.buffer\n try:\n while True:\n kind, opcode, payload = read_message(recv, send_pong)\n if kind == \"close\":\n log(\"daemon WS close\")\n return\n if opcode == 0x1:\n # Desktop StdioConnection expects newline-delimited JSON text\n if not payload.endswith(b\"\\n\"):\n payload += b\"\\n\"\n out.write(payload)\n out.flush()\n else:\n out.write(payload)\n out.flush()\n # Any daemon payload on stdout commits the session: the fallback\n # must never run on already-used stdout even with pristine stdin.\n output_forwarded.set()\n if opcode == 0x1 and _is_daemon_response(payload):\n first_msg.set()\n except Exception as e:\n log(f\"stdout_writer exit: {e}\")\n finally:\n done.set()\n\ndef main() -> int:\n log(f\"start sock={SOCK}\")\n try:\n sock, pending = ws_connect(SOCK)\n except Exception as e:\n log(f\"connect failed: {e}\")\n sys.stderr.write(f\"codex-stdio-to-daemon-ws: {e}\\n\")\n return EXIT_PRE_STDIO\n log(\"connected\")\n connected_at = time.monotonic()\n send_lock = threading.Lock()\n done = threading.Event()\n first_msg = threading.Event()\n output_forwarded = threading.Event()\n\n def _send_timeout() -> float:\n # Same monotonic session budget bounds every send/lock: before the\n # first daemon response the remaining first-RPC budget applies, after\n # it each send gets the connect timeout. Expired budget fails fast so\n # a non-reading peer cannot wedge pong writes or timeout cleanup\n # before any stdin is consumed.\n if not first_msg.is_set():\n return max(0.0, (connected_at + FIRST_MESSAGE_TIMEOUT_S) - time.monotonic())\n return CONNECT_TIMEOUT_S\n\n def send_pong(payload: bytes) -> None:\n locked_send(sock, send_lock, payload, opcode=0xA, timeout=_send_timeout())\n recv = _buffered_reader(sock, pending)\n t = threading.Thread(target=stdout_writer, args=(sock, recv, send_pong, done, first_msg, output_forwarded), daemon=True)\n t.start()\n # Byte-precise commit point: the fallback may run only while zero stdin\n # bytes have left the pipe and no daemon payload was forwarded. Reads use\n # os.read into a userspace line buffer (buffered readline could readahead\n # past the commit point), and every byte read counts \u2014 even a blank line.\n # first_msg tracks the first daemon *response*, not any notification.\n stdin_bytes = 0\n pending_in = b\"\"\n clean_eof = False\n try:\n fd = sys.stdin.fileno()\n while not done.is_set():\n if not first_msg.is_set() and time.monotonic() - connected_at > FIRST_MESSAGE_TIMEOUT_S:\n log(\"first daemon message timed out\")\n break\n ready, _, _ = select.select([fd], [], [], 0.5)\n if not ready:\n continue\n chunk = os.read(fd, 65536)\n if not chunk:\n clean_eof = True\n break\n stdin_bytes += len(chunk)\n pending_in += chunk\n while b\"\\n\" in pending_in:\n line, _, pending_in = pending_in.partition(b\"\\n\")\n line = line.strip()\n if not line:\n continue\n try:\n locked_send(sock, send_lock, line, timeout=_send_timeout())\n except TimeoutError as e:\n log(f\"stdin forward timed out: {e}\")\n break\n except Exception as e:\n log(f\"stdin forward failed: {e}\")\n break\n else:\n continue\n break\n except Exception as e:\n log(f\"stdin loop exit: {e}\")\n # Flush any buffered EOF-tail data before the WS Close: sending Close\n # first would let the daemon discard the tail.\n if clean_eof:\n tail = pending_in.strip()\n if tail and not done.is_set():\n try:\n locked_send(sock, send_lock, tail, timeout=CONNECT_TIMEOUT_S)\n except Exception as e:\n log(f\"eof flush failed: {e}\")\n try:\n locked_send(sock, send_lock, b\"\", opcode=0x8, timeout=CONNECT_TIMEOUT_S)\n except Exception as e:\n log(f\"close send failed: {e}\")\n try:\n try:\n sock.shutdown(socket.SHUT_RDWR)\n except OSError:\n pass\n sock.close()\n except OSError:\n pass\n if clean_eof:\n log(\"exit\")\n return EXIT_SERVED\n # Reader died, first RPC timed out, or stdin broke: fail open only if no\n # stdin byte was ever consumed AND no daemon payload was forwarded. Used\n # stdout forces a mid-session exit so Desktop reconnects instead of\n # falling back to stock Codex on a corrupted stream.\n rc = EXIT_MID_SESSION if (stdin_bytes or output_forwarded.is_set()) else EXIT_PRE_STDIO\n log(f\"exit rc={rc}\")\n return rc\n\nif __name__ == \"__main__\":\n raise SystemExit(main())\n"; diff --git a/test/codex-bridge.test.js b/test/codex-bridge.test.js index 21b6954..49c6e71 100644 --- a/test/codex-bridge.test.js +++ b/test/codex-bridge.test.js @@ -308,6 +308,44 @@ test("server requests stay unanswered until our own turn starts", async () => { } }); +test("codex send ignores approvals naming another thread during our turn", async () => { + const fake = await fakeAppServer({ + ...baseHandlers, + "turn/start": (params, ok, err, send) => { + send({ jsonrpc: "2.0", id: "srv-foreign-thread", method: "item/commandExecution/requestApproval", params: { threadId: "other-thread", command: "rm -rf /" } }); + setTimeout(() => ok({ turn: { id: "turn-20", status: "inProgress", items: [] } }), 50); + }, + }); + try { + const env = { ...process.env, CODEX_HOME: fake.home }; + const outcome = await sendToCodexThread("t-1", "do it", { env }); + assert.equal(outcome.exitCode, 0); + assert.equal(outcome.turnId, "turn-20"); + assert.equal(fake.received.find((m) => m.id === "srv-foreign-thread"), undefined); + } finally { + await fake.close(); + } +}); + +test("codex send ignores approvals naming a Desktop turn during our turn", async () => { + const fake = await fakeAppServer({ + ...baseHandlers, + "turn/start": (params, ok, err, send) => { + send({ jsonrpc: "2.0", id: "srv-foreign-turn", method: "item/commandExecution/requestApproval", params: { threadId: params.threadId, turnId: "desktop-turn-1" } }); + setTimeout(() => ok({ turn: { id: "turn-21", status: "inProgress", items: [] } }), 50); + }, + }); + try { + const env = { ...process.env, CODEX_HOME: fake.home }; + const outcome = await sendToCodexThread("t-1", "do it", { env }); + assert.equal(outcome.exitCode, 0); + assert.equal(outcome.turnId, "turn-21"); + assert.equal(fake.received.find((m) => m.id === "srv-foreign-turn"), undefined); + } finally { + await fake.close(); + } +}); + test("codex send returns once the turn starts; later approval requests are the daemon's to route", async () => { const fake = await fakeAppServer({ ...baseHandlers, diff --git a/test/desktop-shim.test.js b/test/desktop-shim.test.js index caad38c..faaa4c7 100644 --- a/test/desktop-shim.test.js +++ b/test/desktop-shim.test.js @@ -683,3 +683,242 @@ test("login script keeps one CODEX_HOME for GUI env and daemon upkeep", () => { assert.match(renderedEnv, /CODEX_HOME_DIR="\/c"/); assert.match(renderedEnv, /launchctl setenv CODEX_HOME "\$CODEX_HOME_DIR"/); }); + +test("vendored bridge bounds every send/lock with the session budget", () => { + assert.match(BRIDGE_SOURCE, /lock\.acquire\(timeout/); + assert.match(BRIDGE_SOURCE, /send budget expired|send lock timed out/); + assert.match(BRIDGE_SOURCE, /sock\.settimeout\(timeout\)/); + assert.match(BRIDGE_SOURCE, /_send_timeout/); +}); + +test("bridge first-message timer tracks the first response, not any notification", { + skip: !canRunShellBridge && "needs bash + python3", +}, async () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-notif-")); + // An unrelated notification must not disable the first-RPC timer while + // initialize never gets a response. The notification is forwarded (stdout + // becomes used), so the timeout exit is mid-session, never a stock fallback. + const daemon = await fakeWsDaemon(dir, "notif.sock", { + prelude: wsServerFrame(0x1, Buffer.from('{"jsonrpc":"2.0","method":"daemon/event","params":{}}')), + quiet: true, + }); + try { + const started = Date.now(); + const out = await runBridge(writeBridge(dir), dir, { + CODEX_APP_SERVER_SOCK: daemon.socketPath, + CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT: "2", + }, "", { leaveStdinOpen: true }); + const elapsed = Date.now() - started; + assert.equal(out.status, 2, `notification must not satisfy first-RPC; got ${out.status} ${out.stderr}`); + assert.match(out.stdout, /daemon\/event/); + assert.ok(elapsed < 15000, `first-response deadline must fire, took ${elapsed}ms`); + } finally { + await closeDaemon(daemon); + } +}); + +test("bridge never falls back to stock after forwarding daemon output", { + skip: !canRunShellBridge && "needs bash + python3", +}, async () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-used-stdout-")); + // Astra repro: forwarded notification then close with pristine stdin must + // exit mid-session (Desktop reconnects), never pre-stdio stock fallback. + const daemon = await fakeWsDaemon(dir, "used.sock", { + prelude: wsServerFrame(0x1, Buffer.from('{"jsonrpc":"2.0","method":"daemon/event","params":{}}')), + closeAfterMs: 500, + }); + try { + const out = await runBridge(writeBridge(dir), dir, + { CODEX_APP_SERVER_SOCK: daemon.socketPath }, "", { leaveStdinOpen: true }); + assert.match(out.stdout, /daemon\/event/); + assert.equal(out.status, 2, `used stdout must force mid-session, got ${out.status}`); + } finally { + await closeDaemon(daemon); + } +}); + +test("bridge large stdin to a non-reading peer times out instead of hanging", { + skip: !canRunShellBridge && "needs bash + python3", +}, async () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-nonreading-")); + const socketPath = join(dir, "nonreading.sock"); + try { + rmSync(socketPath, { force: true }); + } catch {} + // Upgrade succeeds, then the peer never reads again: small pong/close + // frames would fit the kernel buffer, so fill it with a multi-MB stdin line. + // Unbounded sendall + lock would hang here past the test timeout. + const sockets = new Set(); + const server = createServer((socket) => { + sockets.add(socket); + socket.on("error", () => {}); + socket.on("close", () => sockets.delete(socket)); + let request = Buffer.alloc(0); + let upgraded = false; + socket.on("data", (chunk) => { + if (upgraded) return; // non-reading: drop everything after the upgrade + request = Buffer.concat([request, chunk]); + if (!request.includes("\r\n\r\n")) return; + const keyLine = request.toString("latin1").split("\r\n") + .find((line) => line.toLowerCase().startsWith("sec-websocket-key:")); + const key = keyLine.split(":")[1].trim(); + const accept = createHash("sha1").update(key + WS_GUID).digest().toString("base64"); + upgraded = true; + socket.write(Buffer.from( + `HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: ${accept}\r\n\r\n`, + "latin1", + )); + socket.pause(); // never consume stat frame reads; bridge sends must time out + }); + }); + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(socketPath, resolve); + }); + try { + const bigLine = `{"jsonrpc":"2.0","id":1,"method":"initialize","params":{"pad":"${"x".repeat(4 * 1024 * 1024)}"}}\n`; + const started = Date.now(); + const out = await runBridge(writeBridge(dir), dir, { + CODEX_APP_SERVER_SOCK: socketPath, + CODEX_BRIDGE_CONNECT_TIMEOUT: "2", + CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT: "2", + }, bigLine, { leaveStdinOpen: true }); + const elapsed = Date.now() - started; + assert.equal(out.status, 2, `bounded send must exit mid-session, got ${out.status}`); + assert.ok(elapsed < 15000, `bounded send must not hang, took ${elapsed}ms`); + } finally { + for (const s of sockets) try { s.destroy(); } catch {} + await new Promise((resolve) => server.close(resolve)); + } +}); + +test("bridge rejects HTTP 200 even with a valid accept hash", { + skip: !canRunShellBridge && "needs bash + python3", +}, async () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-200-")); + const socketPath = join(dir, "http200.sock"); + try { + rmSync(socketPath, { force: true }); + } catch {} + const sockets200 = new Set(); + const server = createServer((socket) => { + sockets200.add(socket); + socket.on("error", () => {}); + socket.on("close", () => sockets200.delete(socket)); + let request = Buffer.alloc(0); + socket.on("data", (chunk) => { + request = Buffer.concat([request, chunk]); + if (!request.includes("\r\n\r\n")) return; + const keyLine = request.toString("latin1").split("\r\n") + .find((line) => line.toLowerCase().startsWith("sec-websocket-key:")); + const key = keyLine.split(":")[1].trim(); + const accept = createHash("sha1").update(key + WS_GUID).digest().toString("base64"); + // Valid accept, wrong status: must still fail (require 101 + accept). + socket.write(Buffer.from( + `HTTP/1.1 200 OK\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: ${accept}\r\n\r\n`, + "latin1", + )); + }); + }); + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(socketPath, resolve); + }); + try { + const out = await runBridge(writeBridge(dir), dir, + { CODEX_APP_SERVER_SOCK: socketPath }, "", { leaveStdinOpen: false }); + assert.equal(out.status, 1, `HTTP 200 must fail pre-stdio, got ${out.status}`); + } finally { + for (const s of sockets200) try { s.destroy(); } catch {} + await new Promise((resolve) => server.close(resolve)); + } +}); + +test("bridge flushes buffered EOF-tail data before the WS Close", { + skip: !canRunShellBridge && "needs bash + python3", +}, async () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-eoftail-")); + const socketPath = join(dir, "eoftail.sock"); + try { + rmSync(socketPath, { force: true }); + } catch {} + // Closes the transport as soon as a WS Close frame arrives: a tail sent + // after Close would be lost, so order proves the flush happens first. + const state = { frames: [] }; + const server = createServer((socket) => { + state.socket = socket; + socket.on("error", () => {}); + let request = Buffer.alloc(0); + let upgraded = false; + let buf = Buffer.alloc(0); + socket.on("data", (chunk) => { + if (!upgraded) { + request = Buffer.concat([request, chunk]); + if (!request.includes("\r\n\r\n")) return; + const keyLine = request.toString("latin1").split("\r\n") + .find((line) => line.toLowerCase().startsWith("sec-websocket-key:")); + const key = keyLine.split(":")[1].trim(); + const accept = createHash("sha1").update(key + WS_GUID).digest().toString("base64"); + upgraded = true; + socket.write(Buffer.from( + `HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: ${accept}\r\n\r\n`, + "latin1", + )); + const rest = request.subarray(request.indexOf("\r\n\r\n") + 4); + if (rest.length) buf = Buffer.concat([buf, rest]); + return; + } + buf = Buffer.concat([buf, chunk]); + for (;;) { + if (buf.length < 2) return; + const opcode = buf[0] & 0x0f; + const masked = (buf[1] & 0x80) !== 0; + let len = buf[1] & 0x7f; + let off = 2; + if (len === 126) { + if (buf.length < 4) return; + len = buf.readUInt16BE(2); + off = 4; + } + const maskLen = masked ? 4 : 0; + if (buf.length < off + maskLen + len) return; + let payload = buf.subarray(off + maskLen, off + maskLen + len); + if (masked) { + const mask = buf.subarray(off, off + 4); + payload = Buffer.from(payload.map((b, i) => b ^ mask[i % 4])); + } + state.frames.push({ opcode, payload: Buffer.from(payload) }); + buf = buf.subarray(off + maskLen + len); + if (opcode === 0x8) { + try { socket.destroy(); } catch {} + return; + } + } + }); + }); + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(socketPath, resolve); + }); + try { + // Partial line with no newline at EOF: only the tail flush can deliver it. + const out = await runBridge(writeBridge(dir), dir, + { CODEX_APP_SERVER_SOCK: socketPath }, '{"id":1}', { leaveStdinOpen: false }); + assert.equal(out.status, 0, `clean EOF serves, got ${out.status} ${out.stderr}`); + const dataIdx = state.frames.findIndex((f) => f.opcode === 0x1 && f.payload.toString().includes('"id":1')); + const closeIdx = state.frames.findIndex((f) => f.opcode === 0x8); + assert.ok(dataIdx !== -1, `EOF tail must reach the daemon, saw ${JSON.stringify(state.frames.map((f) => f.opcode))}`); + assert.ok(closeIdx !== -1, "WS Close must still be sent"); + assert.ok(dataIdx < closeIdx, "tail data must precede the WS Close"); + } finally { + try { state.socket?.destroy(); } catch {} + await new Promise((resolve) => server.close(resolve)); + } +}); + +test("vendored bridge pins response-tracked first-RPC, output commit, and strict 101", () => { + assert.match(BRIDGE_SOURCE, /_is_daemon_response/); + assert.match(BRIDGE_SOURCE, /output_forwarded/); + assert.match(BRIDGE_SOURCE, /startswith\(b"HTTP\/1\.1 101"\)/); + assert.match(BRIDGE_SOURCE, /tail.*before.*WS Close|before the WS Close/s); +}); From f8fa0a6fc5c19490116d445c485204e3fe130644 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Tue, 15 Sep 2026 20:59:36 +0000 Subject: [PATCH 2/2] Tighten first-RPC to matching response/error only Co-authored-by: Zack Jackson --- README.md | 2 +- src/core/desktop-shim-bridge.js | 22 +++++++++++++--------- test/desktop-shim.test.js | 28 ++++++++++++++++++++++++++++ 3 files changed, 42 insertions(+), 10 deletions(-) diff --git a/README.md b/README.md index 98589d3..524eb25 100644 --- a/README.md +++ b/README.md @@ -87,7 +87,7 @@ gbot codex desktop-shim uninstall # removes wrapper/bridge/LaunchAgent; Desktop Install copies the scripts out of the package into durable `~/.codex/bin/` paths, so uninstalling the npm package never leaves Desktop pointing into a deleted checkout. The wrapper exports install-time `CODEX_HOME` and derives the socket from it at runtime (`CODEX_APP_SERVER_SOCK`), so custom `CODEX_HOME` layouts work; `gbot` itself resolves the socket the same way (`CODEX_APP_SERVER_SOCK` wins, else `CODEX_HOME`), and `status` reports the effective path and which variable won. -Fail-open runs only before any stdin byte is consumed and no daemon payload was forwarded: the wrapper preflights the daemon socket (3 s deadline) and runs the bridge as a child — never `exec` — falling through to the real standalone `codex` for every non-`app-server` spawn, every preflight failure, and every pre-session bridge failure (exit 1), always on pristine stdio+stdout. The bridge's connect + WebSocket-upgrade handshake runs under one absolute monotonic deadline (`CODEX_BRIDGE_CONNECT_TIMEOUT`, default 10 s) requiring `HTTP/1.1 101` plus a matching `Sec-WebSocket-Accept` (a `200` fails even with a valid hash), and a first-response deadline (`CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT`, default 30 s) bounds connect-to-first-daemon-response covering the first RPC — an unrelated notification never satisfies it — preflight (3 s) + upgrade (10 s) + first RPC (30 s) stays far below a daemon-lock hang, and none of these stages reads stdin. Every send/lock shares the same monotonic session budget (remaining first-RPC budget, else the connect timeout), so a non-reading peer cannot wedge pong writes or timeout cleanup; timeouts close and exit. The commit point is byte-precise on input OR output: the bridge counts every stdin byte read (buffered readahead included, even an unforwarded blank line) and marks stdout used on any forwarded daemon payload, so the fallback only ever runs on truly pristine stdio+stdout. Buffered EOF-tail data flushes before the WS Close. A mid-session bridge failure (exit 2+) makes the wrapper exit promptly so Desktop reconnects; the fallback never runs on half-consumed stdin or already-used stdout. The wrapper never starts the daemon on the spawn path (a wedged daemon lock must not block Desktop) — upkeep belongs to the LaunchAgent login script and install, which share one `CODEX_HOME` for the GUI domain and the daemon they start; if the socket is absent, Desktop simply runs stock Codex until the daemon is started. +Fail-open runs only before any stdin byte is consumed and no daemon payload was forwarded: the wrapper preflights the daemon socket (3 s deadline) and runs the bridge as a child — never `exec` — falling through to the real standalone `codex` for every non-`app-server` spawn, every preflight failure, and every pre-session bridge failure (exit 1), always on pristine stdio+stdout. The bridge's connect + WebSocket-upgrade handshake runs under one absolute monotonic deadline (`CODEX_BRIDGE_CONNECT_TIMEOUT`, default 10 s) requiring `HTTP/1.1 101` plus a matching `Sec-WebSocket-Accept` (a `200` fails even with a valid hash), and a first-response deadline (`CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT`, default 30 s) bounds connect-to-matching-daemon-response covering the first RPC — an unrelated notification or foreign response never satisfies it, only the matching response/error (or timeout) clears it — preflight (3 s) + upgrade (10 s) + first RPC (30 s) stays far below a daemon-lock hang, and none of these stages reads stdin. Every send/lock shares the same monotonic session budget (remaining first-RPC budget, else the connect timeout), so a non-reading peer cannot wedge pong writes or timeout cleanup; timeouts close and exit. The commit point is byte-precise on input OR output: the bridge counts every stdin byte read (buffered readahead included, even an unforwarded blank line) and marks stdout used on any forwarded daemon payload, so the fallback only ever runs on truly pristine stdio+stdout. Buffered EOF-tail data flushes before the WS Close. A mid-session bridge failure (exit 2+) makes the wrapper exit promptly so Desktop reconnects; the fallback never runs on half-consumed stdin or already-used stdout. The wrapper never starts the daemon on the spawn path (a wedged daemon lock must not block Desktop) — upkeep belongs to the LaunchAgent login script and install, which share one `CODEX_HOME` for the GUI domain and the daemon they start; if the socket is absent, Desktop simply runs stock Codex until the daemon is started. Residual risks, stated honestly: a concurrent foreign approval carrying no thread/turn ids inside `gbot`'s own `turn/start` window is indistinguishable from `gbot`'s and will still be refused — don't send to threads Desktop is actively driving. A daemon that speaks framing-valid but semantically unexpected JSON-RPC (unknown methods, id-less responses) is treated as transport; pins are to app-server schema 0.154.0. The bridge trusts the local control socket; a malicious local daemon could hold the session up to the stated budgets, not past them. diff --git a/src/core/desktop-shim-bridge.js b/src/core/desktop-shim-bridge.js index a27f8dd..59b2d0f 100644 --- a/src/core/desktop-shim-bridge.js +++ b/src/core/desktop-shim-bridge.js @@ -9,20 +9,24 @@ // (2) read_message assembles continuation frames; control frames answer inline; // every socket send/lock is bounded by the same monotonic session budget // (remaining first-RPC budget, else the connect timeout) so a non-reading peer -// cannot wedge pong writes or timeout cleanup; timeouts close and exit. +// cannot wedge pong writes or timeout cleanup; timeouts close and exit. The +// wrapper never starts the daemon on the spawn path, so no daemon lock can +// block pre-stdin Desktop. // (3) A first-response deadline (CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT, default -// 30s) bounds connect-to-first-daemon-response covering the first RPC; an -// unrelated notification never disables it. The preflight (3s) + upgrade (10s) -// + first RPC (30s) budget stays far below a daemon-lock hang. Exits are -// 0 served, 1 pre-stdio (wrapper falls back to real Codex on pristine +// 30s) bounds connect-to-matching-daemon-response covering the first RPC; +// unrelated notifications and foreign responses never satisfy it, only the +// matching response/error (or timeout) clears it. The preflight (3s) + +// upgrade (10s) + first RPC (30s) budget stays far below a daemon-lock hang. +// Exits are 0 served, 1 pre-stdio (wrapper falls back to real Codex on pristine // stdio+stdout), 2 mid-session (wrapper exits so Desktop reconnects). // (4) Byte-precise commit on input OR output: os.read into a userspace line // buffer counts every stdin byte (even an unforwarded blank line), and any // forwarded daemon payload marks stdout used, so the fallback only ever runs -// on truly pristine stdio+stdout. Buffered EOF-tail data flushes before the -// WS Close. The reader signals a done event so WS close wakes the select loop; -// no path leaves a thread hung. +// on truly pristine stdio+stdout and never spawns a second Codex after used +// stdout. Buffered EOF-tail data flushes before the WS Close, which drains +// via shutdown+close. The reader signals a done event so WS close wakes the +// select loop; no path leaves a thread hung. // Desktop stdio hangs with stock `codex app-server proxy`, so this custom // bridge is required. Install writes it to the CODEX_HOME bin dir. Never reads // Desktop private pipes (CODEX_APP_TOOLS_PIPE_PATH, /tmp/codex-browser-use/). -export const BRIDGE_SOURCE = "#!/usr/bin/env python3\n\"\"\"Stdio JSONL (Desktop) <-> WebSocket-over-unix (managed app-server daemon).\nReversible: remove CODEX_CLI_PATH wrapper. Does not read Desktop private pipes.\n\"\"\"\nfrom __future__ import annotations\nimport base64, hashlib, json, os, select, socket, struct, sys, threading, time\n\nSOCK = os.environ.get(\n \"CODEX_APP_SERVER_SOCK\",\n os.path.expanduser(\"~/.codex/app-server-control/app-server-control.sock\"),\n)\nLOG = os.environ.get(\"CODEX_STDIO_BRIDGE_LOG\", os.path.expanduser(\"~/.codex/bin/codex-stdio-to-daemon-ws.log\"))\nWS_GUID = \"258EAFA5-E914-47DA-95CA-C5AB0DC85B11\"\n\n# Wrapper contract: 0 served the session, 1 failed before any stdin was\n# consumed and no daemon payload was forwarded (wrapper falls back to real\n# Codex on pristine stdio), 2 failed mid-session (wrapper exits promptly so\n# Desktop reconnects; the fallback must never run on half-consumed stdin or\n# already-used stdout).\nEXIT_SERVED = 0\nEXIT_PRE_STDIO = 1\nEXIT_MID_SESSION = 2\n\nCONNECT_TIMEOUT_S = float(os.environ.get(\"CODEX_BRIDGE_CONNECT_TIMEOUT\", \"10\"))\n# Bound from connect to the first daemon response: a daemon that completes the\n# upgrade but never answers the first RPC must not hold Desktop past this.\nFIRST_MESSAGE_TIMEOUT_S = float(os.environ.get(\"CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT\", \"30\"))\n\ndef log(msg: str) -> None:\n try:\n with open(LOG, \"a\") as f:\n f.write(time.strftime(\"%Y-%m-%dT%H:%M:%SZ\", time.gmtime()) + \" \" + msg + \"\\n\")\n except OSError:\n pass\n\ndef ws_connect(path: str, timeout: float = CONNECT_TIMEOUT_S):\n \"\"\"Connect + HTTP upgrade under one absolute monotonic deadline.\n\n Returns (socket, trailing) where trailing holds bytes coalesced after the\n upgrade headers (TCP may deliver the 101 and the first WS frame together;\n discarding them would lose the first message). Requires HTTP/1.1 101 plus\n a matching Sec-WebSocket-Accept so a 200 (or any non-101) fails fast even\n when the accept hash happens to be valid.\n \"\"\"\n deadline = time.monotonic() + timeout\n s = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)\n s.settimeout(timeout)\n try:\n s.connect(path)\n key = base64.b64encode(os.urandom(16)).decode()\n expected = base64.b64encode(hashlib.sha1((key + WS_GUID).encode()).digest()).decode()\n req = (\n f\"GET /rpc HTTP/1.1\\r\\nHost: localhost\\r\\nUpgrade: websocket\\r\\n\"\n f\"Connection: Upgrade\\r\\nSec-WebSocket-Key: {key}\\r\\nSec-WebSocket-Version: 13\\r\\n\\r\\n\"\n ).encode()\n s.sendall(req)\n buf = b\"\"\n while True:\n if b\"\\r\\n\\r\\n\" in buf:\n break\n remaining = deadline - time.monotonic()\n if remaining <= 0:\n raise TimeoutError(\"WebSocket upgrade timed out\")\n if len(buf) > 65536:\n raise ConnectionError(\"WS upgrade headers too large\")\n s.settimeout(remaining)\n chunk = s.recv(4096)\n if not chunk:\n raise ConnectionError(\"daemon closed during WS upgrade\")\n buf += chunk\n head, _, trailing = buf.partition(b\"\\r\\n\\r\\n\")\n lines = head.split(b\"\\r\\n\")\n if not lines or not lines[0].startswith(b\"HTTP/1.1 101\"):\n raise ConnectionError(f\"WS upgrade failed: {head[:200]!r}\")\n accept = None\n for line in lines[1:]:\n name, sep, value = line.partition(b\":\")\n if sep and name.strip().lower() == b\"sec-websocket-accept\":\n accept = value.strip().decode(\"latin-1\")\n if accept != expected:\n raise ConnectionError(\"WS upgrade accept mismatch\")\n except Exception:\n try:\n s.close()\n except OSError:\n pass\n raise\n s.settimeout(None)\n return s, trailing\n\ndef mask_send(sock: socket.socket, payload: bytes, opcode: int = 1) -> None:\n mask = os.urandom(4)\n ln = len(payload)\n if ln < 126:\n hdr = bytes([0x80 | opcode, 0x80 | ln])\n elif ln < 65536:\n hdr = bytes([0x80 | opcode, 0x80 | 126]) + struct.pack(\"!H\", ln)\n else:\n hdr = bytes([0x80 | opcode, 0x80 | 127]) + struct.pack(\"!Q\", ln)\n sock.sendall(hdr + mask + bytes(b ^ mask[i % 4] for i, b in enumerate(payload)))\n\ndef locked_send(sock: socket.socket, lock: threading.Lock, payload: bytes, opcode: int = 1, timeout: float = CONNECT_TIMEOUT_S) -> None:\n \"\"\"Serialized send bounded by a monotonic budget.\n\n Both the lock wait and the sendall wait share `timeout` seconds. A\n non-reading peer must never wedge pong writes or timeout cleanup: on\n expiry raise TimeoutError so the caller closes and exits without hanging.\n \"\"\"\n if timeout is None:\n timeout = CONNECT_TIMEOUT_S\n if timeout <= 0:\n raise TimeoutError(\"send budget expired\")\n if not lock.acquire(timeout=timeout):\n raise TimeoutError(\"send lock timed out\")\n try:\n sock.settimeout(timeout)\n try:\n mask_send(sock, payload, opcode=opcode)\n finally:\n try:\n sock.settimeout(None)\n except OSError:\n pass\n finally:\n lock.release()\n\ndef _buffered_reader(sock: socket.socket, initial: bytes = b\"\"):\n buf = bytearray(initial)\n def recv(n):\n while len(buf) < n:\n chunk = sock.recv(n - len(buf))\n if not chunk:\n raise ConnectionError(\"eof\")\n buf.extend(chunk)\n out = bytes(buf[:n])\n del buf[:n]\n return out\n return recv\n\ndef read_frame(recv):\n hdr = recv(2)\n fin = bool(hdr[0] & 0x80)\n opcode = hdr[0] & 0x0F\n masked = bool(hdr[1] & 0x80)\n ln = hdr[1] & 0x7F\n if ln == 126:\n ln = struct.unpack(\"!H\", recv(2))[0]\n elif ln == 127:\n ln = struct.unpack(\"!Q\", recv(8))[0]\n mask = recv(4) if masked else b\"\"\n payload = recv(ln)\n if masked:\n payload = bytes(b ^ mask[i % 4] for i, b in enumerate(payload))\n return fin, opcode, payload\n\ndef read_message(recv, send_pong):\n \"\"\"Next complete data message, assembling continuation frames.\n\n Control frames are handled inline (pong answered, close reported) so a\n fragmented message split around a ping still reassembles.\n \"\"\"\n opcode = None\n parts = []\n while True:\n fin, op, payload = read_frame(recv)\n if op == 0x8:\n return (\"close\", None, b\"\")\n if op == 0x9:\n send_pong(payload)\n continue\n if op == 0xA:\n continue\n if op in (0x1, 0x2):\n if opcode is not None:\n raise ConnectionError(\"data frame inside fragmented message\")\n if fin:\n return (\"data\", op, payload)\n opcode = op\n parts.append(payload)\n continue\n if op == 0x0:\n if opcode is None:\n raise ConnectionError(\"stray continuation frame\")\n parts.append(payload)\n if fin:\n return (\"data\", opcode, b\"\".join(parts))\n continue\n raise ConnectionError(f\"unsupported opcode {op}\")\n\ndef _is_daemon_response(payload: bytes) -> bool:\n \"\"\"True only for a JSON-RPC response (initialize answer), not a notification.\n\n The first-message timer must survive unrelated daemon notifications: only\n a message with an id and a result/error (and no method) proves the first\n RPC got a response.\n \"\"\"\n try:\n text = payload.decode(\"utf-8\").strip()\n if text.endswith(\"\\n\"):\n text = text.strip()\n msg = json.loads(text)\n except Exception:\n return False\n return isinstance(msg, dict) and msg.get(\"id\") is not None and \"method\" not in msg and (\"result\" in msg or \"error\" in msg)\n\ndef stdout_writer(sock: socket.socket, recv, send_pong, done: threading.Event, first_msg: threading.Event, output_forwarded: threading.Event) -> None:\n out = sys.stdout.buffer\n try:\n while True:\n kind, opcode, payload = read_message(recv, send_pong)\n if kind == \"close\":\n log(\"daemon WS close\")\n return\n if opcode == 0x1:\n # Desktop StdioConnection expects newline-delimited JSON text\n if not payload.endswith(b\"\\n\"):\n payload += b\"\\n\"\n out.write(payload)\n out.flush()\n else:\n out.write(payload)\n out.flush()\n # Any daemon payload on stdout commits the session: the fallback\n # must never run on already-used stdout even with pristine stdin.\n output_forwarded.set()\n if opcode == 0x1 and _is_daemon_response(payload):\n first_msg.set()\n except Exception as e:\n log(f\"stdout_writer exit: {e}\")\n finally:\n done.set()\n\ndef main() -> int:\n log(f\"start sock={SOCK}\")\n try:\n sock, pending = ws_connect(SOCK)\n except Exception as e:\n log(f\"connect failed: {e}\")\n sys.stderr.write(f\"codex-stdio-to-daemon-ws: {e}\\n\")\n return EXIT_PRE_STDIO\n log(\"connected\")\n connected_at = time.monotonic()\n send_lock = threading.Lock()\n done = threading.Event()\n first_msg = threading.Event()\n output_forwarded = threading.Event()\n\n def _send_timeout() -> float:\n # Same monotonic session budget bounds every send/lock: before the\n # first daemon response the remaining first-RPC budget applies, after\n # it each send gets the connect timeout. Expired budget fails fast so\n # a non-reading peer cannot wedge pong writes or timeout cleanup\n # before any stdin is consumed.\n if not first_msg.is_set():\n return max(0.0, (connected_at + FIRST_MESSAGE_TIMEOUT_S) - time.monotonic())\n return CONNECT_TIMEOUT_S\n\n def send_pong(payload: bytes) -> None:\n locked_send(sock, send_lock, payload, opcode=0xA, timeout=_send_timeout())\n recv = _buffered_reader(sock, pending)\n t = threading.Thread(target=stdout_writer, args=(sock, recv, send_pong, done, first_msg, output_forwarded), daemon=True)\n t.start()\n # Byte-precise commit point: the fallback may run only while zero stdin\n # bytes have left the pipe and no daemon payload was forwarded. Reads use\n # os.read into a userspace line buffer (buffered readline could readahead\n # past the commit point), and every byte read counts \u2014 even a blank line.\n # first_msg tracks the first daemon *response*, not any notification.\n stdin_bytes = 0\n pending_in = b\"\"\n clean_eof = False\n try:\n fd = sys.stdin.fileno()\n while not done.is_set():\n if not first_msg.is_set() and time.monotonic() - connected_at > FIRST_MESSAGE_TIMEOUT_S:\n log(\"first daemon message timed out\")\n break\n ready, _, _ = select.select([fd], [], [], 0.5)\n if not ready:\n continue\n chunk = os.read(fd, 65536)\n if not chunk:\n clean_eof = True\n break\n stdin_bytes += len(chunk)\n pending_in += chunk\n while b\"\\n\" in pending_in:\n line, _, pending_in = pending_in.partition(b\"\\n\")\n line = line.strip()\n if not line:\n continue\n try:\n locked_send(sock, send_lock, line, timeout=_send_timeout())\n except TimeoutError as e:\n log(f\"stdin forward timed out: {e}\")\n break\n except Exception as e:\n log(f\"stdin forward failed: {e}\")\n break\n else:\n continue\n break\n except Exception as e:\n log(f\"stdin loop exit: {e}\")\n # Flush any buffered EOF-tail data before the WS Close: sending Close\n # first would let the daemon discard the tail.\n if clean_eof:\n tail = pending_in.strip()\n if tail and not done.is_set():\n try:\n locked_send(sock, send_lock, tail, timeout=CONNECT_TIMEOUT_S)\n except Exception as e:\n log(f\"eof flush failed: {e}\")\n try:\n locked_send(sock, send_lock, b\"\", opcode=0x8, timeout=CONNECT_TIMEOUT_S)\n except Exception as e:\n log(f\"close send failed: {e}\")\n try:\n try:\n sock.shutdown(socket.SHUT_RDWR)\n except OSError:\n pass\n sock.close()\n except OSError:\n pass\n if clean_eof:\n log(\"exit\")\n return EXIT_SERVED\n # Reader died, first RPC timed out, or stdin broke: fail open only if no\n # stdin byte was ever consumed AND no daemon payload was forwarded. Used\n # stdout forces a mid-session exit so Desktop reconnects instead of\n # falling back to stock Codex on a corrupted stream.\n rc = EXIT_MID_SESSION if (stdin_bytes or output_forwarded.is_set()) else EXIT_PRE_STDIO\n log(f\"exit rc={rc}\")\n return rc\n\nif __name__ == \"__main__\":\n raise SystemExit(main())\n"; +export const BRIDGE_SOURCE = "#!/usr/bin/env python3\n\"\"\"Stdio JSONL (Desktop) <-> WebSocket-over-unix (managed app-server daemon).\nReversible: remove CODEX_CLI_PATH wrapper. Does not read Desktop private pipes.\n\"\"\"\nfrom __future__ import annotations\nimport base64, hashlib, json, os, select, socket, struct, sys, threading, time\n\nSOCK = os.environ.get(\n \"CODEX_APP_SERVER_SOCK\",\n os.path.expanduser(\"~/.codex/app-server-control/app-server-control.sock\"),\n)\nLOG = os.environ.get(\"CODEX_STDIO_BRIDGE_LOG\", os.path.expanduser(\"~/.codex/bin/codex-stdio-to-daemon-ws.log\"))\nWS_GUID = \"258EAFA5-E914-47DA-95CA-C5AB0DC85B11\"\n\n# Wrapper contract: 0 served the session, 1 failed before any stdin was\n# consumed and no daemon payload was forwarded (wrapper falls back to real\n# Codex on pristine stdio), 2 failed mid-session (wrapper exits promptly so\n# Desktop reconnects; the fallback must never run on half-consumed stdin or\n# already-used stdout).\nEXIT_SERVED = 0\nEXIT_PRE_STDIO = 1\nEXIT_MID_SESSION = 2\n\nCONNECT_TIMEOUT_S = float(os.environ.get(\"CODEX_BRIDGE_CONNECT_TIMEOUT\", \"10\"))\n# Bound from connect to the first daemon response: a daemon that completes the\n# upgrade but never answers the first RPC must not hold Desktop past this.\nFIRST_MESSAGE_TIMEOUT_S = float(os.environ.get(\"CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT\", \"30\"))\n\ndef log(msg: str) -> None:\n try:\n with open(LOG, \"a\") as f:\n f.write(time.strftime(\"%Y-%m-%dT%H:%M:%SZ\", time.gmtime()) + \" \" + msg + \"\\n\")\n except OSError:\n pass\n\ndef ws_connect(path: str, timeout: float = CONNECT_TIMEOUT_S):\n \"\"\"Connect + HTTP upgrade under one absolute monotonic deadline.\n\n Returns (socket, trailing) where trailing holds bytes coalesced after the\n upgrade headers (TCP may deliver the 101 and the first WS frame together;\n discarding them would lose the first message). Requires HTTP/1.1 101 plus\n a matching Sec-WebSocket-Accept so a 200 (or any non-101) fails fast even\n when the accept hash happens to be valid.\n \"\"\"\n deadline = time.monotonic() + timeout\n s = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)\n s.settimeout(timeout)\n try:\n s.connect(path)\n key = base64.b64encode(os.urandom(16)).decode()\n expected = base64.b64encode(hashlib.sha1((key + WS_GUID).encode()).digest()).decode()\n req = (\n f\"GET /rpc HTTP/1.1\\r\\nHost: localhost\\r\\nUpgrade: websocket\\r\\n\"\n f\"Connection: Upgrade\\r\\nSec-WebSocket-Key: {key}\\r\\nSec-WebSocket-Version: 13\\r\\n\\r\\n\"\n ).encode()\n s.sendall(req)\n buf = b\"\"\n while True:\n if b\"\\r\\n\\r\\n\" in buf:\n break\n remaining = deadline - time.monotonic()\n if remaining <= 0:\n raise TimeoutError(\"WebSocket upgrade timed out\")\n if len(buf) > 65536:\n raise ConnectionError(\"WS upgrade headers too large\")\n s.settimeout(remaining)\n chunk = s.recv(4096)\n if not chunk:\n raise ConnectionError(\"daemon closed during WS upgrade\")\n buf += chunk\n head, _, trailing = buf.partition(b\"\\r\\n\\r\\n\")\n lines = head.split(b\"\\r\\n\")\n if not lines or not lines[0].startswith(b\"HTTP/1.1 101\"):\n raise ConnectionError(f\"WS upgrade failed: {head[:200]!r}\")\n accept = None\n for line in lines[1:]:\n name, sep, value = line.partition(b\":\")\n if sep and name.strip().lower() == b\"sec-websocket-accept\":\n accept = value.strip().decode(\"latin-1\")\n if accept != expected:\n raise ConnectionError(\"WS upgrade accept mismatch\")\n except Exception:\n try:\n s.close()\n except OSError:\n pass\n raise\n s.settimeout(None)\n return s, trailing\n\ndef mask_send(sock: socket.socket, payload: bytes, opcode: int = 1) -> None:\n mask = os.urandom(4)\n ln = len(payload)\n if ln < 126:\n hdr = bytes([0x80 | opcode, 0x80 | ln])\n elif ln < 65536:\n hdr = bytes([0x80 | opcode, 0x80 | 126]) + struct.pack(\"!H\", ln)\n else:\n hdr = bytes([0x80 | opcode, 0x80 | 127]) + struct.pack(\"!Q\", ln)\n sock.sendall(hdr + mask + bytes(b ^ mask[i % 4] for i, b in enumerate(payload)))\n\ndef locked_send(sock: socket.socket, lock: threading.Lock, payload: bytes, opcode: int = 1, timeout: float = CONNECT_TIMEOUT_S) -> None:\n \"\"\"Serialized send bounded by a monotonic budget.\n\n Both the lock wait and the sendall wait share `timeout` seconds. A\n non-reading peer must never wedge pong writes or timeout cleanup: on\n expiry raise TimeoutError so the caller closes and exits without hanging.\n \"\"\"\n if timeout is None:\n timeout = CONNECT_TIMEOUT_S\n if timeout <= 0:\n raise TimeoutError(\"send budget expired\")\n if not lock.acquire(timeout=timeout):\n raise TimeoutError(\"send lock timed out\")\n try:\n sock.settimeout(timeout)\n try:\n mask_send(sock, payload, opcode=opcode)\n finally:\n try:\n sock.settimeout(None)\n except OSError:\n pass\n finally:\n lock.release()\n\ndef _buffered_reader(sock: socket.socket, initial: bytes = b\"\"):\n buf = bytearray(initial)\n def recv(n):\n while len(buf) < n:\n chunk = sock.recv(n - len(buf))\n if not chunk:\n raise ConnectionError(\"eof\")\n buf.extend(chunk)\n out = bytes(buf[:n])\n del buf[:n]\n return out\n return recv\n\ndef read_frame(recv):\n hdr = recv(2)\n fin = bool(hdr[0] & 0x80)\n opcode = hdr[0] & 0x0F\n masked = bool(hdr[1] & 0x80)\n ln = hdr[1] & 0x7F\n if ln == 126:\n ln = struct.unpack(\"!H\", recv(2))[0]\n elif ln == 127:\n ln = struct.unpack(\"!Q\", recv(8))[0]\n mask = recv(4) if masked else b\"\"\n payload = recv(ln)\n if masked:\n payload = bytes(b ^ mask[i % 4] for i, b in enumerate(payload))\n return fin, opcode, payload\n\ndef read_message(recv, send_pong):\n \"\"\"Next complete data message, assembling continuation frames.\n\n Control frames are handled inline (pong answered, close reported) so a\n fragmented message split around a ping still reassembles.\n \"\"\"\n opcode = None\n parts = []\n while True:\n fin, op, payload = read_frame(recv)\n if op == 0x8:\n return (\"close\", None, b\"\")\n if op == 0x9:\n send_pong(payload)\n continue\n if op == 0xA:\n continue\n if op in (0x1, 0x2):\n if opcode is not None:\n raise ConnectionError(\"data frame inside fragmented message\")\n if fin:\n return (\"data\", op, payload)\n opcode = op\n parts.append(payload)\n continue\n if op == 0x0:\n if opcode is None:\n raise ConnectionError(\"stray continuation frame\")\n parts.append(payload)\n if fin:\n return (\"data\", opcode, b\"\".join(parts))\n continue\n raise ConnectionError(f\"unsupported opcode {op}\")\n\ndef _is_daemon_response(payload: bytes, expect_id=None, known_ids=None) -> bool:\n \"\"\"True only for the matching JSON-RPC response/error, not a notification.\n\n The first-RPC timer must survive unrelated daemon notifications and foreign\n responses: only the response/error whose id matches the client's initialize\n (or, before it is known, any id the client already sent) proves the first\n RPC got an answer. Timeouts clear it via the deadline, never via data.\n \"\"\"\n try:\n text = payload.decode(\"utf-8\").strip()\n if text.endswith(\"\\n\"):\n text = text.strip()\n msg = json.loads(text)\n except Exception:\n return False\n if not isinstance(msg, dict):\n return False\n rid = msg.get(\"id\")\n if rid is None or \"method\" in msg or (\"result\" not in msg and \"error\" not in msg):\n return False\n if expect_id is not None:\n return rid == expect_id\n if known_ids:\n return rid in known_ids\n return False\n\ndef stdout_writer(sock: socket.socket, recv, send_pong, done: threading.Event, first_msg: threading.Event, output_forwarded: threading.Event, client_ids=None, init_state=None) -> None:\n out = sys.stdout.buffer\n try:\n while True:\n kind, opcode, payload = read_message(recv, send_pong)\n if kind == \"close\":\n log(\"daemon WS close\")\n return\n if opcode == 0x1:\n # Desktop StdioConnection expects newline-delimited JSON text\n if not payload.endswith(b\"\\n\"):\n payload += b\"\\n\"\n out.write(payload)\n out.flush()\n else:\n out.write(payload)\n out.flush()\n # Any daemon payload on stdout commits the session: the fallback\n # must never run on already-used stdout even with pristine stdin.\n output_forwarded.set()\n if opcode == 0x1:\n expect = (init_state or {}).get(\"initialize_id\") if isinstance(init_state, dict) else None\n if _is_daemon_response(payload, expect, client_ids or ()):\n first_msg.set()\n except Exception as e:\n log(f\"stdout_writer exit: {e}\")\n finally:\n done.set()\n\ndef main() -> int:\n log(f\"start sock={SOCK}\")\n try:\n sock, pending = ws_connect(SOCK)\n except Exception as e:\n log(f\"connect failed: {e}\")\n sys.stderr.write(f\"codex-stdio-to-daemon-ws: {e}\\n\")\n return EXIT_PRE_STDIO\n log(\"connected\")\n connected_at = time.monotonic()\n send_lock = threading.Lock()\n done = threading.Event()\n first_msg = threading.Event()\n output_forwarded = threading.Event()\n client_ids = set()\n init_state = {\"initialize_id\": None}\n\n def _send_timeout() -> float:\n # Same monotonic session budget bounds every send/lock: before the\n # first daemon response the remaining first-RPC budget applies, after\n # it each send gets the connect timeout. Expired budget fails fast so\n # a non-reading peer cannot wedge pong writes or timeout cleanup\n # before any stdin is consumed.\n if not first_msg.is_set():\n return max(0.0, (connected_at + FIRST_MESSAGE_TIMEOUT_S) - time.monotonic())\n return CONNECT_TIMEOUT_S\n\n def send_pong(payload: bytes) -> None:\n locked_send(sock, send_lock, payload, opcode=0xA, timeout=_send_timeout())\n recv = _buffered_reader(sock, pending)\n t = threading.Thread(target=stdout_writer, args=(sock, recv, send_pong, done, first_msg, output_forwarded, client_ids, init_state), daemon=True)\n t.start()\n # Byte-precise commit point: the fallback may run only while zero stdin\n # bytes have left the pipe and no daemon payload was forwarded. Reads use\n # os.read into a userspace line buffer (buffered readline could readahead\n # past the commit point), and every byte read counts \u2014 even a blank line.\n # first_msg tracks the matching daemon *response/error*, not any notification.\n stdin_bytes = 0\n pending_in = b\"\"\n clean_eof = False\n try:\n fd = sys.stdin.fileno()\n while not done.is_set():\n if not first_msg.is_set() and time.monotonic() - connected_at > FIRST_MESSAGE_TIMEOUT_S:\n log(\"first daemon message timed out\")\n break\n ready, _, _ = select.select([fd], [], [], 0.5)\n if not ready:\n continue\n chunk = os.read(fd, 65536)\n if not chunk:\n clean_eof = True\n break\n stdin_bytes += len(chunk)\n pending_in += chunk\n while b\"\\n\" in pending_in:\n line, _, pending_in = pending_in.partition(b\"\\n\")\n line = line.strip()\n if not line:\n continue\n try:\n try:\n obj = json.loads(line.decode(\"utf-8\"))\n except Exception:\n obj = None\n if isinstance(obj, dict) and obj.get(\"id\") is not None:\n client_ids.add(obj.get(\"id\"))\n if obj.get(\"method\") == \"initialize\" and init_state[\"initialize_id\"] is None:\n init_state[\"initialize_id\"] = obj.get(\"id\")\n except Exception:\n pass\n try:\n locked_send(sock, send_lock, line, timeout=_send_timeout())\n except TimeoutError as e:\n log(f\"stdin forward timed out: {e}\")\n break\n except Exception as e:\n log(f\"stdin forward failed: {e}\")\n break\n else:\n continue\n break\n except Exception as e:\n log(f\"stdin loop exit: {e}\")\n # Flush any buffered EOF-tail data before the WS Close: sending Close\n # first would let the daemon discard the tail.\n if clean_eof:\n tail = pending_in.strip()\n if tail and not done.is_set():\n try:\n locked_send(sock, send_lock, tail, timeout=CONNECT_TIMEOUT_S)\n except Exception as e:\n log(f\"eof flush failed: {e}\")\n try:\n locked_send(sock, send_lock, b\"\", opcode=0x8, timeout=CONNECT_TIMEOUT_S)\n except Exception as e:\n log(f\"close send failed: {e}\")\n try:\n try:\n sock.shutdown(socket.SHUT_RDWR)\n except OSError:\n pass\n sock.close()\n except OSError:\n pass\n if clean_eof:\n log(\"exit\")\n return EXIT_SERVED\n # Reader died, first RPC timed out, or stdin broke: fail open only if no\n # stdin byte was ever consumed AND no daemon payload was forwarded. Used\n # stdout forces a mid-session exit so Desktop reconnects instead of\n # falling back to stock Codex on a corrupted stream.\n rc = EXIT_MID_SESSION if (stdin_bytes or output_forwarded.is_set()) else EXIT_PRE_STDIO\n log(f\"exit rc={rc}\")\n return rc\n\nif __name__ == \"__main__\":\n raise SystemExit(main())\n"; diff --git a/test/desktop-shim.test.js b/test/desktop-shim.test.js index faaa4c7..bdf4639 100644 --- a/test/desktop-shim.test.js +++ b/test/desktop-shim.test.js @@ -921,4 +921,32 @@ test("vendored bridge pins response-tracked first-RPC, output commit, and strict assert.match(BRIDGE_SOURCE, /output_forwarded/); assert.match(BRIDGE_SOURCE, /startswith\(b"HTTP\/1\.1 101"\)/); assert.match(BRIDGE_SOURCE, /tail.*before.*WS Close|before the WS Close/s); + assert.match(BRIDGE_SOURCE, /initialize_id/); + assert.match(BRIDGE_SOURCE, /expect_id/); +}); + +test("bridge first-RPC timer ignores a foreign response id", { + skip: !canRunShellBridge && "needs bash + python3", +}, async () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-foreign-resp-")); + // Daemon answers id 999 while the client asked initialize id 1: only the + // matching response/error/timeout may clear the timer, so it must still fire. + // The foreign response is forwarded (stdout used), hence mid-session exit. + const daemon = await fakeWsDaemon(dir, "foreign.sock", { + prelude: wsServerFrame(0x1, Buffer.from('{"jsonrpc":"2.0","id":999,"result":{}}')), + quiet: true, + }); + try { + const started = Date.now(); + const out = await runBridge(writeBridge(dir), dir, { + CODEX_APP_SERVER_SOCK: daemon.socketPath, + CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT: "2", + }, '{"jsonrpc":"2.0","id":1,"method":"initialize","params":{}}\n', { leaveStdinOpen: true }); + const elapsed = Date.now() - started; + assert.equal(out.status, 2, `foreign response must not satisfy first-RPC; got ${out.status} ${out.stderr}`); + assert.match(out.stdout, /"id":999/); + assert.ok(elapsed < 15000, `matching-response deadline must fire, took ${elapsed}ms`); + } finally { + await closeDaemon(daemon); + } });