diff --git a/gateway-plugin/outbox.py b/gateway-plugin/outbox.py index 0352ef5..9cf336d 100644 --- a/gateway-plugin/outbox.py +++ b/gateway-plugin/outbox.py @@ -35,6 +35,14 @@ _PRUNE_INTERVAL_S = 3600.0 DEFAULT_MAX_ROWS = 5000 +def _frame_thread_id(frame: dict[str, Any]) -> str | None: + """A frame's lane: its ``thread_id``, normalized (absent/blank -> None).""" + tid = frame.get("thread_id") + if not isinstance(tid, str) or not tid.strip(): + return None + return tid + + class Outbox: """Persistent outbox under ``get_hermes_home()/"iris"``. @@ -197,7 +205,12 @@ class Outbox: continue if not isinstance(frame, dict): continue - if thread_id is not None and frame.get("thread_id") != thread_id: + # Exact lane match: a flat-lane history (thread_id=None) must NOT + # include frames that belong to a thread, and vice versa. (The old + # loose filter let auto-threaded messages leak into the flat lane + # on restart, where a delete sent with thread_id=None then matched + # nothing in the outbox and the messages "resurrected" later.) + if _frame_thread_id(frame) != thread_id: continue ftype = frame.get("type") payload = frame.get("payload") @@ -289,6 +302,11 @@ class Outbox: standalone ``message`` frame when present, else the ``message.stop`` 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".) """ if not message_id: return None @@ -296,34 +314,38 @@ class Outbox: rows = self._conn.execute( "SELECT frame FROM outbox WHERE chat_id = ?", (chat_id,) ).fetchall() - msg_frame: dict[str, Any] | None = None - stop_frame: dict[str, Any] | None = None - for r in rows: - try: - frame = json.loads(r["frame"]) - except (json.JSONDecodeError, TypeError): - continue - if not isinstance(frame, dict): - continue - if thread_id is not None and frame.get("thread_id") != thread_id: - continue - payload = frame.get("payload") - if not isinstance(payload, dict) or payload.get("message_id") != message_id: - continue - ftype = frame.get("type") - if ftype == "message": - msg_frame = { - "role": payload.get("role"), - "text": payload.get("text", ""), - "ts": payload.get("ts"), - } - elif ftype == "message.stop": - stop_frame = { - "role": "assistant", - "text": payload.get("final_text", ""), - "ts": payload.get("ts"), - } - return msg_frame or stop_frame + + def scan(lane: str | None) -> dict[str, Any] | None: + msg_frame: dict[str, Any] | None = None + stop_frame: dict[str, Any] | None = None + for r in rows: + try: + frame = json.loads(r["frame"]) + except (json.JSONDecodeError, TypeError): + continue + if not isinstance(frame, dict): + continue + if lane is not None and _frame_thread_id(frame) != lane: + continue + payload = frame.get("payload") + if not isinstance(payload, dict) or payload.get("message_id") != message_id: + continue + ftype = frame.get("type") + if ftype == "message": + msg_frame = { + "role": payload.get("role"), + "text": payload.get("text", ""), + "ts": payload.get("ts"), + } + elif ftype == "message.stop": + stop_frame = { + "role": "assistant", + "text": payload.get("final_text", ""), + "ts": payload.get("ts"), + } + return msg_frame or stop_frame + + return scan(thread_id) or scan(None) def delete_message( self, @@ -337,10 +359,13 @@ class Outbox: ``message.update`` / ``message.stop`` / ``media.offer`` / ``commentary``); all of them are removed so neither ``history`` nor a ``sync`` replay can resurrect the message. The delete is scoped to the - exact lane: a flat-lane delete (``thread_id=None``) matches only frames - with no ``thread_id``, and a thread delete matches only that thread's - frames (a ``message_id`` is unique to one lane, so this is a safety - net, not a filter that drops real frames). Returns the number of rows + exact lane first: a flat-lane delete (``thread_id=None``) matches only + frames with no ``thread_id``, and a thread delete matches only that + thread's frames. When the exact lane matches nothing, the delete falls + back to the ``message_id`` alone (it is a unique uuid4, so it cannot + hit the wrong message) — this keeps deletes working when the request's + lane is stale or missing (e.g. a message the app cached in the flat + lane that the gateway auto-threaded). Returns the number of rows removed (0 when the message is not in the outbox — e.g. already pruned by retention). """ @@ -350,19 +375,24 @@ class Outbox: rows = self._conn.execute( "SELECT cursor, frame FROM outbox WHERE chat_id = ?", (chat_id,) ).fetchall() - cursors: list[int] = [] - for r in rows: - try: - frame = json.loads(r["frame"]) - except (json.JSONDecodeError, TypeError): - continue - if not isinstance(frame, dict): - continue - if frame.get("thread_id") != thread_id: - continue - payload = frame.get("payload") - if isinstance(payload, dict) and payload.get("message_id") == message_id: - cursors.append(int(r["cursor"])) + + def cursors_for(lane: str | None) -> list[int]: + out: list[int] = [] + for r in rows: + try: + frame = json.loads(r["frame"]) + except (json.JSONDecodeError, TypeError): + continue + if not isinstance(frame, dict): + continue + if lane is not None 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) if not cursors: return 0 # One bound-parameter delete per cursor (a message spans only a few