Fix Desktop bridge backpressure, startup deadlines, and approval ownership - #58
Conversation
…mune deadline, foreign-approval silence, strict 101, EOF-before-Close - Bridge commits stdout BEFORE writing (dirty stdout never exits 1); bounded stdout writes/termination; reads blocking after init so healthy idle survives; first-RPC timer cleared only by matching init response. - codex-bridge: foreign thread/turn approvals stay silent; turn-scoped approvals arriving before the turn ack defer until the ack decides. - Wrapper keeps self-fallback refusal, bridge+python3+preflight gating, uninstall reporting; Linux status shell-quotes the export path. - Tests for each must-fix behavior; changeset + README document the bar. Co-authored-by: Zack Jackson <ScriptedAlchemy@users.noreply.github.com>
🦋 Changeset detectedLatest commit: bebed26 The changes in this PR will be included in the next version bump. This PR includes changesets to release 1 package
Not sure what this means? Click here to learn what changesets are. Click here if you're a maintainer who wants to add another changeset to this PR |
…oken) Co-authored-by: Zack Jackson <ScriptedAlchemy@users.noreply.github.com>
…IPE-tolerant chunked flood, destroy paused blackhole socket Co-authored-by: Zack Jackson <ScriptedAlchemy@users.noreply.github.com>
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 492dda2b40
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| // 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.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 # Healthy idle must survive: reads block with no timeout after the\n # handshake. Every send re-arms a bounded timeout around sendall only.\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 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 sock.settimeout(remaining)\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, 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 view = memoryview(data)\n while len(view) > 0:\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\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=CONNECT_TIMEOUT_S)\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"; |
There was a problem hiding this comment.
Make bounded stdout writes nonblocking
When the daemon sends a frame larger than the available stdout pipe capacity after initialization and Desktop stops reading, select() may report the descriptor writable but the subsequent blocking os.write(fd, view) can hang while filling the remainder of the pipe. Because first_msg is already set and stdout_writer cannot reach its finally, the main loop has no deadline that wakes it, defeating the new I/O timeout and potentially wedging Desktop indefinitely. Set the descriptor nonblocking or limit each write to an amount that cannot block before rechecking the deadline.
Useful? React with 👍 / 👎.
| // 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.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 # Healthy idle must survive: reads block with no timeout after the\n # handshake. Every send re-arms a bounded timeout around sendall only.\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 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 sock.settimeout(remaining)\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, 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 view = memoryview(data)\n while len(view) > 0:\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\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=CONNECT_TIMEOUT_S)\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"; |
There was a problem hiding this comment.
Avoid changing the shared socket timeout during sends
During any client send, locked_send temporarily calls sock.settimeout(remaining) on the same socket consumed by stdout_writer. If the reader starts its next recv() during that window, that operation retains the finite timeout even after the sender resets the socket to blocking mode; a successfully initialized but subsequently idle daemon can therefore make the reader time out, set done, and terminate the supposedly idle-safe session. Bound sends with writable polling/nonblocking I/O without changing the timeout observed by the receive thread.
Useful? React with 👍 / 👎.
Problem and result
The Desktop stdio bridge could hang on backpressure or disconnect healthy idle sessions, and a reader still shutting down could allow fallback after committing stdout. This completes the production blocker stack from #56 and the reject of #57 (
48fbcc5). Supersedes those PRs; review only, do not merge.Changes
Local verification
TMPDIR=/tmp node --test test/desktop-shim.test.js: 43 passed, zero failures/skips.TMPDIR=/tmp node --test test/codex-bridge.test.js: 61 passed, zero failures/skips.npm run build: passed, including host load validation.git diff --check: passed.Regression coverage includes partial stdout backpressure, delayed reader shutdown, concurrent idle reads during a bounded send, non-reading socket peer, notification/foreign-response startup deadlines, idle survival, upgrade rejection, EOF ordering, and approval ownership during resume/queue/start.
Limits
Package CI remains queued on GitHub runners. Local verification was on macOS with fake Unix-socket daemons and the built CLI; no new live Desktop or Linux session validation is claimed. Requests without both ownership IDs remain for another Codex client to answer. No npm version bump and no merge.