Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 10 additions & 1 deletion packages/studyloop/src/studyloop/web/routes/session/_ws.py
Original file line number Diff line number Diff line change
Expand Up @@ -161,7 +161,16 @@ async def pty_to_ws() -> None:
{"type": "agent_message", "kind": event.kind, "payload": event.payload}
)
finally:
if not nxt.done():
if nxt.done():
# A pull that completed after the loop's last read -- a drain
# (StopAsyncIteration) landing beside a takeover, or a client
# close that cancelled this coroutine while the pull was
# already done. Read its outcome so asyncio does not log
# "Task exception was never retrieved" from the finalizer;
# a drained stream ending here is expected, not an error.
if not nxt.cancelled():
nxt.exception()
else:
nxt.cancel()
# Only the PTY transport's events() is an async generator, so only
# it has aclose(). The ACP transport deliberately returns a
Expand Down
89 changes: 89 additions & 0 deletions packages/studyloop/tests/test_web_session_ws.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,9 @@

from __future__ import annotations

import asyncio
import gc
import logging
import sys
import time
from pathlib import Path
Expand Down Expand Up @@ -425,3 +428,89 @@ def test_ws_disconnect_releases_active_session(

assert run_async(active.current()) is None
assert stub.end_calls == 1


# ---------------------------------------------------------------------------
# The drained pull future is retrieved on every exit
# ---------------------------------------------------------------------------


class _GatedStub(StubTransport):
"""A transport whose stream drains only when the test says so.

``events()`` yields ``Started`` and then waits on an ``asyncio.Event``
created in the server loop; setting it ends the generator, which is what
a real session's ``end()`` does when it pushes the queue sentinel.
"""

def __init__(self) -> None:
super().__init__()
self.release: asyncio.Event | None = None

async def events(self): # type: ignore[override]
self.release = asyncio.Event()
yield Started(agent="claude")
await self.release.wait()


class TestDrainedPullFuture:
def test_a_drain_that_lands_beside_a_takeover_is_retrieved(
self,
config: SessionConfig,
caplog: pytest.LogCaptureFixture,
monkeypatch: pytest.MonkeyPatch,
) -> None:
"""Regression: ``Task exception was never retrieved ... StopAsyncIteration``.

The pump pulls each transport event as its own future. When a newer
socket takes the consumer slot at the same moment the stream drains,
the pump raised ``_SupersededError`` before reading the drained
future's ``StopAsyncIteration``, and its ``finally`` left a *done*
future neither cancelled nor read -- so asyncio logged the exception
from the future's finalizer. Seen live on 2026-09-23 when a PWA tab
reattached to a session whose transport had already ended.
"""
from studyloop.web.routes.session import _grace

# The pump must wake because the pull completed, never because its
# supersede poll timed out first: that path cancels a still-pending
# pull and never exhibited the leak.
monkeypatch.setattr(_grace, "SUPERSEDE_POLL_S", 5.0)
stub = _GatedStub()
_install_active(stub, config)

with (
TestClient(create_app()) as client,
caplog.at_level(logging.ERROR, logger="asyncio"),
):
portal = client.portal
assert portal is not None
with client.websocket_connect(
"/api/session/ws?study_session_id=study-1",
headers={"Origin": "http://127.0.0.1:8788"},
) as ws:
assert ws.receive_json() == {"type": "started", "agent": "claude"}

def takeover_and_drain() -> None:
# The synchronous half of ``_grace.acquire_consumer``: a
# newer socket has claimed the slot ...
holder = _grace._attachment # pyright: ignore[reportPrivateUsage]
assert holder is not None
holder.superseded = True
# ... and the stream drains in the same loop turn.
assert stub.release is not None
stub.release.set()

portal.call(takeover_and_drain)
frame = ws.receive_json()
assert frame["type"] == "attach_superseded"

# A future whose exception was never read logs from its finalizer;
# force the finalizer so the assertion does not depend on refcount luck.
gc.collect()
leaked = [
record.getMessage()
for record in caplog.records
if record.name == "asyncio" and "never retrieved" in record.getMessage()
]
assert leaked == [], leaked[0]
Loading