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