diff --git a/agent/codex_runtime.py b/agent/codex_runtime.py index bae6ec1951..5d0a6121e3 100644 --- a/agent/codex_runtime.py +++ b/agent/codex_runtime.py @@ -958,8 +958,9 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta def _log_failure(exc: BaseException) -> None: request_body_bytes, exception_chain = _codex_request_failure_details(exc) logger.warning("Codex Responses request failed: serialized_request_body_bytes=%s stream_opened=%s " - "exception_chain=%s model=%s", "unknown" if request_body_bytes is None else request_body_bytes, - str(writer_token["value"] is not None).lower(), exception_chain, getattr(agent, "model", "unknown")) + "exception_chain=%s model=%s attempt=%s", "unknown" if request_body_bytes is None else request_body_bytes, + str(writer_token["value"] is not None).lower(), exception_chain, getattr(agent, "model", "unknown"), + f"{attempt + 1}/{max_stream_retries + 1}") def _codex_stream_created(_raw_stream: Any) -> None: # Claim the delta sink for THIS attempt; a newer attempt supersedes this token. @@ -1069,6 +1070,17 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta return event_stream.final_response raise except _APIConnectionError as exc: + # The SDK wraps every connect/receive failure (``raise APIConnectionError from err``), so the + # raw ``transport_errors`` branch above never sees a pre-stream failure. Before the stream + # opened nothing is billed, so one fresh physical request is safe (#103673); once the writer + # token is claimed the inference may already be billed, so mid-stream failures still raise. + if (attempt < max_stream_retries and writer_token["value"] is None + and isinstance(exc.__cause__, _httpx.TransportError)): + logger.debug( + "Codex Responses pre-stream connect failed (attempt %s/%s); retrying. %s error=%s", + attempt + 1, max_stream_retries + 1, agent._client_log_context(), exc, + ) + continue _log_failure(exc) raise if not agent._interrupt_requested: diff --git a/tests/agent/test_codex_request_transport_diagnostics.py b/tests/agent/test_codex_request_transport_diagnostics.py index 0a65155baf..d05c10b392 100644 --- a/tests/agent/test_codex_request_transport_diagnostics.py +++ b/tests/agent/test_codex_request_transport_diagnostics.py @@ -48,6 +48,7 @@ def test_transport_failure_logs_exact_request_bytes_and_class_chain(caplog): model="gpt-5.6-sol", provider="openai-codex", session_id="", + _client_log_context=lambda: "", ) with caplog.at_level(logging.WARNING, logger="agent.codex_runtime"): @@ -58,6 +59,7 @@ def test_transport_failure_logs_exact_request_bytes_and_class_chain(caplog): assert f"serialized_request_body_bytes={len(request_content)}" in message assert "stream_opened=false" in message assert "exception_chain=APIConnectionError <- RemoteProtocolError" in message + assert "attempt=2/2" in message assert "payload" not in message assert request_content.decode() not in message assert "example.invalid" not in message diff --git a/tests/agent/test_run_agent_codex_responses.py b/tests/agent/test_run_agent_codex_responses.py index bc75572f40..289416f8b7 100644 --- a/tests/agent/test_run_agent_codex_responses.py +++ b/tests/agent/test_run_agent_codex_responses.py @@ -2797,3 +2797,117 @@ def test_run_codex_stream_retired_request_stops_firing_callbacks(monkeypatch): assert streamed == ["keep"] assert "DROPPED" not in streamed + + +def _raise_prestream_transport_error(request): + """Raise the #103673 shape: APIConnectionError <- ReadError <- ReadError.""" + import httpx + + from openai import APIConnectionError + + inner = httpx.ReadError("receive failed", request=request) + mid = httpx.ReadError("receive failed", request=request) + try: + raise mid from inner + except httpx.ReadError as chained: + raise APIConnectionError(request=request) from chained + + +def _completed_create_stream(): + message_item = SimpleNamespace( + type="message", + status="completed", + content=[SimpleNamespace(type="output_text", text="Recovered.")], + ) + usage = SimpleNamespace(input_tokens=10, output_tokens=6, total_tokens=16) + return _FakeCreateStream( + [ + SimpleNamespace(type="response.output_item.done", item=message_item), + SimpleNamespace( + type="response.completed", + response=SimpleNamespace( + status="completed", + usage=usage, + id="resp_prestream_retry_1", + ), + ), + ] + ) + + +def test_run_codex_stream_retries_prestream_apiconnectionerror(monkeypatch): + """Regression test for issue #103673. + + A pre-stream ``APIConnectionError`` wrapping an httpx transport error + (``ReadError`` before the first SSE event, stream never opened) must retry + with a fresh physical request like a raw transport error does, instead of + failing the turn on a transient connect/receive failure. + """ + import httpx + + agent = _build_agent(monkeypatch) + request = httpx.Request( + "POST", + "https://chatgpt.com/backend-api/codex/responses", + content=b'{"model":"gpt-5-codex"}', + ) + calls = {"count": 0} + + def _fake_create(**kwargs): + calls["count"] += 1 + if calls["count"] == 1: + _raise_prestream_transport_error(request) + return _completed_create_stream() + + agent.client = SimpleNamespace(responses=SimpleNamespace(create=_fake_create)) + + response = agent._run_codex_stream(_codex_request_kwargs()) + + assert calls["count"] == 2 + assert response.status == "completed" + assert response.id == "resp_prestream_retry_1" + + +def test_run_codex_stream_prestream_retry_exhaustion_logs_telemetry( + monkeypatch, caplog +): + """Regression test for issue #103673 (observability half). + + When the pre-stream retry is exhausted, the turn still raises, but the + single WARNING must carry the byte count, the stream-open state, the + exception chain, and the attempt count -- without prompt content. + """ + import logging + + import httpx + from openai import APIConnectionError + + agent = _build_agent(monkeypatch) + body = b'{"model":"gpt-5-codex"}' + request = httpx.Request( + "POST", "https://chatgpt.com/backend-api/codex/responses", content=body + ) + calls = {"count": 0} + + def _fake_create(**kwargs): + calls["count"] += 1 + _raise_prestream_transport_error(request) + + agent.client = SimpleNamespace(responses=SimpleNamespace(create=_fake_create)) + + with caplog.at_level(logging.WARNING, logger="agent.codex_runtime"): + with pytest.raises(APIConnectionError): + agent._run_codex_stream(_codex_request_kwargs()) + + assert calls["count"] == 2 + failures = [ + record + for record in caplog.records + if "Codex Responses request failed" in record.message + ] + assert len(failures) == 1 + message = failures[0].message + assert f"serialized_request_body_bytes={len(body)}" in message + assert "stream_opened=false" in message + assert "APIConnectionError <- ReadError <- ReadError" in message + assert "attempt=2/2" in message