From 58f6678e6d01685be026fc05c1bb7eb6f58020bb Mon Sep 17 00:00:00 2001 From: JonthanaHanh <92574114+JonthanaHanh@users.noreply.github.com> Date: Tue, 28 Jul 2026 08:09:17 +0700 Subject: [PATCH] fix(gateway): flush pending messages to disk before shutdown clear (#72680) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit When FTS5 index corruption prevents INSERT INTO messages, the gateway accumulates messages in _pending_messages (memory-only). On shutdown, .clear() discards the only surviving copy — permanent user data loss. Changes: - Add gateway/shutdown_flush.py with flush_pending_to_file() and recover_pending_to_db() for two-phase data preservation. - Patch gateway/run.py: flush runner._pending_messages before clear() in _stop_impl_body, and recover on startup after runner.start(). - Patch gateway/platforms/base.py: flush adapter._pending_messages before clear() in the adapter shutdown path. Recovery behavior: - Reads pending JSON files from ~/.hermes/pending_messages/ - Inserts messages into state.db directly - Per-session isolation: one corrupt session doesn't block others - Successful recovery deletes the flush file - Failed recovery re-saves for next startup retry --- gateway/platforms/base.py | 6 ++ gateway/run.py | 19 ++++ gateway/shutdown_flush.py | 195 ++++++++++++++++++++++++++++++++++++++ 3 files changed, 220 insertions(+) create mode 100644 gateway/shutdown_flush.py diff --git a/gateway/platforms/base.py b/gateway/platforms/base.py index e55c91b4dd18..10ee3e3e109e 100644 --- a/gateway/platforms/base.py +++ b/gateway/platforms/base.py @@ -6073,6 +6073,12 @@ class BasePlatformAdapter(ABC): self._background_tasks.clear() self._expected_cancelled_tasks.clear() self._session_tasks.clear() + # Flush pending messages to disk before clearing (#72680). + try: + from gateway.shutdown_flush import flush_pending_to_file + flush_pending_to_file(self._pending_messages, reason="adapter_shutdown") + except Exception: + pass self._pending_messages.clear() self._active_sessions.clear() for state in list(self._text_debounce_store().values()): diff --git a/gateway/run.py b/gateway/run.py index 3927b7f371a4..2297e264ff2f 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -9968,6 +9968,15 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew self._running_agents_ts.clear() if hasattr(self, "_active_session_leases"): self._active_session_leases.clear() + # Flush pending messages to disk before clearing (#72680). + # When FTS5 corruption prevents message persistence, the + # in-memory _pending_messages dict holds the only surviving + # copy. Clearing without flushing causes permanent data loss. + try: + from gateway.shutdown_flush import flush_pending_to_file + flush_pending_to_file(self._pending_messages, reason="shutdown") + except Exception: + pass self._pending_messages.clear() self._pending_approvals.clear() if hasattr(self, '_busy_ack_ts'): @@ -24490,6 +24499,16 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = success = await runner.start() if not success: return False + # Recover any pending messages flushed during a previous shutdown (#72680). + try: + from gateway.shutdown_flush import recover_pending_to_db + recovered = recover_pending_to_db() + if recovered: + logger.info( + "Recovered %d pending message(s) from shutdown flush", recovered, + ) + except Exception: + pass if runner.should_exit_cleanly: if runner.exit_reason: logger.error("Gateway exiting cleanly: %s", runner.exit_reason) diff --git a/gateway/shutdown_flush.py b/gateway/shutdown_flush.py new file mode 100644 index 000000000000..d12f5b0cc074 --- /dev/null +++ b/gateway/shutdown_flush.py @@ -0,0 +1,195 @@ +"""Flush pending messages to disk before shutdown to prevent data loss. + +When FTS5 index corruption prevents ``INSERT INTO messages``, the gateway +accumulates messages in ``_pending_messages`` (memory-only). On shutdown, +``.clear()`` discards the only surviving copy — permanent user data loss. + +This module provides two hooks: + +1. ``flush_pending_to_file()`` — called BEFORE ``_pending_messages.clear()`` + during shutdown. Serialises any non-empty pending slots to a JSON file + under ``~/.hermes/pending_messages/``. + +2. ``recover_pending_to_db()`` — called AFTER ``runner.start()`` on startup. + Reads flush files, inserts messages directly into state.db, then deletes + the flush file on success. + +See issue #72680 for the full incident report. +""" + +from __future__ import annotations + +import json +import logging +import time +from pathlib import Path +from typing import Any, Dict, Optional + +logger = logging.getLogger(__name__) + +_FLUSH_DIR = Path.home() / ".hermes" / "pending_messages" + + +def _get_flush_dir() -> Path: + _FLUSH_DIR.mkdir(parents=True, exist_ok=True) + return _FLUSH_DIR + + +def flush_pending_to_file( + pending: Dict[str, Any], + *, + reason: str = "shutdown", +) -> int: + """Serialise non-empty ``_pending_messages`` slots to disk. + + Parameters + ---------- + pending: + The adapter or runner ``_pending_messages`` dict. Values may be + ``MessageEvent`` objects (adapter) or plain strings (runner). + reason: + Logged context (``shutdown``, ``restart``, etc.). + + Returns + ------- + int + Number of sessions flushed. + """ + if not pending: + return 0 + + flush_dir = _get_flush_dir() + ts = int(time.time()) + flushed = 0 + + for session_key, value in list(pending.items()): + if value is None: + continue + try: + serialised = _serialise_value(value) + if serialised is None: + continue + safe_key = session_key.replace("/", "_").replace("\\", "_") + path = flush_dir / f"{safe_key}_{ts}.json" + path.write_text( + json.dumps({ + "session_key": session_key, + "reason": reason, + "ts": ts, + "data": serialised, + }, ensure_ascii=False, default=str), + encoding="utf-8", + ) + flushed += 1 + except Exception as exc: + logger.debug( + "Failed to flush pending message for %s: %s", + session_key, exc, + ) + + if flushed: + logger.info( + "Flushed %d pending message(s) to %s (reason=%s)", + flushed, flush_dir, reason, + ) + return flushed + + +def _serialise_value(value: Any) -> Optional[dict]: + """Convert a pending message value to a JSON-serialisable dict.""" + # MessageEvent objects have a .text attribute and other fields + if hasattr(value, "text"): + result: Dict[str, Any] = {"text": getattr(value, "text", "")} + # Preserve additional fields if present + for attr in ("session_id", "platform", "sender_id", "sender_name", + "reply_to", "media", "raw_event"): + val = getattr(value, attr, None) + if val is not None: + try: + json.dumps(val) + result[attr] = val + except (TypeError, ValueError): + result[attr] = str(val) + return result + # Plain string (runner-level _pending_messages) + if isinstance(value, str): + return {"text": value} + # Dict — try direct serialisation + if isinstance(value, dict): + try: + json.dumps(value) + return value + except (TypeError, ValueError): + return {"text": str(value)} + return {"text": str(value)} + + +def recover_pending_to_db( + db_path: Optional[Path] = None, +) -> int: + """Recover flushed pending messages into state.db. + + Reads all ``*.json`` files from the flush directory, inserts messages + into the ``messages`` table, and deletes the file on success. + + Parameters + ---------- + db_path: + Path to state.db. Defaults to ``~/.hermes/state.db``. + + Returns + ------- + int + Number of messages recovered. + """ + import sqlite3 + + flush_dir = _get_flush_dir() + flush_files = sorted(flush_dir.glob("*.json")) + if not flush_files: + return 0 + + if db_path is None: + db_path = Path.home() / ".hermes" / "state.db" + if not db_path.exists(): + logger.warning("state.db not found at %s — skipping recovery", db_path) + return 0 + + recovered = 0 + for path in flush_files: + try: + payload = json.loads(path.read_text(encoding="utf-8")) + session_key = payload.get("session_key", "") + data = payload.get("data", {}) + text = data.get("text", "") + if not text or not session_key: + path.unlink(missing_ok=True) + continue + + conn = sqlite3.connect(str(db_path)) + try: + conn.execute( + "INSERT INTO messages (session_key, role, content, created_at) " + "VALUES (?, ?, ?, ?)", + (session_key, "user", text, payload.get("ts", int(time.time()))), + ) + conn.commit() + recovered += 1 + path.unlink(missing_ok=True) + except Exception as exc: + logger.warning( + "Failed to recover pending message for %s: %s", + session_key, exc, + ) + # Re-save so next startup can retry + finally: + conn.close() + except Exception as exc: + logger.debug("Failed to read flush file %s: %s", path, exc) + path.unlink(missing_ok=True) + + if recovered: + logger.info( + "Recovered %d pending message(s) from shutdown flush", recovered, + ) + return recovered