fix(gateway): flush the FIFO overflow tail to disk at shutdown too (#99882)
Sibling site of the same loss class. The #72680 shutdown flush only serialised the adapter slot (_pending_messages); the FIFO tail parked in SessionState.conversation.queued_events was discarded with the process, so every follow-up queued behind the head at restart time vanished the same way the idle-orphan did. flush_overflow_to_file writes one payload per overflow event in the slot-flush shape (plus seq for arrival order), so the existing recover_pending_to_db startup replay inserts them with no new reader. Wired into _stop_impl beside the slot flush.
This commit is contained in:
@@ -16523,6 +16523,21 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
|
||||
flush_pending_to_file(dict(self._pending_messages), reason="shutdown")
|
||||
except Exception:
|
||||
pass
|
||||
# The FIFO tail lives in SessionState.conversation.queued_events,
|
||||
# not in the slot dict above — flush it too or every follow-up
|
||||
# parked in overflow at restart time is lost (#99882).
|
||||
try:
|
||||
from gateway.shutdown_flush import flush_overflow_to_file
|
||||
flush_overflow_to_file(
|
||||
{
|
||||
_k: list(_v)
|
||||
for _k, _v in dict(getattr(self, "_queued_events", None) or {}).items()
|
||||
if _v
|
||||
},
|
||||
reason="shutdown",
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
# On the real runner these are live SessionState views whose
|
||||
# clear() resets one field per session — never a wholesale dict
|
||||
# swap, so a concurrent writer on another session can't lose its
|
||||
|
||||
@@ -142,6 +142,66 @@ def flush_pending_to_file(
|
||||
return flushed
|
||||
|
||||
|
||||
def flush_overflow_to_file(
|
||||
overflow_by_session: Dict[str, Any],
|
||||
*,
|
||||
reason: str = "shutdown",
|
||||
) -> int:
|
||||
"""Serialise the FIFO overflow tails (``queued_events``) to disk.
|
||||
|
||||
Sibling of :func:`flush_pending_to_file` for the second half of the
|
||||
gateway FIFO (#99882): the adapter slot holds the queue head, and the
|
||||
per-session ``SessionState.conversation.queued_events`` list holds the
|
||||
tail. Shutdown flushed only the slot, so every follow-up parked in
|
||||
overflow at restart time vanished with the process. Each overflow
|
||||
event is written as its own payload in the same shape as a slot flush
|
||||
so ``recover_pending_to_db`` replays them unchanged; a ``seq`` field
|
||||
preserves arrival order within a session.
|
||||
|
||||
Returns the number of events flushed.
|
||||
"""
|
||||
if not overflow_by_session:
|
||||
return 0
|
||||
|
||||
flush_dir = _get_flush_dir()
|
||||
ts = int(time.time())
|
||||
flushed = 0
|
||||
|
||||
for session_key, events in list(overflow_by_session.items()):
|
||||
if not session_key or not events:
|
||||
continue
|
||||
for seq, value in enumerate(list(events)):
|
||||
if value is None:
|
||||
continue
|
||||
try:
|
||||
serialised = _serialise_value(value)
|
||||
if serialised is None:
|
||||
continue
|
||||
_write_payload(
|
||||
flush_dir,
|
||||
{
|
||||
"session_key": session_key,
|
||||
"reason": reason,
|
||||
"ts": ts,
|
||||
"seq": seq,
|
||||
"data": serialised,
|
||||
},
|
||||
)
|
||||
flushed += 1
|
||||
except Exception as exc:
|
||||
logger.debug(
|
||||
"Failed to flush overflow message for %s: %s",
|
||||
session_key, exc,
|
||||
)
|
||||
|
||||
if flushed:
|
||||
logger.info(
|
||||
"Flushed %d queued overflow message(s) to %s (reason=%s)",
|
||||
flushed, flush_dir, reason,
|
||||
)
|
||||
return flushed
|
||||
|
||||
|
||||
# Reason tag for transcript messages dropped by the in-memory pending cap
|
||||
# during live operation (#78182). These payloads carry the full transcript
|
||||
# message dict so they can be replayed verbatim once the DB recovers.
|
||||
|
||||
@@ -11,6 +11,7 @@ import pytest
|
||||
|
||||
from gateway.shutdown_flush import (
|
||||
_serialise_value,
|
||||
flush_overflow_to_file,
|
||||
flush_pending_to_file,
|
||||
recover_pending_to_db,
|
||||
)
|
||||
@@ -168,3 +169,72 @@ def test_get_flush_dir_uses_get_hermes_home(tmp_path, monkeypatch):
|
||||
assert result == tmp_path / "pending_messages"
|
||||
|
||||
|
||||
|
||||
|
||||
# ── FIFO overflow tail durability (#99882) ─────────────────────────────
|
||||
|
||||
|
||||
def _overflow_event(text: str, session_id: str = "20260901_120000_fifo"):
|
||||
event = MagicMock()
|
||||
event.text = text
|
||||
event.session_id = session_id
|
||||
event.platform = "telegram"
|
||||
event.sender_id = "1572286605"
|
||||
event.sender_name = "tester"
|
||||
event.reply_to = None
|
||||
event.media = None
|
||||
event.raw_event = None
|
||||
return event
|
||||
|
||||
|
||||
def test_flush_overflow_writes_one_payload_per_event_in_arrival_order(tmp_path, monkeypatch):
|
||||
"""The FIFO tail (queued_events) must survive shutdown like the slot does.
|
||||
|
||||
Each overflow entry is its own recover_pending_to_db-compatible payload,
|
||||
with ``seq`` recording arrival order inside the session.
|
||||
"""
|
||||
flush_dir = _make_flush_dir(tmp_path)
|
||||
monkeypatch.setattr("gateway.shutdown_flush._get_flush_dir", lambda: flush_dir)
|
||||
|
||||
count = flush_overflow_to_file(
|
||||
{
|
||||
"agent:main:telegram:dm:1": [
|
||||
_overflow_event("follow-up B"),
|
||||
_overflow_event("follow-up C"),
|
||||
],
|
||||
"agent:main:telegram:dm:2": [],
|
||||
"": [_overflow_event("keyless — skipped")],
|
||||
},
|
||||
reason="shutdown",
|
||||
)
|
||||
assert count == 2
|
||||
payloads = sorted(
|
||||
(json.loads(f.read_text(encoding="utf-8")) for f in flush_dir.glob("*.json")),
|
||||
key=lambda p: p["seq"],
|
||||
)
|
||||
assert [p["data"]["text"] for p in payloads] == ["follow-up B", "follow-up C"]
|
||||
assert {p["session_key"] for p in payloads} == {"agent:main:telegram:dm:1"}
|
||||
assert all(p["reason"] == "shutdown" for p in payloads)
|
||||
|
||||
|
||||
def test_flushed_overflow_is_replayed_by_recover_pending_to_db(tmp_path, monkeypatch):
|
||||
"""Round-trip: overflow payloads use the slot-flush shape, so the existing
|
||||
startup recovery inserts them as user rows without any new reader."""
|
||||
flush_dir = _make_flush_dir(tmp_path)
|
||||
monkeypatch.setattr("gateway.shutdown_flush._get_flush_dir", lambda: flush_dir)
|
||||
flush_overflow_to_file({"agent:main:telegram:dm:1": [_overflow_event("orphan-1")]})
|
||||
|
||||
db = MagicMock()
|
||||
recovered = recover_pending_to_db(session_db=db)
|
||||
assert recovered == 1
|
||||
db.append_message.assert_called_once()
|
||||
kwargs = db.append_message.call_args.kwargs
|
||||
assert kwargs["session_id"] == "20260901_120000_fifo"
|
||||
assert kwargs["role"] == "user"
|
||||
assert kwargs["content"] == "orphan-1"
|
||||
assert list(flush_dir.glob("*.json")) == []
|
||||
|
||||
|
||||
def test_flush_overflow_noop_on_empty():
|
||||
assert flush_overflow_to_file({}) == 0
|
||||
assert flush_overflow_to_file({"k": []}) == 0
|
||||
|
||||
Reference in New Issue
Block a user