Follow-up on the cherry-picked #106372 fix (@gaoanze888), which reused the existing 10s _WS_WRITE_TIMEOUT_S as the deadline and only latched _closed. Why a separate, longer constant: _WS_WRITE_TIMEOUT_S is the worker-side wait for the LOOP to run the scheduled send; it deliberately does not mark the transport dead because a GIL-heavy turn stalls the loop for >10s routinely (#48445/#55545). The new clock is different in kind — asyncio.wait_for starts it only once the coroutine actually runs on the loop, and a healthy socket returns from send_text without waiting, so only real kernel backpressure can consume it (verified: a loop stall 3x longer than the deadline does not trip wait_for on 3.11/3.12/3.13 because the wakeup and the timer land in the same iteration and the result wins). Still, tying it to the 10s worker constant would let one test-time monkeypatch flip both semantics, and the issue asks explicitly not to declare the peer dead at the 10s worker wait. 30s is 3x the worker wait and under the client's 45s heartbeat deadline. Why close the socket: latching _closed alone leaves handle_ws parked in receive_text on a half-open peer, so the disconnect teardown (session detach/reap, live-transport unregister, client reconnect) never runs — the frontend keeps waiting on a socket that will never carry another frame. ws.close(1011) unblocks the read loop; websockets bounds the close itself (close_timeout -> abort) when the same full buffer stalls the close frame. Diagnostic now says which side stalled: "socket stalled, loop responsive" versus write()'s "loop stalled". Tests trimmed to the two invariants (≤2 per fix): a stalled send terminates within the deadline, fails both queued writes, latches closed and closes the socket; progress → complete → RPC reply ordering is preserved on a healthy socket. The stalled test uses the issue's own shape (two queued write_async calls) and monkeypatches the deadline with raising=False so on a base without the fix it fails on the symptom (sends never terminate), not on a missing attribute. Refs #106369
72 lines
3.0 KiB
Python
72 lines
3.0 KiB
Python
"""A stalled ``send_text`` must have a bounded lifetime (#106369).
|
|
|
|
``_safe_send_many`` awaits the socket while holding the connection-wide ``_send_lock``. Without a
|
|
deadline, one send parked by socket backpressure trapped every later event and RPC reply behind the
|
|
lock on an apparently open connection, so reconnect recovery never started.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import json
|
|
|
|
from tui_gateway.ws import WSTransport
|
|
|
|
|
|
class _StalledWS:
|
|
"""``send_text`` never completes (kernel backpressure); ``close`` is observable."""
|
|
|
|
def __init__(self) -> None:
|
|
self.closed_with: list[int] = []
|
|
self._release = asyncio.Event()
|
|
|
|
async def send_text(self, line: str) -> None:
|
|
await self._release.wait()
|
|
|
|
async def close(self, code: int = 1000) -> None:
|
|
self.closed_with.append(code)
|
|
self._release.set()
|
|
|
|
|
|
class _FastWS:
|
|
def __init__(self) -> None:
|
|
self.sent: list[str] = []
|
|
|
|
async def send_text(self, line: str) -> None:
|
|
self.sent.append(line)
|
|
|
|
|
|
def test_stalled_send_closes_socket_and_releases_queued_reply(monkeypatch):
|
|
# raising=False: on a base without the deadline the test must fail on the SYMPTOM (sends never terminate).
|
|
monkeypatch.setattr("tui_gateway.ws._WS_SEND_DEADLINE_S", 0.05, raising=False)
|
|
|
|
async def _run() -> None:
|
|
ws = _StalledWS()
|
|
transport = WSTransport(ws, asyncio.get_running_loop(), peer="127.0.0.1:1")
|
|
progress = asyncio.create_task(transport.write_async({"method": "event", "params": {"type": "tool.progress"}}))
|
|
await asyncio.sleep(0) # progress is now inside send_text, holding _send_lock
|
|
reply = asyncio.create_task(transport.write_async({"id": "submit", "result": {"status": "streaming"}}))
|
|
# The loop stays responsive; both sends must still terminate within the deadline (not the 2s cap).
|
|
results = await asyncio.wait_for(asyncio.gather(progress, reply), timeout=2.0)
|
|
assert results == [False, False], "a send that missed the deadline must report failure, not success"
|
|
assert transport.closed, "the transport must latch closed so handle_ws teardown/reconnect can run"
|
|
await asyncio.sleep(0) # let the scheduled close task run
|
|
assert ws.closed_with == [1011], "the stalled socket must be closed, not left half-open"
|
|
|
|
asyncio.run(_run())
|
|
|
|
|
|
def test_progress_then_final_ordering_preserved_on_healthy_socket():
|
|
async def _run() -> None:
|
|
ws = _FastWS()
|
|
transport = WSTransport(ws, asyncio.get_running_loop(), peer="127.0.0.1:1")
|
|
assert await transport.write_async({"method": "event", "params": {"type": "tool.progress"}})
|
|
assert await transport.write_async({"method": "event", "params": {"type": "message.complete"}})
|
|
assert await transport.write_async({"id": "submit", "result": {"status": "done"}})
|
|
assert not transport.closed
|
|
assert [json.loads(s).get("params", {}).get("type", "reply") for s in ws.sent] == [
|
|
"tool.progress", "message.complete", "reply",
|
|
]
|
|
|
|
asyncio.run(_run())
|