diff --git a/hermes_state.py b/hermes_state.py index 4733147521..82daf8e99e 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -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. diff --git a/hermes_state_common.py b/hermes_state_common.py index 4fb7dc848c..e386c1260b 100644 --- a/hermes_state_common.py +++ b/hermes_state_common.py @@ -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 diff --git a/hermes_state_schema.py b/hermes_state_schema.py index 68a888879c..bb735ff72d 100644 --- a/hermes_state_schema.py +++ b/hermes_state_schema.py @@ -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) ) diff --git a/hermes_state_search.py b/hermes_state_search.py index e92869eefe..643c7f7398 100644 --- a/hermes_state_search.py +++ b/hermes_state_search.py @@ -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, diff --git a/tests/state/test_fts_runtime_rebuild.py b/tests/state/test_fts_runtime_rebuild.py index d65094d0fa..e9c14ecd54 100644 --- a/tests/state/test_fts_runtime_rebuild.py +++ b/tests/state/test_fts_runtime_rebuild.py @@ -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()