Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 4 additions & 3 deletions agent/transports/codex_app_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
46 changes: 26 additions & 20 deletions agent/transports/codex_app_server_session.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 ----------

Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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 ----------
Expand All @@ -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:
Expand Down
46 changes: 46 additions & 0 deletions tests/agent/transports/test_codex_app_server_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
137 changes: 131 additions & 6 deletions tests/agent/transports/test_codex_app_server_session.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@

from __future__ import annotations

import threading
import time
from unittest.mock import patch
from typing import Any, Optional
Expand Down Expand Up @@ -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



Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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 ----
Expand Down Expand Up @@ -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]

4 changes: 3 additions & 1 deletion tests/tools/test_terminal_hints.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion tools/terminal_hints.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.",
}


Expand Down
Loading