diff --git a/gateway/run.py b/gateway/run.py index ced859dd7e..9b9e7eb12e 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -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 diff --git a/gateway/shutdown_flush.py b/gateway/shutdown_flush.py index a09b63ee2a..07bcf7e86f 100644 --- a/gateway/shutdown_flush.py +++ b/gateway/shutdown_flush.py @@ -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. diff --git a/tests/gateway/test_shutdown_flush.py b/tests/gateway/test_shutdown_flush.py index efe6f59572..f966ea896d 100644 --- a/tests/gateway/test_shutdown_flush.py +++ b/tests/gateway/test_shutdown_flush.py @@ -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