diff --git a/hermes_state_common.py b/hermes_state_common.py index 8ac4ce37e9..21b276e784 100644 --- a/hermes_state_common.py +++ b/hermes_state_common.py @@ -253,23 +253,27 @@ AUTO_VACUUM_MIN_FREELIST_RATIO = 0.25 # layout 0 (marker absent) with a working inline index until the user opts in. # 1 = v23 external-content layout with a tool-row-excluded trigram # 2 = trigram also excludes structured tool_calls JSON -FTS_STORAGE_VERSION = 2 +# 3 = messages_fts source aligned to a stable projection view +# (``messages_fts_src``): always-truncate tool rows to the prefix, no +# moving high-water boundary. The external-content source now reads +# back EXACTLY what the triggers indexed, so the rank=1 +# 'integrity-check' probe cannot drift from the stored index (the +# recurring fts5 "checksum mismatch" / leaked-token failures). +FTS_STORAGE_VERSION = 3 -# Tool results are often multi-megabyte machine payloads. Index a useful -# prefix for new tool rows instead of tokenizing the entire body while the -# canonical message write holds SQLite's single writer lock. The high-water -# marker lets upgraded databases retain the exact token stream already stored -# for historical rows, so external-content delete/update commands stay valid -# without an eager full-index rebuild. +# Tool results are often multi-megabyte machine payloads. The base FTS index +# stores only a bounded prefix of every tool row; tool rows are skipped by +# default in search, and explicit tool-only search uses a LIKE fallback over +# the full stored content, so no search capability is lost. The projection +# below is STABLE — it depends only on the row being written, never on +# mutable ``state_meta`` markers — which is what keeps the external-content +# integrity checker and the trigger 'delete'/'update' commands in agreement +# with the stored index forever. FTS_TOOL_CONTENT_PREFIX_CHARS = 8_192 -FTS_TOOL_FULL_CONTENT_HIGH_WATER_KEY = "fts_tool_full_content_high_water" def _fts_indexed_content_sql(alias: str) -> str: return f"""CASE WHEN {alias}.role = 'tool' - AND {alias}.id > COALESCE((SELECT CAST(value AS INTEGER) - FROM state_meta - WHERE key = '{FTS_TOOL_FULL_CONTENT_HIGH_WATER_KEY}'), -1) THEN substr(COALESCE({alias}.content, ''), 1, {FTS_TOOL_CONTENT_PREFIX_CHARS}) ELSE {alias}.content END""" @@ -691,12 +695,32 @@ CREATE INDEX IF NOT EXISTS idx_sessions_effective_activity # predicate into a tautology (id > -1 OR id <= -1), i.e. normal operation. # The two state_meta PK probes per write are negligible next to the FTS # insert itself. +# +# messages_fts_src: the base word index no longer reads raw `messages` as its +# external content. Tool rows are indexed as a bounded prefix, so the index +# must read that SAME projection back or FTS5's 'integrity-check' / 'delete' +# commands disagree with the stored tokens and corrupt the index (the +# recurring fts5 checksum-mismatch drift: the projection used to depend on a +# moving state_meta high-water key). The view/trigger/backfill all share the +# one expression in `_fts_indexed_content_sql` — a fixed per-row function +# with no marker lookups — so the boundary can never move again. FTS_SQL = f""" +-- Stable projection the base word index reads and writes through: the view +-- computes EXACTLY what the triggers/backfill insert, so 'rebuild' and the +-- integrity checker always agree with the stored index. +CREATE VIEW IF NOT EXISTS messages_fts_src AS + SELECT id, + CASE WHEN role = 'tool' + THEN substr(COALESCE(content, ''), 1, {FTS_TOOL_CONTENT_PREFIX_CHARS}) + ELSE content END AS content, + tool_name, tool_calls + FROM messages; + CREATE VIRTUAL TABLE IF NOT EXISTS messages_fts USING fts5( content, tool_name, tool_calls, - content='messages', + content='messages_fts_src', content_rowid='id' ); @@ -902,7 +926,9 @@ CREATE VIRTUAL TABLE IF NOT EXISTS messages_fts USING fts5( CREATE TRIGGER IF NOT EXISTS messages_fts_insert AFTER INSERT ON messages BEGIN INSERT INTO messages_fts(rowid, content) VALUES ( new.id, - COALESCE({_FTS_NEW_INDEXED_CONTENT_SQL}, '') + COALESCE(CASE WHEN new.role = 'tool' + THEN substr(COALESCE(new.content, ''), 1, {FTS_TOOL_CONTENT_PREFIX_CHARS}) + ELSE new.content END, '') || ' ' || COALESCE(new.tool_name, '') || ' ' || COALESCE(new.tool_calls, '') ); END; @@ -916,7 +942,9 @@ AFTER UPDATE OF content, tool_name, tool_calls, role ON messages BEGIN DELETE FROM messages_fts WHERE rowid = old.id; INSERT INTO messages_fts(rowid, content) VALUES ( new.id, - COALESCE({_FTS_NEW_INDEXED_CONTENT_SQL}, '') + COALESCE(CASE WHEN new.role = 'tool' + THEN substr(COALESCE(new.content, ''), 1, {FTS_TOOL_CONTENT_PREFIX_CHARS}) + ELSE new.content END, '') || ' ' || COALESCE(new.tool_name, '') || ' ' || COALESCE(new.tool_calls, '') ); END; @@ -932,7 +960,9 @@ CREATE VIRTUAL TABLE IF NOT EXISTS messages_fts_trigram USING fts5( CREATE TRIGGER IF NOT EXISTS messages_fts_trigram_insert AFTER INSERT ON messages BEGIN INSERT INTO messages_fts_trigram(rowid, content) VALUES ( new.id, - COALESCE({_FTS_NEW_INDEXED_CONTENT_SQL}, '') + COALESCE(CASE WHEN new.role = 'tool' + THEN substr(COALESCE(new.content, ''), 1, {FTS_TOOL_CONTENT_PREFIX_CHARS}) + ELSE new.content END, '') || ' ' || COALESCE(new.tool_name, '') || ' ' || COALESCE(new.tool_calls, '') ); END; @@ -946,7 +976,9 @@ AFTER UPDATE OF content, tool_name, tool_calls, role ON messages BEGIN DELETE FROM messages_fts_trigram WHERE rowid = old.id; INSERT INTO messages_fts_trigram(rowid, content) VALUES ( new.id, - COALESCE({_FTS_NEW_INDEXED_CONTENT_SQL}, '') + COALESCE(CASE WHEN new.role = 'tool' + THEN substr(COALESCE(new.content, ''), 1, {FTS_TOOL_CONTENT_PREFIX_CHARS}) + ELSE new.content END, '') || ' ' || COALESCE(new.tool_name, '') || ' ' || COALESCE(new.tool_calls, '') ); END; diff --git a/hermes_state_schema.py b/hermes_state_schema.py index 645556ad67..44a43ee53e 100644 --- a/hermes_state_schema.py +++ b/hermes_state_schema.py @@ -22,7 +22,7 @@ from hermes_startup_watchdog import report_startup_progress from utils import safe_json_loads from hermes_state_common import ( DEFERRED_INDEX_SQL, FTS_CJK_STALE_KEY, FTS_REBUILD_DEFERRAL_KEY, FTS_STALE_KEY, FTS_SQL, - FTS_STORAGE_VERSION, FTS_TOOL_FULL_CONTENT_HIGH_WATER_KEY, FTS_TRIGRAM_SQL, LEGACY_FTS_SQL, + FTS_STORAGE_VERSION, FTS_TOOL_CONTENT_PREFIX_CHARS, FTS_TRIGRAM_SQL, LEGACY_FTS_SQL, LEGACY_FTS_TRIGRAM_SQL, SCHEMA_SQL, SCHEMA_VERSION, _FTS_CJK_TRIGGERS, _FTS_TRIGGERS, _ephemeral_child_sql, _sql_json_extract, fts_rebuild_admission, ) @@ -277,11 +277,17 @@ class SessionSchemaMixin: return len(to_drop) @staticmethod - def _stamp_fts_tool_high_water(cursor: sqlite3.Cursor) -> None: - """Record MAX(messages.id) as the bounded-tool-content high-water mark: rows at or below it keep - their exact stored token stream; newer tool rows index only the prefix (see ``_fts_indexed_content_sql``).""" - high_water = cursor.execute("SELECT COALESCE(MAX(id), 0) FROM messages").fetchone()[0] - cursor.execute(_STATE_META_UPSERT_SQL, (FTS_TOOL_FULL_CONTENT_HIGH_WATER_KEY, str(high_water))) + def _fts_index_is_misaligned_source(cursor: sqlite3.Cursor) -> bool: + """True when ``messages_fts`` is still external-content over the raw + ``messages`` table (FTS_STORAGE_VERSION < 3): its index holds a TRUNCATED + projection for long tool rows that the checker/'delete' commands re-read + as FULL content, a mismatch by construction. Such an index cannot be + repaired in place — it must be 'rebuild'-filled from the aligned + ``messages_fts_src`` view exactly once.""" + row = cursor.execute( + "SELECT sql FROM sqlite_master WHERE type = 'table' AND name = 'messages_fts'" + ).fetchone() + return row is not None and "messages_fts_src" not in (row[0] or "") @staticmethod def _execute_ddl_skipping_settled_triggers(cursor: sqlite3.Cursor, ddl: str) -> None: @@ -330,41 +336,50 @@ class SessionSchemaMixin: if statement.strip(): raise sqlite3.OperationalError("incomplete FTS DDL statement") - def _migrate_bounded_tool_fts_triggers(self, cursor: sqlite3.Cursor, *, legacy: bool) -> None: - """Replace FTS triggers without rebuilding historical indexes. Existing rows keep their - full-content token stream; the durable high-water id makes new tool rows use the bounded - prefix in INSERT and the matching external-content delete/update. One savepoint, so no - concurrent writer lands in a trigger gap. A fresh store has no historical index to migrate; - its FTS family is created later under rebuild admission.""" - if not self._sqlite_table_exists(cursor, "messages_fts"): + def _migrate_misaligned_fts_source(self, cursor: sqlite3.Cursor, *, legacy: bool) -> None: + """Re-point ``messages_fts`` at the stable ``messages_fts_src`` projection view and + rebuild it ONCE (FTS_STORAGE_VERSION 2 -> 3). A v1/v2 base index carries token streams + the raw-``messages`` external-content source cannot read back (truncated long tool + rows, and tool rows whose full content was indexed under an old high-water mark), so + in-place continuity is not achievable — the ONLY valid transition is a full rebuild + from the view, under the shared cross-process rebuild admission. Legacy inline DBs + skip this entirely (their index is self-contained; they still take the DDL on the + optimize path).""" + if legacy or not self._sqlite_table_exists(cursor, "messages_fts"): return - marker = cursor.execute( - "SELECT 1 FROM state_meta WHERE key = ? LIMIT 1", (FTS_TOOL_FULL_CONTENT_HIGH_WATER_KEY,), - ).fetchone() - if marker is not None: + if not self._fts_index_is_misaligned_source(cursor): return - trigram_present = self._sqlite_table_exists(cursor, "messages_fts_trigram") - names = _FTS_BASE_TRIGGERS + (_FTS_TRIGRAM_TRIGGERS if legacy and trigram_present else ()) has_messages = cursor.execute("SELECT 1 FROM messages LIMIT 1").fetchone() is not None - self._fts_tool_prefix_migration_requires_rebuild = bool( - has_messages and self._fts_triggers_missing(cursor, names) - ) - cursor.execute("SAVEPOINT bounded_tool_fts") - try: - self._stamp_fts_tool_high_water(cursor) - for name in names: - cursor.execute(f"DROP TRIGGER IF EXISTS {name}") - if legacy: - self._execute_ddl_script_transactional(cursor, LEGACY_FTS_SQL) - if trigram_present: - self._execute_ddl_script_transactional(cursor, LEGACY_FTS_TRIGRAM_SQL) - else: + if not has_messages: + # Nothing indexed and nothing to index: just swap the shape in place. + cursor.execute("SAVEPOINT fts_align_empty") + try: + for name in _FTS_BASE_TRIGGERS: + cursor.execute(f"DROP TRIGGER IF EXISTS {name}") + cursor.execute("DROP TABLE IF EXISTS messages_fts") self._execute_ddl_script_transactional(cursor, FTS_SQL) - cursor.execute("RELEASE SAVEPOINT bounded_tool_fts") - except BaseException: - cursor.execute("ROLLBACK TO SAVEPOINT bounded_tool_fts") - cursor.execute("RELEASE SAVEPOINT bounded_tool_fts") - raise + cursor.execute(_STATE_META_UPSERT_SQL, ("fts_storage_version", str(FTS_STORAGE_VERSION))) + cursor.execute("RELEASE SAVEPOINT fts_align_empty") + except BaseException: + cursor.execute("ROLLBACK TO SAVEPOINT fts_align_empty") + cursor.execute("RELEASE SAVEPOINT fts_align_empty") + raise + return + self._fts_tool_prefix_migration_requires_rebuild = True + + def do_align() -> None: + self._execute_ddl_script_transactional(cursor, f""" +DROP TRIGGER IF EXISTS messages_fts_insert; +DROP TRIGGER IF EXISTS messages_fts_delete; +DROP TRIGGER IF EXISTS messages_fts_update; +""") + cursor.execute("DROP TABLE IF EXISTS messages_fts") + self._ensure_fts_schema(cursor, "messages_fts", FTS_SQL) + cursor.execute("INSERT INTO messages_fts(messages_fts) VALUES('rebuild')") + cursor.execute(_CLEAR_REBUILD_MARKERS_SQL) + cursor.execute(_STATE_META_UPSERT_SQL, ("fts_storage_version", str(FTS_STORAGE_VERSION))) + + self._run_admitted_startup_rebuild(cursor, do_align) @staticmethod def _sqlite_table_exists(cursor: sqlite3.Cursor, name: str) -> bool: @@ -420,7 +435,6 @@ class SessionSchemaMixin: markers are cleared or the worker would re-insert covered rows (duplicates). ``legacy`` (pre-v23 inline layout) has no external-content 'rebuild' source, so it DELETEs + reinserts the concatenated content the legacy triggers produced.""" - SessionSchemaMixin._stamp_fts_tool_high_water(cursor) tables = ("messages_fts", "messages_fts_trigram") if include_trigram else ("messages_fts",) for tbl in tables: if legacy: @@ -1159,7 +1173,7 @@ class SessionSchemaMixin: cursor, ("messages_fts", "messages_fts_trigram", "messages_fts_cjk"), ) if not self._fts_stale: - self._migrate_bounded_tool_fts_triggers(cursor, legacy=legacy_fts) + self._migrate_misaligned_fts_source(cursor, legacy=legacy_fts) if self._fts_stale: if self._recover_stale_fts(cursor, legacy=legacy_fts): # CJK was detached alongside the base indexes; its ensure path decides when it returns. @@ -1170,8 +1184,10 @@ class SessionSchemaMixin: base_sql, trigram_sql = _FTS_DDL[legacy_fts] # Measure before any DDL. Publishing missing base triggers before rebuild admission lets # another process write through an index whose bootstrap/repair has no owner (#105790). - base_triggers_missing = self._fts_triggers_missing(cursor, _FTS_BASE_TRIGGERS) or getattr( - self, "_fts_tool_prefix_migration_requires_rebuild", False) or "messages_fts" in orphan_repaired + base_triggers_missing = self._fts_triggers_missing(cursor, _FTS_BASE_TRIGGERS) or ( + getattr(self, "_fts_tool_prefix_migration_requires_rebuild", False) + and self._fts_index_is_misaligned_source(cursor) + ) or "messages_fts" in orphan_repaired trigram_triggers_missing = ( self._fts_triggers_missing(cursor, _FTS_TRIGRAM_TRIGGERS) or "messages_fts_trigram" in orphan_repaired ) diff --git a/hermes_state_search.py b/hermes_state_search.py index a497f1f77e..4f797ae5ad 100644 --- a/hermes_state_search.py +++ b/hermes_state_search.py @@ -14,7 +14,7 @@ from typing import Any, Callable, Collection, Dict, List, Optional, Tuple from agent.skill_commands import describe_skill_invocation from hermes_state_common import ( FTS_CJK_STALE_KEY, FTS_SQL, FTS_STALE_KEY, FTS_STORAGE_VERSION, FTS_TOOL_CONTENT_PREFIX_CHARS, - FTS_TOOL_FULL_CONTENT_HIGH_WATER_KEY, FTS_TRIGRAM_EXCLUDED_SOURCES, FTS_TRIGRAM_SQL, + FTS_TRIGRAM_EXCLUDED_SOURCES, FTS_TRIGRAM_SQL, MAX_FTS5_QUERY_CHARS, SCHEMA_VERSION, _FTS_CJK_TRIGGERS, escape_like as _escape_like, fts_rebuild_admission, fts_trigram_session_sql, routed_sessions_setting, ) @@ -222,8 +222,12 @@ class SessionSearchMixin: return {"pending": True, "total": total, "indexed": progress, "percent": min(100, int(100 * progress / total))} # Re-index rows in an id window the index is missing. docsize has one row - # per indexed doc, so the anti-join is exact. Params: (lo, hi) — the base sweep - # takes (hw, prefix_chars, lo, hi): tool rows past the high water index only a prefix. + # per indexed doc, so the anti-join is exact. Params: (lo, hi). + # NOTE: with the aligned projection (FTS_STORAGE_VERSION 3) every writer + # of the index — this sweep, the chunked backfill, and the sync triggers — + # feeds ``messages_fts`` through the ONE stable per-row expression: + # tool rows are truncated to FTS_TOOL_CONTENT_PREFIX_CHARS, everything + # else is verbatim, and nothing consults a moving state_meta marker. _BOUNDARY_SWEEP_SQL = ( "INSERT INTO {table}(rowid, content, tool_name, tool_calls) " "SELECT m.id, m.content, m.tool_name, m.tool_calls FROM messages m WHERE m.id > ? AND m.id <= ? {extra}" @@ -231,7 +235,7 @@ class SessionSearchMixin: ) _BASE_BOUNDARY_SWEEP_SQL = ( "INSERT INTO messages_fts(rowid, content, tool_name, tool_calls) " - "SELECT m.id, CASE WHEN m.role = 'tool' AND m.id > ? THEN substr(COALESCE(m.content, ''), 1, ?) " + "SELECT m.id, CASE WHEN m.role = 'tool' THEN substr(COALESCE(m.content, ''), 1, ?) " "ELSE m.content END, m.tool_name, m.tool_calls FROM messages m WHERE m.id > ? AND m.id <= ? " "AND NOT EXISTS (SELECT 1 FROM messages_fts_docsize d WHERE d.id = m.id)" ) @@ -244,7 +248,9 @@ class SessionSearchMixin: ) _CHUNK_INSERT_SQL = ( "INSERT INTO {table}(rowid, content, tool_name, tool_calls) " - "SELECT id, content, tool_name, tool_calls FROM messages WHERE id > ? AND id <= ?{extra}" + "SELECT id, CASE WHEN role = 'tool' " + f"THEN substr(COALESCE(content, ''), 1, {FTS_TOOL_CONTENT_PREFIX_CHARS}) " + "ELSE content END, tool_name, tool_calls FROM messages WHERE id > ? AND id <= ?{extra}" ) _TRIGRAM_CHUNK_INSERT_SQL = ( "INSERT INTO messages_fts_trigram(rowid, content, tool_name) " @@ -272,13 +278,13 @@ class SessionSearchMixin: def _rebuild_finish(self, prefix: str, sweep_sqls: List[Tuple[str, bool]]) -> None: """Sweep a generous window around the high-water boundary, then clear the markers. - ``(sql, bounded)``: a bounded sweep takes the (hw, prefix_chars) tool-content params first.""" + ``(sql, bounded)``: a bounded sweep takes the tool-content prefix_chars param first.""" def _do(conn): hw_row = _meta_row(conn, f"{prefix}_high_water") if hw_row is not None: hw = int(hw_row[0]) for sql, bounded in sweep_sqls: - params = (hw, FTS_TOOL_CONTENT_PREFIX_CHARS) if bounded else () + params = (FTS_TOOL_CONTENT_PREFIX_CHARS,) if bounded else () conn.execute(sql, (*params, hw - 1000, hw + 1000)) _delete_meta(conn, f"{prefix}_high_water", f"{prefix}_progress") self._execute_write(_do) @@ -478,12 +484,10 @@ class SessionSearchMixin: existing_hw = _meta_row(conn, "fts_rebuild_high_water") if existing_hw is not None and not force: self._reseed_missing_progress(conn) - self.set_meta(FTS_TOOL_FULL_CONTENT_HIGH_WATER_KEY, str(int(existing_hw[0])), cursor=conn) return int(existing_hw[0]) hw = conn.execute("SELECT COALESCE(MAX(id), 0) FROM messages").fetchone()[0] self.set_meta("fts_rebuild_high_water", str(hw), cursor=conn) self.set_meta("fts_rebuild_progress", "0", cursor=conn) - self.set_meta(FTS_TOOL_FULL_CONTENT_HIGH_WATER_KEY, str(hw), cursor=conn) return int(hw) def _repair_optimize_bookkeeping(self) -> None: @@ -1289,11 +1293,6 @@ class SessionSearchMixin: "Deferred in-place FTS rebuild: another process holds the rebuild authority for this state.db.") return 0 with self._lock: - high_water = self._conn.execute("SELECT COALESCE(MAX(id), 0) FROM messages").fetchone()[0] - self._conn.execute( - "INSERT INTO state_meta (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value", - (FTS_TOOL_FULL_CONTENT_HIGH_WATER_KEY, str(high_water)), - ) for tbl in self._present_fts_tables(): try: self._conn.execute(f"INSERT INTO {tbl}({tbl}) VALUES('rebuild')")