diff --git a/agent/transports/codex_app_server.py b/agent/transports/codex_app_server.py index b96228bfca8d..7e8bb696cbb4 100644 --- a/agent/transports/codex_app_server.py +++ b/agent/transports/codex_app_server.py @@ -228,7 +228,7 @@ def request( self._pending[rid] = _Pending(queue=q, method=method) try: self._send({"id": rid, "method": method, "params": params or {}}) - except CodexAppServerTransportError: + except Exception: with self._pending_lock: self._pending.pop(rid, None) raise @@ -314,10 +314,11 @@ def _send(self, obj: dict) -> None: raise CodexAppServerTransportError( "codex app-server stdin not available" ) + payload = (json.dumps(obj) + "\n").encode("utf-8") try: - self._proc.stdin.write((json.dumps(obj) + "\n").encode("utf-8")) + self._proc.stdin.write(payload) self._proc.stdin.flush() - except (BrokenPipeError, ValueError) as exc: + except (OSError, ValueError) as exc: raise CodexAppServerTransportError( f"codex app-server stdin closed unexpectedly: {exc}" ) from exc diff --git a/agent/transports/codex_app_server_session.py b/agent/transports/codex_app_server_session.py index 350b9a9ff280..3c4e5521fa4d 100644 --- a/agent/transports/codex_app_server_session.py +++ b/agent/transports/codex_app_server_session.py @@ -305,6 +305,7 @@ def __init__( self._client: Optional[CodexAppServerClient] = None self._thread_id: Optional[str] = None self._interrupt_event = threading.Event() + self._transport_failed = threading.Event() self._active_turn_id: Optional[str] = None self._active_turn_lock = threading.Lock() # Pending file-change items, keyed by item id. Populated on @@ -417,27 +418,29 @@ def request_steer(self, text: str) -> bool: turn_id = self._active_turn_id thread_id = self._thread_id client = self._client - if not turn_id or not thread_id or client is None: - return False - try: - response = client.request( - "turn/steer", - { - "threadId": thread_id, - "input": [{"type": "text", "text": cleaned}], - "expectedTurnId": turn_id, - }, - timeout=10, + if not turn_id or not thread_id or client is None: + return False + try: + response = client.request( + "turn/steer", + { + "threadId": thread_id, + "input": [{"type": "text", "text": cleaned}], + "expectedTurnId": turn_id, + }, + timeout=10, + ) + except CodexAppServerTransportError: + self._transport_failed.set() + logger.debug("turn/steer transport unavailable", exc_info=True) + return False + except (CodexAppServerError, TimeoutError): + logger.debug("turn/steer rejected for active Codex turn", exc_info=True) + return False + accepted_turn_id = ( + response.get("turnId") if isinstance(response, dict) else None ) - except ( - CodexAppServerError, - CodexAppServerTransportError, - TimeoutError, - ): - logger.debug("turn/steer rejected for active Codex turn", exc_info=True) - return False - accepted_turn_id = response.get("turnId") if isinstance(response, dict) else None - return accepted_turn_id in {None, turn_id} + return accepted_turn_id in {None, turn_id} # ---------- diagnostics ---------- @@ -817,6 +820,7 @@ def run_turn( with self._active_turn_lock: self._active_turn_id = None self._interrupt_event.clear() + result.should_retire = result.should_retire or self._transport_failed.is_set() return result def compact_thread( @@ -1024,6 +1028,7 @@ def compact_thread( ) result.should_retire = True + result.should_retire = result.should_retire or self._transport_failed.is_set() return result # ---------- internals ---------- @@ -1043,6 +1048,7 @@ def _issue_interrupt(self, turn_id: Optional[str]) -> None: except TimeoutError: logger.warning("turn/interrupt timed out") except CodexAppServerTransportError: + self._transport_failed.set() logger.debug("turn/interrupt transport unavailable", exc_info=True) def _handle_server_request(self, req: dict) -> None: diff --git a/tests/agent/transports/test_codex_app_server_runtime.py b/tests/agent/transports/test_codex_app_server_runtime.py index 0136b15c2e62..3a91bd85da99 100644 --- a/tests/agent/transports/test_codex_app_server_runtime.py +++ b/tests/agent/transports/test_codex_app_server_runtime.py @@ -149,6 +149,52 @@ def fail_send(_payload: dict) -> None: assert isinstance(exc.value, RuntimeError) assert client._pending == {} + def test_request_serialization_value_error_clears_pending(self) -> None: + from agent.transports.codex_app_server import ( + CodexAppServerClient, + CodexAppServerTransportError, + ) + + client = object.__new__(CodexAppServerClient) + client._next_id = 1 + client._pending = {} + client._pending_lock = threading.Lock() + client._closed = False + client._proc = type("Proc", (), {"stdin": object()})() + circular: dict = {} + circular["self"] = circular + + with pytest.raises(ValueError, match="Circular reference") as exc: + client.request("turn/start", circular) + + assert not isinstance(exc.value, CodexAppServerTransportError) + assert client._pending == {} + + def test_request_os_write_failure_is_typed_and_clears_pending(self) -> None: + from agent.transports.codex_app_server import ( + CodexAppServerClient, + CodexAppServerTransportError, + ) + + class BrokenStdin: + def write(self, _payload: bytes) -> None: + raise OSError("bad file descriptor") + + def flush(self) -> None: + raise AssertionError("flush should not follow a failed write") + + client = object.__new__(CodexAppServerClient) + client._next_id = 1 + client._pending = {} + client._pending_lock = threading.Lock() + client._closed = False + client._proc = type("Proc", (), {"stdin": BrokenStdin()})() + + with pytest.raises(CodexAppServerTransportError, match="bad file descriptor"): + client.request("turn/start", {}) + + assert client._pending == {} + class TestSpawnEnvIsolation: """The codex spawn must NOT rewrite HOME — codex's shell tool spawns diff --git a/tests/agent/transports/test_codex_app_server_session.py b/tests/agent/transports/test_codex_app_server_session.py index a1a4b2f8a82f..a54d1ccd0b4c 100644 --- a/tests/agent/transports/test_codex_app_server_session.py +++ b/tests/agent/transports/test_codex_app_server_session.py @@ -7,6 +7,7 @@ from __future__ import annotations +import threading import time from unittest.mock import patch from typing import Any, Optional @@ -456,11 +457,85 @@ def test_steer_transport_failure_is_non_fatal(self): session._active_turn_id = "turn-live-123" def fail_write(method, params): - raise CodexAppServerTransportError("stdin closed") + if method == "turn/steer": + raise CodexAppServerTransportError("stdin closed") + if method == "turn/start": + return {"turn": {"id": "turn-fake-001"}} + return {} client._request_handler = fail_write assert session.request_steer("Use Postgres instead") is False + client.queue_notification( + "turn/completed", + threadId="thread-fake-001", + turn={"id": "turn-fake-001", "status": "completed", "error": None}, + ) + + result = session.run_turn("finish", turn_timeout=2.0) + + assert result.should_retire is True + + def test_inflight_steer_failure_precedes_turn_finalization(self): + from agent.transports.codex_app_server import CodexAppServerTransportError + + client = FakeClient() + session = make_session(client) + steer_started = threading.Event() + release_steer = threading.Event() + turn_result = [] + + def handle_request(method, params): + if method == "thread/start": + return { + "thread": {"id": "thread-fake-001"}, + "activePermissionProfile": {"id": "workspace-write"}, + } + if method == "turn/start": + return {"turn": {"id": "turn-fake-001"}} + if method == "turn/steer": + steer_started.set() + assert release_steer.wait(timeout=2.0) + raise CodexAppServerTransportError("stdin closed") + return {} + + client._request_handler = handle_request + run_thread = threading.Thread( + target=lambda: turn_result.append( + session.run_turn( + "start", + turn_timeout=2.0, + notification_poll_timeout=0.01, + ) + ) + ) + run_thread.start() + deadline = time.monotonic() + 2.0 + while session._active_turn_id is None and time.monotonic() < deadline: + time.sleep(0.001) + assert session._active_turn_id == "turn-fake-001" + + steer_result = [] + steer_thread = threading.Thread( + target=lambda: steer_result.append(session.request_steer("change course")) + ) + steer_thread.start() + assert steer_started.wait(timeout=2.0) + client.queue_notification( + "turn/completed", + threadId="thread-fake-001", + turn={"id": "turn-fake-001", "status": "completed", "error": None}, + ) + time.sleep(0.02) + assert run_thread.is_alive() + + release_steer.set() + steer_thread.join(timeout=2.0) + run_thread.join(timeout=2.0) + + assert steer_result == [False] + assert len(turn_result) == 1 + assert turn_result[0].should_retire is True @@ -515,6 +590,45 @@ def test_compact_thread_sends_rpc_and_waits_for_completion(self): assert r.token_usage_last["totalTokens"] == 12 assert r.model_context_window == 200000 + def test_interrupt_transport_failure_retires_compaction_session(self): + from agent.transports.codex_app_server import CodexAppServerTransportError + + client = FakeClient() + session = make_session(client) + client.queue_notification( + "turn/started", + threadId="thread-fake-001", + turn={"id": "compact-turn-1"}, + ) + take_notification = client.take_notification + + def take_and_interrupt(timeout=0.0): + notification = take_notification(timeout) + if notification and notification["method"] == "turn/started": + session.request_interrupt() + return notification + + client.take_notification = take_and_interrupt + + def fail_interrupt(method, params): + if method == "thread/start": + return { + "thread": {"id": "thread-fake-001"}, + "activePermissionProfile": {"id": "workspace-write"}, + } + if method == "thread/compact/start": + return {} + if method == "turn/interrupt": + raise CodexAppServerTransportError("stdin closed") + return {} + + client._request_handler = fail_interrupt + + result = session.compact_thread(turn_timeout=2.0) + + assert result.interrupted is True + assert result.should_retire is True + def test_compact_start_transport_failure_returns_retiring_result(self): from agent.transports.codex_app_server import CodexAppServerTransportError @@ -942,19 +1056,31 @@ def test_dead_subprocess_detected_between_iterations(self): # Stderr-derived auth hint takes precedence over generic message assert r.error and "codex login" in r.error - def test_interrupt_transport_failure_is_non_fatal(self): + def test_interrupt_transport_failure_retires_session(self): from agent.transports.codex_app_server import CodexAppServerTransportError client = FakeClient() session = make_session(client) - session.ensure_started() def fail_write(method, params): - raise CodexAppServerTransportError("stdin closed") + if method == "turn/interrupt": + raise CodexAppServerTransportError("stdin closed") + if method == "thread/start": + return { + "thread": {"id": "thread-fake-001"}, + "activePermissionProfile": {"id": "workspace-write"}, + } + if method == "turn/start": + session.request_interrupt() + return {"turn": {"id": "turn-fake-001"}} + return {} client._request_handler = fail_write - session._issue_interrupt("turn-fake-001") + result = session.run_turn("hi", turn_timeout=2.0) + + assert result.interrupted is True + assert result.should_retire is True # ---- thread/start cross-fill ---- @@ -1028,4 +1154,3 @@ def test_empty_inputs(self): assert _classify_oauth_failure() is None assert _classify_oauth_failure("") is None assert _classify_oauth_failure("", None) is None # type: ignore[arg-type] - diff --git a/tests/tools/test_terminal_hints.py b/tests/tools/test_terminal_hints.py index fe2857f40fe7..3483fa9e33f8 100644 --- a/tests/tools/test_terminal_hints.py +++ b/tests/tools/test_terminal_hints.py @@ -15,7 +15,9 @@ def test_success_never_annotated(self): def test_empty_output_falls_to_exit_code_tier(self): assert "126" in annotate_failure("./run.sh", 126, "") assert "SIGKILL" in annotate_failure("big_job", 137, "") - assert "timeout" in annotate_failure("sleep 999", 124, "") + timeout_hint = annotate_failure("sleep 999", 124, "") + assert "timeout" in timeout_hint + assert "600s" not in timeout_hint def test_unknown_failure_returns_none(self): assert annotate_failure("./x", 1, "some unrecognized error") is None diff --git a/tools/terminal_hints.py b/tools/terminal_hints.py index 7c0176d1512e..12f98af4a11b 100644 --- a/tools/terminal_hints.py +++ b/tools/terminal_hints.py @@ -141,7 +141,7 @@ def _hint_permission_denied(command: str, output: str) -> Optional[str]: _EXIT_CODE_HINTS: dict[int, str] = { 126: "Exit 126: the file was found but is not executable — `chmod +x` it or invoke it via its interpreter (e.g. `bash script.sh`).", 137: "Exit 137: the process was SIGKILLed — usually out-of-memory or an external kill. Reduce memory use or check `dmesg | tail` before retrying.", - 124: "Exit 124: the command hit its timeout. Raise timeout= (foreground max 600s) or run it with background=true and notify_on_complete=true.", + 124: "Exit 124: the command hit its timeout. Raise timeout= up to the configured foreground maximum, or run it with background=true and notify_on_complete=true.", }