fix(state): keep canonical writes available when FTS is corrupt
This commit is contained in:
+69
-10
@@ -60,6 +60,7 @@ from hermes_state_common import ( # noqa: F401 (re-exported for back-compat)
|
||||
DEFERRED_INDEX_SQL,
|
||||
FTS_CJK_STALE_KEY,
|
||||
FTS_SQL,
|
||||
FTS_STALE_KEY,
|
||||
FTS_STORAGE_VERSION,
|
||||
FTS_TRIGRAM_SQL,
|
||||
LEGACY_FTS_SQL,
|
||||
@@ -2542,6 +2543,7 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
|
||||
# incremental FTS merge cadence (see _merge_fts_incrementally).
|
||||
self._fts_usermerge_floor_applied = False
|
||||
self._fts_enabled = False
|
||||
self._fts_stale = False
|
||||
self._trigram_available = False
|
||||
# CJK-bigram index (cjk_unicode61 loadable tokenizer). _fts_cjk_loaded:
|
||||
# extension present on the writer connection; _fts_cjk_available: the
|
||||
@@ -3178,16 +3180,16 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
|
||||
continue
|
||||
# Corrupt FTS shadow tables make every write raise the
|
||||
# malformed/corrupt error class through the FTS sync triggers
|
||||
# while the canonical messages table is intact. The gateway
|
||||
# session store has its own retry queue for transcript
|
||||
# appends (#65637 salvage), but cron and CLI writers call
|
||||
# SessionDB directly — without this, their writes hard-fail
|
||||
# until the next process restart triggers the offline repair.
|
||||
# Rebuild the FTS index in place (once per instance) via
|
||||
# rebuild_fts() and retry the failed write immediately.
|
||||
if not self._try_runtime_fts_rebuild(exc):
|
||||
raise
|
||||
continue
|
||||
# while the canonical messages table is intact. Recover here,
|
||||
# at the shared persistence boundary, so every caller gets the
|
||||
# same guarantee. First try the cheap in-place repair. If that
|
||||
# one-shot path is unavailable or corruption recurs, detach the
|
||||
# derived indexes and retry against the canonical tables.
|
||||
if self._try_runtime_fts_rebuild(exc):
|
||||
continue
|
||||
if self._enter_fts_fail_open(exc):
|
||||
continue
|
||||
raise
|
||||
except sqlite3.Error as exc:
|
||||
# Catch-all for builds that surface 'no more rows available'
|
||||
# as InterfaceError (a sibling of DatabaseError, not a
|
||||
@@ -3291,6 +3293,63 @@ class SessionDB(SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin)
|
||||
)
|
||||
return True
|
||||
|
||||
def _enter_fts_fail_open(self, exc: sqlite3.DatabaseError) -> bool:
|
||||
"""Detach corrupt FTS indexes so canonical writes can continue.
|
||||
|
||||
The stale breadcrumb and trigger removal commit atomically. Its
|
||||
ordering is load-bearing: after triggers are absent, new canonical
|
||||
rows create an index gap of unknown extent, so another process must
|
||||
never reinstall the triggers without first rebuilding every row.
|
||||
"""
|
||||
if not self._fts_enabled or not self._is_fts_write_corruption_error(exc):
|
||||
return False
|
||||
|
||||
try:
|
||||
with self._lock:
|
||||
self._conn.execute("BEGIN IMMEDIATE")
|
||||
try:
|
||||
self._conn.execute(
|
||||
"INSERT INTO state_meta (key, value) VALUES (?, '1') "
|
||||
"ON CONFLICT(key) DO UPDATE SET value = excluded.value",
|
||||
(FTS_STALE_KEY,),
|
||||
)
|
||||
cjk_triggers_present = self._conn.execute(
|
||||
"SELECT 1 FROM sqlite_master WHERE type = 'trigger' "
|
||||
f"AND name IN ({','.join('?' for _ in _FTS_CJK_TRIGGERS)}) "
|
||||
"LIMIT 1",
|
||||
_FTS_CJK_TRIGGERS,
|
||||
).fetchone()
|
||||
if cjk_triggers_present:
|
||||
self._conn.execute(
|
||||
"INSERT INTO state_meta (key, value) VALUES (?, '1') "
|
||||
"ON CONFLICT(key) DO UPDATE SET value = excluded.value",
|
||||
(FTS_CJK_STALE_KEY,),
|
||||
)
|
||||
self._drop_all_fts_triggers(self._conn.cursor())
|
||||
self._conn.commit()
|
||||
except BaseException:
|
||||
self._conn.rollback()
|
||||
raise
|
||||
except sqlite3.Error as detach_exc:
|
||||
logger.error(
|
||||
"Could not detach corrupt FTS indexes; canonical write still "
|
||||
"cannot proceed: %s",
|
||||
detach_exc,
|
||||
)
|
||||
return False
|
||||
|
||||
self._fts_stale = True
|
||||
self._fts_enabled = False
|
||||
self._trigram_available = False
|
||||
self._fts_cjk_available = False
|
||||
logger.error(
|
||||
"state.db FTS indexes remain corrupt (%s); disabled FTS sync and "
|
||||
"retrying the canonical write. Search temporarily uses LIKE until "
|
||||
"a later SessionDB open rebuilds the indexes.",
|
||||
exc,
|
||||
)
|
||||
return True
|
||||
|
||||
def _try_wal_checkpoint(self) -> None:
|
||||
"""Best-effort PASSIVE WAL checkpoint. Never raises.
|
||||
|
||||
|
||||
@@ -549,6 +549,14 @@ _FTS_CJK_TRIGGERS = (
|
||||
FTS_CJK_STALE_KEY = "fts_cjk_stale"
|
||||
|
||||
|
||||
# Durable breadcrumb for a base/trigram FTS index that was detached from the
|
||||
# canonical messages table after runtime corruption. While present, startup
|
||||
# must rebuild the complete index before reinstalling sync triggers: rows may
|
||||
# have been written while those triggers were absent, so merely recreating
|
||||
# them would preserve an unknown index gap.
|
||||
FTS_STALE_KEY = "fts_stale"
|
||||
|
||||
|
||||
# ── Legacy (v22 / inline-content) FTS DDL ──────────────────────────────
|
||||
# Used ONLY to keep an existing pre-v23 install's search working and its
|
||||
# triggers repairable UNTIL the user opts into `hermes db optimize`. This is
|
||||
|
||||
+123
-1
@@ -17,6 +17,7 @@ from hermes_constants import get_hermes_home
|
||||
from hermes_state_common import (
|
||||
DEFERRED_INDEX_SQL,
|
||||
FTS_CJK_STALE_KEY,
|
||||
FTS_STALE_KEY,
|
||||
FTS_SQL,
|
||||
FTS_STORAGE_VERSION,
|
||||
FTS_TRIGRAM_SQL,
|
||||
@@ -24,6 +25,7 @@ from hermes_state_common import (
|
||||
LEGACY_FTS_TRIGRAM_SQL,
|
||||
SCHEMA_SQL,
|
||||
SCHEMA_VERSION,
|
||||
_FTS_CJK_TRIGGERS,
|
||||
_FTS_TRIGGERS,
|
||||
_ephemeral_child_sql,
|
||||
)
|
||||
@@ -115,6 +117,14 @@ class SessionSchemaMixin:
|
||||
self._warn_fts5_unavailable(exc)
|
||||
return False
|
||||
|
||||
def _drop_all_fts_triggers(self, cursor: sqlite3.Cursor) -> None:
|
||||
self._drop_fts_triggers(cursor)
|
||||
for trigger in _FTS_CJK_TRIGGERS:
|
||||
try:
|
||||
cursor.execute(f"DROP TRIGGER IF EXISTS {trigger}")
|
||||
except sqlite3.OperationalError:
|
||||
pass
|
||||
|
||||
@staticmethod
|
||||
def _fts_trigger_count(cursor: sqlite3.Cursor) -> int:
|
||||
placeholders = ",".join("?" for _ in _FTS_TRIGGERS)
|
||||
@@ -336,6 +346,99 @@ class SessionSchemaMixin:
|
||||
return False
|
||||
raise
|
||||
|
||||
def _recover_stale_fts(self, cursor: sqlite3.Cursor, *, legacy: bool) -> bool:
|
||||
"""Atomically rebuild stale base/trigram indexes and resume syncing."""
|
||||
try:
|
||||
trigram_status = self._fts_table_probe(cursor, "messages_fts_trigram")
|
||||
except sqlite3.DatabaseError:
|
||||
# A corrupt vtable may fail even a LIMIT 0 probe. It still needs
|
||||
# to be included in the drop-and-recreate recovery below.
|
||||
trigram_status = True
|
||||
include_trigram = trigram_status is True
|
||||
|
||||
drop_sql = "".join(
|
||||
f"DROP TRIGGER IF EXISTS {trigger};" for trigger in _FTS_TRIGGERS
|
||||
)
|
||||
if include_trigram:
|
||||
drop_sql += "DROP TABLE IF EXISTS messages_fts_trigram;"
|
||||
drop_sql += "DROP VIEW IF EXISTS messages_fts_trigram_src;"
|
||||
drop_sql += "DROP TABLE IF EXISTS messages_fts;"
|
||||
|
||||
if legacy:
|
||||
schema_sql = LEGACY_FTS_SQL
|
||||
if include_trigram:
|
||||
schema_sql += LEGACY_FTS_TRIGRAM_SQL
|
||||
rebuild_sql = schema_sql + """
|
||||
INSERT INTO messages_fts(rowid, content)
|
||||
SELECT id,
|
||||
COALESCE(content, '') || ' ' ||
|
||||
COALESCE(tool_name, '') || ' ' ||
|
||||
COALESCE(tool_calls, '')
|
||||
FROM messages;
|
||||
"""
|
||||
if include_trigram:
|
||||
rebuild_sql += """
|
||||
DELETE FROM messages_fts_trigram;
|
||||
INSERT INTO messages_fts_trigram(rowid, content)
|
||||
SELECT id,
|
||||
COALESCE(content, '') || ' ' ||
|
||||
COALESCE(tool_name, '') || ' ' ||
|
||||
COALESCE(tool_calls, '')
|
||||
FROM messages;
|
||||
"""
|
||||
else:
|
||||
schema_sql = FTS_SQL
|
||||
if include_trigram:
|
||||
schema_sql += FTS_TRIGRAM_SQL
|
||||
rebuild_sql = schema_sql + (
|
||||
"INSERT INTO messages_fts(messages_fts) VALUES('rebuild');"
|
||||
)
|
||||
if include_trigram:
|
||||
rebuild_sql += (
|
||||
"INSERT INTO messages_fts_trigram(messages_fts_trigram) "
|
||||
"VALUES('rebuild');"
|
||||
)
|
||||
rebuild_sql += (
|
||||
"DELETE FROM state_meta WHERE key IN "
|
||||
"('fts_rebuild_high_water', 'fts_rebuild_progress');"
|
||||
)
|
||||
|
||||
# One write transaction closes the dangerous gap: no canonical writer
|
||||
# can slip between the full rebuild and trigger restoration.
|
||||
recovery_sql = (
|
||||
"BEGIN IMMEDIATE;"
|
||||
+ drop_sql
|
||||
+ rebuild_sql
|
||||
+ f"DELETE FROM state_meta WHERE key = '{FTS_STALE_KEY}';"
|
||||
+ "COMMIT;"
|
||||
)
|
||||
try:
|
||||
cursor.executescript(recovery_sql)
|
||||
except sqlite3.DatabaseError as exc:
|
||||
try:
|
||||
self._conn.rollback()
|
||||
except sqlite3.Error:
|
||||
pass
|
||||
# Stale indexes must remain detached even on SQLite builds whose
|
||||
# DDL transaction behavior differs.
|
||||
self._drop_all_fts_triggers(cursor)
|
||||
self._conn.commit()
|
||||
logger.error(
|
||||
"Automatic rebuild of stale FTS indexes failed (%s); "
|
||||
"canonical writes remain enabled with FTS detached.",
|
||||
exc,
|
||||
)
|
||||
return False
|
||||
|
||||
self._fts_stale = False
|
||||
self._fts_enabled = True
|
||||
self._trigram_available = include_trigram
|
||||
logger.warning(
|
||||
"Rebuilt stale state.db FTS indexes from canonical messages and "
|
||||
"restored sync triggers."
|
||||
)
|
||||
return True
|
||||
|
||||
@staticmethod
|
||||
def _parse_schema_columns(schema_sql: str) -> Dict[str, Dict[str, str]]:
|
||||
"""Extract expected columns per table from SCHEMA_SQL.
|
||||
@@ -687,6 +790,14 @@ class SessionSchemaMixin:
|
||||
|
||||
fts5_available = self._sqlite_supports_fts5(cursor)
|
||||
fts_migrations_complete = True
|
||||
self._fts_stale = cursor.execute(
|
||||
"SELECT 1 FROM state_meta WHERE key = ? LIMIT 1",
|
||||
(FTS_STALE_KEY,),
|
||||
).fetchone() is not None
|
||||
if self._fts_stale:
|
||||
# A prior process deliberately detached FTS after corruption.
|
||||
# Keep every FTS writer detached until a full rebuild succeeds.
|
||||
self._drop_all_fts_triggers(cursor)
|
||||
if not fts5_available:
|
||||
# Existing FTS triggers can still fire on messages INSERT/UPDATE
|
||||
# even though the current sqlite runtime cannot read the virtual
|
||||
@@ -1028,7 +1139,18 @@ class SessionSchemaMixin:
|
||||
# its inline triggers exist (via the legacy DDL), and skip the
|
||||
# v23 view/external tables entirely. Fresh installs and opted-in
|
||||
# DBs have no legacy inline FTS, so they get the v23 DDL.
|
||||
if self._db_has_legacy_inline_fts(cursor):
|
||||
legacy_fts = self._db_has_legacy_inline_fts(cursor)
|
||||
if self._fts_stale:
|
||||
if self._recover_stale_fts(cursor, legacy=legacy_fts):
|
||||
# CJK was detached alongside the corrupt base indexes and
|
||||
# has its own stale marker. Its existing ensure path keeps
|
||||
# it offline until its dedicated rebuild.
|
||||
self._ensure_fts_cjk_schema(cursor)
|
||||
else:
|
||||
self._fts_enabled = False
|
||||
self._trigram_available = False
|
||||
self._fts_cjk_available = False
|
||||
elif legacy_fts:
|
||||
triggers_need_repair = (
|
||||
self._fts_trigger_count(cursor) < len(_FTS_TRIGGERS)
|
||||
)
|
||||
|
||||
+237
-81
@@ -20,6 +20,7 @@ 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_TRIGRAM_SQL,
|
||||
MAX_FTS5_QUERY_CHARS,
|
||||
@@ -1459,6 +1460,8 @@ class SessionSearchMixin:
|
||||
def _describe_search_path(self, query: str) -> str:
|
||||
"""Best-effort name of the routing path a query takes (log-only)."""
|
||||
try:
|
||||
if self._fts_stale:
|
||||
return "like_scan_fts_stale"
|
||||
sanitized = self._sanitize_fts5_query(query or "")
|
||||
if not sanitized:
|
||||
return "empty"
|
||||
@@ -1478,6 +1481,221 @@ class SessionSearchMixin:
|
||||
except Exception:
|
||||
return "unknown"
|
||||
|
||||
@staticmethod
|
||||
def _compile_like_boolean_query(
|
||||
query: str,
|
||||
) -> Tuple[str, List[Any], Optional[str]]:
|
||||
"""Compile the supported FTS boolean subset into LIKE predicates.
|
||||
|
||||
Terms within an OR group are ANDed by default, matching FTS5's
|
||||
implicit conjunction. ``NOT`` negates the following term inside that
|
||||
group instead of being discarded, so ``python NOT java`` becomes a
|
||||
positive Python match plus a Java exclusion.
|
||||
"""
|
||||
groups: List[List[Tuple[str, bool]]] = [[]]
|
||||
negate_next = False
|
||||
for raw_token in re.findall(r'"[^"]+"|\S+', query):
|
||||
operator = raw_token.upper()
|
||||
if operator == "OR":
|
||||
if groups[-1]:
|
||||
groups.append([])
|
||||
negate_next = False
|
||||
continue
|
||||
if operator in {"AND", "NEAR"}:
|
||||
continue
|
||||
if operator == "NOT":
|
||||
negate_next = True
|
||||
continue
|
||||
|
||||
term = raw_token.strip('"').strip("*").strip()
|
||||
if term:
|
||||
groups[-1].append((term, negate_next))
|
||||
negate_next = False
|
||||
|
||||
compiled_groups: List[str] = []
|
||||
params: List[Any] = []
|
||||
snippet_term: Optional[str] = None
|
||||
for group in groups:
|
||||
if not group or not any(not negated for _, negated in group):
|
||||
continue
|
||||
clauses: List[str] = []
|
||||
for term, negated in group:
|
||||
escaped = (
|
||||
term.replace("\\", "\\\\")
|
||||
.replace("%", "\\%")
|
||||
.replace("_", "\\_")
|
||||
)
|
||||
clause = (
|
||||
"(COALESCE(m.content, '') LIKE ? ESCAPE '\\' OR "
|
||||
"COALESCE(m.tool_name, '') LIKE ? ESCAPE '\\' OR "
|
||||
"COALESCE(m.tool_calls, '') LIKE ? ESCAPE '\\')"
|
||||
)
|
||||
clauses.append(f"NOT {clause}" if negated else clause)
|
||||
params.extend([f"%{escaped}%"] * 3)
|
||||
if snippet_term is None and not negated:
|
||||
snippet_term = term
|
||||
compiled_groups.append(f"({' AND '.join(clauses)})")
|
||||
|
||||
return " OR ".join(compiled_groups), params, snippet_term
|
||||
|
||||
def _search_messages_like_fallback(
|
||||
self,
|
||||
query: str,
|
||||
*,
|
||||
source_filter: Optional[List[str]],
|
||||
exclude_sources: Optional[List[str]],
|
||||
role_filter: Optional[List[str]],
|
||||
limit: int,
|
||||
offset: int,
|
||||
sort: Optional[str],
|
||||
include_inactive: bool,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Search canonical messages while derived FTS state is stale."""
|
||||
predicate, params, snippet_term = self._compile_like_boolean_query(query)
|
||||
if not predicate or snippet_term is None:
|
||||
return []
|
||||
|
||||
where = [f"({predicate})"]
|
||||
if not include_inactive:
|
||||
where.append("(m.active = 1 OR m.compacted = 1)")
|
||||
if source_filter is not None:
|
||||
where.append(f"s.source IN ({','.join('?' for _ in source_filter)})")
|
||||
params.extend(source_filter)
|
||||
if exclude_sources is not None:
|
||||
where.append(
|
||||
f"s.source NOT IN ({','.join('?' for _ in exclude_sources)})"
|
||||
)
|
||||
params.extend(exclude_sources)
|
||||
if role_filter:
|
||||
where.append(f"m.role IN ({','.join('?' for _ in role_filter)})")
|
||||
params.extend(role_filter)
|
||||
|
||||
order = (
|
||||
"ASC"
|
||||
if isinstance(sort, str) and sort.strip().lower() == "oldest"
|
||||
else "DESC"
|
||||
)
|
||||
sql = f"""
|
||||
SELECT m.id, m.session_id, m.role,
|
||||
substr(m.content, max(1, instr(m.content, ?) - 40), 120) AS snippet,
|
||||
m.content, m.timestamp, m.tool_name,
|
||||
s.source, s.model, s.started_at AS session_started
|
||||
FROM messages m
|
||||
JOIN sessions s ON s.id = m.session_id
|
||||
WHERE {' AND '.join(where)}
|
||||
ORDER BY m.timestamp {order}, m.id {order}
|
||||
LIMIT ? OFFSET ?
|
||||
"""
|
||||
with self._read_ctx() as conn:
|
||||
rows = conn.execute(
|
||||
sql, [snippet_term, *params, limit, offset]
|
||||
).fetchall()
|
||||
return [dict(row) for row in rows]
|
||||
|
||||
def _refresh_fts_stale_state(self) -> None:
|
||||
"""Observe fail-open initiated by another process sharing state.db."""
|
||||
if self._fts_stale or not self._fts_enabled:
|
||||
return
|
||||
try:
|
||||
with self._read_ctx() as conn:
|
||||
stale = conn.execute(
|
||||
"SELECT 1 FROM state_meta WHERE key = ? LIMIT 1",
|
||||
(FTS_STALE_KEY,),
|
||||
).fetchone()
|
||||
except sqlite3.Error:
|
||||
return
|
||||
if stale is not None:
|
||||
self._fts_stale = True
|
||||
self._fts_enabled = False
|
||||
self._trigram_available = False
|
||||
self._fts_cjk_available = False
|
||||
|
||||
def _finalize_search_matches(
|
||||
self,
|
||||
matches: List[Dict[str, Any]],
|
||||
result_fields: Optional[Collection[str]] = None,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Attach neighboring messages and trim full content from results.
|
||||
|
||||
Context (1 message before + after each match) is only loaded when
|
||||
the selected result projection consumes it. Each query takes its
|
||||
own fresh read transaction via _read_ctx, so we never hold a lock
|
||||
across N sequential queries.
|
||||
"""
|
||||
context_matches = (
|
||||
matches if result_fields is None or "context" in result_fields else ()
|
||||
)
|
||||
for match in context_matches:
|
||||
try:
|
||||
with self._read_ctx() as conn:
|
||||
ctx_cursor = conn.execute(
|
||||
"""WITH target AS (
|
||||
SELECT session_id, timestamp, id
|
||||
FROM messages
|
||||
WHERE id = ?
|
||||
)
|
||||
SELECT role, content
|
||||
FROM (
|
||||
SELECT m.id, m.timestamp, m.role, m.content
|
||||
FROM messages m
|
||||
JOIN target t ON t.session_id = m.session_id
|
||||
WHERE (m.timestamp < t.timestamp)
|
||||
OR (m.timestamp = t.timestamp AND m.id < t.id)
|
||||
ORDER BY m.timestamp DESC, m.id DESC
|
||||
LIMIT 1
|
||||
)
|
||||
UNION ALL
|
||||
SELECT role, content
|
||||
FROM messages
|
||||
WHERE id = ?
|
||||
UNION ALL
|
||||
SELECT role, content
|
||||
FROM (
|
||||
SELECT m.id, m.timestamp, m.role, m.content
|
||||
FROM messages m
|
||||
JOIN target t ON t.session_id = m.session_id
|
||||
WHERE (m.timestamp > t.timestamp)
|
||||
OR (m.timestamp = t.timestamp AND m.id > t.id)
|
||||
ORDER BY m.timestamp ASC, m.id ASC
|
||||
LIMIT 1
|
||||
)""",
|
||||
(match["id"], match["id"]),
|
||||
)
|
||||
context_msgs = []
|
||||
for row in ctx_cursor.fetchall():
|
||||
decoded = self._decode_content(row["content"])
|
||||
if isinstance(decoded, list):
|
||||
text_parts = [
|
||||
part.get("text", "")
|
||||
for part in decoded
|
||||
if isinstance(part, dict)
|
||||
and part.get("type") == "text"
|
||||
]
|
||||
text = " ".join(t for t in text_parts if t).strip()
|
||||
preview = text or "[multimodal content]"
|
||||
elif isinstance(decoded, str):
|
||||
preview = decoded
|
||||
else:
|
||||
preview = ""
|
||||
context_msgs.append(
|
||||
{"role": row["role"], "content": preview[:200]}
|
||||
)
|
||||
match["context"] = context_msgs
|
||||
except Exception:
|
||||
match["context"] = []
|
||||
|
||||
# Remove full content from result (snippet is enough, saves tokens)
|
||||
for match in matches:
|
||||
match.pop("content", None)
|
||||
|
||||
if result_fields is not None:
|
||||
matches = [
|
||||
{field: match[field] for field in result_fields if field in match}
|
||||
for match in matches
|
||||
]
|
||||
|
||||
return matches
|
||||
|
||||
def _search_messages_impl(
|
||||
self,
|
||||
query: str,
|
||||
@@ -1523,9 +1741,6 @@ class SessionSearchMixin:
|
||||
"""
|
||||
result_fields = self._search_message_fields(fields)
|
||||
|
||||
if not self._fts_enabled:
|
||||
return []
|
||||
|
||||
if not query or not query.strip():
|
||||
return []
|
||||
|
||||
@@ -1533,6 +1748,24 @@ class SessionSearchMixin:
|
||||
if not query:
|
||||
return []
|
||||
|
||||
self._refresh_fts_stale_state()
|
||||
if self._fts_stale:
|
||||
matches = self._search_messages_like_fallback(
|
||||
query,
|
||||
source_filter=source_filter,
|
||||
exclude_sources=exclude_sources,
|
||||
role_filter=role_filter,
|
||||
limit=limit,
|
||||
offset=offset,
|
||||
sort=sort,
|
||||
include_inactive=include_inactive,
|
||||
)
|
||||
return self._finalize_search_matches(
|
||||
matches, result_fields=result_fields
|
||||
)
|
||||
if not self._fts_enabled:
|
||||
return []
|
||||
|
||||
# Normalise sort. Anything not in the allowed set falls back to None
|
||||
# (FTS5 rank-only) so callers can pass through user input without
|
||||
# validation.
|
||||
@@ -1972,84 +2205,7 @@ class SessionSearchMixin:
|
||||
if tri_matches:
|
||||
matches = tri_matches
|
||||
|
||||
# Add surrounding context (1 message before + after each match) only
|
||||
# when the selected result projection consumes it. Each query takes
|
||||
# its own fresh read transaction via _read_ctx, so we never hold a
|
||||
# lock across N sequential queries.
|
||||
context_matches = (
|
||||
matches if result_fields is None or "context" in result_fields else ()
|
||||
)
|
||||
for match in context_matches:
|
||||
try:
|
||||
with self._read_ctx() as conn:
|
||||
ctx_cursor = conn.execute(
|
||||
"""WITH target AS (
|
||||
SELECT session_id, timestamp, id
|
||||
FROM messages
|
||||
WHERE id = ?
|
||||
)
|
||||
SELECT role, content
|
||||
FROM (
|
||||
SELECT m.id, m.timestamp, m.role, m.content
|
||||
FROM messages m
|
||||
JOIN target t ON t.session_id = m.session_id
|
||||
WHERE (m.timestamp < t.timestamp)
|
||||
OR (m.timestamp = t.timestamp AND m.id < t.id)
|
||||
ORDER BY m.timestamp DESC, m.id DESC
|
||||
LIMIT 1
|
||||
)
|
||||
UNION ALL
|
||||
SELECT role, content
|
||||
FROM messages
|
||||
WHERE id = ?
|
||||
UNION ALL
|
||||
SELECT role, content
|
||||
FROM (
|
||||
SELECT m.id, m.timestamp, m.role, m.content
|
||||
FROM messages m
|
||||
JOIN target t ON t.session_id = m.session_id
|
||||
WHERE (m.timestamp > t.timestamp)
|
||||
OR (m.timestamp = t.timestamp AND m.id > t.id)
|
||||
ORDER BY m.timestamp ASC, m.id ASC
|
||||
LIMIT 1
|
||||
)""",
|
||||
(match["id"], match["id"]),
|
||||
)
|
||||
context_msgs = []
|
||||
for r in ctx_cursor.fetchall():
|
||||
raw = r["content"]
|
||||
decoded = self._decode_content(raw)
|
||||
# Multimodal context: render a compact text-only
|
||||
# summary for search previews.
|
||||
if isinstance(decoded, list):
|
||||
text_parts = [
|
||||
p.get("text", "") for p in decoded
|
||||
if isinstance(p, dict) and p.get("type") == "text"
|
||||
]
|
||||
text = " ".join(t for t in text_parts if t).strip()
|
||||
preview = text or "[multimodal content]"
|
||||
elif isinstance(decoded, str):
|
||||
preview = decoded
|
||||
else:
|
||||
preview = ""
|
||||
context_msgs.append(
|
||||
{"role": r["role"], "content": preview[:200]}
|
||||
)
|
||||
match["context"] = context_msgs
|
||||
except Exception:
|
||||
match["context"] = []
|
||||
|
||||
# Remove full content from result (snippet is enough, saves tokens)
|
||||
for match in matches:
|
||||
match.pop("content", None)
|
||||
|
||||
if result_fields is not None:
|
||||
matches = [
|
||||
{field: match[field] for field in result_fields if field in match}
|
||||
for match in matches
|
||||
]
|
||||
|
||||
return matches
|
||||
return self._finalize_search_matches(matches, result_fields=result_fields)
|
||||
|
||||
def _search_unindexed_gap(
|
||||
self,
|
||||
|
||||
@@ -7,16 +7,24 @@ intact. Before this fix the gateway swallowed the failure at debug level and
|
||||
the in-memory session advanced while disk silently fell behind — surfacing
|
||||
later as "Persisted transcript lagged live cached history" amnesia.
|
||||
|
||||
The fix: ``_execute_write`` detects the malformed-image class, performs a
|
||||
one-shot in-place FTS rebuild (FTS5 ``'rebuild'`` command — index rewritten
|
||||
from canonical rows, no messages touched), and retries the failed write.
|
||||
The fix: ``_execute_write`` first attempts a one-shot in-place FTS rebuild.
|
||||
If corruption persists, it records a durable stale marker, detaches the FTS
|
||||
sync triggers, and retries the canonical write. Search degrades to ``LIKE``
|
||||
until a later open atomically rebuilds the index and restores the triggers.
|
||||
"""
|
||||
|
||||
import sqlite3
|
||||
|
||||
import pytest
|
||||
|
||||
from hermes_state import SessionDB
|
||||
from hermes_state import (
|
||||
FTS_STALE_KEY,
|
||||
LEGACY_FTS_SQL,
|
||||
LEGACY_FTS_TRIGRAM_SQL,
|
||||
SCHEMA_SQL,
|
||||
SessionDB,
|
||||
_FTS_TRIGGERS,
|
||||
)
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
@@ -55,6 +63,26 @@ def _message_contents(db_path):
|
||||
return [r[0] for r in rows]
|
||||
|
||||
|
||||
def _meta_value(db_path, key):
|
||||
raw = sqlite3.connect(str(db_path))
|
||||
row = raw.execute(
|
||||
"SELECT value FROM state_meta WHERE key = ?", (key,)
|
||||
).fetchone()
|
||||
raw.close()
|
||||
return None if row is None else row[0]
|
||||
|
||||
|
||||
def _base_fts_triggers(db_path):
|
||||
raw = sqlite3.connect(str(db_path))
|
||||
rows = raw.execute(
|
||||
"SELECT name FROM sqlite_master WHERE type = 'trigger' "
|
||||
f"AND name IN ({','.join('?' for _ in _FTS_TRIGGERS)})",
|
||||
_FTS_TRIGGERS,
|
||||
).fetchall()
|
||||
raw.close()
|
||||
return {row[0] for row in rows}
|
||||
|
||||
|
||||
class TestRuntimeFtsRebuild:
|
||||
def test_corruption_error_classification_covers_both_sqlite_messages(self):
|
||||
"""SQLite's message for a corrupt FTS index varies by version: older
|
||||
@@ -153,18 +181,185 @@ class TestRuntimeFtsRebuild:
|
||||
assert any(">>>" in (r.get("snippet") or "") for r in results)
|
||||
|
||||
|
||||
def test_rebuild_is_one_shot_per_instance(self, db, tmp_path):
|
||||
def test_second_corruption_fails_open_and_rebuilds_on_reopen(
|
||||
self, db, tmp_path
|
||||
):
|
||||
if not db._fts_enabled:
|
||||
pytest.skip("FTS5 unavailable in this build")
|
||||
db_path = tmp_path / "state.db"
|
||||
db.create_session("s1", source="test")
|
||||
db.append_message("s1", "user", "seed")
|
||||
_corrupt_fts(tmp_path / "state.db")
|
||||
_corrupt_fts(db_path)
|
||||
db.append_message("s1", "user", "first heal") # consumes the one shot
|
||||
assert db._fts_runtime_rebuild_attempted is True
|
||||
|
||||
# Corrupt again: the guard must NOT loop — the write now propagates.
|
||||
_corrupt_fts(tmp_path / "state.db")
|
||||
with pytest.raises(sqlite3.DatabaseError):
|
||||
db.append_message("s1", "user", "second corruption")
|
||||
# A second corruption must not strand the canonical transcript. The
|
||||
# derived indexes are detached and marked stale instead of looping.
|
||||
_corrupt_fts(db_path)
|
||||
db.append_message("s1", "user", "second corruption")
|
||||
assert _message_contents(db_path) == [
|
||||
"seed",
|
||||
"first heal",
|
||||
"second corruption",
|
||||
]
|
||||
assert db._fts_stale is True
|
||||
assert _meta_value(db_path, FTS_STALE_KEY) == "1"
|
||||
assert _base_fts_triggers(db_path) == set()
|
||||
|
||||
# Search remains available from canonical rows while FTS is stale.
|
||||
results = db.search_messages("second corruption")
|
||||
assert results
|
||||
assert any("second corruption" in row["snippet"] for row in results)
|
||||
|
||||
# A later open atomically rebuilds all canonical rows before triggers
|
||||
# return, then clears the durable breadcrumb.
|
||||
db.close()
|
||||
reopened = SessionDB(db_path=db_path)
|
||||
try:
|
||||
assert reopened._fts_stale is False
|
||||
assert _meta_value(db_path, FTS_STALE_KEY) is None
|
||||
assert _base_fts_triggers(db_path) == set(_FTS_TRIGGERS)
|
||||
results = reopened.search_messages("second corruption")
|
||||
assert results
|
||||
finally:
|
||||
reopened.close()
|
||||
|
||||
def test_failed_in_place_rebuild_fails_open(self, db, tmp_path, monkeypatch):
|
||||
if not db._fts_enabled:
|
||||
pytest.skip("FTS5 unavailable in this build")
|
||||
db_path = tmp_path / "state.db"
|
||||
db.create_session("s1", source="test")
|
||||
db.append_message("s1", "user", "seed")
|
||||
_corrupt_fts(db_path)
|
||||
|
||||
def _failed_rebuild():
|
||||
raise sqlite3.DatabaseError("rebuild could not read corrupt FTS")
|
||||
|
||||
monkeypatch.setattr(db, "rebuild_fts", _failed_rebuild)
|
||||
db.append_message("s1", "user", "canonical survives")
|
||||
|
||||
assert _message_contents(db_path)[-1] == "canonical survives"
|
||||
assert _meta_value(db_path, FTS_STALE_KEY) == "1"
|
||||
assert _base_fts_triggers(db_path) == set()
|
||||
|
||||
def test_stale_search_preserves_not_semantics(self, db, tmp_path, monkeypatch):
|
||||
if not db._fts_enabled:
|
||||
pytest.skip("FTS5 unavailable in this build")
|
||||
db_path = tmp_path / "state.db"
|
||||
db.create_session("s1", source="test")
|
||||
db.append_message("s1", "user", "python language guide")
|
||||
db.append_message("s1", "user", "python java interoperability")
|
||||
_corrupt_fts(db_path)
|
||||
|
||||
monkeypatch.setattr(
|
||||
db,
|
||||
"rebuild_fts",
|
||||
lambda: (_ for _ in ()).throw(
|
||||
sqlite3.DatabaseError("rebuild could not read corrupt FTS")
|
||||
),
|
||||
)
|
||||
db.append_message("s1", "user", "canonical write survives")
|
||||
assert db._fts_stale is True
|
||||
|
||||
results = db.search_messages("python NOT java")
|
||||
snippets = [row["snippet"] for row in results]
|
||||
assert any("python language guide" in snippet for snippet in snippets)
|
||||
assert all("java" not in snippet for snippet in snippets)
|
||||
|
||||
def test_existing_peer_observes_fail_open_marker(
|
||||
self, db, tmp_path, monkeypatch
|
||||
):
|
||||
if not db._fts_enabled:
|
||||
pytest.skip("FTS5 unavailable in this build")
|
||||
db_path = tmp_path / "state.db"
|
||||
db.create_session("s1", source="test")
|
||||
db.append_message("s1", "user", "seed")
|
||||
peer = SessionDB(db_path=db_path)
|
||||
try:
|
||||
_corrupt_fts(db_path)
|
||||
|
||||
def _failed_rebuild():
|
||||
raise sqlite3.DatabaseError("rebuild failed")
|
||||
|
||||
monkeypatch.setattr(db, "rebuild_fts", _failed_rebuild)
|
||||
db.append_message("s1", "user", "visible through canonical search")
|
||||
|
||||
assert peer._fts_stale is False
|
||||
results = peer.search_messages("canonical search")
|
||||
assert peer._fts_stale is True
|
||||
assert results
|
||||
finally:
|
||||
peer.close()
|
||||
|
||||
def test_failed_startup_rebuild_keeps_fts_detached(
|
||||
self, db, tmp_path, monkeypatch
|
||||
):
|
||||
if not db._fts_enabled:
|
||||
pytest.skip("FTS5 unavailable in this build")
|
||||
db_path = tmp_path / "state.db"
|
||||
db.create_session("s1", source="test")
|
||||
db.append_message("s1", "user", "seed")
|
||||
_corrupt_fts(db_path)
|
||||
monkeypatch.setattr(
|
||||
db,
|
||||
"rebuild_fts",
|
||||
lambda: (_ for _ in ()).throw(sqlite3.DatabaseError("still corrupt")),
|
||||
)
|
||||
db.append_message("s1", "user", "before restart")
|
||||
db.close()
|
||||
|
||||
monkeypatch.setattr(
|
||||
SessionDB,
|
||||
"_recover_stale_fts",
|
||||
lambda self, cursor, legacy: False,
|
||||
)
|
||||
reopened = SessionDB(db_path=db_path)
|
||||
try:
|
||||
assert reopened._fts_stale is True
|
||||
assert _meta_value(db_path, FTS_STALE_KEY) == "1"
|
||||
assert _base_fts_triggers(db_path) == set()
|
||||
reopened.append_message("s1", "user", "after failed recovery")
|
||||
assert _message_contents(db_path)[-1] == "after failed recovery"
|
||||
assert reopened.search_messages("failed recovery")
|
||||
finally:
|
||||
reopened.close()
|
||||
|
||||
def test_legacy_inline_fts_fails_open_and_recovers(self, tmp_path, monkeypatch):
|
||||
db_path = tmp_path / "legacy-state.db"
|
||||
raw = sqlite3.connect(str(db_path))
|
||||
raw.executescript(SCHEMA_SQL)
|
||||
try:
|
||||
raw.executescript(LEGACY_FTS_SQL + LEGACY_FTS_TRIGRAM_SQL)
|
||||
except sqlite3.OperationalError as exc:
|
||||
raw.close()
|
||||
pytest.skip(f"required FTS tokenizer unavailable: {exc}")
|
||||
raw.commit()
|
||||
raw.close()
|
||||
|
||||
legacy = SessionDB(db_path=db_path)
|
||||
try:
|
||||
assert legacy._db_has_legacy_inline_fts(legacy._conn.cursor())
|
||||
legacy.create_session("s1", source="test")
|
||||
legacy.append_message("s1", "user", "legacy seed")
|
||||
_corrupt_fts(db_path)
|
||||
monkeypatch.setattr(
|
||||
legacy,
|
||||
"rebuild_fts",
|
||||
lambda: (_ for _ in ()).throw(
|
||||
sqlite3.DatabaseError("legacy rebuild failed")
|
||||
),
|
||||
)
|
||||
legacy.append_message("s1", "user", "legacy canonical survives")
|
||||
assert _message_contents(db_path)[-1] == "legacy canonical survives"
|
||||
assert _meta_value(db_path, FTS_STALE_KEY) == "1"
|
||||
finally:
|
||||
legacy.close()
|
||||
|
||||
recovered = SessionDB(db_path=db_path)
|
||||
try:
|
||||
assert recovered._fts_stale is False
|
||||
assert _meta_value(db_path, FTS_STALE_KEY) is None
|
||||
assert recovered.search_messages("canonical survives")
|
||||
finally:
|
||||
recovered.close()
|
||||
|
||||
|
||||
Reference in New Issue
Block a user