diff --git a/agent/chat_completion_helpers.py b/agent/chat_completion_helpers.py index 4b95950db1..b4b96c611b 100644 --- a/agent/chat_completion_helpers.py +++ b/agent/chat_completion_helpers.py @@ -3424,11 +3424,14 @@ class _StreamingCall(StreamingWaitMonitor): self.agent._log_stream_retry(kind="exhausted", error=e, attempt=max_retries + 1, max_attempts=max_retries + 1, mid_tool_call=False, diag=self.clients.diag) # Empty stream: "connection failed" would send users chasing network issues. - _what = ("Provider returned malformed streaming data after" if _is_stream_parse_err - else "Provider returned an empty response stream after" if _is_empty_stream - else "Connection to provider failed after") - self.agent._buffer_diagnostic_status( - f"❌ {_what} {max_retries + 1} attempts. The provider may be experiencing issues — try again in a moment.") + if _is_stream_parse_err or _is_empty_stream: + _what = ("Provider returned malformed streaming data after" if _is_stream_parse_err + else "Provider returned an empty response stream after") + self.agent._buffer_diagnostic_status( + f"❌ {_what} {max_retries + 1} attempts. The provider may be experiencing issues — try again in a moment.") + else: + from agent.stream_diag import buffer_connect_exhausted_notice + buffer_connect_exhausted_notice(self.agent, e, attempts=max_retries + 1, base_url=self.agent.base_url) else: self._maybe_disable_streaming(e) logger.exception("Streaming failed before delivery: %s", e) diff --git a/agent/codex_runtime.py b/agent/codex_runtime.py index 09009f5996..3e37bf39f2 100644 --- a/agent/codex_runtime.py +++ b/agent/codex_runtime.py @@ -17,6 +17,7 @@ from typing import Any, Callable, Dict, List from agent.stream_single_writer import claim_stream_writer, stream_writer_is_current from agent.transports.hermes_tools_mcp_server import HERMES_TOOLS_MCP_SERVER_NAME from agent.sdk_transform_bypass import bypass_sdk_request_transform +from agent.stream_diag import buffer_connect_exhausted_notice from agent.usage_anchor import set_usage_anchor logger = logging.getLogger(__name__) @@ -1040,6 +1041,11 @@ def run_codex_stream(agent, api_kwargs: dict, client: Any = None, on_first_delta "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}") + if writer_token["value"] is None: + # No stream ever opened: the user gets one line naming host/attempts/size (#97548). + buffer_connect_exhausted_notice( + agent, exc, attempts=attempt + 1, + base_url=getattr(active_client, "base_url", None) or getattr(agent, "base_url", "")) def _codex_stream_created(_raw_stream: Any) -> None: # Claim the delta sink for THIS attempt; a newer attempt supersedes this token. diff --git a/agent/stream_diag.py b/agent/stream_diag.py index 3bc1d56b53..cb1b42e25c 100644 --- a/agent/stream_diag.py +++ b/agent/stream_diag.py @@ -155,8 +155,67 @@ def emit_stream_drop( pass +# Above this size a refused/reset connect is as likely a body-size limit on the endpoint or a proxy +# in front of it as an outage, so the notice says so (#97548: an ~829 KB request failed twice +# before the stream opened while short chats went through). +LARGE_REQUEST_HINT_BYTES = 256 * 1024 + + +def _failed_request(error: BaseException) -> Any: + """The buffered ``httpx.Request`` carried by the exception chain (SDK wrapper or httpx error), else None.""" + link: Optional[BaseException] = error + for _ in range(8): + if link is None: + return None + try: + request = getattr(link, "request", None) + except RuntimeError: # httpx raises when the error was built without a request + request = None + if request is not None: + return request + link = link.__cause__ or link.__context__ + return None + + +def _request_body_bytes(request: Any) -> Optional[int]: + content = getattr(request, "content", None) + if isinstance(content, str): + return len(content.encode("utf-8")) + if isinstance(content, (bytes, bytearray, memoryview)): + return len(content) + return None + + +def connect_exhausted_notice(error: BaseException, *, attempts: int, base_url: Any) -> str: + """The one user-facing line for \"no stream event ever arrived and the connect retries are spent\": + names the host, the attempt count and the serialized request size so a body-size limit is + distinguishable from an outage without reading agent.log.""" + from urllib.parse import urlparse + request = _failed_request(error) + host = urlparse(str(getattr(request, "url", None) or base_url or "")).hostname or "the endpoint" + size = _request_body_bytes(request) + line = f"❌ Could not open a stream to {host} after {attempts} attempt{'s' if attempts != 1 else ''}" + if size is None: + return line + "; the endpoint looks unreachable — try again in a moment." + line += f" (request {max(1, round(size / 1024))} KB)" + if size >= LARGE_REQUEST_HINT_BYTES: + return line + "; the endpoint or a proxy in front of it may reject requests this large." + return line + "; the endpoint looks unreachable — try again in a moment." + + +def buffer_connect_exhausted_notice(agent: Any, error: BaseException, *, attempts: int, base_url: Any) -> None: + """Buffer :func:`connect_exhausted_notice` once per turn: the outer retry/fallback loop re-enters the + stream call several times, and every pass exhausting the same budget must not add another copy.""" + text = connect_exhausted_notice(error, attempts=attempts, base_url=base_url) + if any(str(msg) == text for _kind, msg in getattr(agent, "_retry_status_buffer", None) or ()): + return + agent._buffer_diagnostic_status(text) + + __all__ = [ "STREAM_DIAG_HEADERS", + "connect_exhausted_notice", + "buffer_connect_exhausted_notice", "stream_diag_init", "stream_diag_capture_response", "flatten_exception_chain", diff --git a/tests/agent/test_codex_request_transport_diagnostics.py b/tests/agent/test_codex_request_transport_diagnostics.py index fa1e272d8e..c1226af3b2 100644 --- a/tests/agent/test_codex_request_transport_diagnostics.py +++ b/tests/agent/test_codex_request_transport_diagnostics.py @@ -52,6 +52,7 @@ def test_transport_failure_logs_exact_request_bytes_and_class_chain(caplog): provider="openai-codex", session_id="", _client_log_context=lambda: "", + _buffer_diagnostic_status=lambda message: None, ) with caplog.at_level(logging.WARNING, logger="agent.codex_runtime"): diff --git a/tests/agent/test_run_agent_codex_responses.py b/tests/agent/test_run_agent_codex_responses.py index 7de957f44b..f9b635a658 100644 --- a/tests/agent/test_run_agent_codex_responses.py +++ b/tests/agent/test_run_agent_codex_responses.py @@ -2958,6 +2958,31 @@ def test_run_codex_stream_prestream_retry_exhaustion_logs_telemetry( assert "attempt=2/2" in message +def test_run_codex_stream_prestream_exhaustion_buffers_one_user_line_with_host_attempts_size(monkeypatch): + """#97548: when the pre-stream connect retries are spent the user gets ONE line naming the + endpoint host, the attempt count and the serialized request size (agent.log was the only place + those lived), and re-entering the stream call from the outer retry loop does not add a copy.""" + import httpx + from openai import APIConnectionError + + agent = _build_agent(monkeypatch) + body = b'{"model":"gpt-5-codex","input":"' + b"x" * (829 * 1024) + b'"}' + request = httpx.Request("POST", "https://api.example.com/backend-api/codex/responses", content=body) + agent.client = SimpleNamespace(responses=SimpleNamespace( + create=lambda **kwargs: _raise_prestream_transport_error(request))) + + for _outer_retry in range(2): + with pytest.raises(APIConnectionError): + agent._run_codex_stream(_codex_request_kwargs()) + + lines = [str(msg) for _kind, msg in agent._retry_status_buffer] + assert len(lines) == 1, lines + assert "api.example.com" in lines[0] + assert "after 2 attempts" in lines[0] + assert f"request {round(len(body) / 1024)} KB" in lines[0] + assert "reject requests this large" in lines[0] + + def _codex_truncated_tool_call_response(): """``status=incomplete`` (max_output_tokens) whose function_call item was cut mid-arguments and settled as ``completed`` — the self-hosted /v1/responses shape from #91770.""" diff --git a/website/docs/reference/faq.md b/website/docs/reference/faq.md index 83c9c1a555..9d96c5474d 100644 --- a/website/docs/reference/faq.md +++ b/website/docs/reference/faq.md @@ -230,6 +230,12 @@ To isolate the source: See [Security](../user-guide/security.md) for Hermes' documented execution controls and [Providers](../integrations/providers.md) for provider configuration. +#### "Could not open a stream to `` after N attempts (request X KB)" + +**Meaning:** every connect attempt to that endpoint failed before a single stream event arrived, so nothing was billed; the normal retry/fallback chain still runs afterwards. The line names the host actually contacted, how many attempts were made, and the serialized request size — the three things that separate an outage from a request-size limit. + +**Solution:** if the request is large (hundreds of KB — long coding sessions reach this once the context grows) and short new chats work, the endpoint or a proxy in front of it is likely rejecting bodies that size: raise its body limit, or run `/compress` to shrink the context. If the request is small, the endpoint is unreachable — check the `base_url`, then retry with `/retry`. `logs/agent.log` records the exception chain for each attempt. + #### `/model` only shows one provider / can't switch providers **Cause:** `/model` (inside a chat session) can only switch between providers you've **already configured**. If you've only set up OpenRouter, that's all `/model` will show.