fix(relay): a raising socket write returns the result dict, not an exception
The socket can die BETWEEN _request_response's 'is None' liveness guard
and the actual write: the reader's finally hasn't cleared _ws yet, so
_send raises ConnectionClosed straight into callers whose contract is a
result dict (RelayAdapter.send consumes it with no try — only the
cosmetic typing lanes wrap the call). No liveness check can close this
window; it has to be caught at the write.
Convert the raise to {'success': False, 'error': ...} like every other
failed send, log the traceback at debug (the returned string alone
can't distinguish an ordinary dead socket from a defect in the frame
building above), and rely on the existing finally to drop the pending
entry. CancelledError is a BaseException, so cancellation still
propagates.
Same disposition as the equivalent guard in PR #82238; regression test
drives a socket whose write raises while the reader is still parked,
and fails without the except clause.
This commit is contained in:
@@ -812,6 +812,16 @@ class WebSocketRelayTransport:
|
||||
return await asyncio.wait_for(fut, timeout=self._outbound_timeout_s)
|
||||
except asyncio.TimeoutError:
|
||||
return {"success": False, "error": "relay outbound timed out"}
|
||||
except Exception as exc: # noqa: BLE001 - a dead socket is a failed send, not a raise
|
||||
# No `is None` check can close the window where the socket dies
|
||||
# BETWEEN the liveness guard above and the actual write — the
|
||||
# reader's finally hasn't cleared _ws yet, so _send raises
|
||||
# ConnectionClosed straight into callers whose contract is a
|
||||
# result dict (RelayAdapter.send consumes it with no try).
|
||||
# Report it like every other failed send. CancelledError is a
|
||||
# BaseException, so cancellation still propagates.
|
||||
logger.debug("relay %s send failed", frame_type, exc_info=True)
|
||||
return {"success": False, "error": f"relay send failed: {exc}"}
|
||||
finally:
|
||||
self._pending.pop(request_id, None)
|
||||
|
||||
|
||||
@@ -326,3 +326,48 @@ async def test_send_after_drop_with_reconnect_disabled_fails_fast():
|
||||
)
|
||||
assert result["success"] is False
|
||||
assert t._pending == {}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_send_raising_socket_returns_error_dict():
|
||||
"""The socket can die BETWEEN the `_ws is None` liveness guard and the
|
||||
actual write (the reader's finally hasn't cleared the handle yet). The
|
||||
write then raises ConnectionClosed — but send_outbound's contract is a
|
||||
result dict, and RelayAdapter.send consumes it with no try. The raise
|
||||
must be converted to {"success": False, ...}, with no future left in
|
||||
_pending."""
|
||||
|
||||
class _RaisingWS:
|
||||
"""Send raises (already dead); the reader hasn't noticed yet."""
|
||||
|
||||
def __init__(self):
|
||||
self.reader_release = asyncio.Event()
|
||||
|
||||
async def send(self, data):
|
||||
raise ConnectionClosedError(None, None)
|
||||
|
||||
def __aiter__(self):
|
||||
return self
|
||||
|
||||
async def __anext__(self):
|
||||
await self.reader_release.wait()
|
||||
raise ConnectionClosedError(None, None)
|
||||
|
||||
async def close(self):
|
||||
pass
|
||||
|
||||
t = WebSocketRelayTransport("ws://unused", "discord", "bot1", outbound_timeout_s=5.0)
|
||||
fake = _RaisingWS()
|
||||
t._ws = fake
|
||||
t._reader = asyncio.create_task(t._read_loop())
|
||||
await asyncio.sleep(0)
|
||||
|
||||
result = await asyncio.wait_for(
|
||||
t.send_outbound({"op": "send_message", "text": "hi"}), timeout=2.0
|
||||
)
|
||||
assert result["success"] is False
|
||||
assert "relay send failed" in result["error"]
|
||||
assert t._pending == {}
|
||||
|
||||
fake.reader_release.set()
|
||||
await t._reader
|
||||
|
||||
Reference in New Issue
Block a user