diff --git a/app/core/clients/proxy_websocket.py b/app/core/clients/proxy_websocket.py index 962e0c5cf1..1a9d6a8f5f 100644 --- a/app/core/clients/proxy_websocket.py +++ b/app/core/clients/proxy_websocket.py @@ -78,6 +78,9 @@ rf"|{_LIVE_CALL_UUID_CORE})" ) _LIVE_CALL_ID_PATTERN = re.compile(rf"{REALTIME_LIVE_CALL_ID_ROUTE_REGEX}\Z") +UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE = "upstream_websocket_liveness_timeout" +_WEBSOCKETS_KEEPALIVE_TIMEOUT_REASON = "keepalive ping timeout" +_AIOHTTP_HEARTBEAT_TIMEOUT_PREFIX = "No PONG received after " class RealtimeWebSocketProtocol(StrEnum): @@ -99,16 +102,19 @@ class _UpstreamWebSocketPolicy: preserve_close_semantics: bool +# Responses turns may be silent at the application layer for minutes, but a +# healthy transport still answers ping control frames. Keep both watchdogs on: +# disabling them turns a black-holed VPN route into a multi-hour request stall. _RESPONSES_WEBSOCKET_POLICY = _UpstreamWebSocketPolicy( operation="responses websocket", include_responses_beta=True, archive_payloads=True, - enable_routed_heartbeat=False, + enable_routed_heartbeat=True, retry_handshake_status=True, preserve_handshake_status=False, credential_safe_connect_errors=False, retry_routed_network_errors=True, - enable_direct_ping_timeout=False, + enable_direct_ping_timeout=True, preserve_close_semantics=False, ) _LIVE_SIDEBAND_WEBSOCKET_POLICY = _UpstreamWebSocketPolicy( @@ -176,6 +182,8 @@ def __init__(self, message: str, *, error_code: str) -> None: def _websocket_transport_error_code(exc: BaseException, *, uses_proxy: bool) -> str: + if _is_websocket_liveness_timeout(exc): + return UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE return process_network_error_code( exc, fallback="upstream_unavailable", @@ -183,17 +191,64 @@ def _websocket_transport_error_code(exc: BaseException, *, uses_proxy: bool) -> ) +def is_account_neutral_websocket_error_code(error_code: str | None) -> bool: + """Return whether transport provenance rules out an account-health penalty.""" + + # These failures occur below the selected account's application protocol. + # They follow an ambiguous send, so relay owners must fail rather than + # replay while leaving the account eligible for unrelated requests. Keep + # the compatibility keepalive code here as long as adapters can emit it. + return error_code in { + PROCESS_NETWORK_UNAVAILABLE_CODE, + UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + "upstream_keepalive_timeout", + } + + +def _is_websocket_liveness_timeout(exc: BaseException) -> bool: + if isinstance(exc, ConnectionClosedError): + # websockets emits this locally-sent 1011 when its own ping watchdog + # expires. A peer may acknowledge it, leaving both close frames on the + # exception; send-first ordering still proves the marker came from our + # watchdog without trusting a peer that sends the same code and reason. + return ( + exc.sent is not None + and int(exc.sent.code) == 1011 + and exc.sent.reason == _WEBSOCKETS_KEEPALIVE_TIMEOUT_REASON + and (exc.rcvd is None or exc.rcvd_then_sent is False) + ) + # aiohttp surfaces its heartbeat watchdog through WSMsgType.ERROR with a + # ServerTimeoutError carrying this library-defined prefix. + return isinstance(exc, aiohttp.ServerTimeoutError) and str(exc).startswith(_AIOHTTP_HEARTBEAT_TIMEOUT_PREFIX) + + +def _aiohttp_stored_liveness_exception(websocket: Any) -> Exception | None: + # When aiohttp's heartbeat expires between receive() calls, no waiter is + # available for WSMsgType.ERROR. aiohttp stores the timeout instead and the + # next receive returns CLOSED, so every post-connect path must consult it. + exception_getter = getattr(websocket, "exception", None) + if not callable(exception_getter): + return None + exception = exception_getter() + return exception if isinstance(exception, Exception) and _is_websocket_liveness_timeout(exception) else None + + def _relay_receive_error_code(error_code: str) -> str | None: - """Expose only account-neutral process failures across the adapter boundary.""" + """Expose account-neutral transport failures across the adapter boundary.""" # Relay owners map an absent code to their established stream_incomplete # contract. Leaking the adapter's generic fallback would bypass that path. - return error_code if error_code in {PROCESS_NETWORK_UNAVAILABLE_CODE, "upstream_keepalive_timeout"} else None + return error_code if is_account_neutral_websocket_error_code(error_code) else None def _is_keepalive_timeout_close(exc: ConnectionClosedError) -> bool: """Classify peer/proxy heartbeat failures without exposing socket details.""" + # Treat the legacy text marker as trusted only when this endpoint initiated + # the close. A peer can send the same public code and reason, so peer-first + # ordering must retain ordinary close/error semantics. + if exc.sent is None or (exc.rcvd is not None and exc.rcvd_then_sent is not False): + return False reason = _close_reason_from_exception(exc) return "keepalive ping timeout" in f"{exc} {reason or ''}".lower() @@ -282,11 +337,12 @@ async def receive(self) -> UpstreamWebSocketMessage: ) error_code = _websocket_transport_error_code(exc, uses_proxy=self._uses_proxy) await _rotate_after_websocket_network_failure(error_code) - relay_error_code = ( - "upstream_keepalive_timeout" - if _is_keepalive_timeout_close(exc) - else _relay_receive_error_code(error_code) - ) + relay_error_code = _relay_receive_error_code(error_code) + if relay_error_code is None and _is_keepalive_timeout_close(exc): + # Prefer the stable, provenance-checked watchdog code above. + # This text fallback preserves compatibility with keepalive + # failures whose exception shape lacks the local-send marker. + relay_error_code = "upstream_keepalive_timeout" # ConnectionClosedError describes an incomplete close handshake, # not generic transport provenance. Let Responses relay owners map # it to stream_incomplete while live relays preserve received closes. @@ -357,7 +413,8 @@ async def send_text(self, text: str) -> None: if asyncio.iscoroutine(result): await result except Exception as exc: - await _raise_websocket_send_error(exc, endpoint_id=self._endpoint_id, uses_proxy=True) + classification_exc = _aiohttp_stored_liveness_exception(self._websocket) or exc + await _raise_websocket_send_error(classification_exc, endpoint_id=self._endpoint_id, uses_proxy=True) async def send_bytes(self, data: bytes) -> None: try: @@ -365,27 +422,43 @@ async def send_bytes(self, data: bytes) -> None: if asyncio.iscoroutine(result): await result except Exception as exc: - await _raise_websocket_send_error(exc, endpoint_id=self._endpoint_id, uses_proxy=True) + classification_exc = _aiohttp_stored_liveness_exception(self._websocket) or exc + await _raise_websocket_send_error(classification_exc, endpoint_id=self._endpoint_id, uses_proxy=True) async def receive(self) -> UpstreamWebSocketMessage: try: msg = await self._websocket.receive() except Exception as exc: - error_code = _websocket_transport_error_code(exc, uses_proxy=True) + classification_exc = _aiohttp_stored_liveness_exception(self._websocket) or exc + error_code = _websocket_transport_error_code(classification_exc, uses_proxy=True) await _rotate_after_websocket_network_failure(error_code) return UpstreamWebSocketMessage( kind="error", - error=codex_transport_error_message("websocket receive", self._endpoint_id, exc), + error=codex_transport_error_message("websocket receive", self._endpoint_id, classification_exc), error_code=_relay_receive_error_code(error_code), ) if msg.type in (aiohttp.WSMsgType.CLOSE, aiohttp.WSMsgType.CLOSING, aiohttp.WSMsgType.CLOSED): + liveness_exception = _aiohttp_stored_liveness_exception(self._websocket) + if liveness_exception is not None: + return UpstreamWebSocketMessage( + kind="error", + close_code=_aiohttp_ws_close_code(self._websocket, msg), + error=codex_transport_error_message( + "websocket receive", + self._endpoint_id, + liveness_exception, + ), + error_code=UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + ) return UpstreamWebSocketMessage( kind="close", close_code=_aiohttp_ws_close_code(self._websocket, msg), close_reason=_aiohttp_ws_close_reason(msg), ) if msg.type == aiohttp.WSMsgType.ERROR: - exception = msg.data if isinstance(msg.data, BaseException) else None + exception = ( + msg.data if isinstance(msg.data, Exception) else _aiohttp_stored_liveness_exception(self._websocket) + ) error_code = ( _websocket_transport_error_code(exception, uses_proxy=True) if exception is not None @@ -845,11 +918,8 @@ async def _connect_upstream_websocket( settings.upstream_websocket_proxy_env() if hasattr(settings, "upstream_websocket_proxy_env") else os.environ ) proxy_url = resolve_websocket_proxy_from_env(url, proxy_env) if settings.upstream_websocket_trust_env else None - # Long Responses turns can spend minutes without application frames, - # so that existing transport keeps its own watchdog disabled. Live - # sideband traffic uses ping/pong liveness rather than an application- - # frame idle timeout because WebRTC media may remain healthy while the - # sideband itself is silent. + # Ping/pong control frames verify transport liveness without treating valid + # application-frame silence as an idle response. ping_timeout = ( settings.proxy_downstream_websocket_idle_timeout_seconds if policy.enable_direct_ping_timeout else None ) diff --git a/app/modules/proxy/_service/http_bridge/request_submit.py b/app/modules/proxy/_service/http_bridge/request_submit.py index 83dc8d56ff..a84ea4aec6 100644 --- a/app/modules/proxy/_service/http_bridge/request_submit.py +++ b/app/modules/proxy/_service/http_bridge/request_submit.py @@ -38,7 +38,11 @@ from app.core.clients.proxy import codex_control_request as core_codex_control_request # noqa: F401 from app.core.clients.proxy import compact_responses as core_compact_responses # noqa: F401 from app.core.clients.proxy import transcribe_audio as core_transcribe_audio # noqa: F401 -from app.core.clients.proxy_websocket import UpstreamWebSocketTransportError +from app.core.clients.proxy_websocket import ( + UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + UpstreamWebSocketTransportError, + is_account_neutral_websocket_error_code, +) from app.core.errors import ( openai_error, ) @@ -250,6 +254,26 @@ async def _send_http_bridge_request_text_with_archive_id( reset_request_id(token) +async def _settle_claimed_http_bridge_liveness_failure( + service: Any, + session: "_HTTPBridgeSession", + *, + error_message: str, +) -> None: + """Finish the pending-deque settlement claimed beside a failed send.""" + + if session.liveness_settlement_owner != "send": + raise RuntimeError("HTTP bridge liveness settlement started without the send claim") + async with session.lifecycle_lock: + await service._fail_http_bridge_reader_and_maybe_retire( + session, + error_code=UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + error_message=error_message, + penalize_account=False, + force_retire=True, + ) + + def _text_with_account_installation_id(text_data: str, codex_installation_id: str | None) -> str: payload = json.loads(text_data) if not isinstance(payload, dict): @@ -1160,7 +1184,7 @@ async def _submit_http_bridge_request_with_handoff( upstream_send_started = True try: await _send_http_bridge_request_text_with_archive_id(session, request_state, text_data) - except BaseException: + except BaseException as exc: request_state.recovery_attempt_dispatched = True # Publish retirement while lifecycle ownership is still # held; a gate waiter must never reuse an ambiguously sent @@ -1168,6 +1192,15 @@ async def _submit_http_bridge_request_with_handoff( session.closed = True session.upstream_control.reconnect_requested = True session.upstream_control.retire_after_drain = True + if ( + isinstance(exc, UpstreamWebSocketTransportError) + and exc.error_code == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE + ): + # Only this narrow claim, not ``closed``, tells the + # reader that the submitter will settle siblings. + # Keep it inside lifecycle_lock with the failing + # send so the reader cannot observe an ownership gap. + session.claim_liveness_settlement() raise request_state.recovery_attempt_dispatched = True session.last_used_at = _service_time().monotonic() @@ -1248,31 +1281,54 @@ async def _submit_http_bridge_request_with_handoff( # handed to the kernel. Never reconnect-and-resend from this path; # only failures proven to precede dispatch may be replayed. error_code = exc.error_code if isinstance(exc, UpstreamWebSocketTransportError) else "stream_incomplete" - account_neutral = error_code == "proxy_network_unavailable" - await self._cleanup_http_bridge_submit_interruption( - session, - request_state=request_state, - gate_acquired=gate_acquired, - request_enqueued=request_enqueued, - counted_in_queue=True, - admission_waiter_registered=admission_waiter_registered, - ) - await self._fail_pending_websocket_requests( - account=session.account, - account_id_value=session.account.id, - pending_requests=deque([request_state]), - pending_lock=anyio.Lock(), - error_code=error_code, - error_message=str(exc) or "Upstream websocket closed before response.completed", - api_key=None, - response_create_gate=session.response_create_gate, - penalize_account=not account_neutral, - ) - session.closed = True - try: - await session.upstream.close() - except Exception: - logger.debug("Failed to close HTTP bridge upstream websocket after send failure", exc_info=True) + # Liveness expiry and local network loss are transport failures, + # not evidence against the selected account. Keep this in sync + # with the reader path's shared provenance classification. + account_neutral = is_account_neutral_websocket_error_code(error_code) + if error_code == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE: + # The sender claimed ownership beside the failing send while + # holding lifecycle_lock. It therefore owns the entire session + # deque, including older in-flight requests; settling only this + # request would strand its siblings after the reader yields. + # Publish the cleanup task before the first await after the + # claim. Shielding it makes cancellation wait for settlement, + # so the claim can never outlive its exactly-once owner. + settlement_task = asyncio.create_task( + _settle_claimed_http_bridge_liveness_failure( + self, + session, + error_message=str(exc) or "Upstream websocket liveness failed", + ), + name="http-bridge-liveness-send-settlement", + ) + _, settlement_cancellation = await _await_task_deferring_cancellation(settlement_task) + if settlement_cancellation is not None: + raise settlement_cancellation + else: + await self._cleanup_http_bridge_submit_interruption( + session, + request_state=request_state, + gate_acquired=gate_acquired, + request_enqueued=request_enqueued, + counted_in_queue=True, + admission_waiter_registered=admission_waiter_registered, + ) + await self._fail_pending_websocket_requests( + account=session.account, + account_id_value=session.account.id, + pending_requests=deque([request_state]), + pending_lock=anyio.Lock(), + error_code=error_code, + error_message=str(exc) or "Upstream websocket closed before response.completed", + api_key=None, + response_create_gate=session.response_create_gate, + penalize_account=not account_neutral, + ) + session.closed = True + try: + await session.upstream.close() + except Exception: + logger.debug("Failed to close HTTP bridge upstream websocket after send failure", exc_info=True) # Always raise 502 so the client can retry with # previous_response_id intact. Returning 400 # previous_response_not_found causes the client to drop diff --git a/app/modules/proxy/_service/http_bridge/upstream_events.py b/app/modules/proxy/_service/http_bridge/upstream_events.py index 0444b55df3..6e8c1f47ef 100644 --- a/app/modules/proxy/_service/http_bridge/upstream_events.py +++ b/app/modules/proxy/_service/http_bridge/upstream_events.py @@ -29,7 +29,12 @@ from app.core.clients.proxy import codex_control_request as core_codex_control_request # noqa: F401 from app.core.clients.proxy import compact_responses as core_compact_responses # noqa: F401 from app.core.clients.proxy import transcribe_audio as core_transcribe_audio # noqa: F401 -from app.core.clients.proxy_websocket import UpstreamWebSocketMessage, UpstreamWebSocketTransportError +from app.core.clients.proxy_websocket import ( + UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + UpstreamWebSocketMessage, + UpstreamWebSocketTransportError, + is_account_neutral_websocket_error_code, +) from app.core.errors import response_failed_event from app.core.openai.models import OpenAIEvent from app.core.openai.parsing import parse_sse_event_payload @@ -1033,13 +1038,13 @@ async def _relay_http_bridge_upstream_messages( session.last_upstream_close_generation += 1 session.last_upstream_close_code = message.close_code retried = False - # A process-network receive failure does not prove that the - # upstream rejected response.create. Do not replay ordinary - # requests in that ambiguous case: the first request may - # still be executing and replay could duplicate work, billing, - # or tool side effects. Clean websocket closes remain eligible - # for the bounded pre-created retry path below. - if message.error_code != "proxy_network_unavailable": + # Account-neutral transport failures do not prove that the + # upstream rejected response.create. The request may still be + # executing, so replay could duplicate work, billing, or tool + # side effects. Clean closes remain eligible for the bounded + # pre-created retry circuit maintained by the session. + account_neutral = is_account_neutral_websocket_error_code(message.error_code) + if not account_neutral: retried = await self._retry_http_bridge_precreated_request(session) if retried: continue @@ -1049,6 +1054,15 @@ async def _relay_http_bridge_upstream_messages( else None ) async with session.lifecycle_lock: + if ( + session.liveness_settlement_owner == "send" + and message.error_code == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE + ): + # A submitter publishes this dedicated claim beside the + # failing send while holding lifecycle_lock. ``closed`` + # alone is only an admission/retirement state and must + # never suppress settlement of still-pending siblings. + break await self._fail_http_bridge_reader_and_maybe_retire( session, error_code=message.error_code or "stream_incomplete", @@ -1061,16 +1075,15 @@ async def _relay_http_bridge_upstream_messages( else "websocket_transport_error" ), penalize_account=( - message.error_code != "proxy_network_unavailable" - and message.error_code != "upstream_keepalive_timeout" - and not ( - message.kind == "close" - and _classify_upstream_close( - message.close_code, - response_events_seen=response_events_seen, - ) - == "clean" - ) + not account_neutral and not (message.kind == "close" and close_classification == "clean") + ), + **( + # An admission waiter must not inherit a socket whose + # heartbeat already proved it dead. Other failures + # preserve the existing deferred-retirement handoff. + {"force_retire": True} + if message.error_code == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE + else {} ), ) break @@ -1084,18 +1097,27 @@ async def _relay_http_bridge_upstream_messages( exc_info=True, ) error_code = exc.error_code if isinstance(exc, UpstreamWebSocketTransportError) else "stream_incomplete" - account_neutral = error_code in {"proxy_network_unavailable", "upstream_keepalive_timeout"} + account_neutral = is_account_neutral_websocket_error_code(error_code) async with session.lifecycle_lock: - await self._fail_http_bridge_reader_and_maybe_retire( - session, - error_code=error_code, - error_message=( - str(exc) - if isinstance(exc, UpstreamWebSocketTransportError) - else "HTTP bridge upstream reader crashed before response.completed" - ), - penalize_account=not account_neutral, - ) + if not ( + session.liveness_settlement_owner == "send" + and error_code == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE + ): + # Match the message path above when receive() raises while + # a concurrent send failure already owns settlement. + await self._fail_http_bridge_reader_and_maybe_retire( + session, + error_code=error_code, + error_message=( + str(exc) + if isinstance(exc, UpstreamWebSocketTransportError) + else "HTTP bridge upstream reader crashed before response.completed" + ), + penalize_account=not account_neutral, + # Preserve ordinary crash handoff behavior, but never hand + # a heartbeat-expired socket to an admission waiter. + **({"force_retire": True} if error_code == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE else {}), + ) finally: await _cancel_http_bridge_reader_child( wakeup_task, diff --git a/app/modules/proxy/_service/support.py b/app/modules/proxy/_service/support.py index 134f10abca..56287be5ae 100644 --- a/app/modules/proxy/_service/support.py +++ b/app/modules/proxy/_service/support.py @@ -1030,6 +1030,13 @@ class _HTTPBridgeSession: last_upstream_close_code: int | None = None last_upstream_close_generation: int = 0 closed: bool = False + # ``closed`` rejects new admissions but is written by many unrelated + # retirement paths; it never proves that a sender owns pending settlement. + # Only the submitter may claim this, while holding ``lifecycle_lock``, when + # its own send reports a liveness timeout. The reader remains the default + # settlement owner for every other close, including an already-closed + # session whose still-running transport later loses heartbeat liveness. + liveness_settlement_owner: Literal["send"] | None = None # Set when the session proved silent/wedged (reattached stream with # response events but no ``response.created``, or repeated eventless # timeouts). A quarantined session must never be selected for reuse or @@ -1049,6 +1056,17 @@ class _HTTPBridgeSession: upstream_proxy_fallback_used: bool | None = None upstream_proxy_fail_closed_reason: str | None = None + def claim_liveness_settlement(self) -> bool: + """Claim whole-deque settlement for a liveness-failed submitter. + + The caller must hold ``lifecycle_lock`` across the failing send and + this synchronous claim so the reader cannot settle the same deque. + """ + + if self.liveness_settlement_owner is None: + self.liveness_settlement_owner = "send" + return self.liveness_settlement_owner == "send" + def _complete_http_bridge_handoff( session: _HTTPBridgeSession, diff --git a/app/modules/proxy/_service/websocket/mixin.py b/app/modules/proxy/_service/websocket/mixin.py index 50308cbc78..371654853b 100644 --- a/app/modules/proxy/_service/websocket/mixin.py +++ b/app/modules/proxy/_service/websocket/mixin.py @@ -61,6 +61,7 @@ UpstreamWebSocket, UpstreamWebSocketTransportError, filter_inbound_websocket_headers, + is_account_neutral_websocket_error_code, ) from app.core.errors import ( OpenAIErrorEnvelope, @@ -1099,7 +1100,14 @@ async def _process_upstream_websocket_transport_end( replay_refusal_reasons: list[str] = [] replay_request_state = None message_error_code = getattr(message, "error_code", None) - if message_error_code != "proxy_network_unavailable": + # A classified local transport failure says nothing about whether an + # already-sent response.create was accepted. Keep it account-neutral and + # terminal: replay here could duplicate work, billing, or tool side effects. + account_neutral = is_account_neutral_websocket_error_code(message_error_code) + if account_neutral: + if any(state.last_downstream_sequence_number is not None for state in reader_owned): + replay_refusal_reasons.append("sequenced_downstream_frame") + else: replay_request_state = await _pop_replayable_precreated_websocket_request_state( reader_owned, pending_lock=anyio.Lock(), @@ -1131,7 +1139,7 @@ async def _process_upstream_websocket_transport_end( client_send_lock=client_send_lock, response_create_gate=response_create_gate, downstream_activity=downstream_activity, - penalize_account=message_error_code != "proxy_network_unavailable", + penalize_account=not account_neutral, suppress_sequenced_downstream_errors=sequenced_downstream_replay_refused, ) # A terminal receive can race the outer session cleanup, especially when @@ -2349,7 +2357,7 @@ def take_reader_replay_request_state() -> _WebSocketRequestState | None: client_send_lock=client_send_lock, response_create_gate=response_create_gate, downstream_activity=downstream_activity, - penalize_account=exc.error_code != "proxy_network_unavailable", + penalize_account=not is_account_neutral_websocket_error_code(exc.error_code), ), name="proxy-websocket-finalization-transport-send-failure", ) @@ -2373,7 +2381,7 @@ def take_reader_replay_request_state() -> _WebSocketRequestState | None: client_send_lock=client_send_lock, response_create_gate=response_create_gate, downstream_activity=downstream_activity, - penalize_account=exc.error_code != "proxy_network_unavailable", + penalize_account=not is_account_neutral_websocket_error_code(exc.error_code), suppress_sequenced_downstream_errors=sequenced_downstream_replay_refused, ) if sequenced_downstream_replay_refused: diff --git a/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/.openspec.yaml b/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/.openspec.yaml new file mode 100644 index 0000000000..e08b5f89a2 --- /dev/null +++ b/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/.openspec.yaml @@ -0,0 +1,2 @@ +schema: spec-driven +created: 2026-08-03 diff --git a/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/design.md b/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/design.md new file mode 100644 index 0000000000..15dcabe9fc --- /dev/null +++ b/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/design.md @@ -0,0 +1,69 @@ +## Context + +Responses upstream WebSockets use two existing transports: `websockets` for direct egress and aiohttp for routed egress. Both libraries already support ping/pong liveness detection, and the same connection policy already enables it for Realtime live sideband traffic. The Responses policy disables both mechanisms because application-frame silence is valid during long turns. That distinction is unnecessary for ping/pong: control-frame replies prove transport liveness without requiring an application event. + +After a VPN disconnect, an established TCP connection can be black-holed without an immediate DNS, route, or socket exception. Downstream keepalives then keep the client attached while the upstream request remains pending. Once a request frame has been sent, however, the proxy cannot know whether upstream accepted it before connectivity was lost, so transparent replay could duplicate work or side effects. + +## Goals / Non-Goals + +**Goals:** + +- Bound silent upstream Responses WebSocket failure detection on both direct and routed egress. +- Preserve a stable classification from transport adapter through the direct WebSocket and HTTP bridge owners. +- Settle pending work once, without replaying an ambiguously delivered request or penalizing a healthy account. +- Ensure subsequent client retries open a fresh upstream connection and therefore use the current host route. + +**Non-Goals:** + +- Add a host route watcher, VPN-specific integration, background recovery coordinator, or proactive socket migration. +- Change the long Responses request budget or application-event idle timeout. +- Automatically resume a turn whose upstream acceptance is unknown. +- Add a configuration setting, dependency, persistence change, or operator-facing UI. + +## Decisions + +### Reuse transport ping/pong support and the existing timeout + +Enable `heartbeat` for routed aiohttp Responses sockets and `ping_timeout` for direct `websockets` Responses sockets. Both values come from `proxy_downstream_websocket_idle_timeout_seconds`, matching the existing live-sideband policy and keeping the fix zero-config. + +Application-level watchdogs were rejected because valid Responses turns may be silent for minutes and downstream synthetic keepalives are not evidence of upstream health. A host-network watcher was rejected because it is platform-specific, cannot reliably enumerate every network transition, and duplicates transport-layer failure detection. + +### Give liveness expiry a narrow stable classification + +Map only library-specific ping/pong timeout signals to `upstream_websocket_liveness_timeout`: the `websockets` locally sent close with reason `keepalive ping timeout`, and aiohttp's `ServerTimeoutError` produced by its heartbeat watchdog. Ordinary upstream closes and other receive exceptions keep their current behavior. + +A shared code predicate identifies account-neutral WebSocket failures so relay owners do not duplicate string comparisons as new neutral conditions are added. + +### Fail closed after ambiguous delivery + +Both relay owners treat the classified liveness timeout like a post-send network failure: no transparent replay, no account-health write, exact-once pending-request settlement, and retirement of the affected socket. The downstream error remains retryable at the client boundary, where a fresh client connection can safely establish a new upstream route under the client's existing retry semantics. + +The HTTP bridge cannot infer settlement ownership from `session.closed`: that +flag also rejects admission after continuity-persistence and other failures +that settle only the submitting request. A submitter claims whole-deque +liveness settlement explicitly while holding the session lifecycle lock around +the failing send. The reader skips its normal settlement only when that claim +exists; a later liveness expiry on an otherwise closed session still settles +every pending sibling. + +Once published, the send claim owns the whole pending deque. The submitter +therefore starts settlement as a shielded child before its next await and +defers caller cancellation until that child completes; cancellation cannot +leave a permanent claim whose reader has already yielded. + +For example, if `response.create` was written before a VPN route disappears and no pong returns, codex-lb emits a terminal `upstream_websocket_liveness_timeout` failure for that pending request and closes or retires the upstream session. It does not resend `response.create` on another account or socket. + +## Risks / Trade-offs + +- [Some intermediaries do not answer WebSocket pings correctly] → Use the established, configurable timeout already deployed for live sideband traffic; operators can adjust the existing value if needed. +- [Detection is not instantaneous] → A bounded delay is preferable to false positives and is far shorter than the multi-hour Responses request budget. +- [The client must retry the interrupted turn] → This avoids duplicate model work and tool side effects when delivery status is unknowable. +- [Library wording could change] → Pin direct detection to the concrete close code/reason emitted by the installed `websockets` API and drive a real no-pong expiry in integration coverage, in addition to adapter unit tests. + +## Migration Plan + +No data or configuration migration is required. Deploy normally; rollback restores the previous policy flags and classification behavior. + +## Open Questions + +None. diff --git a/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/proposal.md b/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/proposal.md new file mode 100644 index 0000000000..2d90716341 --- /dev/null +++ b/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/proposal.md @@ -0,0 +1,27 @@ +## Why + +An established upstream Responses WebSocket can remain silently black-holed after a host network transition such as VPN disconnection. Because Responses sockets currently disable the transports' existing ping/pong liveness checks, downstream keepalives can keep a conversation waiting until the much longer request timeout instead of terminating promptly so the client can reconnect. + +## What Changes + +- Enable the existing WebSocket transport liveness checks for direct and routed upstream Responses connections, using the current downstream WebSocket idle-timeout setting as the liveness budget. +- Classify transport-detected ping/pong timeouts with a stable internal error code. +- Treat a post-send liveness timeout as account-neutral and terminal for pending work, without transparent replay, because upstream request acceptance is ambiguous. +- Retire the affected upstream socket so a later client retry establishes a fresh network route. +- Add transport, direct WebSocket relay, and HTTP bridge regression coverage for the liveness and settlement invariants. + +## Capabilities + +### New Capabilities + +None. + +### Modified Capabilities + +- `responses-api-compat`: Require bounded upstream Responses WebSocket liveness detection and safe handling of liveness timeouts across direct WebSocket and HTTP bridge clients. + +## Impact + +- Affects the shared upstream WebSocket adapters and both Responses relay owners. +- Reuses existing aiohttp heartbeat, websockets ping timeout, and configuration; no new dependency, setting, API, schema, migration, or dashboard surface is introduced. +- A stalled conversation terminates after the configured liveness budget and relies on the downstream client to retry on a fresh connection. diff --git a/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/specs/responses-api-compat/spec.md b/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/specs/responses-api-compat/spec.md new file mode 100644 index 0000000000..7d21be2d96 --- /dev/null +++ b/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/specs/responses-api-compat/spec.md @@ -0,0 +1,77 @@ +## ADDED Requirements + +### Requirement: Responses upstream websocket liveness is bounded + +The proxy MUST configure direct and routed upstream Responses WebSocket transports with finite ping/pong liveness detection derived from `proxy_downstream_websocket_idle_timeout_seconds`. When an established Responses WebSocket is terminated because its transport did not receive the required pong, the adapter MUST classify the failure as `upstream_websocket_liveness_timeout`. Direct WebSocket and HTTP bridge relay owners MUST treat that failure as account neutral, MUST NOT transparently replay a pending request whose delivery is ambiguous, MUST finalize its pending request ownership exactly once, and MUST retire the affected upstream socket so a later client retry opens a fresh connection. An HTTP bridge reader MUST suppress its own pending-deque settlement only when a concurrent submitter explicitly claimed liveness-settlement ownership under the session lifecycle lock; `session.closed` alone MUST NOT suppress settlement. + +#### Scenario: Direct Responses websocket loses pong liveness + +- **GIVEN** a direct upstream Responses WebSocket has been established +- **WHEN** the `websockets` keepalive watchdog terminates it after a pong timeout +- **THEN** the pending request fails with `upstream_websocket_liveness_timeout` +- **AND** the request is not transparently replayed +- **AND** the selected account receives no failure-health signal +- **AND** the affected upstream socket is retired + +#### Scenario: Routed Responses websocket loses pong liveness + +- **GIVEN** a routed upstream Responses WebSocket has been established for an HTTP bridge or direct WebSocket client +- **WHEN** the aiohttp heartbeat watchdog terminates it after a pong timeout +- **THEN** the pending request fails with `upstream_websocket_liveness_timeout` +- **AND** the request is not transparently replayed +- **AND** the selected account receives no failure-health signal +- **AND** the affected upstream socket is retired + +#### Scenario: Long turn remains healthy through control frames + +- **GIVEN** a Responses turn emits no application event within the liveness interval +- **WHEN** the upstream WebSocket continues replying to transport pings +- **THEN** the proxy keeps the upstream socket open +- **AND** the existing Responses request budget remains authoritative for the turn + +#### Scenario: Closed bridge without a sender claim later loses pong liveness + +- **GIVEN** an HTTP bridge session has multiple pending requests +- **AND** a separate submit failure marks the session closed without claiming liveness-settlement ownership +- **WHEN** the still-running upstream transport later expires its heartbeat +- **THEN** the reader settles every pending request with `upstream_websocket_liveness_timeout` +- **AND** the selected account receives no failure-health signal + +#### Scenario: Claimed bridge settlement survives submitter cancellation + +- **GIVEN** an HTTP bridge submitter claims liveness-settlement ownership after its send fails +- **WHEN** the submitter is cancelled before whole-deque settlement completes +- **THEN** settlement continues until every pending sibling is finalized exactly once +- **AND** the submitter cancellation is preserved after settlement completes + +## MODIFIED Requirements + +### Requirement: Upstream websocket drops penalize affected accounts +When an upstream websocket closes while one or more streamed response requests are pending and have not reached a terminal event, the proxy MUST record a transient upstream error for the account before signaling failure for those pending requests, except when the close carries a classified process-wide network failure or upstream WebSocket liveness timeout. A classified process-wide network failure or upstream WebSocket liveness timeout MUST remain account neutral and use its classified error code. For other closes, the proxy MUST surface `stream_incomplete` to affected pending requests except when a direct Responses WebSocket request has already successfully emitted a finite integer `sequence_number`. For that sequenced direct-WebSocket case, the proxy MUST record the request outcome as `stream_incomplete` without emitting a synthetic terminal frame under the active response id, then MUST close the downstream WebSocket with code 1011. + +#### Scenario: websocket closes before pending responses complete + +- **GIVEN** a streamed response request is pending on an upstream websocket +- **AND** the direct downstream response has not emitted a numeric sequence, or the request uses another transport +- **WHEN** the websocket closes before a terminal response event is observed +- **AND** the close does not carry a classified process-wide network failure or upstream WebSocket liveness timeout +- **THEN** the pending request fails with `stream_incomplete` +- **AND** the account receives a transient upstream failure signal for routing + +#### Scenario: sequenced direct websocket closes before completion + +- **GIVEN** a direct Responses WebSocket request has successfully emitted a finite integer `sequence_number` +- **WHEN** the upstream websocket closes before a terminal response event is observed +- **AND** the close does not carry a classified process-wide network failure or upstream WebSocket liveness timeout +- **THEN** the request is recorded as failed with `stream_incomplete` +- **AND** no synthetic terminal frame is emitted under the active response id +- **AND** the downstream WebSocket closes with code 1011 +- **AND** the account receives a transient upstream failure signal for routing + +#### Scenario: websocket liveness timeout remains account neutral + +- **GIVEN** a streamed response request is pending on an upstream websocket +- **WHEN** its transport reports `upstream_websocket_liveness_timeout` +- **THEN** the pending request fails with that classified error code +- **AND** the account receives no failure-health signal +- **AND** the request is not transparently replayed diff --git a/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/tasks.md b/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/tasks.md new file mode 100644 index 0000000000..6da3ddfffa --- /dev/null +++ b/openspec/changes/archive/2026-08-04-recover-responses-websocket-liveness/tasks.md @@ -0,0 +1,24 @@ +## 1. Transport liveness + +- [x] 1.1 Enable the existing finite heartbeat and ping timeout for routed and direct Responses WebSocket connections. +- [x] 1.2 Classify aiohttp heartbeat expiry and websockets keepalive expiry with the stable liveness-timeout code and a shared account-neutral predicate. + +## 2. Relay safety + +- [x] 2.1 Make direct Responses WebSocket liveness failures terminal, non-replayable, account-neutral, and fully settled. +- [x] 2.2 Apply the same no-replay, account-neutral, forced-retirement behavior to HTTP bridge upstream readers. + +## 3. Regression coverage + +- [x] 3.1 Cover direct and routed transport policy values and library-specific liveness classification. +- [x] 3.2 Cover direct WebSocket and HTTP bridge no-replay, account-health, settlement, and retirement invariants. + +## 4. Verification + +- [x] 4.1 Run focused tests, formatting, lint, type, architecture, and strict OpenSpec validation checks. + +## 5. Maintainer review follow-up + +- [x] 5.1 Replace the HTTP bridge reader's overloaded `closed` ownership guard with an explicit submitter liveness-settlement claim and cover closed-session sibling settlement. +- [x] 5.2 Drive a real installed-library no-pong timeout so integration coverage pins the `websockets` watchdog shape used by classification. +- [x] 5.3 Shield a submitter's claimed bridge liveness settlement from cancellation and cover the claim-to-settlement window. diff --git a/openspec/specs/responses-api-compat/context.md b/openspec/specs/responses-api-compat/context.md index 3eb8c29881..9a62f360da 100644 --- a/openspec/specs/responses-api-compat/context.md +++ b/openspec/specs/responses-api-compat/context.md @@ -32,6 +32,9 @@ See `openspec/specs/responses-api-compat/spec.md` for normative requirements. - `/v1/responses/compact` is supported only when the upstream implements it. - `prompt_cache_key` affinity on OpenAI-style routes is intentionally bounded by a dashboard-managed freshness window, unlike durable backend `session_id` or dashboard sticky-thread routing. - Codex-native direct websocket `/backend-api/codex/responses` treats upstream `previous_response_id` as an ephemeral anchor. If that anchor goes stale, the proxy must mask raw `previous_response_not_found` details and emit a sanitized `codex_previous_response_stale` classifier so compatible Codex clients can soft-reset and retry without `previous_response_id`. +- Upstream Responses WebSockets use transport ping/pong control frames to detect a black-holed connection without confusing valid application-event silence with an idle turn. Direct and routed connections reuse `proxy_downstream_websocket_idle_timeout_seconds` for this zero-config liveness budget. +- A post-send liveness timeout is delivery-ambiguous. It remains account-neutral, is never transparently replayed, and retires the affected upstream socket so a client retry opens a fresh route without risking duplicated model work or tool side effects. +- HTTP bridge settlement ownership is explicit: `closed` rejects new work but does not imply that a submitter owns existing siblings. Only a liveness-failed send claims whole-deque settlement under the lifecycle lock; otherwise the reader remains responsible for settling pending requests when the transport dies. ## Fast Mode and Service Tiers @@ -114,6 +117,7 @@ when upstream reports a different actual tier. - **Codex websocket reconnects:** Reconnect continuity now depends on the client replaying the accepted `x-codex-turn-state`; generated turn-state is emitted on accept for backend Codex routes and echoed back when the client already supplies one. - **Codex websocket stale previous-response anchors:** Direct backend Codex websocket stale-anchor failures are surfaced as `response.failed` / `codex_previous_response_stale` without the raw upstream code or missing `resp_...` id; OpenAI-compatible `/v1/responses` websocket clients continue to receive generic `stream_incomplete` masking. - **Websocket handshake forbidden/not-found:** Auto transport now fails loud on `403` / `404` instead of silently hiding the websocket regression behind HTTP fallback. +- **Upstream websocket stops answering pings:** Pending direct-WebSocket and HTTP-bridge work fails with `upstream_websocket_liveness_timeout`; the account remains healthy and the request is not replayed because upstream acceptance is unknown. - **Invalid request payloads:** Return 4xx with `invalid_request_error`. ## Error Envelope Mapping (Reference) @@ -174,5 +178,6 @@ OpenSpec change first. - Post-deploy: monitor `capacity_exhausted_active_sessions`, Codex-session bridge reuse/evict counts, websocket handshake 403/404 rates after the narrower auto-fallback policy, and backend Codex HTTP vs websocket cache-ratio gaps. - When tracing compact incidents, confirm that request logs and upstream logs show direct `/codex/responses/compact` usage without surrogate `/codex/responses` fallback. - Post-deploy: monitor `no_accounts`, `stream_incomplete`, and `upstream_unavailable`. +- Post-deploy: monitor `upstream_websocket_liveness_timeout`; recurring failures indicate a host route, VPN, proxy, or intermediary that black-holes established WebSockets. - Post-deploy: monitor `codex_previous_response_stale` on `/backend-api/codex/responses`; recurring spikes mean clients are still relying on stale upstream anchors and should perform the documented full-context retry without `previous_response_id`. - Websocket/Codex CLI tier verification runbook: `openspec/specs/responses-api-compat/ops.md` diff --git a/openspec/specs/responses-api-compat/spec.md b/openspec/specs/responses-api-compat/spec.md index 577c772514..328f750b84 100644 --- a/openspec/specs/responses-api-compat/spec.md +++ b/openspec/specs/responses-api-compat/spec.md @@ -83,15 +83,59 @@ The default compact request budget MUST be at least 180 seconds, and the default - **THEN** `compact_request_budget_seconds` is at least 180 seconds - **AND** `stream_idle_timeout_seconds` is at least 600 seconds +### Requirement: Responses upstream websocket liveness is bounded + +The proxy MUST configure direct and routed upstream Responses WebSocket transports with finite ping/pong liveness detection derived from `proxy_downstream_websocket_idle_timeout_seconds`. When an established Responses WebSocket is terminated because its transport did not receive the required pong, the adapter MUST classify the failure as `upstream_websocket_liveness_timeout`. Direct WebSocket and HTTP bridge relay owners MUST treat that failure as account neutral, MUST NOT transparently replay a pending request whose delivery is ambiguous, MUST finalize its pending request ownership exactly once, and MUST retire the affected upstream socket so a later client retry opens a fresh connection. An HTTP bridge reader MUST suppress its own pending-deque settlement only when a concurrent submitter explicitly claimed liveness-settlement ownership under the session lifecycle lock; `session.closed` alone MUST NOT suppress settlement. + +#### Scenario: Direct Responses websocket loses pong liveness + +- **GIVEN** a direct upstream Responses WebSocket has been established +- **WHEN** the `websockets` keepalive watchdog terminates it after a pong timeout +- **THEN** the pending request fails with `upstream_websocket_liveness_timeout` +- **AND** the request is not transparently replayed +- **AND** the selected account receives no failure-health signal +- **AND** the affected upstream socket is retired + +#### Scenario: Routed Responses websocket loses pong liveness + +- **GIVEN** a routed upstream Responses WebSocket has been established for an HTTP bridge or direct WebSocket client +- **WHEN** the aiohttp heartbeat watchdog terminates it after a pong timeout +- **THEN** the pending request fails with `upstream_websocket_liveness_timeout` +- **AND** the request is not transparently replayed +- **AND** the selected account receives no failure-health signal +- **AND** the affected upstream socket is retired + +#### Scenario: Long turn remains healthy through control frames + +- **GIVEN** a Responses turn emits no application event within the liveness interval +- **WHEN** the upstream WebSocket continues replying to transport pings +- **THEN** the proxy keeps the upstream socket open +- **AND** the existing Responses request budget remains authoritative for the turn + +#### Scenario: Closed bridge without a sender claim later loses pong liveness + +- **GIVEN** an HTTP bridge session has multiple pending requests +- **AND** a separate submit failure marks the session closed without claiming liveness-settlement ownership +- **WHEN** the still-running upstream transport later expires its heartbeat +- **THEN** the reader settles every pending request with `upstream_websocket_liveness_timeout` +- **AND** the selected account receives no failure-health signal + +#### Scenario: Claimed bridge settlement survives submitter cancellation + +- **GIVEN** an HTTP bridge submitter claims liveness-settlement ownership after its send fails +- **WHEN** the submitter is cancelled before whole-deque settlement completes +- **THEN** settlement continues until every pending sibling is finalized exactly once +- **AND** the submitter cancellation is preserved after settlement completes + ### Requirement: Upstream websocket drops penalize affected accounts -When an upstream websocket closes while one or more streamed response requests are pending and have not reached a terminal event, the proxy MUST record a transient upstream error for the account before signaling failure for those pending requests, except when the close carries a classified process-wide network failure. A classified process-wide network failure MUST remain account neutral and use its network error code. For other closes, the proxy MUST surface `stream_incomplete` to affected pending requests except when a direct Responses WebSocket request has already successfully emitted a finite integer `sequence_number`. For that sequenced direct-WebSocket case, the proxy MUST record the request outcome as `stream_incomplete` without emitting a synthetic terminal frame under the active response id, then MUST close the downstream WebSocket with code 1011. +When an upstream websocket closes while one or more streamed response requests are pending and have not reached a terminal event, the proxy MUST record a transient upstream error for the account before signaling failure for those pending requests, except when the close carries a classified process-wide network failure or upstream WebSocket liveness timeout. A classified process-wide network failure or upstream WebSocket liveness timeout MUST remain account neutral and use its classified error code. For other closes, the proxy MUST surface `stream_incomplete` to affected pending requests except when a direct Responses WebSocket request has already successfully emitted a finite integer `sequence_number`. For that sequenced direct-WebSocket case, the proxy MUST record the request outcome as `stream_incomplete` without emitting a synthetic terminal frame under the active response id, then MUST close the downstream WebSocket with code 1011. #### Scenario: websocket closes before pending responses complete - **GIVEN** a streamed response request is pending on an upstream websocket - **AND** the direct downstream response has not emitted a numeric sequence, or the request uses another transport - **WHEN** the websocket closes before a terminal response event is observed -- **AND** the close does not carry a classified process-wide network failure +- **AND** the close does not carry a classified process-wide network failure or upstream WebSocket liveness timeout - **THEN** the pending request fails with `stream_incomplete` - **AND** the account receives a transient upstream failure signal for routing @@ -99,11 +143,20 @@ When an upstream websocket closes while one or more streamed response requests a - **GIVEN** a direct Responses WebSocket request has successfully emitted a finite integer `sequence_number` - **WHEN** the upstream websocket closes before a terminal response event is observed +- **AND** the close does not carry a classified process-wide network failure or upstream WebSocket liveness timeout - **THEN** the request is recorded as failed with `stream_incomplete` - **AND** no synthetic terminal frame is emitted under the active response id - **AND** the downstream WebSocket closes with code 1011 - **AND** the account receives a transient upstream failure signal for routing +#### Scenario: websocket liveness timeout remains account neutral + +- **GIVEN** a streamed response request is pending on an upstream websocket +- **WHEN** its transport reports `upstream_websocket_liveness_timeout` +- **THEN** the pending request fails with that classified error code +- **AND** the account receives no failure-health signal +- **AND** the request is not transparently replayed + ### Requirement: Single HTTP bridge previous-response misses recover or fail closed When an HTTP bridge session receives an anonymous upstream `previous_response_not_found` error for a single pending follow-up request, the service MUST treat the error as an internal continuity-loss signal. It MUST either recover through the existing previous-response rebind path or rewrite the error to a retryable continuity failure instead of forwarding the raw upstream invalid-request error. diff --git a/tests/integration/test_proxy_websocket_responses.py b/tests/integration/test_proxy_websocket_responses.py index db39eeb9a7..b7c11cebcb 100644 --- a/tests/integration/test_proxy_websocket_responses.py +++ b/tests/integration/test_proxy_websocket_responses.py @@ -1,6 +1,8 @@ from __future__ import annotations import asyncio +import base64 +import hashlib import json import logging import threading @@ -16,11 +18,16 @@ from httpx import Headers from sqlalchemy import select from starlette.websockets import WebSocketDisconnect +from websockets.asyncio.client import connect as websocket_connect import app.modules.proxy.api as proxy_api_module import app.modules.proxy.service as proxy_module from app.core import shutdown as shutdown_state from app.core.auth.refresh import RefreshError +from app.core.clients.proxy_websocket import ( + UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + WebsocketsUpstreamWebSocket, +) from app.core.utils.request_id import get_request_id from app.db.models import Account, AccountStatus, ApiKeyUsageReservation, RequestLog from app.db.session import SessionLocal @@ -38,6 +45,58 @@ _REAL_WRITE_REQUEST_LOG = proxy_module.ProxyService._write_request_log +@pytest.mark.asyncio +async def test_real_websockets_keepalive_expiry_preserves_liveness_classification() -> None: + async def accept_without_answering_frames( + reader: asyncio.StreamReader, + writer: asyncio.StreamWriter, + ) -> None: + try: + request = await reader.readuntil(b"\r\n\r\n") + websocket_key = next( + line.split(b":", 1)[1].strip() + for line in request.split(b"\r\n") + if line.lower().startswith(b"sec-websocket-key:") + ) + accept = base64.b64encode(hashlib.sha1(websocket_key + b"258EAFA5-E914-47DA-95CA-C5AB0DC85B11").digest()) + writer.write( + b"HTTP/1.1 101 Switching Protocols\r\n" + b"Upgrade: websocket\r\n" + b"Connection: Upgrade\r\n" + b"Sec-WebSocket-Accept: " + accept + b"\r\n\r\n" + ) + await writer.drain() + # Read and discard every frame. In particular, never answer pings, + # so the production client watchdog must terminate the connection. + while await reader.read(65536): + pass + finally: + writer.close() + try: + await writer.wait_closed() + except ConnectionError: + pass + + server = await asyncio.start_server(accept_without_answering_frames, "127.0.0.1", 0) + async with server: + assert server.sockets + port = server.sockets[0].getsockname()[1] + async with websocket_connect( + f"ws://127.0.0.1:{port}", + ping_interval=0.05, + ping_timeout=0.05, + close_timeout=0.05, + proxy=None, + ) as connection: + upstream = WebsocketsUpstreamWebSocket(connection) + message = await asyncio.wait_for(upstream.receive(), timeout=1.0) + + # This assertion pins the actual close code/reason emitted by the installed + # websockets watchdog to the adapter's stable, account-neutral classifier. + assert message.kind == "error" + assert message.error_code == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE + + def _assert_previous_response_not_found_error(error: dict[str, object]) -> None: assert error["code"] == proxy_module.PREVIOUS_RESPONSE_NOT_FOUND_CODE assert error["message"] == proxy_module.PREVIOUS_RESPONSE_NOT_FOUND_MESSAGE @@ -107,6 +166,7 @@ def __init__(self, messages: list[_FakeUpstreamMessage]) -> None: self.archived_receive_request_ids: list[str | None] = [] self.archived_receive_texts: list[str] = [] self.closed = False + self.closed_event = threading.Event() self._messages: asyncio.Queue[_FakeUpstreamMessage] = asyncio.Queue() for message in messages: self._messages.put_nowait(message) @@ -130,6 +190,7 @@ def archive_received(self, message: _FakeUpstreamMessage) -> None: async def close(self) -> None: self.closed = True + self.closed_event.set() class _SequencedUpstreamWebSocket(_FakeUpstreamWebSocket): @@ -6201,6 +6262,9 @@ async def fake_resolve_previous_response_owner( for call in log_calls ) assert any(call["status"] == "success" and call["request_id"] == "resp_ws_inflight" for call in log_calls) + # TestClient runs the ASGI task in a worker thread. Wait for its owned + # cancellation cleanup instead of racing that thread on the plain flag. + assert fake_upstream.closed_event.wait(timeout=1.0) assert fake_upstream.closed is True diff --git a/tests/unit/test_proxy_http_bridge.py b/tests/unit/test_proxy_http_bridge.py index 68aede235e..ec32daa972 100644 --- a/tests/unit/test_proxy_http_bridge.py +++ b/tests/unit/test_proxy_http_bridge.py @@ -25,6 +25,7 @@ from app.core.auth.refresh import RefreshError from app.core.clients.proxy import CODEX_RESPONSES_LITE_WEBSOCKET_METADATA_KEY, ProxyResponseError from app.core.clients.proxy_websocket import ( + UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, CodexUpstreamWebSocket, UpstreamWebSocket, UpstreamWebSocketMessage, @@ -9762,6 +9763,31 @@ async def stubborn_child() -> None: assert cleanup_tasks == set() +@pytest.mark.asyncio +async def test_await_cancelled_task_allows_outer_cancellation_cleanup_to_finish() -> None: + cleanup_finished = asyncio.Event() + + async def cancelled_owner() -> None: + child = asyncio.create_task(asyncio.Event().wait()) + try: + await asyncio.Event().wait() + finally: + await proxy_service._await_cancelled_task( + child, + timeout_seconds=1.0, + label="outer cancellation child", + ) + cleanup_finished.set() + + owner = asyncio.create_task(cancelled_owner()) + await asyncio.sleep(0) + owner.cancel() + + with pytest.raises(asyncio.CancelledError): + await owner + assert cleanup_finished.is_set() + + @pytest.mark.asyncio async def test_await_cancelled_task_defers_stubborn_child_cleanup() -> None: child_cancelled = asyncio.Event() @@ -21039,6 +21065,336 @@ async def test_retry_http_bridge_fresh_hard_request_excludes_silent_account( release_lease.assert_awaited_once_with(old_lease) +@pytest.mark.asyncio +async def test_http_bridge_liveness_timeout_is_neutral_not_replayed_and_forces_retirement( + monkeypatch: pytest.MonkeyPatch, +) -> None: + service = proxy_service.ProxyService(cast(Any, nullcontext())) + request_state = proxy_service._WebSocketRequestState( + request_id="req-bridge-liveness-timeout", + model="gpt-5.4", + service_tier=None, + reasoning_effort=None, + api_key_reservation=None, + started_at=time.monotonic(), + awaiting_response_created=True, + request_text='{"type":"response.create","model":"gpt-5.4","input":"hello"}', + transport="http", + ) + session = _make_bridge_session( + key_value="bridge-liveness-timeout", + pending_requests=deque([request_state]), + queued_request_count=1, + ) + session.admission_waiter_count = 1 + session.upstream = cast( + UpstreamWebSocket, + SimpleNamespace( + receive=AsyncMock( + return_value=UpstreamWebSocketMessage( + kind="error", + error="Upstream websocket liveness failed", + error_code=UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + ) + ), + close=AsyncMock(), + ), + ) + monkeypatch.setattr(proxy_service, "get_settings", lambda: _make_app_settings()) + retry_precreated = AsyncMock(return_value=True) + fail_pending = AsyncMock() + retire = AsyncMock() + monkeypatch.setattr(service, "_retry_http_bridge_precreated_request", retry_precreated) + monkeypatch.setattr(service, "_fail_pending_websocket_requests", fail_pending) + monkeypatch.setattr(service, "_retire_stale_pending_http_bridge_session", retire) + + await service._relay_http_bridge_upstream_messages(session) + + retry_precreated.assert_not_awaited() + fail_pending.assert_awaited_once() + fail_pending_args = fail_pending.await_args + assert fail_pending_args is not None + assert fail_pending_args.kwargs["error_code"] == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE + assert fail_pending_args.kwargs["penalize_account"] is False + retire.assert_awaited_once_with( + session, + detail=UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + response_events_seen=0, + ) + assert session.queued_request_count == 0 + assert session.closed is True + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + "cancel_submitter", + [False, True], + ids=["normal", "cancelled-after-claim"], +) +async def test_http_bridge_liveness_send_receive_race_settles_request_once( + monkeypatch: pytest.MonkeyPatch, + cancel_submitter: bool, +) -> None: + class RacingLivenessUpstream: + def __init__(self) -> None: + self.send_started = asyncio.Event() + self.receive_returned = asyncio.Event() + self.close = AsyncMock() + + async def send_text(self, _text: str) -> None: + self.send_started.set() + await self.receive_returned.wait() + # Let the reader queue on lifecycle_lock before the submitter + # publishes its send-side failure ownership and releases the lock. + await asyncio.sleep(0) + raise UpstreamWebSocketTransportError( + "Codex upstream websocket send failed: heartbeat expired", + error_code=UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + ) + + async def receive(self) -> UpstreamWebSocketMessage: + await self.send_started.wait() + self.receive_returned.set() + return UpstreamWebSocketMessage( + kind="error", + error="Codex upstream websocket receive failed: heartbeat expired", + error_code=UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + ) + + service = proxy_service.ProxyService(cast(Any, nullcontext())) + sibling_queue: asyncio.Queue[str | None] = asyncio.Queue() + sibling_state = proxy_service._WebSocketRequestState( + request_id="req-bridge-liveness-race-sibling", + model="gpt-5.4", + service_tier=None, + reasoning_effort=None, + api_key_reservation=None, + started_at=time.monotonic(), + response_id="resp-bridge-liveness-race-sibling", + event_queue=sibling_queue, + transport="http", + skip_request_log=True, + ) + request_queue: asyncio.Queue[str | None] = asyncio.Queue() + request_state = proxy_service._WebSocketRequestState( + request_id="req-bridge-liveness-race", + model="gpt-5.4", + service_tier=None, + reasoning_effort=None, + api_key_reservation=None, + started_at=time.monotonic(), + awaiting_response_created=True, + event_queue=request_queue, + request_text='{"type":"response.create","model":"gpt-5.4","input":"hello"}', + transport="http", + skip_request_log=True, + ) + upstream = RacingLivenessUpstream() + session = _make_bridge_session( + key_value="bridge-liveness-race", + pending_requests=deque([sibling_state]), + queued_request_count=1, + ) + session.upstream = cast(UpstreamWebSocket, upstream) + service._http_bridge_sessions[session.key] = session + fail_pending = AsyncMock(wraps=service._fail_pending_websocket_requests) + retire = AsyncMock() + retry_precreated = AsyncMock(return_value=True) + settlement_started = asyncio.Event() + allow_settlement = asyncio.Event() + fail_reader = service._fail_http_bridge_reader_and_maybe_retire + + async def controlled_fail_reader( + controlled_session: proxy_service._HTTPBridgeSession, + **kwargs: Any, + ) -> None: + if cancel_submitter: + settlement_started.set() + await allow_settlement.wait() + await fail_reader(controlled_session, **kwargs) + + monkeypatch.setattr(proxy_service, "get_settings", lambda: _make_app_settings()) + monkeypatch.setattr(service, "_fail_pending_websocket_requests", fail_pending) + monkeypatch.setattr(service, "_retire_stale_pending_http_bridge_session", retire) + monkeypatch.setattr(service, "_retry_http_bridge_precreated_request", retry_precreated) + monkeypatch.setattr(service, "_fail_http_bridge_reader_and_maybe_retire", controlled_fail_reader) + + reader_task = asyncio.create_task(service._relay_http_bridge_upstream_messages(session)) + submit_task = asyncio.create_task( + service._submit_http_bridge_request( + session, + request_state=request_state, + text_data=request_state.request_text or "{}", + queue_limit=8, + ) + ) + try: + if cancel_submitter: + await asyncio.wait_for(settlement_started.wait(), timeout=1.0) + assert session.liveness_settlement_owner == "send" + submit_task.cancel() + await asyncio.sleep(0) + assert submit_task.done() is False + allow_settlement.set() + with pytest.raises(asyncio.CancelledError): + await submit_task + else: + with pytest.raises(ProxyResponseError) as exc_info: + await submit_task + assert exc_info.value.payload["error"]["code"] == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE + await asyncio.wait_for(reader_task, timeout=1.0) + finally: + allow_settlement.set() + if not submit_task.done(): + submit_task.cancel() + if not reader_task.done(): + reader_task.cancel() + await asyncio.gather(submit_task, reader_task, return_exceptions=True) + + fail_pending.assert_awaited_once() + failure_call = fail_pending.await_args + assert failure_call is not None + assert failure_call.kwargs["penalize_account"] is False + retry_precreated.assert_not_awaited() + for event_queue in (sibling_queue, request_queue): + terminal_event = await asyncio.wait_for(event_queue.get(), timeout=0.1) + assert terminal_event is not None + assert UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE in terminal_event + assert await asyncio.wait_for(event_queue.get(), timeout=0.1) is None + assert request_state.replay_count == 0 + assert session.pending_requests == deque() + assert session.queued_request_count == 0 + assert session.closed is True + assert session.liveness_settlement_owner == "send" + retire.assert_awaited_once_with( + session, + detail=UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + response_events_seen=0, + ) + + +@pytest.mark.asyncio +async def test_http_bridge_closed_without_liveness_claim_still_settles_pending_siblings( + monkeypatch: pytest.MonkeyPatch, +) -> None: + class DelayedLivenessUpstream: + def __init__(self) -> None: + self.receive_started = asyncio.Event() + self.release_liveness = asyncio.Event() + self.send_text = AsyncMock() + self.close = AsyncMock() + + async def receive(self) -> UpstreamWebSocketMessage: + self.receive_started.set() + await self.release_liveness.wait() + return UpstreamWebSocketMessage( + kind="error", + error="Codex upstream websocket receive failed: heartbeat expired", + error_code=UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + ) + + def pending_sibling(request_id: str) -> proxy_service._WebSocketRequestState: + return proxy_service._WebSocketRequestState( + request_id=request_id, + model="gpt-5.6-sol", + service_tier=None, + reasoning_effort=None, + api_key_reservation=None, + started_at=time.monotonic(), + response_id=f"resp-{request_id}", + event_queue=asyncio.Queue(), + transport="http", + skip_request_log=True, + ) + + service = proxy_service.ProxyService(cast(Any, nullcontext())) + siblings = [pending_sibling("sibling-one"), pending_sibling("sibling-two")] + upstream = DelayedLivenessUpstream() + session = _make_bridge_session( + key_value="closed-without-liveness-claim", + pending_requests=deque(siblings), + queued_request_count=len(siblings), + ) + session.upstream = cast(UpstreamWebSocket, upstream) + session.durable_session_id = "durable-closed-without-liveness-claim" + session.durable_owner_epoch = 2 + service._http_bridge_sessions[session.key] = session + record_recovery_attempt = AsyncMock(return_value=None) + service._durable_bridge = cast( + Any, + SimpleNamespace( + lookup_retry_circuit=AsyncMock(return_value=None), + record_recovery_attempt=record_recovery_attempt, + release_live_session=AsyncMock(return_value=None), + ), + ) + third_request = proxy_service._WebSocketRequestState( + request_id="third-submit", + model="gpt-5.6-sol", + service_tier=None, + reasoning_effort=None, + api_key_reservation=None, + started_at=time.monotonic(), + awaiting_response_created=True, + request_text='{"type":"response.create","model":"gpt-5.6-sol","input":"hi"}', + fresh_upstream_request_text='{"type":"response.create","model":"gpt-5.6-sol","input":"hi"}', + fresh_upstream_request_is_retry_safe=True, + transport="http", + skip_request_log=True, + ) + fail_pending = AsyncMock(wraps=service._fail_pending_websocket_requests) + retire = AsyncMock() + retry_precreated = AsyncMock(return_value=True) + monkeypatch.setattr(proxy_service, "get_settings", lambda: _make_app_settings()) + monkeypatch.setattr(service, "_fail_pending_websocket_requests", fail_pending) + monkeypatch.setattr(service, "_retire_stale_pending_http_bridge_session", retire) + monkeypatch.setattr(service, "_retry_http_bridge_precreated_request", retry_precreated) + + reader_task = asyncio.create_task(service._relay_http_bridge_upstream_messages(session)) + try: + await asyncio.wait_for(upstream.receive_started.wait(), timeout=1.0) + with pytest.raises(ProxyResponseError) as exc_info: + await service._submit_http_bridge_request( + session, + request_state=third_request, + text_data=third_request.request_text or "{}", + queue_limit=8, + ) + + assert exc_info.value.payload["error"]["code"] == "bridge_continuity_persistence_failed" + record_recovery_attempt.assert_awaited_once() + assert session.closed is True + assert session.liveness_settlement_owner is None + assert session.pending_requests == deque(siblings) + + upstream.release_liveness.set() + await asyncio.wait_for(reader_task, timeout=1.0) + finally: + if not reader_task.done(): + reader_task.cancel() + await asyncio.gather(reader_task, return_exceptions=True) + + fail_pending.assert_awaited_once() + failure_call = fail_pending.await_args + assert failure_call is not None + assert failure_call.kwargs["error_code"] == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE + assert failure_call.kwargs["penalize_account"] is False + retry_precreated.assert_not_awaited() + for sibling in siblings: + assert sibling.event_queue is not None + terminal_event = await asyncio.wait_for(sibling.event_queue.get(), timeout=0.1) + assert terminal_event is not None + assert UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE in terminal_event + assert await asyncio.wait_for(sibling.event_queue.get(), timeout=0.1) is None + assert session.pending_requests == deque() + assert session.queued_request_count == 0 + retire.assert_awaited_once_with( + session, + detail=UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + response_events_seen=0, + ) + + @pytest.mark.asyncio async def test_http_bridge_retry_send_network_failure_is_neutral_and_not_replayed( monkeypatch: pytest.MonkeyPatch, diff --git a/tests/unit/test_proxy_utils.py b/tests/unit/test_proxy_utils.py index bc84f388bf..7bd4bb817d 100644 --- a/tests/unit/test_proxy_utils.py +++ b/tests/unit/test_proxy_utils.py @@ -37,6 +37,7 @@ from app.core.balancer.types import UpstreamError from app.core.clients.proxy import _build_upstream_headers, filter_inbound_headers from app.core.clients.proxy_websocket import ( + UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, CodexUpstreamWebSocket, UpstreamWebSocket, UpstreamWebSocketTransportError, @@ -27504,10 +27505,172 @@ async def connect(*_args: object, **_kwargs: object): "response.created", "response.failed", ] - assert downstream.close_calls == [(1011, "upstream replay requires a fresh request")] + assert downstream.close_calls[0] == (1011, "upstream replay requires a fresh request") handle_stream_error.assert_not_awaited() +@pytest.mark.asyncio +async def test_proxy_responses_websocket_liveness_race_awaits_reader_settlement(monkeypatch): + request_logs = _RequestLogsRecorder() + service = proxy_service.ProxyService(_repo_factory(request_logs)) + settings = _make_proxy_settings() + settings.stream_idle_timeout_seconds = 300.0 + settings.proxy_downstream_websocket_idle_timeout_seconds = 120.0 + monkeypatch.setattr(proxy_service, "get_settings_cache", lambda: _SettingsCache(settings)) + monkeypatch.setattr(proxy_service, "get_settings", lambda: settings) + handle_stream_error = AsyncMock() + monkeypatch.setattr(proxy_service.ProxyService, "_handle_stream_error", handle_stream_error) + + request_texts = [ + json.dumps( + { + "type": "response.create", + "model": "gpt-5.4", + "instructions": "", + "input": [{"role": "user", "content": label}], + "stream": True, + }, + separators=(",", ":"), + ) + for label in ("first", "second") + ] + + class RacingDownstreamWebSocket: + def __init__(self) -> None: + self.request_index = 0 + self.first_created = asyncio.Event() + self.done = asyncio.Event() + self.sent_text: list[str] = [] + + async def receive(self) -> dict[str, object]: + if self.request_index == 0: + self.request_index = 1 + return {"type": "websocket.receive", "text": request_texts[0]} + if self.request_index == 1: + await self.first_created.wait() + self.request_index = 2 + return {"type": "websocket.receive", "text": request_texts[1]} + await self.done.wait() + return {"type": "websocket.disconnect"} + + async def send_text(self, text: str) -> None: + self.sent_text.append(text) + payload = json.loads(text) + if payload.get("type") == "response.created": + self.first_created.set() + if sum(json.loads(item).get("type") == "response.failed" for item in self.sent_text) == 2: + self.done.set() + + async def send_bytes(self, _data: bytes) -> None: + return None + + async def close(self, code: int = 1000, reason: str | None = None) -> None: + del code, reason + self.done.set() + + settlement_started = asyncio.Event() + allow_settlement = asyncio.Event() + + class RacingUpstreamWebSocket: + def __init__(self) -> None: + self.send_count = 0 + self.receive_count = 0 + self.second_send_started = asyncio.Event() + self.closed = False + + async def send_text(self, _text: str) -> None: + self.send_count += 1 + if self.send_count == 2: + self.second_send_started.set() + await settlement_started.wait() + raise UpstreamWebSocketTransportError( + "Codex upstream websocket send failed: heartbeat expired", + error_code=UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + ) + + async def send_bytes(self, _data: bytes) -> None: + return None + + async def receive(self) -> SimpleNamespace: + self.receive_count += 1 + if self.receive_count == 1: + return SimpleNamespace( + kind="text", + text=json.dumps( + { + "type": "response.created", + "response": {"id": "resp_liveness_race", "status": "in_progress"}, + }, + separators=(",", ":"), + ), + data=None, + close_code=None, + error=None, + error_code=None, + ) + await self.second_send_started.wait() + return SimpleNamespace( + kind="error", + text=None, + data=None, + close_code=1011, + error="Codex upstream websocket receive failed: heartbeat expired", + error_code=UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + ) + + async def close(self) -> None: + self.closed = True + + released_request_ids: list[str] = [] + + async def controlled_release(request_state: proxy_service._WebSocketRequestState) -> None: + if not settlement_started.is_set(): + settlement_started.set() + await allow_settlement.wait() + released_request_ids.append(request_state.request_id) + + downstream = RacingDownstreamWebSocket() + upstream = RacingUpstreamWebSocket() + account = _make_account("acc_ws_liveness_race") + + async def connect(*_args: object, **_kwargs: object): + return account, upstream + + monkeypatch.setattr(proxy_service.ProxyService, "_connect_proxy_websocket", connect) + monkeypatch.setattr(service, "_release_websocket_request_state_reservation", controlled_release) + + proxy_task = asyncio.create_task( + service.proxy_responses_websocket( + cast(WebSocket, downstream), + {}, + codex_session_affinity=False, + openai_cache_affinity=False, + api_key=None, + ) + ) + try: + await asyncio.wait_for(settlement_started.wait(), timeout=1.0) + assert proxy_task.done() is False + allow_settlement.set() + await asyncio.wait_for(proxy_task, timeout=1.0) + finally: + allow_settlement.set() + if not proxy_task.done(): + proxy_task.cancel() + await asyncio.gather(proxy_task, return_exceptions=True) + + emitted = [json.loads(text) for text in downstream.sent_text] + failures = [payload for payload in emitted if payload.get("type") == "response.failed"] + assert len(failures) == 2 + assert {payload["response"]["error"]["code"] for payload in failures} == {UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE} + assert len(released_request_ids) == 2 + assert len(set(released_request_ids)) == 2 + assert upstream.send_count == 2 + assert upstream.closed is True + handle_stream_error.assert_not_awaited() + assert len(request_logs.calls) == 2 + + @pytest.mark.asyncio async def test_stream_api_key_settlement_detaches_and_closes_repo(monkeypatch): started = asyncio.Event() @@ -28373,12 +28536,21 @@ async def close(self) -> None: @pytest.mark.asyncio -async def test_relay_upstream_websocket_network_failure_is_neutral_and_not_replayed(monkeypatch): +@pytest.mark.parametrize( + "error_code", + ["proxy_network_unavailable", UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE], + ids=["process-network", "liveness-timeout"], +) +async def test_relay_upstream_websocket_account_neutral_failure_is_not_replayed( + monkeypatch: pytest.MonkeyPatch, + error_code: str, +) -> None: request_logs = _RequestLogsRecorder() service = proxy_service.ProxyService(_repo_factory(request_logs)) handle_stream_error = AsyncMock() + release_reservation = AsyncMock() monkeypatch.setattr(service, "_handle_stream_error", handle_stream_error) - monkeypatch.setattr(service, "_release_websocket_request_state_reservation", AsyncMock()) + monkeypatch.setattr(service, "_release_websocket_request_state_reservation", release_reservation) class _FakeDownstreamWebSocket: def __init__(self) -> None: @@ -28390,19 +28562,22 @@ async def send_text(self, text: str) -> None: async def close(self, code: int = 1000, reason: str | None = None) -> None: del code, reason - class _NetworkFailureUpstream: + class _AccountNeutralFailureUpstream: + def __init__(self) -> None: + self.closed = False + async def receive(self) -> SimpleNamespace: return SimpleNamespace( kind="error", text=None, data=None, close_code=None, - error="Codex upstream websocket receive failed via proxy endpoint ep_1: OSError", - error_code="proxy_network_unavailable", + error="Upstream websocket liveness failed", + error_code=error_code, ) async def close(self) -> None: - return None + self.closed = True request_state = proxy_service._WebSocketRequestState( request_id="ws_req_network_failure", @@ -28418,10 +28593,11 @@ async def close(self) -> None: pending_requests = deque([request_state]) upstream_control = proxy_service._WebSocketUpstreamControl() downstream = _FakeDownstreamWebSocket() + upstream = _AccountNeutralFailureUpstream() await service._relay_upstream_websocket_messages( cast(WebSocket, downstream), - cast(proxy_service.UpstreamWebSocket, _NetworkFailureUpstream()), + cast(proxy_service.UpstreamWebSocket, upstream), account=_make_account("acc_ws_network_failure"), account_id_value="acc_ws_network_failure", pending_requests=pending_requests, @@ -28439,8 +28615,85 @@ async def close(self) -> None: assert upstream_control.reconnect_requested is True assert list(pending_requests) == [] handle_stream_error.assert_not_awaited() + release_reservation.assert_awaited_once_with(request_state) + # Upstream now retires every terminal receive before the reader exits; the + # liveness distinction controls replay and account health, not ownership. + assert upstream.closed is True terminal = json.loads(downstream.sent_text[-1]) - assert terminal["response"]["error"]["code"] == "proxy_network_unavailable" + assert terminal["response"]["error"]["code"] == error_code + + +@pytest.mark.asyncio +async def test_relay_upstream_websocket_liveness_timeout_preserves_sequenced_failure_contract( + monkeypatch: pytest.MonkeyPatch, +) -> None: + service = proxy_service.ProxyService(_repo_factory(_RequestLogsRecorder())) + fail_pending = AsyncMock() + monkeypatch.setattr(service, "_fail_pending_websocket_requests", fail_pending) + + class _DownstreamWebSocket: + def __init__(self) -> None: + self.close_calls: list[tuple[int, str | None]] = [] + + async def close(self, code: int = 1000, reason: str | None = None) -> None: + self.close_calls.append((code, reason)) + + class _LivenessTimeoutUpstream: + def __init__(self) -> None: + self.closed = False + + async def receive(self) -> SimpleNamespace: + return SimpleNamespace( + kind="error", + text=None, + data=None, + close_code=1011, + error="Upstream websocket liveness failed", + error_code=UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, + ) + + async def close(self) -> None: + self.closed = True + + request_state = proxy_service._WebSocketRequestState( + request_id="ws_req_sequenced_liveness_timeout", + model="gpt-5.1", + service_tier=None, + reasoning_effort=None, + api_key_reservation=None, + started_at=time.monotonic(), + request_text='{"type":"response.create","model":"gpt-5.1","input":"hi"}', + response_create_sent_at=0.0, + awaiting_response_created=False, + last_downstream_sequence_number=1, + ) + pending_requests = deque([request_state]) + downstream = _DownstreamWebSocket() + upstream = _LivenessTimeoutUpstream() + + await service._relay_upstream_websocket_messages( + cast(WebSocket, downstream), + cast(proxy_service.UpstreamWebSocket, upstream), + account=_make_account("acc_ws_sequenced_liveness_timeout"), + account_id_value="acc_ws_sequenced_liveness_timeout", + pending_requests=pending_requests, + pending_lock=anyio.Lock(), + client_send_lock=anyio.Lock(), + api_key=None, + upstream_control=proxy_service._WebSocketUpstreamControl(), + response_create_gate=asyncio.Semaphore(1), + proxy_request_budget_seconds=5.0, + stream_idle_timeout_seconds=5.0, + downstream_activity=proxy_service._DownstreamWebSocketActivity(), + ) + + fail_pending.assert_awaited_once() + fail_pending_args = fail_pending.await_args + assert fail_pending_args is not None + assert fail_pending_args.kwargs["penalize_account"] is False + assert fail_pending_args.kwargs["suppress_sequenced_downstream_errors"] is True + assert upstream.closed is True + assert downstream.close_calls[0] == (1011, "upstream replay requires a fresh request") @pytest.mark.asyncio diff --git a/tests/unit/test_proxy_websocket_client.py b/tests/unit/test_proxy_websocket_client.py index 23b2511b63..b0079f87c4 100644 --- a/tests/unit/test_proxy_websocket_client.py +++ b/tests/unit/test_proxy_websocket_client.py @@ -19,6 +19,7 @@ from app.core.clients.codex import CodexTransportError, CodexWebSocketResult from app.core.clients.proxy import ProxyResponseError, is_confirmed_pre_dispatch_transport_error from app.core.clients.proxy_websocket import ( + UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE, CodexUpstreamWebSocket, RealtimeWebSocketProtocol, UpstreamWebSocketTransportError, @@ -148,6 +149,9 @@ async def recv(self) -> tuple[bytes, int]: async def receive(self) -> object: return b'{"type":"response.completed"}' + def exception(self) -> BaseException | None: + return None + async def close(self, *, code: int = 1000, message: bytes = b"") -> None: del code, message self.closed = True @@ -228,6 +232,121 @@ async def recv(self): assert message.error is None +@pytest.mark.asyncio +async def test_direct_adapter_classifies_keepalive_timeout() -> None: + class Connection: + async def recv(self): + raise ConnectionClosedError(None, Close(1011, "keepalive ping timeout")) + + websocket = WebsocketsUpstreamWebSocket(cast(Any, Connection())) + + message = await websocket.receive() + + assert message.kind == "error" + assert message.error_code == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE + + +@pytest.mark.asyncio +async def test_direct_adapter_classifies_keepalive_timeout_after_close_ack() -> None: + class Connection: + async def recv(self): + raise ConnectionClosedError( + Close(1000, "acknowledged"), + Close(1011, "keepalive ping timeout"), + False, + ) + + websocket = WebsocketsUpstreamWebSocket(cast(Any, Connection())) + + message = await websocket.receive() + + assert message.kind == "error" + assert message.error_code == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE + + +@pytest.mark.asyncio +async def test_direct_adapter_does_not_trust_peer_keepalive_timeout_marker() -> None: + class Connection: + async def recv(self): + raise ConnectionClosedError( + Close(1011, "keepalive ping timeout"), + Close(1011, "keepalive ping timeout"), + True, + ) + + websocket = WebsocketsUpstreamWebSocket(cast(Any, Connection())) + + message = await websocket.receive() + + assert message.kind == "error" + assert message.error_code is None + + +@pytest.mark.asyncio +async def test_routed_adapter_classifies_heartbeat_timeout() -> None: + websocket = CodexUpstreamWebSocket( + _FakeCodexErrorWebSocket(aiohttp.ServerTimeoutError("No PONG received after 60.0 seconds")) + ) + + message = await websocket.receive() + + assert message.kind == "error" + assert message.error_code == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE + + +@pytest.mark.asyncio +async def test_routed_adapter_classifies_heartbeat_timeout_stored_between_receive_calls() -> None: + heartbeat_timeout = aiohttp.ServerTimeoutError("No PONG received after 60.0 seconds") + + class ClosedWebSocket(_FakeCodexWebSocket): + async def receive(self) -> aiohttp.WSMessage: + return aiohttp.WSMessage(aiohttp.WSMsgType.CLOSED, None, None) + + def exception(self) -> BaseException | None: + return heartbeat_timeout + + websocket = CodexUpstreamWebSocket(ClosedWebSocket()) + + message = await websocket.receive() + + assert message.kind == "error" + assert message.error_code == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE + + +@pytest.mark.asyncio +@pytest.mark.parametrize( + "payload", + ["request", b"request"], + ids=["text", "bytes"], +) +async def test_routed_adapter_send_preserves_stored_heartbeat_timeout( + payload: str | bytes, +) -> None: + heartbeat_timeout = aiohttp.ServerTimeoutError("No PONG received after 60.0 seconds") + + class ClosedWebSocket(_FakeCodexWebSocket): + async def send_str(self, data: str) -> None: + del data + raise RuntimeError("Cannot write to closing transport") + + async def send_bytes(self, data: bytes) -> None: + del data + raise RuntimeError("Cannot write to closing transport") + + def exception(self) -> BaseException | None: + return heartbeat_timeout + + websocket = CodexUpstreamWebSocket(ClosedWebSocket()) + + with pytest.raises(UpstreamWebSocketTransportError) as exc_info: + if isinstance(payload, str): + await websocket.send_text(payload) + else: + await websocket.send_bytes(payload) + + assert exc_info.value.error_code == UPSTREAM_WEBSOCKET_LIVENESS_TIMEOUT_CODE + + @pytest.mark.asyncio async def test_codex_responses_websocket_closes_owned_client_when_context_exit_fails(): class _FailingContext: @@ -266,6 +385,7 @@ async def fake_websocket_connect(url: str, **kwargs): lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=False, ), @@ -294,7 +414,7 @@ async def fake_websocket_connect(url: str, **kwargs): assert kwargs["proxy"] is None assert kwargs["open_timeout"] == 7.0 assert "ping_interval" not in kwargs - assert kwargs["ping_timeout"] is None + assert kwargs["ping_timeout"] == 120.0 assert kwargs["max_size"] == 4321 assert "subprotocols" not in kwargs additional_headers = cast(dict[str, str], kwargs["additional_headers"]) @@ -328,6 +448,7 @@ async def recv(self) -> str: lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=False, ), @@ -366,6 +487,7 @@ async def test_connect_responses_websocket_routed_codex_call_preserves_size_limi lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=False, ), @@ -390,6 +512,7 @@ async def test_connect_responses_websocket_routed_codex_call_preserves_size_limi assert call["route"] is route assert call["timeout"] == 7.0 assert call["max_msg_size"] == 4321 + assert call["heartbeat"] == 120.0 assert "max_size" not in call assert "protocols" not in call assert websocket.response_header("x-codex-turn-state") == "turn-routed" @@ -666,6 +789,7 @@ async def test_connect_responses_websocket_routed_transport_error_maps_proxy_err lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=False, ), @@ -702,6 +826,7 @@ async def test_connect_responses_websocket_routed_pre_dispatch_failure_carries_p upstream_connect_timeout_seconds=7.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=False, + proxy_downstream_websocket_idle_timeout_seconds=120.0, ), ) @@ -752,6 +877,7 @@ async def test_connect_responses_websocket_routed_tls_verification_failure_is_no upstream_connect_timeout_seconds=7.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=False, + proxy_downstream_websocket_idle_timeout_seconds=120.0, ), ) @@ -802,6 +928,7 @@ async def fake_websocket_connect(url: str, **kwargs): lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=False, ), @@ -837,6 +964,7 @@ async def fake_websocket_connect(url: str, **kwargs): lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=False, ), @@ -902,6 +1030,7 @@ async def fake_websocket_connect(url: str, **kwargs): lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=False, ), @@ -938,6 +1067,7 @@ async def fake_websocket_connect(url: str, **kwargs): lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=True, ), @@ -976,6 +1106,7 @@ async def fake_websocket_connect(url: str, **kwargs): lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=True, ), @@ -1023,6 +1154,7 @@ async def test_connect_responses_websocket_sanitizes_ws_error_payload(monkeypatc lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=True, ), @@ -1078,6 +1210,7 @@ async def fake_websocket_connect(url: str, **kwargs): lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=True, ), @@ -1120,6 +1253,7 @@ async def fake_websocket_connect(url: str, **kwargs): lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=True, ), @@ -1164,6 +1298,7 @@ async def fake_websocket_connect(url: str, **kwargs): lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=True, ), @@ -1208,6 +1343,7 @@ async def fake_websocket_connect(url: str, **kwargs): lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=True, ), @@ -1250,6 +1386,7 @@ def upstream_websocket_proxy_env(self): lambda: _Settings( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=True, ), @@ -1293,6 +1430,7 @@ def upstream_websocket_proxy_env(self): lambda: _Settings( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=True, ), @@ -1329,6 +1467,7 @@ async def fake_websocket_connect(url: str, **kwargs): lambda: SimpleNamespace( upstream_base_url="http://chatgpt.local/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=True, ), @@ -1387,6 +1526,7 @@ async def upstream_handler(connection): lambda: SimpleNamespace( upstream_base_url=f"http://127.0.0.1:{upstream_port}/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=True, ), @@ -1425,6 +1565,7 @@ async def fake_websocket_connect(url: str, **kwargs): lambda: SimpleNamespace( upstream_base_url="http://chatgpt.local/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=True, ), @@ -1472,6 +1613,7 @@ async def fake_websocket_connect(url: str, **kwargs): lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=True, ), @@ -1506,6 +1648,7 @@ async def fake_websocket_connect(url: str, **kwargs): lambda: SimpleNamespace( upstream_base_url="https://chatgpt.com/backend-api", upstream_connect_timeout_seconds=7.0, + proxy_downstream_websocket_idle_timeout_seconds=120.0, max_sse_event_bytes=4321, upstream_websocket_trust_env=True, ),