fix(realtime): treat any client receive failure as a client hangup

client_ack_messages classified a websockets ConnectionClosed raised by
the client socket as the backend closing, so bidirectional_forward kept
waiting on the upstream instead of ending the session. Starlette clients
raise WebSocketDisconnect, but the realtime test client in
tests/llm_translation/realtime raises websockets.exceptions.ConnectionClosed,
which hung test_openai_realtime_simple.py until the run was killed.

Only the receive_text call now maps every exception to
CLIENT_DISCONNECTED; the loop body keeps ConnectionClosed as
BACKEND_CLOSED, since the backend socket is the only websockets socket
touched there.
This commit is contained in:
mateo-berri 2026-09-04 19:52:48 -07:00
parent 85d45fbb4b
commit da9dbdba96
2 changed files with 25 additions and 1 deletions

View file

@ -1297,13 +1297,22 @@ class RealTimeStreaming:
item["content"] = new_content
return item
async def _receive_client_message(self) -> str | None:
try:
return await self.websocket.receive_text()
except Exception as e: # noqa: BLE001 # whatever the client socket raises, the client is gone
verbose_logger.debug("Client disconnected: %s", e)
return None
async def client_ack_messages(self) -> ClientLoopExit:
import websockets
client_event: _ClientEventFrame
try:
while True:
message = await self.websocket.receive_text()
message = await self._receive_client_message()
if message is None:
return ClientLoopExit.CLIENT_DISCONNECTED
## GUARDRAIL: intercept conversation.item.create for text-based injection.
guardrail_turn_detection_injected = False

View file

@ -3307,3 +3307,18 @@ async def test_client_hanging_up_first_ends_the_session_without_a_relayed_close(
assert session.logging.logged_sessions == ((),)
assert session.logging.logged_failures == ()
client_ws.close.assert_not_awaited()
@pytest.mark.asyncio
async def test_client_hanging_up_with_a_websockets_close_is_not_mistaken_for_the_backend_closing():
client_ws: Final = _client_ws_that_never_sends()
client_ws.receive_text = AsyncMock(side_effect=ConnectionClosed(None, None))
backend_ws: Final = MagicMock()
backend_ws.recv = AsyncMock(side_effect=_wait_forever)
session: Final = _relay_session(client_ws, backend_ws)
await session.run()
assert session.logging.logged_sessions == ((),)
assert session.logging.logged_failures == ()
client_ws.close.assert_not_awaited()