Skip to content

Commit 2883588

Browse files
Merge PR #40 (fix/ws-drained-pull-future): read a done pull future on every WS pump exit
Merge commit, not a fast-forward, by necessity: #36, #40, #42 and #41 were each cut from main @ c1a28de in parallel, so once #36 landed no other tip descends from main's head. A rebase would change the SHAs the PR body cites (RED 8bcf7d9, GREEN b33ac04) and the repository ruleset refuses the force-push a rebased branch needs. The merge keeps every cited commit intact. Files touched are disjoint from the other three.
2 parents 3b8a753 + b33ac04 commit 2883588

2 files changed

Lines changed: 99 additions & 1 deletion

File tree

‎packages/studyloop/src/studyloop/web/routes/session/_ws.py‎

Lines changed: 10 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -161,7 +161,16 @@ async def pty_to_ws() -> None:
161161
{"type": "agent_message", "kind": event.kind, "payload": event.payload}
162162
)
163163
finally:
164-
if not nxt.done():
164+
if nxt.done():
165+
# A pull that completed after the loop's last read -- a drain
166+
# (StopAsyncIteration) landing beside a takeover, or a client
167+
# close that cancelled this coroutine while the pull was
168+
# already done. Read its outcome so asyncio does not log
169+
# "Task exception was never retrieved" from the finalizer;
170+
# a drained stream ending here is expected, not an error.
171+
if not nxt.cancelled():
172+
nxt.exception()
173+
else:
165174
nxt.cancel()
166175
# Only the PTY transport's events() is an async generator, so only
167176
# it has aclose(). The ACP transport deliberately returns a

‎packages/studyloop/tests/test_web_session_ws.py‎

Lines changed: 89 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,9 @@
1414

1515
from __future__ import annotations
1616

17+
import asyncio
18+
import gc
19+
import logging
1720
import sys
1821
import time
1922
from pathlib import Path
@@ -425,3 +428,89 @@ def test_ws_disconnect_releases_active_session(
425428

426429
assert run_async(active.current()) is None
427430
assert stub.end_calls == 1
431+
432+
433+
# ---------------------------------------------------------------------------
434+
# The drained pull future is retrieved on every exit
435+
# ---------------------------------------------------------------------------
436+
437+
438+
class _GatedStub(StubTransport):
439+
"""A transport whose stream drains only when the test says so.
440+
441+
``events()`` yields ``Started`` and then waits on an ``asyncio.Event``
442+
created in the server loop; setting it ends the generator, which is what
443+
a real session's ``end()`` does when it pushes the queue sentinel.
444+
"""
445+
446+
def __init__(self) -> None:
447+
super().__init__()
448+
self.release: asyncio.Event | None = None
449+
450+
async def events(self): # type: ignore[override]
451+
self.release = asyncio.Event()
452+
yield Started(agent="claude")
453+
await self.release.wait()
454+
455+
456+
class TestDrainedPullFuture:
457+
def test_a_drain_that_lands_beside_a_takeover_is_retrieved(
458+
self,
459+
config: SessionConfig,
460+
caplog: pytest.LogCaptureFixture,
461+
monkeypatch: pytest.MonkeyPatch,
462+
) -> None:
463+
"""Regression: ``Task exception was never retrieved ... StopAsyncIteration``.
464+
465+
The pump pulls each transport event as its own future. When a newer
466+
socket takes the consumer slot at the same moment the stream drains,
467+
the pump raised ``_SupersededError`` before reading the drained
468+
future's ``StopAsyncIteration``, and its ``finally`` left a *done*
469+
future neither cancelled nor read -- so asyncio logged the exception
470+
from the future's finalizer. Seen live on 2026-09-23 when a PWA tab
471+
reattached to a session whose transport had already ended.
472+
"""
473+
from studyloop.web.routes.session import _grace
474+
475+
# The pump must wake because the pull completed, never because its
476+
# supersede poll timed out first: that path cancels a still-pending
477+
# pull and never exhibited the leak.
478+
monkeypatch.setattr(_grace, "SUPERSEDE_POLL_S", 5.0)
479+
stub = _GatedStub()
480+
_install_active(stub, config)
481+
482+
with (
483+
TestClient(create_app()) as client,
484+
caplog.at_level(logging.ERROR, logger="asyncio"),
485+
):
486+
portal = client.portal
487+
assert portal is not None
488+
with client.websocket_connect(
489+
"/api/session/ws?study_session_id=study-1",
490+
headers={"Origin": "http://127.0.0.1:8788"},
491+
) as ws:
492+
assert ws.receive_json() == {"type": "started", "agent": "claude"}
493+
494+
def takeover_and_drain() -> None:
495+
# The synchronous half of ``_grace.acquire_consumer``: a
496+
# newer socket has claimed the slot ...
497+
holder = _grace._attachment # pyright: ignore[reportPrivateUsage]
498+
assert holder is not None
499+
holder.superseded = True
500+
# ... and the stream drains in the same loop turn.
501+
assert stub.release is not None
502+
stub.release.set()
503+
504+
portal.call(takeover_and_drain)
505+
frame = ws.receive_json()
506+
assert frame["type"] == "attach_superseded"
507+
508+
# A future whose exception was never read logs from its finalizer;
509+
# force the finalizer so the assertion does not depend on refcount luck.
510+
gc.collect()
511+
leaked = [
512+
record.getMessage()
513+
for record in caplog.records
514+
if record.name == "asyncio" and "never retrieved" in record.getMessage()
515+
]
516+
assert leaked == [], leaked[0]

0 commit comments

Comments
 (0)