fix(relay): reader without a socket settles pending waiters instead of asserting
_read_loop opened with 'assert self._ws is not None' — an exit that escaped BEFORE the finally that fails pending futures, contradicting the 'fails pending on ANY exit path' invariant the hardening commit established. Production currently assigns _ws before scheduling the reader, so this was latent, but any future lifecycle change hitting it would strand every in-flight waiter for the full outbound timeout with only an AssertionError in the logs. Turn it into a guarded early-return INSIDE the try: the reader logs the lifecycle bug and unwinds through the same finally as every other exit, settling all waiters. Regression test drives _read_loop with _ws=None against a registered pending future.
This commit is contained in:
@@ -822,13 +822,20 @@ class WebSocketRelayTransport:
|
||||
await self._ws.send(json.dumps(frame) + "\n")
|
||||
|
||||
async def _read_loop(self) -> None:
|
||||
assert self._ws is not None
|
||||
# Bind the socket this reader serves: the finally below must only
|
||||
# clear _ws if it still points at THIS socket (a supervisor re-dial
|
||||
# may have already installed a fresh one by the time we unwind).
|
||||
ws = self._ws
|
||||
buf = ""
|
||||
try:
|
||||
if ws is None:
|
||||
# Scheduled without a socket (a lifecycle bug, not a normal
|
||||
# path). The old `assert` here escaped BEFORE the finally
|
||||
# existed to fail pending futures — the one exit that could
|
||||
# still strand waiters for the full outbound timeout. Fall
|
||||
# through to the finally instead; it settles them all.
|
||||
logger.error("relay ws read loop started with no socket")
|
||||
return
|
||||
try:
|
||||
async for chunk in self._ws:
|
||||
buf += chunk if isinstance(chunk, str) else chunk.decode("utf-8")
|
||||
|
||||
@@ -263,6 +263,24 @@ async def _run_reader_to_exit(t: WebSocketRelayTransport, fake: _DroppingWS) ->
|
||||
await t._reader
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_read_loop_without_socket_still_fails_pending():
|
||||
"""If the reader is ever scheduled with no socket (lifecycle bug), it must
|
||||
still settle in-flight waiters on its way out — the old `assert` escaped
|
||||
before the fail-pending cleanup and left them to the full 30s timeout."""
|
||||
t = WebSocketRelayTransport("ws://unused", "discord", "bot1", outbound_timeout_s=30.0)
|
||||
loop = asyncio.get_running_loop()
|
||||
fut: asyncio.Future = loop.create_future()
|
||||
t._pending["rid"] = fut
|
||||
t._ws = None
|
||||
|
||||
await t._read_loop() # must not raise
|
||||
|
||||
assert fut.done()
|
||||
assert fut.result() == {"success": False, "error": "relay transport connection lost"}
|
||||
assert t._pending == {}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_send_after_terminal_4401_revocation_fails_fast():
|
||||
"""A terminal 4401 revocation deliberately arms NO reconnect supervisor,
|
||||
|
||||
Reference in New Issue
Block a user