diff --git a/agent/agent_runtime_helpers.py b/agent/agent_runtime_helpers.py index 1d7e287c30bb1..46af2f9bb1223 100644 --- a/agent/agent_runtime_helpers.py +++ b/agent/agent_runtime_helpers.py @@ -25,7 +25,7 @@ from agent.credential_pool import ( STATUS_EXHAUSTED, credential_pool_matches_provider, resolve_runtime_pool_key ) -from agent.error_classifier import FailoverReason +from agent.error_classifier import TRANSPORT_ERROR_TYPES, FailoverReason from agent.turn_context import drop_stale_api_content from utils import base_url_host_matches, base_url_hostname, env_var_enabled, atomic_json_write logger = logging.getLogger(__name__) @@ -1195,10 +1195,16 @@ def _load_primary_pool(): # Transient transport failures worth one more attempt with a rebuilt client / connection pool. -_TRANSIENT_TRANSPORT_ERRORS = frozenset({ - "ReadTimeout", "ConnectTimeout", "PoolTimeout", "ConnectError", "RemoteProtocolError", - "APIConnectionError", "APITimeoutError", -}) +# Derived from the canonical classifier so the gate cannot drift from "what counts as a +# transport fault": a hand-maintained copy omitted ``ReadError`` — the shape a mid-response +# ``[Errno 104] Connection reset by peer`` takes on routes that read the body themselves — +# so the one connection the reset poisoned was never retired and every reset skipped the +# rebuild. This widens WHICH failures get the single rebuild, not how many attempts anything +# gets: recovery still fires at most once per API-call block, after the classifier already +# called the error retryable and the normal retry budget is spent. +# ``PoolTimeout`` stays a local addition: it is a fault of the pool this function rebuilds, +# which is narrower than the classifier's question. +_TRANSIENT_TRANSPORT_ERRORS = TRANSPORT_ERROR_TYPES | {"PoolTimeout"} _INLINE_REASONING_PATTERNS = tuple( re.compile(rf"<{tag}>(.*?)", re.DOTALL | re.IGNORECASE) for tag in ("think", "thinking", "thought", "reasoning", "REASONING_SCRATCHPAD") diff --git a/agent/error_classifier.py b/agent/error_classifier.py index 7bf3b8716f15d..c3cc34aa2b88e 100644 --- a/agent/error_classifier.py +++ b/agent/error_classifier.py @@ -285,7 +285,11 @@ def billing_unverified(self) -> bool: # SSL names keep provider-wrapped SSL errors (chain lost) as transport, not # unknown; OpenAI SDK errors are not subclasses of Python builtins. -_TRANSPORT_ERROR_TYPES = frozenset({ +# Public: this is the canonical answer to "is this exception type a transport fault?". +# ``agent/agent_runtime_helpers.py`` derives its client-rebuild gate from it so the two +# cannot drift (they did: ``ReadError`` — the fleet's most common reset shape — was +# classified as transport here but missing from the rebuild gate). +TRANSPORT_ERROR_TYPES = frozenset({ "ReadTimeout", "ConnectTimeout", "PoolTimeout", "ConnectError", "RemoteProtocolError", "ConnectionError", "ConnectionResetError", "ConnectionAbortedError", "BrokenPipeError", "TimeoutError", "ReadError", "ServerDisconnectedError", @@ -571,7 +575,7 @@ def _by_transport(c: _Ctx) -> Optional[Verdict]: # network call): as ``unknown`` it would burn every retry instantly. if c.error_type == "RuntimeError" and "consecutive stale attempts" in msg and "aborting this call" in msg: return _v(_R.timeout, **_ABORT_FALLBACK) - transport = c.error_type in _TRANSPORT_ERROR_TYPES or isinstance(c.error, (TimeoutError, ConnectionError, OSError)) + transport = c.error_type in TRANSPORT_ERROR_TYPES or isinstance(c.error, (TimeoutError, ConnectionError, OSError)) return _V_TIMEOUT if transport else None diff --git a/gateway/run.py b/gateway/run.py index f59ae79a82717..3569ef48541bf 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -391,6 +391,25 @@ def _gateway_surface_passes_raw_text(platform: Any) -> bool: r"cannot\s+connect", r"failed\s+to\s+establish", r"could\s+not\s+connect") _GATEWAY_CONNECTION_ERROR_RE = re.compile("(" + "|".join(_CONNECTION_ERROR_MARKERS) + ")", re.IGNORECASE) +# An ESTABLISHED connection died mid-transfer. Says nothing about whether the endpoint is up — +# in the reported incident an earlier call in the same turn had already been answered by it. +_CONNECTION_INTERRUPTED_MARKERS = ( + r"connection\s+reset", r"connection\s+aborted", r"errno\s+104", r"errno\s+103", + r"broken\s+pipe", r"server\s+disconnected", r"peer\s+closed\s+connection", + r"connection\s+was\s+closed", r"network\s+connection\s+lost", r"unexpected\s+eof", + r"incomplete\s+chunked\s+read", r"response\s+ended\s+prematurely", r"socket\s+hang\s+up", + r"(?:\w+\.)?remoteprotocolerror", r"(?:\w+\.)?readerror") +_GATEWAY_CONNECTION_INTERRUPTED_RE = re.compile( + "(" + "|".join(_CONNECTION_INTERRUPTED_MARKERS) + ")", re.IGNORECASE) + +# Nothing accepted the connection / no path to the host: "the endpoint is not up" IS the diagnosis. +_ENDPOINT_UNREACHABLE_MARKERS = ( + r"connection\s+refused", r"actively\s+refused", r"winerror\s+10061", r"errno\s+111", + r"no\s+route\s+to\s+host", r"network\s+is\s+unreachable", r"cannot\s+connect", + r"failed\s+to\s+establish", r"could\s+not\s+connect", r"(?:\w+\.)?connect\s*(?:error|timeout)") +_GATEWAY_ENDPOINT_UNREACHABLE_RE = re.compile( + "(" + "|".join(_ENDPOINT_UNREACHABLE_MARKERS) + ")", re.IGNORECASE) + _GATEWAY_SECRET_PATTERNS = ( re.compile(r"\bsk-[A-Za-z0-9][A-Za-z0-9_\-]{12,}\b"), re.compile(r"\bgh[pousr]_[A-Za-z0-9_]{20,}\b"), re.compile(r"\bxapp-\d+-[A-Za-z0-9\-]{20,}\b"), @@ -600,14 +619,32 @@ def _format_exec_approval_fallback( + ", ".join(choices[:-1]) + f", or {choices[-1]}.") # Ordered: auth beats policy beats rate-limit beats connection; first match wins. +# +# The three connection rows are NOT interchangeable, and collapsing them is what made a +# mid-response TCP reset read as "your model server is down" (2026-09-06) while the same turn +# had already been answered by that endpoint: +# * interrupted — an established connection died mid-transfer. Whether the endpoint is up is +# unknown from this alone, so we must not guess. +# * unreachable — nothing accepted the connection at all; "not running / unreachable" is the +# actual diagnosis, and this is the case that wording was written for (#86570). +# * ambiguous — connection-shaped but the cause was flattened away (an SDK-wrapped +# ``APIConnectionError: Connection error.`` keeps neither). Name both +# possibilities; assert neither. _PROVIDER_ERROR_REPLIES = ( (_GATEWAY_AUTH_ERROR_RE, "⚠️ Provider authentication failed. Check the configured credentials; " "raw provider details are in the gateway logs."), (_GATEWAY_PROVIDER_POLICY_RE, "⚠️ The model provider rejected the request. I kept the raw provider " "error out of chat; check gateway logs for details or try rephrasing."), (_GATEWAY_RATE_LIMIT_RE, "⏱️ The model provider is rate-limiting requests. Please wait a moment and try again."), - (_GATEWAY_CONNECTION_ERROR_RE, "⚠️ The model server is not responding — it looks like the configured " - "model endpoint is not running or is unreachable.")) + (_GATEWAY_CONNECTION_INTERRUPTED_RE, "⚠️ The connection to the model provider was interrupted before the " + "reply arrived, and the retries hit the same problem. Please try " + "again; the transport details are in the gateway logs."), + (_GATEWAY_ENDPOINT_UNREACHABLE_RE, "⚠️ The model server is not responding — it looks like the configured " + "model endpoint is not running or is unreachable."), + (_GATEWAY_CONNECTION_ERROR_RE, "⚠️ The model request could not be completed over the network after " + "retries — the connection was either never established or dropped before " + "the reply arrived. Please try again; if it keeps happening, check that the " + "configured model endpoint is reachable. Details are in the gateway logs.")) def _gateway_provider_error_reply(text: str) -> str: @@ -689,7 +726,14 @@ def _prepare_gateway_status_message(platform: Any, event_type: str, message: str ): return None if _looks_like_gateway_provider_error(text): - return _gateway_provider_error_reply(text) + # The status and the turn's final response are the SAME failure: the agent's terminal + # paths (agent/turn_recovery.py) emit the envelope through status_callback and then + # return it as final_response, and _sanitize_gateway_final_response maps both onto the + # same reply here — so delivering this one posts the identical bubble twice + # (2026-09-06: three connection resets, one warning shown twice). Chat surfaces get the + # provider-failure category from the final response only; the raw diagnostic still goes + # to the logs and to the programmatic surfaces above. + return None return text diff --git a/tests/gateway/test_connection_error_reply_wording.py b/tests/gateway/test_connection_error_reply_wording.py new file mode 100644 index 0000000000000..a9f8f417f6501 --- /dev/null +++ b/tests/gateway/test_connection_error_reply_wording.py @@ -0,0 +1,139 @@ +"""Connection failures must be described accurately on chat surfaces. + +A mid-response `[Errno 104] Connection reset by peer` is NOT evidence that the +configured endpoint is down — the same turn had already completed an earlier call +against it. Telling the user "the endpoint is not running" sends them debugging a +server that is running (reported 2026-09-06). + +The refusal case (nothing listening: ECONNREFUSED / WinError 10061 / no route) IS a +"your endpoint is not up" diagnosis and must keep it, as must the auth, policy and +rate-limit classifications. +""" + +import pytest + +from gateway.config import Platform +from gateway.run import ( + _gateway_provider_error_reply, + _sanitize_gateway_final_response, +) + +CHAT_PLATFORMS = [Platform.TELEGRAM, "slack", "feishu"] + +# Envelopes the agent's terminal path actually produces for a mid-transfer drop: an +# established connection died, which says nothing about whether the endpoint is up. +INTERRUPTED_ENVELOPES = [ + "API call failed after 3 retries: httpx.ReadError: [Errno 104] Connection reset by peer", + "API call failed after 3 retries: ConnectionResetError: [Errno 104] Connection reset by peer", + "API call failed after 3 retries: httpx.RemoteProtocolError: peer closed connection " + "without sending complete message body", + "API call failed after 3 retries: httpx.ReadError: server disconnected without sending a response", +] + +# Envelopes that really do mean "nothing is listening at the configured endpoint". +UNREACHABLE_ENVELOPES = [ + "API call failed after 3 retries: httpx.ConnectError: [Errno 111] Connection refused", + "API call failed after 3 retries: ConnectionError: [WinError 10061] No connection could " + "be made because the target machine actively refused it", + "API call failed after 3 retries: httpx.ConnectError: [Errno 113] No route to host", +] + +# Connection-shaped, but the SDK flattened the cause away — neither diagnosis is supported. +AMBIGUOUS_ENVELOPES = [ + "API call failed after 3 retries: openai.APIConnectionError: Connection error.", + "❌ API failed after 3 retries — openai.APIConnectionError: Connection error.", +] + +# Wording that ASSERTS the endpoint is down / not started. +_DOWN_DIAGNOSIS_CLAIMS = ("is not running", "not responding", "is unreachable", "not started") + + +def _claims_endpoint_is_down(reply: str) -> bool: + return any(claim in reply.lower() for claim in _DOWN_DIAGNOSIS_CLAIMS) + + +@pytest.mark.parametrize("envelope", INTERRUPTED_ENVELOPES) +def test_interrupted_connection_is_not_diagnosed_as_a_dead_endpoint(envelope): + """A reset/dropped response must be reported as an interruption, not a dead server.""" + reply = _gateway_provider_error_reply(envelope) + + assert reply + assert not _claims_endpoint_is_down(reply), reply + + +@pytest.mark.parametrize("envelope", UNREACHABLE_ENVELOPES) +def test_refused_connection_keeps_the_endpoint_down_diagnosis(envelope): + """A refused/unroutable connect is exactly the case the old wording was written for.""" + reply = _gateway_provider_error_reply(envelope) + + assert _claims_endpoint_is_down(reply), reply + + +@pytest.mark.parametrize("envelope", AMBIGUOUS_ENVELOPES) +def test_ambiguous_connection_error_asserts_neither_cause(envelope): + """An SDK-flattened ``Connection error.`` supports no diagnosis — so make none.""" + reply = _gateway_provider_error_reply(envelope) + + assert not _claims_endpoint_is_down(reply), reply + # It still has to be actionable: the user is told what to check, not what is broken. + assert "reachable" in reply.lower(), reply + + +def test_the_three_connection_causes_are_distinct_categories(): + """The causes must not collapse onto one message again.""" + interrupted = {_gateway_provider_error_reply(e) for e in INTERRUPTED_ENVELOPES} + unreachable = {_gateway_provider_error_reply(e) for e in UNREACHABLE_ENVELOPES} + ambiguous = {_gateway_provider_error_reply(e) for e in AMBIGUOUS_ENVELOPES} + + assert len(interrupted) == len(unreachable) == len(ambiguous) == 1 + assert len(interrupted | unreachable | ambiguous) == 3 + + +@pytest.mark.parametrize( + "envelope, expected_marker", + [ + ( + "API call failed after 3 retries: HTTP 401 Unauthorized: incorrect api key provided", + "authentication", + ), + ( + "API call failed after 3 retries: HTTP 429: rate limit exceeded for this model", + "rate-limiting", + ), + ( + "API call failed after 3 retries: HTTP 400: request blocked under the provider " + "safety policy", + "rejected", + ), + ], + ids=["auth", "rate_limit", "policy"], +) +def test_other_provider_error_classifications_are_preserved(envelope, expected_marker): + """Auth beats policy beats rate-limit beats connection — that ordering still holds.""" + reply = _gateway_provider_error_reply(envelope) + + assert expected_marker in reply.lower(), reply + + +@pytest.mark.parametrize("platform", CHAT_PLATFORMS) +def test_reset_final_response_is_sanitized_and_secret_free(platform): + """The reset envelope still goes through redaction + the safe-category rewrite.""" + raw = ( + "API call failed after 3 retries: httpx.ReadError: [Errno 104] Connection reset " + "by peer (Authorization: Bearer sk-ABCDEF0123456789abcdef0123)" + ) + + sanitized = _sanitize_gateway_final_response(platform, raw) + + assert "sk-ABCDEF" not in sanitized + assert "Errno 104" not in sanitized + assert not _claims_endpoint_is_down(sanitized), sanitized + assert sanitized.strip() + + +@pytest.mark.parametrize("platform", ["local", "api_server", "webhook"]) +def test_programmatic_surfaces_keep_the_raw_connection_error(platform): + """CLI/API consumers still need the bottom exception, not a chat-safe category.""" + raw = "API call failed after 3 retries: httpx.ReadError: [Errno 104] Connection reset by peer" + + assert _sanitize_gateway_final_response(platform, raw) == raw diff --git a/tests/gateway/test_local_model_connection_reply.py b/tests/gateway/test_local_model_connection_reply.py index ae52e318a71ca..1bfe9aac5162f 100644 --- a/tests/gateway/test_local_model_connection_reply.py +++ b/tests/gateway/test_local_model_connection_reply.py @@ -11,8 +11,8 @@ class TestGatewayConnectionErrorReply: def test_connection_error_strings_produce_specific_reply(self): + """A connect that was REFUSED/unroutable is the local-endpoint-down case.""" samples = [ - "openai.APIConnectionError", "httpx.ConnectError: connection refused", "ConnectionError: [WinError 10061] No connection could be made", "Errno 111 Connection refused", @@ -24,6 +24,21 @@ def test_connection_error_strings_produce_specific_reply(self): assert "not responding" in reply.lower(), text assert "not running or is unreachable" in reply, text + def test_cause_free_connection_error_points_at_reachability_without_asserting_it(self): + """A bare ``APIConnectionError`` kept no cause: it may equally be a dropped reply. + + Still a provider-error envelope, and still points the user at endpoint + reachability — but it no longer *states* that the endpoint is down (issue #15: + the same wording was shown for a mid-response TCP reset from a live endpoint). + """ + text = "openai.APIConnectionError" + + assert _looks_like_gateway_provider_error(text) + reply = _gateway_provider_error_reply(text) + + assert "reachable" in reply.lower() + assert "not running or is unreachable" not in reply + def test_broad_connection_phrases_still_map_once_classified(self): """Reply selector keeps the full phrase set; the gate does not.""" for text in ( diff --git a/tests/gateway/test_telegram_noise_filter.py b/tests/gateway/test_telegram_noise_filter.py index 3c7b292bf95ac..f97729106205e 100644 --- a/tests/gateway/test_telegram_noise_filter.py +++ b/tests/gateway/test_telegram_noise_filter.py @@ -255,8 +255,14 @@ def test_chat_gateways_drop_interrupt_sentinel(platform): assert _sanitize_gateway_final_response("local", sentinel) == sentinel -def test_telegram_status_sanitizes_raw_provider_security_errors(): - """Provider policy/security bodies should be replaced before chat delivery.""" +def test_telegram_status_never_delivers_raw_provider_security_errors(): + """Provider policy/security bodies must not reach chat through the status lane. + + The terminal failure envelope is delivered ONCE, as the turn's final response + (test_telegram_final_response_sanitizes_raw_provider_errors below) — the status + copy is dropped so the user does not get the same warning twice. Whatever the + delivery decision, the raw body/request id must never survive it. + """ raw = ( "❌ API failed after 3 retries — HTTP 400: request blocked because " "Operation contains cybersecurity risk. request_id=req_123" @@ -264,11 +270,9 @@ def test_telegram_status_sanitizes_raw_provider_security_errors(): sanitized = _prepare_gateway_status_message(Platform.TELEGRAM, "lifecycle", raw) - assert sanitized is not None - assert "provider rejected" in sanitized.lower() - assert "cybersecurity risk" not in sanitized.lower() - assert "HTTP 400" not in sanitized - assert "req_123" not in sanitized + assert sanitized is None + # The programmatic surfaces still get the full diagnostic. + assert _prepare_gateway_status_message("local", "lifecycle", raw) == raw def test_telegram_final_response_sanitizes_raw_provider_errors(): diff --git a/tests/gateway/test_terminal_api_failure_single_notice.py b/tests/gateway/test_terminal_api_failure_single_notice.py new file mode 100644 index 0000000000000..5694e237b4110 --- /dev/null +++ b/tests/gateway/test_terminal_api_failure_single_notice.py @@ -0,0 +1,238 @@ +"""A terminal provider failure must produce exactly ONE user-visible chat bubble. + +Reported 2026-09-06: three `httpx.ReadError: [Errno 104] Connection reset by peer` +attempts exhausted the retry budget on a Telegram turn and the user saw the *same* +sanitized connection warning twice. + +The duplicate is structural, not platform-specific: ``max_retries_exhausted_result`` +emits a terminal status through ``status_callback`` AND returns a ``final_response`` +carrying the same failure envelope. On chat surfaces the gateway rewrites both through +the same provider-error reply table, so both bubbles render identically. + +These tests drive the REAL terminal path (real ``AIAgent``, real classifier, real +``httpx.ReadError``) and pipe its two outputs through the REAL gateway filters, so they +assert the delivery contract rather than any particular wording. +""" + +import httpx +import pytest +from unittest.mock import MagicMock, patch + +from agent.error_classifier import classify_api_error +from agent.turn_recovery import max_retries_exhausted_result +from gateway.config import Platform +from gateway.run import ( + _prepare_gateway_status_message, + _sanitize_gateway_final_response, +) +from run_agent import AIAgent + +CHAT_PLATFORMS = [Platform.TELEGRAM, "slack", "feishu", "irc"] +PROGRAMMATIC_SURFACES = ["local", "api_server", "webhook"] + + +def _make_agent(): + """A minimal real AIAgent (no network, no live endpoint probing).""" + with ( + patch("model_tools.get_tool_definitions", return_value=[]), + patch("model_tools.check_toolset_requirements", return_value={}), + patch("agent.process_bootstrap.OpenAI"), + patch( + "agent.context_compressor.get_model_context_length", + return_value=200_000, + ), + ): + agent = AIAgent( + api_key="test-key-12345678", + base_url="https://my-llm.example.com/v1", + provider="custom", + model="test-model", + quiet_mode=True, + skip_context_files=True, + skip_memory=True, + ) + agent.client = MagicMock() + agent._persist_session = MagicMock() + return agent + + +def _run_terminal_failure(agent, api_error, *, is_rate_limited=False): + """Run the real retries-exhausted terminal path; return (statuses, final_response).""" + statuses: list[tuple[str, str]] = [] + agent.status_callback = lambda kind, message: statuses.append((kind, message)) + classified = classify_api_error(api_error) + result = max_retries_exhausted_result( + agent, + api_error, + classified, + max_retries=3, + is_rate_limited=is_rate_limited, + error_msg=str(api_error).lower(), + api_kwargs=None, + api_messages=[{"role": "user", "content": "hi"}], + messages=[{"role": "user", "content": "hi"}], + conversation_history=None, + api_call_count=2, + approx_tokens=51_987, + provider="openai-codex", + base_url="https://example.invalid/backend-api/codex", + model="test-model", + ) + return statuses, result["final_response"] + + +def _delivered_bubbles(platform, statuses, final_response): + """Everything this platform would actually render for the failed turn.""" + bubbles = [ + prepared + for kind, message in statuses + if (prepared := _prepare_gateway_status_message(platform, kind, message)) + ] + final = _sanitize_gateway_final_response(platform, final_response) + if final: + bubbles.append(final) + return bubbles + + +def _reset_error(): + """The exact live failure shape: a mid-response read reset.""" + return httpx.ReadError("[Errno 104] Connection reset by peer") + + +@pytest.mark.parametrize("platform", CHAT_PLATFORMS) +def test_connection_reset_terminal_failure_renders_one_bubble(platform): + """The reported shape: retries exhausted on a connection reset.""" + agent = _make_agent() + statuses, final_response = _run_terminal_failure(agent, _reset_error()) + + # The agent still reports the failure on both channels — that is its contract. + assert statuses, "terminal failure must still emit a status for local/CLI surfaces" + assert final_response + + bubbles = _delivered_bubbles(platform, statuses, final_response) + + assert len(bubbles) == 1, f"user saw {len(bubbles)} bubbles: {bubbles}" + assert "reset by peer" not in bubbles[0].lower() + + +@pytest.mark.parametrize( + "api_error", + [ + httpx.ReadError("[Errno 104] Connection reset by peer"), + httpx.RemoteProtocolError("peer closed connection without sending complete message body"), + httpx.ConnectError("[Errno 111] Connection refused"), + RuntimeError("HTTP 500: internal server error"), + ], + ids=["read_reset", "remote_protocol", "connect_refused", "http_500"], +) +def test_every_terminal_failure_shape_renders_one_bubble(api_error): + """No terminal provider failure envelope may be delivered twice.""" + agent = _make_agent() + statuses, final_response = _run_terminal_failure(agent, api_error) + + bubbles = _delivered_bubbles(Platform.TELEGRAM, statuses, final_response) + + assert len(bubbles) == 1, f"user saw {len(bubbles)} bubbles: {bubbles}" + + +def test_rate_limited_terminal_failure_renders_one_bubble(): + """The rate-limit terminal variant must keep the same single-bubble contract.""" + agent = _make_agent() + error = RuntimeError("HTTP 429: too many requests, please retry after 20s") + statuses, final_response = _run_terminal_failure(agent, error, is_rate_limited=True) + + bubbles = _delivered_bubbles(Platform.TELEGRAM, statuses, final_response) + + assert len(bubbles) == 1, f"user saw {len(bubbles)} bubbles: {bubbles}" + assert "rate" in bubbles[0].lower() or "limit" in bubbles[0].lower() + + +def test_non_retryable_terminal_failure_renders_one_bubble(): + """The sibling terminal path (non-retryable 4xx) has the same two outputs.""" + from agent.turn_recovery import nonretryable_client_error_result + + agent = _make_agent() + statuses: list[tuple[str, str]] = [] + agent.status_callback = lambda kind, message: statuses.append((kind, message)) + api_error = RuntimeError("HTTP 401: Incorrect API key provided: sk-live_abcdefghij0123456789") + api_error.status_code = 401 + + result = nonretryable_client_error_result( + agent, + api_error, + classify_api_error(api_error), + status_code=401, + api_kwargs=None, + api_messages=[{"role": "user", "content": "hi"}], + messages=[{"role": "user", "content": "hi"}], + conversation_history=None, + api_call_count=1, + approx_tokens=100, + provider="openai-codex", + base_url="https://example.invalid/backend-api/codex", + model="test-model", + ) + + bubbles = _delivered_bubbles(Platform.TELEGRAM, statuses, result["final_response"]) + + assert len(bubbles) == 1, f"user saw {len(bubbles)} bubbles: {bubbles}" + assert "sk-live_" not in bubbles[0] + + +@pytest.mark.parametrize("surface", PROGRAMMATIC_SURFACES) +def test_programmatic_surfaces_keep_both_diagnostic_channels(surface): + """CLI/API/webhook consumers keep the raw status stream AND the raw final.""" + agent = _make_agent() + statuses, final_response = _run_terminal_failure(agent, _reset_error()) + + prepared = [ + _prepare_gateway_status_message(surface, kind, message) + for kind, message in statuses + ] + + assert [p for p in prepared if p] == [m for _, m in statuses] + assert _sanitize_gateway_final_response(surface, final_response) == final_response + + +@pytest.mark.parametrize("platform", CHAT_PLATFORMS) +def test_progress_and_unrelated_warnings_still_reach_chat(platform): + """Suppressing the duplicate must not silence ordinary work/warning statuses.""" + keepers = [ + ("lifecycle", "still on it"), + ("lifecycle", "⏳ Working — 3 min"), + ("warn", "⚠️ The command you approved touches files outside the workspace."), + ("warn", "🛑 Tool budget exhausted for this turn."), + ] + + for kind, message in keepers: + assert _prepare_gateway_status_message(platform, kind, message) == message + + +@pytest.mark.parametrize("platform", CHAT_PLATFORMS) +def test_context_overflow_warning_still_reaches_chat(platform): + """The context-overflow advisory is a real user warning, not a provider envelope.""" + from agent.conversation_compression import ( + CONTEXT_OVERFLOW_BLOCKED_WARNING_TEMPLATE, + ) + + message = CONTEXT_OVERFLOW_BLOCKED_WARNING_TEMPLATE.format( + tokens=85_000, threshold=72_000, reason="ineffective" + ) + + assert _prepare_gateway_status_message(platform, "warn", message) == message + + +@pytest.mark.parametrize("platform", CHAT_PLATFORMS) +def test_terminal_failure_status_never_leaks_provider_detail(platform): + """Whatever the delivery decision, raw provider bodies/secrets never reach chat.""" + raw = ( + "❌ API failed after 3 retries — HTTP 401 Unauthorized: " + "Authorization: Bearer sk-ABCDEF0123456789abcdef0123 request_id=req_9" + ) + + prepared = _prepare_gateway_status_message(platform, "lifecycle", raw) + + if prepared is not None: + assert "sk-ABCDEF" not in prepared + assert "req_9" not in prepared + assert "HTTP 401" not in prepared diff --git a/tests/run_agent/test_readerror_reset_recovery.py b/tests/run_agent/test_readerror_reset_recovery.py new file mode 100644 index 0000000000000..8948d39cd61de --- /dev/null +++ b/tests/run_agent/test_readerror_reset_recovery.py @@ -0,0 +1,242 @@ +"""Reset recovery, proven against a real socket that really resets the connection. + +Live evidence (2026-09-06): 64 of 67 failed provider attempts surfaced as a bare +``httpx.ReadError: [Errno 104] Connection reset by peer``. The Responses/Codex route +consumes the SSE body itself, so a mid-response reset escapes the OpenAI SDK +un-wrapped — the retry loop sees ``ReadError``, not ``APIConnectionError``. + +``try_recover_primary_transport`` is the one place that retires the poisoned httpx +connection pool and rebuilds the client before giving up, and it is gated on a typed +allow-list. If that list disagrees with the canonical transport classifier, the exact +error the fleet actually hits skips the rebuild. + +This test injects a REAL reset from a local listener (``SO_LINGER 0`` → RST): the +exception is produced by httpx, not constructed, and the recovery it drives goes on to +make a REAL successful request through the rebuilt client. No provider credentials. +""" + +import json +import socket +import struct +import threading +import time + +import httpx +import openai +import pytest +from unittest.mock import MagicMock, patch + +from agent.error_classifier import classify_api_error +from run_agent import AIAgent + +_COMPLETION_BODY = json.dumps( + { + "id": "chatcmpl-recovered", + "object": "chat.completion", + "created": 1, + "model": "test-model", + "choices": [ + { + "index": 0, + "message": {"role": "assistant", "content": "recovered"}, + "finish_reason": "stop", + } + ], + "usage": {"prompt_tokens": 1, "completion_tokens": 1, "total_tokens": 2}, + } +).encode() + +# Generous enough that a loaded runner still has the client parked on the read when the +# RST lands; the whole harness pays it once per injected fault. +_ABORT_DELAY_SECONDS = 0.25 + + +class ResettingEndpoint: + """Local HTTP listener that aborts an ARMED response mid-stream with a TCP RST. + + Arming is explicit rather than "reset the first request" so the injected fault always + lands on the request the test is measuring, whatever else happens to touch the port. + """ + + def __init__(self) -> None: + self._sock = socket.socket() + self._sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + self._sock.bind(("127.0.0.1", 0)) + self._sock.listen(16) + self.port = self._sock.getsockname()[1] + self.connections = 0 + self.resets = 0 + self._armed = threading.Event() + self._stop = threading.Event() + self._thread = threading.Thread(target=self._serve, daemon=True) + self._thread.start() + + @property + def base_url(self) -> str: + return f"http://127.0.0.1:{self.port}/v1" + + def arm(self) -> None: + """Reset the next model request; every request after it is served normally.""" + self._armed.set() + + def _serve(self) -> None: + while not self._stop.is_set(): + try: + conn, _ = self._sock.accept() + except OSError: + return + self.connections += 1 + try: + request = conn.recv(65536) + if b"/completions" in request and self._armed.is_set(): + self._armed.clear() + self.resets += 1 + self._abort(conn) + else: + self._respond(conn) + except OSError: + pass + finally: + try: + conn.close() + except OSError: + pass + + def _abort(self, conn: socket.socket) -> None: + """Open the response, emit one frame, then RST the connection mid-body.""" + conn.sendall( + b"HTTP/1.1 200 OK\r\nContent-Type: text/event-stream\r\n\r\n" + b'data: {"id":"c","object":"chat.completion.chunk","created":1,' + b'"model":"test-model","choices":[{"index":0,"delta":{"content":"hi"}}]}\n\n' + ) + # Let the client drain that frame and block on the next read first. Resetting while + # the frame is still in flight races: the client can consume it and see the close as + # a clean end-of-body, which is a *different* failure (silent truncation), not the + # mid-read reset under test. + time.sleep(_ABORT_DELAY_SECONDS) + # SO_LINGER with a zero timeout makes close() send RST instead of FIN. + conn.setsockopt(socket.SOL_SOCKET, socket.SO_LINGER, struct.pack("ii", 1, 0)) + + def _respond(self, conn: socket.socket) -> None: + conn.sendall( + b"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: " + + str(len(_COMPLETION_BODY)).encode() + + b"\r\nConnection: close\r\n\r\n" + + _COMPLETION_BODY + ) + + def close(self) -> None: + self._stop.set() + try: + self._sock.close() + except OSError: + pass + + +@pytest.fixture +def endpoint(): + server = ResettingEndpoint() + try: + yield server + finally: + server.close() + + +def _make_agent(base_url, *, provider="custom"): + with ( + patch("model_tools.get_tool_definitions", return_value=[]), + patch("model_tools.check_toolset_requirements", return_value={}), + patch("agent.process_bootstrap.OpenAI"), + patch("agent.context_compressor.get_model_context_length", return_value=200_000), + ): + agent = AIAgent( + api_key="test-key-12345678", + base_url=base_url, + provider=provider, + model="test-model", + quiet_mode=True, + skip_context_files=True, + skip_memory=True, + ) + return agent + + +def _provoke_real_reset(endpoint, client) -> BaseException: + """Drive a real request against the resetting endpoint and return what was raised.""" + endpoint.arm() + with pytest.raises(BaseException) as excinfo: # noqa: PT011 - the type IS the assertion + stream = client.chat.completions.create( + model="test-model", + messages=[{"role": "user", "content": "hello"}], + stream=True, + ) + for _chunk in stream: + pass + return excinfo.value + + +def test_local_fault_injection_reproduces_the_reported_error_shape(endpoint): + """Sanity: the harness really produces the live bottom exception, unmocked.""" + client = openai.OpenAI(api_key="k", base_url=endpoint.base_url, max_retries=0) + + error = _provoke_real_reset(endpoint, client) + + assert isinstance(error, httpx.ReadError) + assert "reset by peer" in str(error).lower() + assert classify_api_error(error).retryable is True + + +def test_reset_drives_client_rebuild_and_a_successful_follow_up(endpoint): + """The reported failure must reach the transport-recovery path and recover.""" + agent = _make_agent(endpoint.base_url) + poisoned = openai.OpenAI(api_key="test-key-12345678", base_url=endpoint.base_url, max_retries=0) + agent.client = poisoned + retired: list = [] + agent._retire_shared_openai_client = lambda client, reason=None: retired.append(client) + + error = _provoke_real_reset(endpoint, poisoned) + assert isinstance(error, httpx.ReadError) + assert endpoint.resets == 1 + + with patch("agent.agent_runtime_helpers.time.sleep"): + recovered = agent._try_recover_primary_transport(error, retry_count=3, max_retries=3) + + assert recovered is True, "a real connection reset must get the one-shot client rebuild" + assert retired == [poisoned], "the poisoned pool must be retired, not left checked out" + assert agent.client is not poisoned + + # Credentials, model and route are unchanged by recovery. + assert agent.model == "test-model" + assert agent.provider == "custom" + assert str(agent.client.base_url).rstrip("/") == endpoint.base_url.rstrip("/") + assert agent.client.api_key == "test-key-12345678" + + # The rebuilt client is a working client, not just a new object. + response = agent.client.chat.completions.create( + model="test-model", messages=[{"role": "user", "content": "again"}] + ) + assert response.choices[0].message.content == "recovered" + assert endpoint.connections >= 2 + + +def test_recovery_is_bounded_and_respects_client_ownership(endpoint): + """Recovery never fires once fallback owns the client, and never for aggregators.""" + error = httpx.ReadError("[Errno 104] Connection reset by peer") + + on_fallback = _make_agent(endpoint.base_url) + on_fallback.client = MagicMock() + on_fallback._fallback_activated = True + assert on_fallback._try_recover_primary_transport(error, retry_count=3, max_retries=3) is False + + aggregator = _make_agent("https://openrouter.ai/api/v1", provider="openrouter") + aggregator.client = MagicMock() + assert aggregator._try_recover_primary_transport(error, retry_count=3, max_retries=3) is False + + +def test_non_transport_errors_still_skip_the_rebuild(endpoint): + """The fix must not turn the typed allow-list into "retry everything".""" + agent = _make_agent(endpoint.base_url) + agent.client = MagicMock() + + for error in (ValueError("bad request shape"), KeyError("model")): + assert agent._try_recover_primary_transport(error, retry_count=3, max_retries=3) is False