From 31f97611123b4ea05a0fc5d0d2c49c7583de4816 Mon Sep 17 00:00:00 2001 From: GISCE Bot Date: Mon, 3 Aug 2026 14:36:36 +0000 Subject: [PATCH] fix: cancel journal read before closing stream --- src/github_agent_bridge/backend.py | 7 +++++++ tests/test_backend.py | 25 +++++++++++++++++++++++++ 2 files changed, 32 insertions(+) diff --git a/src/github_agent_bridge/backend.py b/src/github_agent_bridge/backend.py index 89988e4..09c1e19 100644 --- a/src/github_agent_bridge/backend.py +++ b/src/github_agent_bridge/backend.py @@ -296,6 +296,8 @@ async def _session_stream_events(db: str | Path, job_id: int, *, after_id: int | async def _journal_stream_events(unit: str, *, shutdown_event: asyncio.Event | None = None): try: stream = stream_journal_lines(unit) + line_task: asyncio.Task | None = None + shutdown_task: asyncio.Task | None = None try: while shutdown_event is None or not shutdown_event.is_set(): line_task = asyncio.create_task(anext(stream)) @@ -322,6 +324,11 @@ async def _journal_stream_events(unit: str, *, shutdown_event: asyncio.Event | N return yield _sse_event("journal_line", {"unit": unit, "line": line}) finally: + pending_tasks = [task for task in (line_task, shutdown_task) if task is not None and not task.done()] + for task in pending_tasks: + task.cancel() + if pending_tasks: + await asyncio.gather(*pending_tasks, return_exceptions=True) await stream.aclose() except FileNotFoundError: yield _sse_event("journal_error", {"unit": unit, "error": "journalctl_not_found"}) diff --git a/tests/test_backend.py b/tests/test_backend.py index da8e4f8..12265e9 100644 --- a/tests/test_backend.py +++ b/tests/test_backend.py @@ -958,6 +958,31 @@ async def stream_until_shutdown(): asyncio.run(stream_until_shutdown()) +def test_dashboard_journal_stream_cancels_pending_read_when_client_disconnects(monkeypatch): + closed = False + + async def fake_stream_journal_lines(unit): + nonlocal closed + try: + await asyncio.sleep(60) + yield "unreachable" + finally: + closed = True + + async def disconnect_stream(): + shutdown = asyncio.Event() + monkeypatch.setattr("github_agent_bridge.backend.stream_journal_lines", fake_stream_journal_lines) + stream = _journal_stream_events("github-agent-bridge-dashboard.service", shutdown_event=shutdown) + pending_chunk = asyncio.create_task(anext(stream)) + await asyncio.sleep(0) + pending_chunk.cancel() + with pytest.raises(asyncio.CancelledError): + await pending_chunk + assert closed is True + + asyncio.run(disconnect_stream()) + + def test_dashboard_requires_auth_by_default(tmp_path): db = tmp_path / "bridge.sqlite3" JobQueue(db)