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).
This commit is contained in:
1 parent
d801a18db5
commit
29d0c1a73f
2 files changed
+33
-16
No files matched your search
@@ -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
|
||||
|
||||
+14
-10
@@ -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
|
||||
|
||||
Reference in new issue
Block a user