Fix flaky message deletion: exact lane match in outbox history + message_id fallback on delete
A flat-lane history (thread_id=None) returned frames from ALL threads, so auto-threaded messages leaked into the flat lane on restart. Deleting them from the flat lane then sent thread_id=None, which matched nothing in outbox.delete_message/message_info (exact lane match) -> removed=0, no session-store purge, and the messages resurrected from the outbox on the next app restart. - history: exact lane match (flat lane shows only flat-lane frames, per docs/06 §6.3) - delete_message / message_info: fall back to the unique message_id (uuid4) when the exact lane matches nothing, so deletes with a stale/missing thread_id still remove the frames and the purge finds its row
This commit is contained in:
1 parent
ca622d3a39
commit
a4e4a4ea63
1 file changed
+76
-46
+76
-46
@@ -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
|
||||
|
||||
Reference in new issue
Block a user