fix(state): align messages_fts external content with its indexed projection

`messages_fts` declared `content='messages'` while its triggers indexed only a
bounded prefix of every long tool row, with the boundary held in a
`fts_tool_full_content_high_water` state_meta marker that the migration, the
rebuild seeding and the in-place rebuild path all re-stamped. FTS5's strict
integrity check re-reads the external content source and compares it with the
stored token stream, so the two could never agree: one long tool row is enough to
make

    INSERT INTO messages_fts(messages_fts, rank) VALUES('integrity-check', 1)

fail with `fts5: checksum mismatch for table "messages_fts"`. The delete/update
triggers re-evaluated the *new* marker, so they sent full content for a row whose
index held a prefix and left tokens behind that survived deleting the row.

Fix the class by removing the moving part. `messages_fts` now reads a view,
`messages_fts_src`, that computes exactly what the writers index - tool rows
truncated to FTS_TOOL_CONTENT_PREFIX_CHARS, everything else verbatim - through a
fixed per-row expression with no state_meta lookups, shared by the triggers, the
boundary sweep and the chunked insert.

- `messages_fts_src` view added; `messages_fts` external content points at it
- triggers, boundary sweep and chunked backfill all read that one projection
- `_stamp_fts_tool_high_water`, the marker seeding in `_seed_fts_rebuild_markers`
  and the in-place rebuild stamp are gone, with
  FTS_TOOL_FULL_CONTENT_HIGH_WATER_KEY
- FTS_STORAGE_VERSION 2 -> 3: an existing index is re-pointed and rebuilt ONCE by
  `_migrate_misaligned_fts_source` under the shared cross-process rebuild
  admission, because a v2 index holds token streams its old source cannot read
  back; an empty index swaps shape in place, and legacy inline DBs are untouched
- search behaviour is unchanged: tool rows were already prefix-indexed, and
  explicit tool search already used the stored-content LIKE path
This commit is contained in:
Kelly Griffin
2026-09-17 10:31:30 -04:00
committed by Teknium
parent 1b083b85a3
commit 42e97f3808
3 changed files with 118 additions and 71 deletions

View File

@@ -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;

View File

@@ -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
)

View File

@@ -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')")