From 29d0c1a73f015c1f651e7e81ab05e869b84f01b4 Mon Sep 17 00:00:00 2001 From: ARIA Date: Sun, 23 Aug 2026 14:45:09 +0200 Subject: [PATCH] Fix two gateway test failures: outbox lane scoping + SSE teardown race - outbox: delete_message/message_info now match the exact lane first (a flat-lane delete/lookup with thread_id=None sees only frames with no thread_id) and fall back to the message_id across all lanes only when the exact lane matches nothing. Previously lane=None meant 'any lane' in the first pass, so a flat-lane delete also removed same-id frames from threads (test expected 3 removed, got 4). - http_server: the SSE live loop skipped queued frames when stop() set sub.closed before the handler thread reached the loop (descheduled under load between the initial hello/status writes and the loop). The loop now drains frames queued before the close, so the status{restarting} teardown broadcast always reaches the client before EOF (test_disconnect_broadcasts_status_restarting was flaky ~70% under CPU load). --- gateway-plugin/http_server.py | 25 +++++++++++++++++++------ gateway-plugin/outbox.py | 24 ++++++++++++++---------- 2 files changed, 33 insertions(+), 16 deletions(-) diff --git a/gateway-plugin/http_server.py b/gateway-plugin/http_server.py index a2e495d..e9751fc 100644 --- a/gateway-plugin/http_server.py +++ b/gateway-plugin/http_server.py @@ -523,12 +523,25 @@ class HttpServer: for snap in self._adapter.todo_snapshot_frames(): self._write_sse(handler, "frame", None, snap.to_json()) # 3. Live frames (cursor=None frames have no id). - while not sub.closed.is_set(): - try: - item = sub.q.get(timeout=SSE_HEARTBEAT_S) - except queue.Empty: - self._write_raw(handler, ": hb\n\n") - continue + while True: + if sub.closed.is_set(): + # stop() can land between the initial writes above and + # this loop (the handler thread is descheduled under + # load): drain the frames queued before the close — e.g. + # the status{restarting} teardown broadcast — so the + # client sees them before EOF instead of losing them to + # the closed check. + try: + item = sub.q.get_nowait() + except queue.Empty: + reason = "stopped" + break + else: + try: + item = sub.q.get(timeout=SSE_HEARTBEAT_S) + except queue.Empty: + self._write_raw(handler, ": hb\n\n") + continue if item is _STOP: reason = "stopped" break diff --git a/gateway-plugin/outbox.py b/gateway-plugin/outbox.py index 9cf336d..0e0f98e 100644 --- a/gateway-plugin/outbox.py +++ b/gateway-plugin/outbox.py @@ -303,10 +303,11 @@ class Outbox: frame of a streamed reply -- or ``None`` when the message is not in the outbox (e.g. already pruned by retention). - The lane is matched exactly first; when that finds nothing the lookup - falls back to the ``message_id`` alone (it is a unique uuid4), so a - stale/missing ``thread_id`` on the request still resolves the row. - (``lane=None`` in the scan means "any lane".) + The lane is matched exactly first (a flat-lane lookup, ``thread_id + = None``, sees only frames with no ``thread_id``); when that finds + nothing the lookup falls back to the ``message_id`` alone across all + lanes (it is a unique uuid4), so a stale/missing ``thread_id`` on the + request still resolves the row. """ if not message_id: return None @@ -315,7 +316,7 @@ class Outbox: "SELECT frame FROM outbox WHERE chat_id = ?", (chat_id,) ).fetchall() - def scan(lane: str | None) -> dict[str, Any] | None: + def scan(lane: str | None, exact: bool) -> dict[str, Any] | None: msg_frame: dict[str, Any] | None = None stop_frame: dict[str, Any] | None = None for r in rows: @@ -325,7 +326,7 @@ class Outbox: continue if not isinstance(frame, dict): continue - if lane is not None and _frame_thread_id(frame) != lane: + if exact and _frame_thread_id(frame) != lane: continue payload = frame.get("payload") if not isinstance(payload, dict) or payload.get("message_id") != message_id: @@ -345,7 +346,7 @@ class Outbox: } return msg_frame or stop_frame - return scan(thread_id) or scan(None) + return scan(thread_id, exact=True) or scan(None, exact=False) def delete_message( self, @@ -376,7 +377,7 @@ class Outbox: "SELECT cursor, frame FROM outbox WHERE chat_id = ?", (chat_id,) ).fetchall() - def cursors_for(lane: str | None) -> list[int]: + def cursors_for(lane: str | None, exact: bool) -> list[int]: out: list[int] = [] for r in rows: try: @@ -385,14 +386,17 @@ class Outbox: continue if not isinstance(frame, dict): continue - if lane is not None and _frame_thread_id(frame) != lane: + if exact and _frame_thread_id(frame) != lane: continue payload = frame.get("payload") if isinstance(payload, dict) and payload.get("message_id") == message_id: out.append(int(r["cursor"])) return out - cursors = cursors_for(thread_id) or cursors_for(None) + # Exact lane first (a flat-lane delete must not reach into + # threads); fall back to the message_id across all lanes only + # when the exact lane matches nothing (stale/missing thread_id). + cursors = cursors_for(thread_id, exact=True) or cursors_for(thread_id, exact=False) if not cursors: return 0 # One bound-parameter delete per cursor (a message spans only a few