diff --git a/hermes_state.py b/hermes_state.py index 4f7b6453461a..fd19a173ebb3 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -1486,7 +1486,8 @@ def repair_state_db_schema(db_path: Path, *, backup: bool = True) -> Dict[str, A VALUES ('delete', old.id, old.content, old.tool_name, old.tool_calls); END; -CREATE TRIGGER IF NOT EXISTS messages_fts_cjk_update AFTER UPDATE ON messages +CREATE TRIGGER IF NOT EXISTS messages_fts_cjk_update +AFTER UPDATE OF content, tool_name, tool_calls, role ON messages WHEN (old.content IS NOT new.content OR old.tool_name IS NOT new.tool_name OR old.tool_calls IS NOT new.tool_calls diff --git a/hermes_state_common.py b/hermes_state_common.py index 19dbaa39d78e..6c62bd99e93e 100644 --- a/hermes_state_common.py +++ b/hermes_state_common.py @@ -357,7 +357,11 @@ def _ephemeral_child_sql(alias: str = "s") -> str: VALUES ('delete', old.id, old.content, old.tool_name, old.tool_calls); END; -CREATE TRIGGER IF NOT EXISTS messages_fts_update AFTER UPDATE ON messages +-- UPDATE OF skips the trigger entirely for non-content column writes +-- (status/compacted/observed/etc.), which is stronger than the WHEN gate +-- alone and avoids FTS I/O saturation on large state.db (#68858 / #73639). +CREATE TRIGGER IF NOT EXISTS messages_fts_update +AFTER UPDATE OF content, tool_name, tool_calls ON messages WHEN (old.content IS NOT new.content OR old.tool_name IS NOT new.tool_name OR old.tool_calls IS NOT new.tool_calls) @@ -425,7 +429,8 @@ def _ephemeral_child_sql(alias: str = "s") -> str: VALUES ('delete', old.id, old.content, old.tool_name, old.tool_calls); END; -CREATE TRIGGER IF NOT EXISTS messages_fts_trigram_update AFTER UPDATE ON messages +CREATE TRIGGER IF NOT EXISTS messages_fts_trigram_update +AFTER UPDATE OF content, tool_name, tool_calls, role ON messages WHEN (old.content IS NOT new.content OR old.tool_name IS NOT new.tool_name OR old.tool_calls IS NOT new.tool_calls @@ -486,7 +491,8 @@ def _ephemeral_child_sql(alias: str = "s") -> str: DELETE FROM messages_fts WHERE rowid = old.id; END; -CREATE TRIGGER IF NOT EXISTS messages_fts_update AFTER UPDATE ON messages BEGIN +CREATE TRIGGER IF NOT EXISTS messages_fts_update +AFTER UPDATE OF content, tool_name, tool_calls ON messages BEGIN DELETE FROM messages_fts WHERE rowid = old.id; INSERT INTO messages_fts(rowid, content) VALUES ( new.id, @@ -513,7 +519,8 @@ def _ephemeral_child_sql(alias: str = "s") -> str: DELETE FROM messages_fts_trigram WHERE rowid = old.id; END; -CREATE TRIGGER IF NOT EXISTS messages_fts_trigram_update AFTER UPDATE ON messages BEGIN +CREATE TRIGGER IF NOT EXISTS messages_fts_trigram_update +AFTER UPDATE OF content, tool_name, tool_calls ON messages BEGIN DELETE FROM messages_fts_trigram WHERE rowid = old.id; INSERT INTO messages_fts_trigram(rowid, content) VALUES ( new.id, diff --git a/hermes_state_schema.py b/hermes_state_schema.py index 15513a7e7337..b15c22eb0d40 100644 --- a/hermes_state_schema.py +++ b/hermes_state_schema.py @@ -16,6 +16,7 @@ from hermes_constants import get_hermes_home from hermes_state_common import ( DEFERRED_INDEX_SQL, + FTS_CJK_STALE_KEY, FTS_SQL, FTS_STORAGE_VERSION, FTS_TRIGRAM_SQL, @@ -56,6 +57,146 @@ def _fts_trigger_count(cursor: sqlite3.Cursor) -> int: ).fetchone() return int(row[0] if not isinstance(row, sqlite3.Row) else row[0]) + + @staticmethod + def _fts_update_trigger_needs_narrowing(sql: Optional[str]) -> bool: + """True when trigger SQL is missing AFTER UPDATE OF (still broad).""" + if not sql: + return False + # Collapse whitespace so multi-line DDL still matches. + compact = " ".join(sql.split()).upper() + # Already narrowed. + if "AFTER UPDATE OF " in compact: + return False + # Broad UPDATE trigger that we still need to replace. + return "AFTER UPDATE ON " in compact + + def _migrate_broad_fts_update_triggers(self, cursor: sqlite3.Cursor) -> int: + """Replace broad AFTER UPDATE FTS triggers with AFTER UPDATE OF variants. + + ``CREATE TRIGGER IF NOT EXISTS`` will not replace an existing broad + trigger, so installs that already created ``AFTER UPDATE ON messages`` + would keep firing on every messages row touch (status/compaction + writes included). Inspect ``sqlite_master``, drop any still-broad + UPDATE triggers, and re-apply the current DDL constants. + + No FTS rebuild: content correctness was already gated by WHEN clauses + on modern installs; OF only skips unnecessary trigger evaluation. + + Returns the number of triggers dropped (0 when already converged). + """ + import re as _re + + # CJK is a v23-only surface. Decide the layout before selecting + # destructive candidates so the legacy branch never drops a trigger + # it does not recreate. + legacy_layout = self._db_has_legacy_inline_fts(cursor) + update_names = ( + "messages_fts_update", + "messages_fts_trigram_update", + ) + if not legacy_layout and hasattr(self, "_ensure_fts_cjk_schema"): + update_names += ("messages_fts_cjk_update",) + placeholders = ", ".join("?" for _ in update_names) + rows = cursor.execute( + "SELECT name, sql FROM sqlite_master " + f"WHERE type = 'trigger' AND name IN ({placeholders})", + update_names, + ).fetchall() + to_drop = [] + for row in rows: + name = row[0] if not isinstance(row, sqlite3.Row) else row["name"] + sql = row[1] if not isinstance(row, sqlite3.Row) else row["sql"] + if self._fts_update_trigger_needs_narrowing(sql): + to_drop.append(name) + if not to_drop: + return 0 + + for name in to_drop: + # Trigger names are from our fixed allowlist, not user input. + if not _re.fullmatch(r"[A-Za-z0-9_]+", name): + continue + cursor.execute(f"DROP TRIGGER IF EXISTS {name}") + + # Re-apply current DDL so CREATE TRIGGER installs the OF variants. + # Choose legacy vs v23 the same way _init_schema does. + if legacy_layout: + self._ensure_fts_schema(cursor, "messages_fts", LEGACY_FTS_SQL) + self._ensure_fts_schema( + cursor, "messages_fts_trigram", LEGACY_FTS_TRIGRAM_SQL + ) + else: + self._ensure_fts_schema(cursor, "messages_fts", FTS_SQL) + self._ensure_fts_schema( + cursor, "messages_fts_trigram", FTS_TRIGRAM_SQL + ) + # CJK triggers live on the host SessionDB; only recreate one that + # this migration actually dropped. ``_ensure_fts_cjk_schema`` is + # documented never-raises and soft-fails OperationalError by + # clearing availability — raise-path handling alone is not + # enough. After ensure, require a narrowed CJK UPDATE trigger or + # durable quarantine (stale breadcrumb + unavailable). + if "messages_fts_cjk_update" in to_drop: + try: + self._ensure_fts_cjk_schema(cursor) + except Exception: + self._quarantine_cjk_after_update_of_migration(cursor) + logger.exception( + "CJK FTS re-ensure after UPDATE OF migration failed" + ) + raise + if not self._cjk_update_trigger_is_narrowed(cursor): + self._quarantine_cjk_after_update_of_migration(cursor) + logger.warning( + "CJK FTS UPDATE trigger missing or still broad after " + "UPDATE OF migration; marked stale and unavailable" + ) + + logger.info( + "Migrated %d broad FTS UPDATE trigger(s) to AFTER UPDATE OF " + "(no rebuild required)", + len(to_drop), + ) + return len(to_drop) + + def _cjk_update_trigger_is_narrowed(self, cursor: sqlite3.Cursor) -> bool: + """True when messages_fts_cjk_update exists with AFTER UPDATE OF.""" + row = cursor.execute( + "SELECT sql FROM sqlite_master " + "WHERE type = 'trigger' AND name = ?", + ("messages_fts_cjk_update",), + ).fetchone() + if not row: + return False + sql = row[0] if not isinstance(row, sqlite3.Row) else row["sql"] + return not self._fts_update_trigger_needs_narrowing(sql) + + def _quarantine_cjk_after_update_of_migration( + self, cursor: sqlite3.Cursor + ) -> None: + """Fail-closed after dropping CJK UPDATE during OF migration. + + Clears availability, persists ``fts_cjk_stale``, and drops any + residual broad/partial CJK UPDATE trigger so a later open cannot + ``CREATE TRIGGER IF NOT EXISTS`` a gap without rebuild. + """ + self._fts_cjk_available = False + try: + self.set_meta(FTS_CJK_STALE_KEY, "1", cursor=cursor) + except Exception: + logger.debug( + "Could not persist CJK FTS stale breadcrumb", + exc_info=True, + ) + try: + cursor.execute("DROP TRIGGER IF EXISTS messages_fts_cjk_update") + except Exception: + logger.debug( + "Could not drop residual CJK UPDATE trigger after quarantine", + exc_info=True, + ) + + @staticmethod def _rebuild_fts_indexes( cursor: sqlite3.Cursor, @@ -853,6 +994,11 @@ def _init_schema(self): # the surfaces above and gated on the loadable tokenizer: self._ensure_fts_cjk_schema(cursor) + # Replace any pre-existing broad AFTER UPDATE triggers with + # AFTER UPDATE OF variants. IF NOT EXISTS cannot rewrite them. + if getattr(self, "_fts_enabled", False): + self._migrate_broad_fts_update_triggers(cursor) + self._conn.commit() def _backfill_gateway_metadata_from_sessions_json( diff --git a/tests/test_fts_update_of_narrowing.py b/tests/test_fts_update_of_narrowing.py new file mode 100644 index 000000000000..2a8e0d504f0e --- /dev/null +++ b/tests/test_fts_update_of_narrowing.py @@ -0,0 +1,283 @@ +"""FTS UPDATE OF narrowing + migration (#73639 retargeted onto split SessionDB).""" + +from __future__ import annotations + +import sqlite3 +import tempfile +from pathlib import Path + +import pytest + +from hermes_state import SessionDB +from hermes_state_common import FTS_CJK_STALE_KEY +from hermes_state_schema import SessionSchemaMixin + + +def _trigger_sql(conn: sqlite3.Connection, name: str) -> str | None: + row = conn.execute( + "SELECT sql FROM sqlite_master WHERE type='trigger' AND name=?", + (name,), + ).fetchone() + return row[0] if row else None + + +def _assert_nonindexed_updates_bypass_missing_fts_target( + db: SessionDB, message_id: int +) -> None: + """Prove the UPDATE OF gate, not the trigger's content-change WHEN.""" + db._conn.execute("DROP TABLE messages_fts") + assert _trigger_sql(db._conn, "messages_fts_update") is not None + + db._conn.execute( + "UPDATE messages SET active = 0, compacted = 1, observed = 1 " + "WHERE id = ?", + (message_id,), + ) + with pytest.raises(sqlite3.OperationalError, match=r"no such table.*messages_fts"): + db._conn.execute( + "UPDATE messages SET content = 'changed' WHERE id = ?", + (message_id,), + ) + + +def _install_legacy_inline_base_fts(db: SessionDB) -> None: + """Replace v23 FTS with the broad inline shape shipped by v11..v22.""" + db._drop_fts_triggers(db._conn) + db._conn.executescript( + """ + DROP TABLE IF EXISTS messages_fts; + DROP TABLE IF EXISTS messages_fts_trigram; + DROP VIEW IF EXISTS messages_fts_trigram_src; + + CREATE VIRTUAL TABLE messages_fts USING fts5(content); + CREATE TRIGGER messages_fts_insert AFTER INSERT ON messages BEGIN + INSERT INTO messages_fts(rowid, content) VALUES ( + new.id, + COALESCE(new.content, '') || ' ' || + COALESCE(new.tool_name, '') || ' ' || + COALESCE(new.tool_calls, '') + ); + END; + CREATE TRIGGER messages_fts_delete AFTER DELETE ON messages BEGIN + DELETE FROM messages_fts WHERE rowid = old.id; + END; + CREATE TRIGGER messages_fts_update AFTER UPDATE ON messages BEGIN + DELETE FROM messages_fts WHERE rowid = old.id; + INSERT INTO messages_fts(rowid, content) VALUES ( + new.id, + COALESCE(new.content, '') || ' ' || + COALESCE(new.tool_name, '') || ' ' || + COALESCE(new.tool_calls, '') + ); + END; + """ + ) + + +def test_fresh_db_installs_update_of_triggers(tmp_path: Path): + db = SessionDB(db_path=tmp_path / "state.db") + try: + sql = _trigger_sql(db._conn, "messages_fts_update") + assert sql is not None + compact = " ".join(sql.split()).upper() + assert "AFTER UPDATE OF " in compact + assert "CONTENT" in compact + assert "TOOL_NAME" in compact + assert "TOOL_CALLS" in compact + + tri = _trigger_sql(db._conn, "messages_fts_trigram_update") + if tri: # trigram may be unavailable on some builds + tcompact = " ".join(tri.split()).upper() + assert "AFTER UPDATE OF " in tcompact + finally: + db.close() + + +def test_migrate_replaces_broad_update_trigger(tmp_path: Path): + path = tmp_path / "state.db" + db = SessionDB(db_path=path) + try: + # Force a broad trigger the way older installs had it. + db._conn.execute("DROP TRIGGER IF EXISTS messages_fts_update") + db._conn.execute( + """ + CREATE TRIGGER messages_fts_update AFTER UPDATE ON messages + BEGIN + SELECT 1; + END + """ + ) + db._conn.commit() + before = _trigger_sql(db._conn, "messages_fts_update") + assert "AFTER UPDATE OF" not in " ".join(before.split()).upper() + + dropped = db._migrate_broad_fts_update_triggers(db._conn) + db._conn.commit() + assert dropped >= 1 + + after = _trigger_sql(db._conn, "messages_fts_update") + assert after is not None + compact = " ".join(after.split()).upper() + assert "AFTER UPDATE OF " in compact + finally: + db.close() + + +def test_needs_narrowing_helper(): + assert SessionSchemaMixin._fts_update_trigger_needs_narrowing( + "CREATE TRIGGER t AFTER UPDATE ON messages BEGIN SELECT 1; END" + ) + assert not SessionSchemaMixin._fts_update_trigger_needs_narrowing( + "CREATE TRIGGER t AFTER UPDATE OF content ON messages BEGIN SELECT 1; END" + ) + assert not SessionSchemaMixin._fts_update_trigger_needs_narrowing(None) + + +def test_v23_status_only_update_bypasses_fts_trigger_body(tmp_path: Path): + db = SessionDB(db_path=tmp_path / "state.db") + try: + sid = "s1" + db.create_session(sid, source="test") + mid = db.append_message(sid, role="user", content="hello searchable") + _assert_nonindexed_updates_bypass_missing_fts_target(db, mid) + finally: + db.close() + + +def test_legacy_status_only_update_bypasses_migrated_fts_trigger(tmp_path: Path): + db = SessionDB(db_path=tmp_path / "state.db") + try: + sid = "s1" + db.create_session(sid, source="test") + mid = db.append_message(sid, role="user", content="legacy searchable") + _install_legacy_inline_base_fts(db) + + assert db._db_has_legacy_inline_fts(db._conn) + assert db._migrate_broad_fts_update_triggers(db._conn) >= 1 + assert "AFTER UPDATE OF" in " ".join( + _trigger_sql(db._conn, "messages_fts_update").split() + ).upper() + + _assert_nonindexed_updates_bypass_missing_fts_target(db, mid) + finally: + db.close() + + +def test_cjk_ensure_failure_marks_unavailable_and_propagates( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +): + path = tmp_path / "state.db" + db = SessionDB(db_path=path) + try: + db._conn.execute("DROP TRIGGER IF EXISTS messages_fts_cjk_update") + db._conn.execute( + "CREATE TRIGGER messages_fts_cjk_update " + "AFTER UPDATE ON messages BEGIN SELECT 1; END" + ) + db._fts_cjk_available = True + + def _fail_cjk_ensure(_cursor): + raise sqlite3.DatabaseError("injected CJK ensure failure") + + monkeypatch.setattr(db, "_ensure_fts_cjk_schema", _fail_cjk_ensure) + + with pytest.raises(sqlite3.DatabaseError, match="injected CJK ensure failure"): + db._migrate_broad_fts_update_triggers(db._conn) + assert db._fts_cjk_available is False + assert _trigger_sql(db._conn, "messages_fts_cjk_update") is None + + # The DROP is autocommitted. Fail-closed therefore needs a durable + # breadcrumb that other processes can observe, not just an instance + # flag on the SessionDB whose initialization is aborting. + with sqlite3.connect(path) as observer: + stale = observer.execute( + "SELECT value FROM state_meta WHERE key = ?", + (FTS_CJK_STALE_KEY,), + ).fetchone() + assert stale == ("1",) + assert _trigger_sql(observer, "messages_fts_cjk_update") is None + finally: + db.close() + + +def test_cjk_soft_fail_ensure_without_raise_marks_stale( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +): + """Production ensure swallows OperationalError — still must quarantine.""" + path = tmp_path / "state.db" + db = SessionDB(db_path=path) + try: + db._conn.execute("DROP TRIGGER IF EXISTS messages_fts_cjk_update") + db._conn.execute( + "CREATE TRIGGER messages_fts_cjk_update " + "AFTER UPDATE ON messages BEGIN SELECT 1; END" + ) + db._fts_cjk_available = True + + def _soft_fail_cjk_ensure(_cursor): + # Mirrors real _ensure_fts_cjk_schema OperationalError path: + # clear availability, do not raise, do not recreate triggers. + db._fts_cjk_available = False + + monkeypatch.setattr(db, "_ensure_fts_cjk_schema", _soft_fail_cjk_ensure) + + dropped = db._migrate_broad_fts_update_triggers(db._conn) + assert dropped >= 1 + assert db._fts_cjk_available is False + assert _trigger_sql(db._conn, "messages_fts_cjk_update") is None + + with sqlite3.connect(path) as observer: + stale = observer.execute( + "SELECT value FROM state_meta WHERE key = ?", + (FTS_CJK_STALE_KEY,), + ).fetchone() + assert stale == ("1",) + assert _trigger_sql(observer, "messages_fts_cjk_update") is None + finally: + db.close() + + +def test_cjk_broad_trigger_is_restored_as_update_of( + tmp_path: Path, monkeypatch: pytest.MonkeyPatch +): + db = SessionDB(db_path=tmp_path / "state.db") + try: + db._conn.execute("DROP TRIGGER IF EXISTS messages_fts_cjk_update") + db._conn.execute( + "CREATE TRIGGER messages_fts_cjk_update " + "AFTER UPDATE ON messages BEGIN SELECT 1; END" + ) + + def _restore_cjk_update_trigger(cursor): + cursor.execute( + "CREATE TRIGGER messages_fts_cjk_update " + "AFTER UPDATE OF content, tool_name, tool_calls ON messages " + "BEGIN SELECT 1; END" + ) + + monkeypatch.setattr(db, "_ensure_fts_cjk_schema", _restore_cjk_update_trigger) + + assert db._migrate_broad_fts_update_triggers(db._conn) == 1 + cjk_sql = _trigger_sql(db._conn, "messages_fts_cjk_update") + assert cjk_sql is not None + assert "AFTER UPDATE OF" in " ".join(cjk_sql.split()).upper() + finally: + db.close() + + +def test_legacy_migration_does_not_drop_cjk_trigger(tmp_path: Path): + db = SessionDB(db_path=tmp_path / "state.db") + try: + _install_legacy_inline_base_fts(db) + db._conn.execute("DROP TRIGGER IF EXISTS messages_fts_cjk_update") + db._conn.execute( + "CREATE TRIGGER messages_fts_cjk_update " + "AFTER UPDATE ON messages BEGIN SELECT 1; END" + ) + + assert db._migrate_broad_fts_update_triggers(db._conn) >= 1 + cjk_sql = _trigger_sql(db._conn, "messages_fts_cjk_update") + assert cjk_sql is not None + assert "AFTER UPDATE OF" not in " ".join(cjk_sql.split()).upper() + finally: + db.close()