Merge commit 'ca94646b7eb1687e78ff65dfddc6e7f955c5693a' into HEAD
This commit is contained in:
@@ -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:
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user