diff --git a/.changeset/potato-mustfix-stack.md b/.changeset/potato-mustfix-stack.md new file mode 100644 index 0000000..d774a60 --- /dev/null +++ b/.changeset/potato-mustfix-stack.md @@ -0,0 +1,5 @@ +--- +"grok-bot-cli": patch +--- + +Stack the must-fix set from the Codex reject of the Act-On-PR bridge and the unfinished follow-up accept bar: commit stdout BEFORE writing (dirty stdout never exits 1 / never fail-opens to stock Codex) with bounded stdout writes and termination, keep healthy idle reads untimed (separate absolute handshake/first-RPC deadline from mid-session write/lock budgets), clear the init timer only on the matching initialize response/error (notifications never satisfy it), stay silent on foreign serverRequest approvals (deferring turn-scoped ones until the turn id is known), require strict HTTP/1.1 101 plus accept, flush EOF-tail data before the WS Close, and keep the wrapper hardening (self-fallback refusal, bridge+python3+preflight gating, uninstall reporting, Linux status wording with shell-quoted path). diff --git a/README.md b/README.md index 1f80b56..47f1410 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 committed to stdout: 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 preflight and upgrade leave stdin untouched. Reads after a successful init block with no timeout so healthy idle sessions survive silence; every mid-session socket send and stdout write runs under its own write budget (`CODEX_BRIDGE_IO_TIMEOUT`, default 30 s), so a wedged or non-reading peer cannot hang Desktop. The commit point is byte-precise on input OR output, committed before the stdout write: the bridge counts every stdin byte read (buffered readahead included, even an unforwarded blank line) and marks stdout used before writing any forwarded daemon payload via a bounded write, 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 dirty 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: requests without both matching thread and turn IDs remain unanswered; a Codex client must handle those requests. 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-silent, so `gbot` cannot reject Desktop's approval. A request naming a turn that arrives before the acknowledgment supplies `gbot`'s turn id waits instead of being classified foreign, and is answered only on a match. Requests with missing thread or turn IDs remain unanswered because ownership cannot be established. `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..c98d569 100644 --- a/src/core/codex-bridge.js +++ b/src/core/codex-bridge.js @@ -282,11 +282,43 @@ 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, + // Approvals that name a turn while our own turn id is still unknown + // (arriving before — or in the same socket batch as — the turn/start + // acknowledgment) wait here instead of being classified foreign: once + // the acknowledgment supplies the turn id, matching ones are answered + // and the rest stay silent. + deferred: [], + // Adopt the acknowledged turn id, answering deferred requests that + // prove to be ours. Runs inside the connection closure so it can send. + _adoptTurn(turnId) { + client.expectedTurnId = turnId; + for (const entry of client.deferred.splice(0)) { + const params = entry.params && typeof entry.params === "object" && !Array.isArray(entry.params) + ? entry.params + : null; + const threadId = params ? (params.threadId ?? params.thread_id) : null; + const namedTurnId = params + ? (params.turnId ?? params.turn_id ?? (params.turn && params.turn.id)) + : null; + if (threadId === client.expectedThreadId && namedTurnId === turnId) { + entry.answered = true; + sendJson({ + jsonrpc: "2.0", + id: entry.id, + error: { code: -32601, message: "gbot codex does not answer " + entry.method + "; configure approval_policy on the daemon" }, + }); + } + } + }, request(method, params) { const id = nextId++; return new Promise((res, rej) => { @@ -327,8 +359,29 @@ export function connectCodexAppServer(path, { timeoutMs = 15000 } = {}) { const onMessage = (msg) => { if (msg.id != null && msg.method) { - const answered = client.answerServerRequests; - refused.push({ id: msg.id, method: msg.method, params: msg.params, answered }); + let answered = client.answerServerRequests; + let defer = false; + // Approval ownership: only answer requests for the turn gbot itself + // started. A request naming another thread — or naming a turn that is + // not ours — is recorded foreign-silent, never rejected. A request + // naming a turn while our turn id is still unknown is deferred, not + // classified: the turn/start acknowledgment decides it. + 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 && params.turn.id); + if (threadId !== client.expectedThreadId || turnId == null) answered = false; + else if (client.expectedTurnId == null) { answered = false; defer = true; } + else if (turnId !== client.expectedTurnId) answered = false; + } else answered = false; + } + const entry = { id: msg.id, method: msg.method, params: msg.params, answered }; + refused.push(entry); + if (defer) { + client.deferred.push(entry); + return; + } if (!answered) return; sendJson({ jsonrpc: "2.0", @@ -893,11 +946,15 @@ 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. Requests naming a turn + // that arrives before the acknowledgment supplies our turn id wait in + // client.deferred and are adopted only on a match. Requests missing + // thread or turn IDs remain unanswered because ownership is unknown. + client.expectedThreadId = threadId; + client.expectedTurnId = null; client.answerServerRequests = true; let turn; try { @@ -924,8 +981,18 @@ async function sendToCodexThreadInner(threadId, text, { env, envelope, whenBusy { delivery: "unknown", reason: "bad-response", threadId, envelope }, ); } + client._adoptTurn(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 && 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..7b3481e 100644 --- a/src/core/desktop-shim-bridge.js +++ b/src/core/desktop-shim-bridge.js @@ -2,21 +2,33 @@ // 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 shares one monotonic budget (remaining first-RPC +// budget, else the mid-session write budget) so a non-reading peer cannot +// wedge pong writes or timeout cleanup; timeouts close and exit. Reads wait +// with no timeout after the handshake so healthy idle sessions survive. +// 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-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, committed BEFORE the stdout +// write: os.read into a userspace line buffer counts every stdin byte (even +// an unforwarded blank line), and any daemon payload marks stdout used before +// it is written via a bounded stdout write, so the fallback only ever runs +// 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 with a bounded join. 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 while stdio was still\n# pristine -- no stdin byte consumed and no stdout frame committed (wrapper\n# falls back to real Codex), 2 failed mid-session (wrapper exits promptly so\n# Desktop reconnects; the fallback must never run on half-consumed stdin or\n# dirty 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 matching first daemon response: a daemon that\n# completes the upgrade but never answers the first RPC must not hold Desktop\n# past this. Only the matching initialize response/error (or the deadline\n# itself) clears it; notifications never do.\nFIRST_MESSAGE_TIMEOUT_S = float(os.environ.get(\"CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT\", \"30\"))\n# Bound on every mid-session socket send and stdout write: a wedged or\n# non-reading peer must not hang Desktop past this. Reads after a successful\n# init block with no timeout so healthy idle sessions survive silence.\nIO_TIMEOUT_S = float(os.environ.get(\"CODEX_BRIDGE_IO_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.settimeout(max(0.0, deadline - time.monotonic()))\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 lines[0].split(b\" \", 2)[:2] != [b\"HTTP/1.1\", b\"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 # Keep one socket mode across reader and writer threads. Reads wait\n # indefinitely in select; sends use their own absolute deadline.\n s.setblocking(False)\n return s, trailing\n\ndef mask_send(sock: socket.socket, payload: bytes, opcode: int = 1, deadline=None) -> 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 view = memoryview(hdr + mask + bytes(b ^ mask[i % 4] for i, b in enumerate(payload)))\n while view:\n remaining = deadline - time.monotonic()\n if remaining <= 0:\n raise TimeoutError(\"socket write timed out\")\n _, writable, _ = select.select([], [sock], [], remaining)\n if not writable:\n raise TimeoutError(\"socket write timed out\")\n try:\n n = sock.send(view)\n except BlockingIOError:\n continue\n if not n:\n raise ConnectionError(\"socket closed during write\")\n view = view[n:]\n\ndef locked_send(sock: socket.socket, lock: threading.Lock, payload: bytes, opcode: int = 1, timeout: float = CONNECT_TIMEOUT_S) -> None:\n \"\"\"Serialized send sharing one monotonic budget across lock wait + send.\n\n A 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 deadline = time.monotonic() + timeout\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 remaining = deadline - time.monotonic()\n if remaining <= 0:\n raise TimeoutError(\"send budget expired\")\n mask_send(sock, payload, opcode=opcode, deadline=deadline)\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 select.select([sock], [], [])\n try:\n chunk = sock.recv(n - len(buf))\n except BlockingIOError:\n continue\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, never 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 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 _bounded_stdout_write(data: bytes, timeout: float = IO_TIMEOUT_S) -> None:\n \"\"\"Write one daemon frame to stdout without hanging on a non-reading peer.\n\n Waits for writability under the mid-session write budget and writes via\n os.write in a loop, so a Desktop that stopped reading cannot hang the\n bridge and then exit 1 on dirty stdout.\n \"\"\"\n if timeout is None:\n timeout = IO_TIMEOUT_S\n if timeout <= 0:\n raise TimeoutError(\"stdout budget expired\")\n fd = sys.stdout.fileno()\n deadline = time.monotonic() + timeout\n was_blocking = os.get_blocking(fd)\n os.set_blocking(fd, False)\n try:\n view = memoryview(data)\n while view:\n remaining = deadline - time.monotonic()\n if remaining <= 0:\n raise TimeoutError(\"stdout write timed out\")\n _, writable, _ = select.select([], [fd], [], remaining)\n if not writable:\n raise TimeoutError(\"stdout write timed out\")\n try:\n n = os.write(fd, view)\n except BlockingIOError:\n continue\n view = view[n:]\n finally:\n os.set_blocking(fd, was_blocking)\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 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 # Commit BEFORE writing: once a daemon byte is about to hit\n # stdout, the session is dirty and the fallback must never run,\n # even if the write itself blocks or fails.\n output_forwarded.set()\n _bounded_stdout_write(payload)\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 first_deadline = connected_at + FIRST_MESSAGE_TIMEOUT_S\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 # Before the first matching daemon response the remaining first-RPC\n # budget bounds every send; after it each send gets the mid-session\n # write budget. Idle reads are never timed out.\n if not first_msg.is_set():\n return max(0.0, first_deadline - time.monotonic())\n return IO_TIMEOUT_S\n\n def _cleanup_timeout() -> float:\n # Timeout cleanup stays within the expired session budget: once the\n # first-RPC deadline has passed, close directly instead of opening\n # fresh windows that let a short budget take several timeout periods.\n if not first_msg.is_set():\n return max(0.0, first_deadline - 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 committed to stdout.\n # Reads use os.read into a userspace line buffer (buffered readline could\n # readahead past the commit point), and every byte read counts -- even a\n # blank line. first_msg tracks the matching daemon response/error, never\n # a 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() >= first_deadline:\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=_cleanup_timeout())\n except Exception as e:\n log(f\"eof flush failed: {e}\")\n try:\n locked_send(sock, send_lock, b\"\", opcode=0x8, timeout=_cleanup_timeout())\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 # Bound termination: wake the reader so no thread is left hung on recv.\n t.join(timeout=0.5)\n # A reader still exiting could commit output after this check. In that\n # case pristine stdio cannot be proven, so never permit fallback.\n if t.is_alive():\n output_forwarded.set()\n if clean_eof:\n log(\"exit\")\n return EXIT_SERVED\n # Reader died, first RPC timed out, post-upgrade I/O timed out, or stdin\n # broke: fail open only while stdio is still pristine -- no stdin byte\n # ever consumed AND no daemon payload ever committed to stdout. Committed\n # output 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/src/core/desktop-shim.js b/src/core/desktop-shim.js index c35da42..7f79ea5 100644 --- a/src/core/desktop-shim.js +++ b/src/core/desktop-shim.js @@ -86,11 +86,12 @@ export function renderWrapperScript({ realPath, bridgePath, wrapperLogPath, code # daemon's control socket. Reversible: \`gbot codex desktop-shim uninstall\`. # The daemon is kept up by the LaunchAgent login script, not here: this hot path # never starts the daemon itself, so a wedged daemon lock cannot hang -# Desktop. Fail-open runs ONLY before any stdin is consumed: an absent socket -# fails the preflight and a pre-session bridge failure (exit 1) falls through +# Desktop. Fail-open runs ONLY while stdio is still pristine (no stdin +# consumed, no stdout written): an absent socket fails the preflight and a +# pre-session bridge failure (exit 1) falls through # to the real standalone codex on pristine stdio. A mid-session bridge failure # (exit 2+) exits promptly so Desktop reconnects; the fallback never runs on -# half-consumed stdin. Bridge exit 0 means it served the session. Never touches +# half-consumed stdin or dirty stdout. Bridge exit 0 means it served the session. Never touches # Desktop binaries or its private tool pipe. set -u REAL="\${CODEX_DESKTOP_WRAPPER_REAL:-${realPath}}" @@ -107,7 +108,11 @@ echo "$(ts) argv: $*" >>"$LOG" 2>/dev/null || true if [[ ! -x "$REAL" ]]; then FALLBACK="$(command -v codex 2>/dev/null || true)" if [[ -n "$FALLBACK" && -x "$FALLBACK" ]]; then - REAL="$FALLBACK" + if [[ "$FALLBACK" -ef "$0" ]] 2>/dev/null || { [[ -n "\${CODEX_CLI_PATH:-}" ]] && [[ "$FALLBACK" -ef "\${CODEX_CLI_PATH}" ]] 2>/dev/null; }; then + echo "$(ts) refusing self fallback $FALLBACK" >>"$LOG" 2>/dev/null || true + else + REAL="$FALLBACK" + fi fi fi @@ -128,7 +133,7 @@ preflight() { } if [[ "$has_app_server" -eq 1 && "$has_daemon" -eq 0 && "$has_proxy" -eq 0 && "$has_generate" -eq 0 ]]; then - if [[ -x "$REAL" && -f "$BRIDGE" ]] && command -v python3 >/dev/null 2>&1; then + if [[ -f "$BRIDGE" ]] && command -v python3 >/dev/null 2>&1; then if preflight; then echo "$(ts) rewrite -> stdio-ws bridge" >>"$LOG" 2>/dev/null || true python3 "$BRIDGE" @@ -339,9 +344,10 @@ export function uninstallDesktopShim({ const warnings = []; const removed = []; for (const path of [paths.wrapperPath, paths.bridgePath, paths.envScriptPath]) { + const existed = fileExists(path); try { if (rmSync(path, { force: true }) === undefined && fileExists(path)) warnings.push(`could not remove ${path}`); - else if (!fileExists(path)) removed.push(path); + else if (existed && !fileExists(path)) removed.push(path); } catch (error) { warnings.push(`could not remove ${path} (${error instanceof Error ? error.message : String(error)})`); } diff --git a/src/core/format.js b/src/core/format.js index 8021bfb..d46a6c1 100644 --- a/src/core/format.js +++ b/src/core/format.js @@ -106,6 +106,11 @@ export function formatCodexThread(t) { + "\n " + stripTerminalControls(t.cwd ?? "") + preview; } +/** Shell-safe single-quoting for copy-pasteable export lines (spaces, quotes, $). */ +export function shellQuote(value) { + return "'" + String(value).replace(/'/g, "'\\''") + "'"; +} + export function formatDesktopShimStatus(s) { const mark = (ok) => (ok ? "yes" : "no"); const cliLine = s.platform === "darwin" @@ -125,7 +130,11 @@ export function formatDesktopShimStatus(s) { if (!s.installed) { lines.push("shim: not installed — run `gbot codex desktop-shim install` (Desktop keeps stock behavior until then)"); } else if (!s.wrapperPointsAtShim) { - lines.push("shim: installed but CODEX_CLI_PATH does not point at it — reinstall or relaunch ChatGPT.app after login"); + lines.push( + s.platform === "darwin" + ? "shim: installed but CODEX_CLI_PATH does not point at it — reinstall or relaunch ChatGPT.app after login" + : "shim: installed but CODEX_CLI_PATH does not point at it — export CODEX_CLI_PATH=" + shellQuote(s.wrapperPath), + ); } else { lines.push("shim: active — Desktop app-server spawns bridge onto the managed daemon"); } diff --git a/test/codex-bridge.test.js b/test/codex-bridge.test.js index 21b6954..f0cfe46 100644 --- a/test/codex-bridge.test.js +++ b/test/codex-bridge.test.js @@ -28,10 +28,15 @@ async function fakeAppServer(handlers) { const socketPath = join(home, "app-server-control", "app-server-control.sock"); const received = []; const sockets = new Set(); + let resolveDisconnected; + const disconnected = new Promise((resolve) => { resolveDisconnected = resolve; }); const server = createServer(); server.on("upgrade", (req, socket) => { sockets.add(socket); - socket.on("close", () => sockets.delete(socket)); + socket.on("close", () => { + sockets.delete(socket); + resolveDisconnected(); + }); socket.write( "HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\n" + "Sec-WebSocket-Accept: " + websocketAccept(req.headers["sec-websocket-key"]) + "\r\n\r\n", @@ -61,6 +66,7 @@ async function fakeAppServer(handlers) { return { home, received, + disconnected, close: () => new Promise((resolve) => { for (const sock of sockets) sock.destroy(); server.close(resolve); @@ -243,7 +249,7 @@ test("codex send refuses server approval requests and fails with guidance", asyn const fake = await fakeAppServer({ ...baseHandlers, "turn/start": (params, ok, err, send) => { - send({ jsonrpc: "2.0", id: "srv-1", method: "item/commandExecution/requestApproval", params: { threadId: params.threadId, command: "rm -rf /" } }); + send({ jsonrpc: "2.0", id: "srv-1", method: "item/commandExecution/requestApproval", params: { threadId: params.threadId, turnId: "turn-10", command: "rm -rf /" } }); setTimeout(() => ok({ turn: { id: "turn-10", status: "inProgress", items: [] } }), 50); }, }); @@ -513,7 +519,7 @@ test("codex send preserves turn and thread ids with accepted delivery on refusal const fake = await fakeAppServer({ ...baseHandlers, "turn/start": (params, ok, err, send) => { - send({ jsonrpc: "2.0", id: "srv-1", method: "item/commandExecution/requestApproval", params: { threadId: params.threadId } }); + send({ jsonrpc: "2.0", id: "srv-1", method: "item/commandExecution/requestApproval", params: { threadId: params.threadId, turnId: "turn-10" } }); setTimeout(() => ok({ turn: { id: "turn-10", status: "inProgress", items: [] } }), 20); }, }); @@ -637,7 +643,7 @@ test("codex send emits structured JSON errors with delivery and ids", async () = const refusing = await fakeAppServer({ ...baseHandlers, "turn/start": (params, ok, err, send) => { - send({ jsonrpc: "2.0", id: "srv-1", method: "item/commandExecution/requestApproval", params: { threadId: params.threadId } }); + send({ jsonrpc: "2.0", id: "srv-1", method: "item/commandExecution/requestApproval", params: { threadId: params.threadId, turnId: "turn-10" } }); setTimeout(() => ok({ turn: { id: "turn-10", status: "inProgress", items: [] } }), 20); }, }); @@ -1235,3 +1241,87 @@ test("formatCodexStatus names the private-stdio case and keeps the unknown line" assert.match(formatCodexStatus({ ...daemon, desktopAttached: "private-stdio" }), /codex app-server daemon start/); assert.match(formatCodexStatus({ ...daemon, desktopAttached: "unknown" }), /desktop attached: unknown \(not observable from the socket\)/); }); + +test("codex send stays silent on foreign approvals inside the turn window", async () => { + const fake = await fakeAppServer({ + ...baseHandlers, + "turn/start": (params, ok, err, send) => { + // Another thread's approval and a Desktop turn's approval arrive inside + // gbot's own turn/start window: both must stay silent, never -32601. + send({ jsonrpc: "2.0", id: "srv-foreign-thread", method: "item/commandExecution/requestApproval", params: { threadId: "other", command: "rm -rf /" } }); + send({ jsonrpc: "2.0", id: "srv-foreign-turn", method: "item/commandExecution/requestApproval", params: { threadId: params.threadId, turnId: "turn-desktop" } }); + send({ id: "srv-unscoped", method: "item/commandExecution/requestApproval", params: {} }); + send({ id: "srv-thread-only", method: "item/commandExecution/requestApproval", params: { threadId: params.threadId } }); + send({ id: "srv-turn-only", method: "item/commandExecution/requestApproval", params: { turnId: "turn-9" } }); + setTimeout(() => ok({ turn: { id: "turn-9", status: "inProgress", items: [] } }), 20); + }, + }); + const env = { ...process.env, CODEX_HOME: fake.home }; + try { + const outcome = await sendToCodexThread("t-1", "go", { env }); + assert.equal(outcome.delivery, "accepted"); + assert.equal(outcome.turnId, "turn-9"); + assert.equal(outcome.exitCode, 0); + await fake.disconnected; + assert.equal(fake.received.find((m) => m.id === "srv-foreign-thread"), undefined); + assert.equal(fake.received.find((m) => m.id === "srv-foreign-turn"), undefined); + for (const id of ["srv-unscoped", "srv-thread-only", "srv-turn-only"]) { + assert.equal(fake.received.find((m) => m.id === id), undefined, `${id} must stay unanswered`); + } + } finally { + await fake.close(); + } +}); + +test("codex send still refuses our own approval that arrives before the turn ack", async () => { + const fake = await fakeAppServer({ + ...baseHandlers, + "turn/start": (params, ok, err, send) => { + // Our own turn's approval arrives before — or in the same batch as — + // the turn/start acknowledgment: it waits for the ack, then refuses. + send({ jsonrpc: "2.0", id: "srv-early-own", method: "item/commandExecution/requestApproval", params: { threadId: params.threadId, turnId: "turn-12" } }); + setTimeout(() => ok({ turn: { id: "turn-12", status: "inProgress", items: [] } }), 20); + }, + }); + const env = { ...process.env, CODEX_HOME: fake.home }; + try { + const outcome = await sendToCodexThread("t-1", "go", { env }); + assert.equal(outcome.delivery, "accepted"); + assert.equal(outcome.reason, "approval-refused"); + assert.equal(outcome.turnId, "turn-12"); + await fake.disconnected; + const refusal = fake.received.find((m) => m.id === "srv-early-own"); + assert.equal(refusal.error.code, -32601); + assert.equal(refusal.result, undefined); + } finally { + await fake.close(); + } +}); + +for (const whenBusy of ["reject", "queue"]) { + test(`source send leaves resume and ${whenBusy} approvals unanswered`, async () => { + const fake = await fakeAppServer({ + ...baseHandlers, + "thread/resume": (params, ok, err, send) => { + send({ id: "resume-approval", method: "item/commandExecution/requestApproval", params: { threadId: params.threadId, turnId: "desktop-turn" } }); + ok({ thread: { id: params.threadId, status: { type: "active", activeFlags: [] } } }); + }, + "thread/queue/add": (params, ok, err, send) => { + send({ id: "queue-approval", method: "item/commandExecution/requestApproval", params: { threadId: params.threadId, turnId: "desktop-turn" } }); + ok({ queuedSubmission: { id: "q-1" } }); + }, + }); + try { + const result = await sendToCodexThread("t-1", "later", { + whenBusy, + env: { CODEX_HOME: fake.home, GROK_BOT_CODEX_EXPERIMENTAL: "1" }, + }); + await fake.disconnected; + assert.equal(result.delivery, whenBusy === "queue" ? "queued" : "rejected"); + assert.ok(!fake.received.some((msg) => msg.id === "resume-approval" || msg.id === "queue-approval")); + assert.ok(!fake.received.some((msg) => msg.method === "turn/start")); + } finally { + await fake.close(); + } + }); +} diff --git a/test/desktop-shim.test.js b/test/desktop-shim.test.js index caad38c..9a102ba 100644 --- a/test/desktop-shim.test.js +++ b/test/desktop-shim.test.js @@ -9,6 +9,7 @@ import { join } from "node:path"; import test from "node:test"; import { BRIDGE_SOURCE } from "../src/core/desktop-shim-bridge.js"; +import { decodeFrame, encodeFrame } from "../src/core/codex-bridge.js"; import { defaultPaths, desktopShimStatus, @@ -20,6 +21,7 @@ import { shouldBridge, uninstallDesktopShim, } from "../src/core/desktop-shim.js"; +import { formatDesktopShimStatus } from "../src/core/format.js"; test("shouldBridge only rewrites bare app-server spawns", () => { assert.equal(shouldBridge(["app-server"]), true); @@ -418,12 +420,15 @@ const WS_GUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"; /** Unmasked server-to-client WS frame. */ function wsServerFrame(opcode, payload, fin = true) { - const head = Buffer.alloc(payload.length < 126 ? 2 : 4); + const head = Buffer.alloc(payload.length < 126 ? 2 : payload.length < 65536 ? 4 : 10); head[0] = (fin ? 0x80 : 0) | opcode; if (payload.length < 126) head[1] = payload.length; - else { + else if (payload.length < 65536) { head[1] = 126; head.writeUInt16BE(payload.length, 2); + } else { + head[1] = 127; + head.writeBigUInt64BE(BigInt(payload.length), 2); } return Buffer.concat([head, payload]); } @@ -505,6 +510,10 @@ const runBridge = (bridgePath, dir, env, input, { leaveStdinOpen = false } = {}) child.stderr.on("data", (chunk) => { stderr += chunk; }); + child.stdin.on("error", (err) => { + // The bridge can time out while the parent still has queued input. + if (err.code !== "EPIPE") reject(err); + }); const timer = setTimeout(() => { child.kill("SIGKILL"); reject(new Error(`bridge hung: stdout=${stdout} stderr=${stderr}`)); @@ -677,9 +686,515 @@ test("vendored bridge pins the first-RPC deadline and byte-precise commit", () = assert.match(BRIDGE_SOURCE, /first_msg/); }); +test("vendored bridge pins the must-fix set: commit-before-write, strict 101, EOF ordering", () => { + // Commit-before-write: stdout is marked used BEFORE the bounded write, so a + // crash between commit and flush can never exit 1 on dirty stdout. + const commitAt = BRIDGE_SOURCE.indexOf("output_forwarded.set()"); + const writeAt = BRIDGE_SOURCE.indexOf("_bounded_stdout_write(payload)"); + assert.ok(commitAt !== -1 && writeAt !== -1 && commitAt < writeAt, "commit precedes the stdout write"); + assert.match(BRIDGE_SOURCE, /rc = EXIT_MID_SESSION if \(stdin_bytes or output_forwarded\.is_set\(\)\)/); + // Healthy idle: reads block with no timeout after the handshake; only sends + // re-arm a bounded timeout. + assert.match(BRIDGE_SOURCE, /s\.setblocking\(False\)/); + assert.match(BRIDGE_SOURCE, /CODEX_BRIDGE_IO_TIMEOUT/); + // Init immunity: only the matching response/error clears the first-RPC timer. + assert.match(BRIDGE_SOURCE, /_is_daemon_response/); + assert.match(BRIDGE_SOURCE, /notifications never do|never satisfy it|never a notification/i); + // Strict upgrade: HTTP/1.1 101 required, so a 200 fails even with a valid hash. + assert.match(BRIDGE_SOURCE, /\[b"HTTP\/1\.1", b"101"\]/); + // EOF-tail flushes before the WS Close. + const tailAt = BRIDGE_SOURCE.indexOf("tail = pending_in.strip()"); + const closeAt = BRIDGE_SOURCE.indexOf('opcode=0x8, timeout=_cleanup_timeout()'); + assert.ok(tailAt !== -1 && closeAt !== -1 && tailAt < closeAt, "EOF tail flushes before Close"); +}); + test("login script keeps one CODEX_HOME for GUI env and daemon upkeep", () => { const renderedEnv = renderEnvScript({ codexHome: "/c", envLogPath: "/l", realPath: "/r", wrapperPath: "/w" }); assert.match(renderedEnv, /export CODEX_HOME="\$CODEX_HOME_DIR"/); assert.match(renderedEnv, /CODEX_HOME_DIR="\/c"/); assert.match(renderedEnv, /launchctl setenv CODEX_HOME "\$CODEX_HOME_DIR"/); }); + +// ---- must-fix stack: failing-then-passing behavioral guards ---- + +test("bridge dirty stdout never falls back: greet-then-hold exits 2, not 1", { + skip: !canRunShellBridge && "needs bash + python3", +}, async () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-dirty-")); + // The daemon greets with a bare notification (not the matching init + // response) then holds silently: stdout is dirty, yet the first-RPC + // deadline must still fire — and must exit mid-session, never pre-stdio. + const daemon = await fakeWsDaemon(dir, "greet.sock", { + prelude: wsServerFrame(0x1, Buffer.from('{"greet":true}')), + 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, `dirty stdout must exit mid-session, got ${out.status} ${out.stderr}`); + assert.match(out.stdout, /"greet"/); + assert.ok(elapsed < 15000, `first-RPC deadline must fire on dirty stdout, took ${elapsed}ms`); + } finally { + await closeDaemon(daemon); + } +}); + +test("bridge healthy idle survives silence past the write budget after init", { + skip: !canRunShellBridge && "needs bash + python3", +}, async () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-idle-")); + const socketPath = join(dir, "idle.sock"); + try { + rmSync(socketPath, { force: true }); + } catch {} + // Answers initialize with the matching id, then stays silent for 3 s + // (past the 1 s mid-session write budget) before a late frame: reads must + // block with no timeout once init succeeded. + const server = createServer((socket) => { + socket.on("error", () => {}); + let request = Buffer.alloc(0); + let upgraded = false; + let rest = Buffer.alloc(0); + let replied = false; + 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"); + 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", + )); + upgraded = true; + return; + } + rest = Buffer.concat([rest, chunk]); + for (;;) { + const frame = decodeFrame(rest); + if (!frame) return; + rest = frame.rest; + if (frame.opcode === 0x8) return; + if (frame.opcode !== 0x1 || replied) continue; + let msg; + try { + msg = JSON.parse(frame.payload.toString()); + } catch { + continue; + } + if (msg.method === "initialize" && msg.id != null) { + replied = true; + socket.write(wsServerFrame(0x1, Buffer.from(JSON.stringify({ jsonrpc: "2.0", id: msg.id, result: { ok: true } })))); + setTimeout(() => { + try { + socket.write(wsServerFrame(0x1, Buffer.from('{"late":true}'))); + } catch {} + }, 3000).unref(); + setTimeout(() => { + try { + socket.destroy(); + } catch {} + }, 3500).unref(); + } + } + }); + }); + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(socketPath, resolve); + }); + try { + const started = Date.now(); + const out = await runBridge(writeBridge(dir), dir, { + CODEX_APP_SERVER_SOCK: socketPath, + CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT: "10", + CODEX_BRIDGE_IO_TIMEOUT: "1", + }, '{"jsonrpc":"2.0","id":7,"method":"initialize","params":{}}\n', { leaveStdinOpen: true }); + const elapsed = Date.now() - started; + assert.ok(out.stdout.includes('"late":true'), `late frame after 3 s idle must survive, got: ${out.stdout}`); + assert.ok(elapsed >= 3000, `bridge must have waited out the idle gap, took ${elapsed}ms`); + assert.equal(out.status, 2, `expected mid-session exit after the daemon went away, got ${out.status}`); + } finally { + await new Promise((resolve) => server.close(resolve)); + } +}); + +test("bridge notification storm never satisfies the first-RPC deadline", { + skip: !canRunShellBridge && "needs bash + python3", +}, async () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-storm-")); + const storm = []; + for (let i = 0; i < 20; i++) { + storm.push(wsServerFrame(0x1, Buffer.from(JSON.stringify({ jsonrpc: "2.0", method: "ping", params: { n: i } })))); + } + storm.push(wsServerFrame(0x1, Buffer.from(JSON.stringify({ jsonrpc: "2.0", id: 999, result: { foreign: true } })))); + const daemon = await fakeWsDaemon(dir, "storm.sock", { prelude: Buffer.concat(storm), 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-dirtied stdout must exit mid-session, got ${out.status} ${out.stderr}`); + assert.match(out.stdout, /"method":"ping"/); + assert.ok(elapsed < 15000, `first-RPC deadline must fire through the storm, took ${elapsed}ms`); + } finally { + await closeDaemon(daemon); + } +}); + +/** Raw listener answering the upgrade with a caller-chosen status line (valid accept). */ +async function statusLineDaemon(dir, name, statusLine) { + const socketPath = join(dir, name); + try { + rmSync(socketPath, { force: true }); + } catch {} + const server = createServer((socket) => { + socket.on("error", () => {}); + 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"); + socket.write(Buffer.from( + `${statusLine}\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); + }); + return { server, socketPath }; +} + +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 daemon = await statusLineDaemon(dir, "ok200.sock", "HTTP/1.1 200 OK"); + try { + const started = Date.now(); + const out = await runBridge(writeBridge(dir), dir, { CODEX_APP_SERVER_SOCK: daemon.socketPath }, ""); + assert.equal(out.status, 1); + assert.ok(Date.now() - started < 15000, "200 fails fast instead of hanging"); + } finally { + await new Promise((resolve) => daemon.server.close(resolve)); + } +}); + +test("bridge rejects a non-HTTP/1.1 101 upgrade", { + skip: !canRunShellBridge && "needs bash + python3", +}, async () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-101-")); + const daemon = await statusLineDaemon(dir, "old101.sock", "HTTP/1.0 101 Switching Protocols"); + try { + const started = Date.now(); + const out = await runBridge(writeBridge(dir), dir, { CODEX_APP_SERVER_SOCK: daemon.socketPath }, ""); + assert.equal(out.status, 1); + assert.ok(Date.now() - started < 15000, "non-1.1 101 fails fast instead of serving"); + } finally { + await new Promise((resolve) => daemon.server.close(resolve)); + } +}); + +test("bridge flushes EOF-tail data before the WS Close", { + skip: !canRunShellBridge && "needs bash + python3", +}, async () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-order-")); + const socketPath = join(dir, "order.sock"); + try { + rmSync(socketPath, { force: true }); + } catch {} + const frames = []; + let resolveClosed; + const closed = new Promise((resolve) => { + resolveClosed = resolve; + }); + const server = createServer((socket) => { + socket.on("error", () => {}); + let request = Buffer.alloc(0); + let upgraded = false; + let rest = 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"); + 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.write(wsServerFrame(0x1, Buffer.from('{"jsonrpc":"2.0","id":1,"result":{}}'))); + upgraded = true; + return; + } + rest = Buffer.concat([rest, chunk]); + for (;;) { + const frame = decodeFrame(rest); + if (!frame) return; + rest = frame.rest; + frames.push({ opcode: frame.opcode, payload: frame.payload.toString() }); + } + }); + socket.on("close", () => resolveClosed()); + }); + 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 }, + '{"jsonrpc":"2.0","id":1,"method":"initialize","params":{}}\nTAIL-NO-NEWLINE'); + assert.equal(out.status, 0, out.stderr); + await Promise.race([ + closed, + new Promise((_, reject) => setTimeout(() => reject(new Error("daemon never saw the session end")), 5000)), + ]); + const tailIdx = frames.findIndex((f) => f.opcode === 0x1 && f.payload === "TAIL-NO-NEWLINE"); + const closeIdx = frames.findIndex((f) => f.opcode === 0x8); + assert.ok(tailIdx !== -1, `tail flushed, got ${JSON.stringify(frames)}`); + assert.ok(closeIdx !== -1, "close sent"); + assert.ok(tailIdx < closeIdx, "EOF tail flushes before Close"); + } finally { + await new Promise((resolve) => server.close(resolve)); + } +}); + +test("bridge send to a non-reading peer times out instead of hanging", { + skip: !canRunShellBridge && "needs bash + python3", +}, async () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-blackhole-")); + const socketPath = join(dir, "blackhole.sock"); + try { + rmSync(socketPath, { force: true }); + } catch {} + let peer; + const server = createServer((socket) => { + peer = socket; + socket.on("error", () => {}); + let request = Buffer.alloc(0); + let upgraded = false; + socket.on("data", (chunk) => { + if (upgraded) return; + 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"); + 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", + )); + upgraded = true; + // Blackhole: never read again, so the client's send buffer fills. + socket.removeAllListeners("data"); + socket.pause(); + }); + }); + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(socketPath, resolve); + }); + try { + const line = `{"jsonrpc":"2.0","id":1,"method":"initialize","params":{},"pad":"${"x".repeat(30000)}"}\n`; + const started = Date.now(); + const out = await runBridge(writeBridge(dir), dir, { + CODEX_APP_SERVER_SOCK: socketPath, + CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT: "5", + }, line.repeat(300)); + const elapsed = Date.now() - started; + assert.equal(out.status, 2, `wedged send must exit mid-session, got ${out.status} ${out.stderr}`); + assert.ok(elapsed < 20000, `send must time out instead of hanging, took ${elapsed}ms`); + } finally { + peer?.destroy(); + await new Promise((resolve) => server.close(resolve)); + } +}); + +test("wrapper refuses self fallback and gates the bridge without REAL", () => { + const script = renderWrapperScript({ + bridgeLogPath: "/b/bridge.log", + bridgePath: "/b/bridge.py", + codexHome: "/b", + realPath: "/r/codex", + wrapperLogPath: "/b/w.log", + }); + // A `command -v codex` hit that resolves to the wrapper itself (via $0 or + // CODEX_CLI_PATH) must not exec-loop; it fails with no runnable codex. + assert.match(script, /-ef "\$0"/); + assert.match(script, /refusing self fallback/); + // Bridge attempt needs only the bridge file + python3 + preflight: REAL + // only gates passthrough/fallthrough, never the bridge attempt. + assert.match(script, /if \[\[ -f "\$BRIDGE" \]\] && command -v python3/); + assert.doesNotMatch(script, /-x "\$REAL" && -f "\$BRIDGE"/); +}); + +test("uninstall reports only paths that existed before delete", () => { + const home = mkdtempSync(join(tmpdir(), "gbot-shim-uninstall-home-")); + const codexHome = mkdtempSync(join(tmpdir(), "gbot-shim-uninstall-codex-")); + const env = { CODEX_HOME: codexHome, HOME: home }; + const runner = () => ({ status: 0 }); + const empty = uninstallDesktopShim({ env, home, platform: "linux", runner }); + assert.deepEqual(empty.removed, [], "nothing installed, nothing reported removed"); +}); + +test("linux status quotes the export path for spaces", () => { + const home = mkdtempSync(join(tmpdir(), "gbot-shim-quote-home-")); + const codexHome = mkdtempSync(join(tmpdir(), "gbot-shim-quote-codex-")); + const installed = installDesktopShim({ + env: { CODEX_HOME: codexHome, HOME: home }, + home, + platform: "linux", + runner: () => ({ status: 0 }), + }); + const status = desktopShimStatus({ env: { CODEX_HOME: codexHome, HOME: home }, home, platform: "linux" }); + assert.equal(status.wrapperPointsAtShim, false); + const text = formatDesktopShimStatus({ ...status, wrapperPath: "/home/First Last/.codex/bin/codex-desktop-to-daemon" }); + assert.match(text, /export CODEX_CLI_PATH='\/home\/First Last\/.codex\/bin\/codex-desktop-to-daemon'/); + assert.equal(installed.exitCode, 0); +}); + +// Execute the installed Python helpers against real socket pairs and pipes. +for (const [name, script] of [ + ["a reader still alive at shutdown cannot authorize fallback", ` +a, b = socket.socketpair() +r, w = os.pipe() +sys.stdin = os.fdopen(r, "r") +g = bridge["main"].__globals__ +g["CONNECT_TIMEOUT_S"] = 0.01 +g["FIRST_MESSAGE_TIMEOUT_S"] = 0.01 +g["ws_connect"] = lambda path: (a, b"") +def delayed_writer(*args): + time.sleep(2) +g["stdout_writer"] = delayed_writer +try: + assert bridge["main"]() == bridge["EXIT_MID_SESSION"] +finally: + b.close() + os.close(w) +`], + ["stdout backpressure expires after a partial write", ` +r, w = os.pipe() +old_stdout = sys.stdout +sys.stdout = os.fdopen(w, "w") +started = time.monotonic() +try: + try: + bridge["_bounded_stdout_write"](b"x" * (4 * 1024 * 1024), 0.2) + raise AssertionError("full pipe did not time out") + except TimeoutError: + assert time.monotonic() - started < 1 + assert os.read(r, 1) == b"x", "must exercise a partial write" +finally: + sys.stdout.close() + sys.stdout = old_stdout + os.close(r) +`], + ["socket send budget does not time out a concurrent idle reader", ` +a, b = socket.socketpair() +a.setblocking(False) +a.setsockopt(socket.SOL_SOCKET, socket.SO_SNDBUF, 4096) +lock = threading.Lock() +errors = [] +def send(): + try: + bridge["locked_send"](a, lock, b"x" * 1000000, timeout=0.2) + except TimeoutError: + pass + except Exception as e: + errors.append(e) +t = threading.Thread(target=send) +t.start() +while not lock.locked(): + time.sleep(0.001) +time.sleep(0.05) +late = threading.Timer(0.5, lambda: b.sendall(b"z")) +late.start() +try: + assert bridge["_buffered_reader"](a)(1) == b"z" + t.join(1) + assert not t.is_alive() + assert not errors, errors +finally: + late.join() + a.close() + b.close() +`], +]) { + test(`bridge ${name}`, { skip: !canRunShellBridge && "needs python3" }, () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-python-")); + const out = spawnSync("python3", ["-c", ` +import os, runpy, socket, sys, threading, time +bridge = runpy.run_path(sys.argv[1]) +${script} +`, writeBridge(dir)], { encoding: "utf8", timeout: 3000 }); + assert.equal(out.error, undefined, out.error?.message); + assert.equal(out.status, 0, out.stderr); + }); +} + +test("wrapper exits without fallback while Desktop leaves stdout unread", { + skip: !canRunShellBridge && "needs bash + python3", +}, async () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-backpressure-")); + const payload = Buffer.from(JSON.stringify({ method: "notice", params: { text: "x".repeat(2_000_000) } })); + const daemon = await fakeWsDaemon(dir, "output.sock", { prelude: encodeFrame(1, payload), quiet: true }); + const { wrapperPath } = stageWrapper(dir, { bridgePath: writeBridge(dir) }); + const child = spawn("bash", [wrapperPath, "app-server"], { + env: { + ...process.env, + CODEX_APP_SERVER_SOCK: daemon.socketPath, + CODEX_BRIDGE_FIRST_MESSAGE_TIMEOUT: "5", + CODEX_BRIDGE_IO_TIMEOUT: "0.2", + }, + }); + let stderr = ""; + child.stderr.on("data", (chunk) => { stderr += chunk; }); + const timer = setTimeout(() => child.kill("SIGKILL"), 3000); + try { + // Observe process exit before draining stdout: a blocked write must not + // keep the process alive until Desktop resumes reading. + const [code, signal] = await once(child, "exit"); + assert.equal(signal, null, stderr); + assert.equal(code, 2, stderr); + let output = ""; + for await (const chunk of child.stdout) output += chunk; + assert.ok(output.length > 0 && output.length < payload.length); + assert.doesNotMatch(output, /FAKE-REAL/); + } finally { + clearTimeout(timer); + child.kill("SIGKILL"); + child.stdin.destroy(); + child.stdout.destroy(); + await closeDaemon(daemon); + } +}); + +test("bridge rejects a status token beginning with 101", { + skip: !canRunShellBridge && "needs bash + python3", +}, async () => { + const dir = mkdtempSync(join(tmpdir(), "gbot-shim-1010-")); + const daemon = await statusLineDaemon(dir, "invalid.sock", "HTTP/1.1 1010 Invalid"); + try { + const out = await runBridge(writeBridge(dir), dir, { CODEX_APP_SERVER_SOCK: daemon.socketPath }, ""); + assert.equal(out.status, 1, out.stderr); + } finally { + await new Promise((resolve) => daemon.server.close(resolve)); + } +});