diff --git a/hermes_state_search.py b/hermes_state_search.py index 7abfe3801f..0074f43eb7 100644 --- a/hermes_state_search.py +++ b/hermes_state_search.py @@ -16,53 +16,49 @@ 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_TRIGRAM_SQL, - MAX_FTS5_QUERY_CHARS, - SCHEMA_VERSION, - _FTS_CJK_TRIGGERS, - escape_like as _escape_like, - fts_rebuild_admission, + FTS_CJK_STALE_KEY, FTS_SQL, FTS_STALE_KEY, FTS_STORAGE_VERSION, FTS_TRIGRAM_SQL, + MAX_FTS5_QUERY_CHARS, SCHEMA_VERSION, _FTS_CJK_TRIGGERS, + escape_like as _escape_like, fts_rebuild_admission, ) # Keep the pre-split logger identity so log filtering/capture is unchanged. logger = logging.getLogger("hermes_state") -# Characters FTS5's query grammar rejects outside a quoted phrase. Anything -# missing here reaches MATCH raw and raises, which the execute site swallows -# into zero results. Assembled through re.escape so the backslash cannot be -# eaten as a regex escape. ``%`` is deliberately excluded: the CJK LIKE -# fallback needs it as a literal (that path escapes wildcards itself). +# Characters FTS5's query grammar rejects outside a quoted phrase; anything +# missing here reaches MATCH raw and raises (swallowed into zero results). +# re.escape so the backslash is not eaten. ``%`` is deliberately excluded: the +# CJK LIKE fallback needs it literal (that path escapes wildcards itself). _FTS5_SPECIAL_CHARS = '+{}():"^@/#&|~[]<>,;!?$=\\\'' _FTS5_SPECIAL_RE = re.compile(f"[{re.escape(_FTS5_SPECIAL_CHARS)}]") _FTS_OPERATORS = frozenset({"AND", "OR", "NOT"}) +_LIKE_SKIP_TOKENS = _FTS_OPERATORS | {"NEAR"} +_LIKE_TOKEN_RE = re.compile(r'"[^"]+"|\S+') # Column list shared by every search route (snippet + metadata, never content). -_SEARCH_SELECT_TAIL = ( - "m.timestamp, m.tool_name, s.source, s.model, s.started_at AS session_started" -) +_SEARCH_SELECT_TAIL = "m.timestamp, m.tool_name, s.source, s.model, s.started_at AS session_started" _LIKE_SNIPPET_SQL = "substr(m.content, max(1, instr(m.content, ?) - 40), 120) AS snippet" _LIKE_ANY_COLUMN_SQL = ( "(m.content LIKE ? ESCAPE '\\' OR m.tool_name LIKE ? ESCAPE '\\' " "OR m.tool_calls LIKE ? ESCAPE '\\')" ) +_LIKE_COALESCED_COLUMN_SQL = ( + "(COALESCE(m.content, '') LIKE ? ESCAPE '\\' OR " + "COALESCE(m.tool_name, '') LIKE ? ESCAPE '\\' OR " + "COALESCE(m.tool_calls, '') LIKE ? ESCAPE '\\')" +) +# ``sort`` -> ORDER BY for the FTS routes; anything unknown is rank-only so +# callers can pass user input through. +_FTS_ORDER_BY = {"newest": "ORDER BY m.timestamp DESC, rank", "oldest": "ORDER BY m.timestamp ASC, rank"} def _meta_row(conn, key: str) -> Optional[sqlite3.Row]: """Point-read one ``state_meta`` row (``None`` when absent).""" - return conn.execute( - "SELECT value FROM state_meta WHERE key = ?", (key,) - ).fetchone() + return conn.execute("SELECT value FROM state_meta WHERE key = ?", (key,)).fetchone() def _delete_meta(conn, *keys: str) -> None: - conn.execute( - f"DELETE FROM state_meta WHERE key IN ({','.join('?' for _ in keys)})", keys - ) + conn.execute(f"DELETE FROM state_meta WHERE key IN ({','.join('?' for _ in keys)})", keys) def _quote_fts_tokens(raw_query: str) -> str: @@ -74,14 +70,23 @@ def _quote_fts_tokens(raw_query: str) -> str: ) +def _like_params(term: str) -> List[str]: + """One ``%term%`` bind per column of ``_LIKE_ANY_COLUMN_SQL``.""" + return [f"%{_escape_like(term)}%"] * 3 + + +def _flatten_text(decoded: Any) -> str: + """Multimodal part list -> joined text (or the placeholder); str passes + through; anything else is ''.""" + if isinstance(decoded, list): + parts = [p.get("text", "") for p in decoded if isinstance(p, dict) and p.get("type") == "text"] + return " ".join(t for t in parts if t).strip() or "[multimodal content]" + return decoded if isinstance(decoded, str) else "" + + def _search_filter_clauses( - where: List[str], - params: list, - *, - include_inactive: bool, - source_filter: Optional[List[str]], - exclude_sources: Optional[List[str]], - role_filter: Optional[List[str]], + where: List[str], params: list, *, include_inactive: bool, source_filter: Optional[List[str]], + exclude_sources: Optional[List[str]], role_filter: Optional[List[str]], ) -> None: """Append the visibility/source/role predicates every search route shares. @@ -105,22 +110,12 @@ class SessionSearchMixin: """See module docstring — mixin for SessionDB (Search cluster).""" _SEARCH_MESSAGE_RESULT_FIELDS = ( - "id", - "session_id", - "role", - "snippet", - "timestamp", - "tool_name", - "source", - "model", - "session_started", - "context", + "id", "session_id", "role", "snippet", "timestamp", "tool_name", "source", "model", + "session_started", "context", ) @classmethod - def _search_message_fields( - cls, fields: Optional[Collection[str]] - ) -> Optional[Tuple[str, ...]]: + def _search_message_fields(cls, fields: Optional[Collection[str]]) -> Optional[Tuple[str, ...]]: """Validate and canonically order an optional result projection.""" if fields is None: return None @@ -130,34 +125,36 @@ class SessionSearchMixin: unknown = requested.difference(cls._SEARCH_MESSAGE_RESULT_FIELDS) if unknown: raise ValueError(f"unknown search result field(s): {', '.join(sorted(unknown))}") - return tuple( - field for field in cls._SEARCH_MESSAGE_RESULT_FIELDS if field in requested - ) + return tuple(field for field in cls._SEARCH_MESSAGE_RESULT_FIELDS if field in requested) def _try_incremental_merge_fts(self) -> None: - """Run one bounded FTS5 merge pass without failing the completed write.""" + """Run one bounded FTS5 merge pass without failing the completed write. + + The canonical write is already committed: no maintenance failure (even + the bare SystemError CPython's sqlite3 layer can raise under + cross-thread errmsg scrambling) may make the caller replay an + ambiguous, possibly-durable write. + """ if not self._fts_enabled: return try: - self._merge_fts_incrementally( - max_pages=self._FTS_MERGE_MAX_PAGES_PER_INDEX - ) + self._merge_fts_incrementally(max_pages=self._FTS_MERGE_MAX_PAGES_PER_INDEX) except Exception as exc: # noqa: BLE001 - post-commit maintenance - # The canonical write is already committed. No maintenance failure - # — including the bare SystemError CPython's sqlite3 layer can raise - # under cross-thread errmsg scrambling — may escape and make the - # caller replay an ambiguous, possibly-durable write. logger.warning("FTS incremental merge failed after commit: %s", exc) + # ── Deferred rebuild engine (base + CJK backfills) ───────────────────── + def fts_rebuild_status(self) -> Optional[Dict[str, Any]]: """Deferred-rebuild progress ``{"pending", "total", "indexed", - "percent"}``, or None when no rebuild is pending. - - Reads state_meta via the pooled reader rather than get_meta() (which - takes self._lock) so search_messages never blocks on the writer lock. - """ + "percent"}``, or None when no rebuild is pending. Reads via the pooled + reader (not get_meta, which takes self._lock) so search_messages never + blocks on the writer lock.""" return self._rebuild_status("fts_rebuild") + def fts_cjk_rebuild_status(self) -> Optional[Dict[str, Any]]: + """CJK-index backfill progress, or None when none is pending.""" + return self._rebuild_status("fts_cjk_rebuild") + def _rebuild_status(self, prefix: str) -> Optional[Dict[str, Any]]: rows = self._read_all( "SELECT key, value FROM state_meta WHERE key IN (?, ?)", @@ -174,6 +171,21 @@ class SessionSearchMixin: pct = min(100, int(100 * progress / total)) return {"pending": True, "total": total, "indexed": progress, "percent": pct} + # Re-index rows in an id window the index is missing. docsize has one row + # per indexed doc, so the anti-join is exact. + _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}" + "AND NOT EXISTS (SELECT 1 FROM {table}_docsize d WHERE d.id = m.id)" + ) + _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}" + ) + def _fts_rebuild_finish(self) -> None: """Finalize the deferred rebuild: boundary sweep + clear markers. @@ -185,21 +197,17 @@ class SessionSearchMixin: """ sweeps = [self._BOUNDARY_SWEEP_SQL.format(table="messages_fts", extra="")] if self._trigram_available: - sweeps.append(self._BOUNDARY_SWEEP_SQL.format( - table="messages_fts_trigram", extra="AND m.role <> 'tool' " - )) + sweeps.append(self._BOUNDARY_SWEEP_SQL.format(table="messages_fts_trigram", extra="AND m.role <> 'tool' ")) self._rebuild_finish("fts_rebuild", sweeps) logger.info("Deferred FTS rebuild complete — all messages indexed.") - # Re-index rows in an id window the index is missing. docsize has one row - # per indexed doc, so the anti-join is exact. - _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}" - "AND NOT EXISTS (SELECT 1 FROM {table}_docsize d WHERE d.id = m.id)" - ) + def _fts_cjk_rebuild_finish(self) -> None: + """Boundary sweep + clear the cjk markers; index becomes servable.""" + self._rebuild_finish("fts_cjk_rebuild", [ + self._BOUNDARY_SWEEP_SQL.format(table="messages_fts_cjk", extra="AND m.role <> 'tool' ") + ]) + self._fts_cjk_available = True + logger.info("CJK FTS index backfill complete — serving CJK search.") def _rebuild_finish(self, prefix: str, sweep_sqls: List[str]) -> None: """Sweep a generous window around the high-water boundary, then clear @@ -213,91 +221,6 @@ class SessionSearchMixin: _delete_meta(conn, f"{prefix}_high_water", f"{prefix}_progress") self._execute_write(_do) - def _fts_teardown_trash_step(self) -> bool: - """Tear down one chunk of a demoted v22 FTS shadow table; True while - work remains. - - Trash tables are PLAIN tables (their vtable parent was demoted), so - chunked DELETE + final DROP involve no FTS5 machinery. Integer - single-column-key tables are drained with a high-water marker so each - chunk's scan is bounded (re-scanning from the start was O(n²) on - large tables). Compound-key tables cannot use a scalar high-water - comparison and keep the chunked ``LIMIT`` delete — they are small by - construction. - """ - with self._lock: - trash = [ - r[0] for r in self._conn.execute( - "SELECT name FROM sqlite_master WHERE type = 'table' " - "AND name LIKE ? ESCAPE '\\'", - (self._FTS_TRASH_PREFIX.replace("_", "\\_") + "%",), - ).fetchall() - ] - if not trash: - return False - - tbl = trash[0] - - def _do(conn): - pk_info = [ - (r[1], (r[2] or "").upper()) - for r in conn.execute(f"PRAGMA table_info({tbl})") - if r[5] > 0 - ] - pk_cols = [name for name, _typ in pk_info] - key = ", ".join(pk_cols) if pk_cols else "rowid" - - if len(pk_cols) == 1 and (not pk_info or pk_info[0][1] == "INTEGER"): - # High-water drain. The marker is read/written in the same - # BEGIN IMMEDIATE as the DELETE, so concurrent callers claim - # disjoint key ranges. Only integer PKs can anchor the numeric - # comparison — the config shadow table (TEXT pk) falls through - # to the chunked delete below. - marker_key = f"fts_teardown_{tbl}_progress" - row = _meta_row(conn, marker_key) - high_water = int(row[0]) if row is not None else 0 - - # Claim the chunk's upper bound: the LAST row of the - # LIMIT window, so a full chunk is deleted per step. - upper_rows = conn.execute( - f"SELECT {key} FROM {tbl} WHERE {key} > ? " - f"ORDER BY {key} LIMIT {self._FTS_REBUILD_CHUNK_ROWS}", - (high_water,), - ).fetchall() - if not upper_rows: - # Drained — the DROP is cheap now. - conn.execute(f"DROP TABLE IF EXISTS {tbl}") - _delete_meta(conn, marker_key) - logger.info("Old FTS shadow table %s torn down.", tbl) - return True - - upper = upper_rows[-1][0] - cur = conn.execute( - f"DELETE FROM {tbl} WHERE {key} > ? AND {key} <= ?", - (high_water, upper), - ) - if cur.rowcount > 0: - self.set_meta(marker_key, str(upper), cursor=conn) - return True - - # Compound-key or rowid trash table: chunked delete (small tables, - # quadratic re-scan is not a concern). - cur = conn.execute( - f"DELETE FROM {tbl} WHERE ({key}) IN " - f"(SELECT {key} FROM {tbl} LIMIT {self._FTS_REBUILD_CHUNK_ROWS})" - ) - if cur.rowcount == 0: - # Empty — the DROP is cheap now. - conn.execute(f"DROP TABLE IF EXISTS {tbl}") - logger.info("Old FTS shadow table %s torn down.", tbl) - return True # re-check: more trash tables / chunks may remain - - try: - return bool(self._execute_write(_do)) - except sqlite3.OperationalError as exc: - logger.debug("FTS trash teardown chunk failed (will retry): %s", exc) - return True - def fts_rebuild_step(self) -> bool: """Backfill one chunk of the deferred FTS rebuild; True while work remains. Safe from any process: chunks are claimed atomically inside @@ -307,24 +230,23 @@ class SessionSearchMixin: return False inserts = [self._CHUNK_INSERT_SQL.format(table="messages_fts", extra="")] if self._trigram_available: - inserts.append(self._CHUNK_INSERT_SQL.format( - table="messages_fts_trigram", extra=" AND role <> 'tool'" - )) + inserts.append(self._CHUNK_INSERT_SQL.format(table="messages_fts_trigram", extra=" AND role <> 'tool'")) return self._rebuild_step( - "fts_rebuild", inserts, - fail_msg="FTS rebuild chunk failed (will retry): %s", + "fts_rebuild", inserts, fail_msg="FTS rebuild chunk failed (will retry): %s", finish=self._fts_rebuild_finish, ) - _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}" - ) + def fts_cjk_rebuild_step(self) -> bool: + """Backfill one chunk of the CJK index. True while work remains.""" + if not self._fts_enabled or not self._fts_cjk_loaded: + return False + return self._rebuild_step( + "fts_cjk_rebuild", + [self._CHUNK_INSERT_SQL.format(table="messages_fts_cjk", extra=" AND role <> 'tool'")], + fail_msg="CJK FTS rebuild chunk failed (will retry): %s", finish=self._fts_cjk_rebuild_finish, + ) - def _rebuild_step( - self, prefix: str, insert_sqls: List[str], *, fail_msg: str, finish - ) -> bool: + def _rebuild_step(self, prefix: str, insert_sqls: List[str], *, fail_msg: str, finish) -> bool: """Shared chunk engine for the base and CJK deferred backfills.""" high_water_raw = self.get_meta(f"{prefix}_high_water") if high_water_raw is None: @@ -333,27 +255,21 @@ class SessionSearchMixin: chunk = self._FTS_REBUILD_CHUNK_ROWS def _do(conn): - # Re-read progress inside the write transaction (BEGIN IMMEDIATE - # is already held by _execute_write) — this is the claim: two - # workers can't read the same progress value concurrently. + # Re-reading progress inside the BEGIN IMMEDIATE held by + # _execute_write IS the claim: two workers cannot read the same value. row = _meta_row(conn, f"{prefix}_progress") if row is None: return False # finished (or cleared) by another process progress = int(row[0]) if progress >= high_water: return False - - # The chunk upper bound is an id, not a row count, so gaps from - # deleted rows don't shrink chunks below the claimed range. + # Upper bound is an id, not a row count, so deleted-row gaps don't + # shrink chunks below the claimed range. upper = min(progress + chunk, high_water) for sql in insert_sqls: conn.execute(sql, (progress, upper)) - # Publish progress in the same transaction as the rows it - # covers — crash-atomic: either both land or neither does. - conn.execute( - "UPDATE state_meta SET value = ? WHERE key = ?", - (str(upper), f"{prefix}_progress"), - ) + # Progress lands in the same transaction as its rows (crash-atomic). + conn.execute("UPDATE state_meta SET value = ? WHERE key = ?", (str(upper), f"{prefix}_progress")) return upper < high_water try: @@ -368,32 +284,70 @@ class SessionSearchMixin: return False return bool(more) - def fts_cjk_rebuild_status(self) -> Optional[Dict[str, Any]]: - """CJK-index backfill progress, or None when none is pending.""" - return self._rebuild_status("fts_cjk_rebuild") + def _fts_teardown_trash_step(self) -> bool: + """Tear down one chunk of a demoted v22 FTS shadow table; True while + work remains. - def fts_cjk_rebuild_step(self) -> bool: - """Backfill one chunk of the CJK index. True while work remains.""" - if not self._fts_enabled or not self._fts_cjk_loaded: + Trash tables are PLAIN tables (their vtable parent was demoted), so no + FTS5 machinery is involved. Integer single-column-key tables drain with + a high-water marker so each chunk's scan is bounded (restarting the + scan was O(n²)); compound-key tables cannot use a scalar high-water and + keep the chunked ``LIMIT`` delete — they are small by construction. + """ + with self._lock: + trash = [ + r[0] for r in self._conn.execute( + "SELECT name FROM sqlite_master WHERE type = 'table' " + "AND name LIKE ? ESCAPE '\\'", + (self._FTS_TRASH_PREFIX.replace("_", "\\_") + "%",), + ).fetchall() + ] + if not trash: return False - return self._rebuild_step( - "fts_cjk_rebuild", - [self._CHUNK_INSERT_SQL.format( - table="messages_fts_cjk", extra=" AND role <> 'tool'" - )], - fail_msg="CJK FTS rebuild chunk failed (will retry): %s", - finish=self._fts_cjk_rebuild_finish, - ) + tbl = trash[0] - def _fts_cjk_rebuild_finish(self) -> None: - """Boundary sweep + clear the cjk markers; index becomes servable.""" - self._rebuild_finish("fts_cjk_rebuild", [ - self._BOUNDARY_SWEEP_SQL.format( - table="messages_fts_cjk", extra="AND m.role <> 'tool' " + def _do(conn): + pk_info = [(r[1], (r[2] or "").upper()) for r in conn.execute(f"PRAGMA table_info({tbl})") if r[5] > 0] + pk_cols = [name for name, _typ in pk_info] + key = ", ".join(pk_cols) if pk_cols else "rowid" + if len(pk_cols) == 1 and (not pk_info or pk_info[0][1] == "INTEGER"): + # High-water drain; the marker is read/written in the same + # BEGIN IMMEDIATE as the DELETE so concurrent callers claim + # disjoint ranges. Only INTEGER pks anchor the comparison (the + # TEXT-pk config shadow table falls through below). + marker_key = f"fts_teardown_{tbl}_progress" + row = _meta_row(conn, marker_key) + high_water = int(row[0]) if row is not None else 0 + # Claim the LAST row of the LIMIT window so a full chunk goes per step. + upper_rows = conn.execute( + f"SELECT {key} FROM {tbl} WHERE {key} > ? " + f"ORDER BY {key} LIMIT {self._FTS_REBUILD_CHUNK_ROWS}", + (high_water,), + ).fetchall() + if not upper_rows: + conn.execute(f"DROP TABLE IF EXISTS {tbl}") + _delete_meta(conn, marker_key) + logger.info("Old FTS shadow table %s torn down.", tbl) + return True + upper = upper_rows[-1][0] + cur = conn.execute(f"DELETE FROM {tbl} WHERE {key} > ? AND {key} <= ?", (high_water, upper)) + if cur.rowcount > 0: + self.set_meta(marker_key, str(upper), cursor=conn) + return True + cur = conn.execute( + f"DELETE FROM {tbl} WHERE ({key}) IN " + f"(SELECT {key} FROM {tbl} LIMIT {self._FTS_REBUILD_CHUNK_ROWS})" ) - ]) - self._fts_cjk_available = True - logger.info("CJK FTS index backfill complete — serving CJK search.") + if cur.rowcount == 0: + conn.execute(f"DROP TABLE IF EXISTS {tbl}") + logger.info("Old FTS shadow table %s torn down.", tbl) + return True # re-check: more trash tables / chunks may remain + + try: + return bool(self._execute_write(_do)) + except sqlite3.OperationalError as exc: + logger.debug("FTS trash teardown chunk failed (will retry): %s", exc) + return True def _fts_cjk_reset_if_stale(self) -> None: """From-scratch rebuild of a stale cjk index (triggers were dropped, @@ -410,15 +364,11 @@ class SessionSearchMixin: conn.execute(f"DROP TRIGGER IF EXISTS {trig}") conn.execute("DROP TABLE IF EXISTS messages_fts_cjk") conn.execute("DROP VIEW IF EXISTS messages_fts_cjk_src") - _delete_meta( - conn, FTS_CJK_STALE_KEY, "fts_cjk_rebuild_high_water", "fts_cjk_rebuild_progress" - ) + _delete_meta(conn, FTS_CJK_STALE_KEY, "fts_cjk_rebuild_high_water", "fts_cjk_rebuild_progress") return True - was_stale = self._execute_write(_do) - if was_stale: - # Recreate outside the write transaction — _ensure_fts_cjk_schema - # uses executescript(), which implicitly commits and must not run - # inside _execute_write's BEGIN IMMEDIATE. + if self._execute_write(_do): + # Recreate OUTSIDE the write transaction: _ensure_fts_cjk_schema uses + # executescript(), which implicitly commits. with self._lock: self._ensure_fts_cjk_schema(self._conn) self._conn.commit() @@ -426,42 +376,29 @@ class SessionSearchMixin: def _fts_external_index_empty_with_messages(self, conn) -> bool: """True when the base FTS table indexes nothing while ``messages`` has rows (the post-demote empty-index shape). Caller holds ``self._lock``. - Healthy and mid-backfill installs never match.""" + Healthy and mid-backfill installs never match. docsize is the + authoritative "is this rowid indexed" surface for external-content + FTS5; EXISTS not COUNT(*) because this runs on every writable open.""" try: - has_msg = conn.execute( - "SELECT EXISTS(SELECT 1 FROM messages)" - ).fetchone()[0] - if not has_msg: + if not conn.execute("SELECT EXISTS(SELECT 1 FROM messages)").fetchone()[0]: return False - # docsize is the authoritative "is this rowid indexed" surface for - # external-content FTS5 (probing the vtable is unreliable across - # builds). EXISTS not COUNT(*): this runs on every writable open - # and COUNT(*) is a full b-tree scan. - has_fts = conn.execute( - "SELECT EXISTS(SELECT 1 FROM messages_fts_docsize)" - ).fetchone()[0] - return not has_fts + return not conn.execute("SELECT EXISTS(SELECT 1 FROM messages_fts_docsize)").fetchone()[0] except sqlite3.OperationalError: - # Table absent / FTS disabled mid-init — not this failure class. - return False + return False # table absent / FTS disabled mid-init — not this failure class def _fts_index_known_empty(self, conn) -> bool: """True when the base external-content index holds no rows; a missing table counts as empty (the schema ensure that follows creates it).""" try: - n = conn.execute( - "SELECT COUNT(*) FROM messages_fts_docsize" - ).fetchone()[0] - return int(n) == 0 + return int(conn.execute("SELECT COUNT(*) FROM messages_fts_docsize").fetchone()[0]) == 0 except sqlite3.OperationalError: return True def _reset_fts_index_to_empty(self, conn) -> None: """Truncate the v23 external-content tables via FTS5 ``'delete-all'``. - A plain ``DELETE`` is O(rows) on external-content FTS5 (minutes on a - large index, holding the write lock) and corrupts the index when - indexed rows have diverged from ``messages`` — precisely the shape + A plain DELETE is O(rows) on external-content FTS5 and corrupts the + index when indexed rows diverged from ``messages`` — exactly the shape this repair handles. The backfill worker replays its id range with no anti-join, so a replay from zero is only safe once the index is known empty; this is how a partially indexed DB gets there. @@ -472,59 +409,43 @@ class SessionSearchMixin: except sqlite3.OperationalError: pass # table absent — already an empty surface + def _reseed_missing_progress(self, conn) -> None: + """high_water without progress: fts_rebuild_step reads missing progress + as "done by another process" and optimize would no-op then stamp. + Reset a partially indexed DB to a known-empty surface (the chunk + worker replays without an anti-join) and re-seed progress.""" + if _meta_row(conn, "fts_rebuild_progress") is None: + if not self._fts_index_known_empty(conn): + self._reset_fts_index_to_empty(conn) + self.set_meta("fts_rebuild_progress", "0", cursor=conn) + def _seed_fts_rebuild_markers(self, conn, *, force: bool = False) -> int: """Write ``fts_rebuild_high_water`` / ``fts_rebuild_progress`` for a - full backfill; returns the high-water id. - - Without ``force`` and with high_water already set, only repairs a - missing progress key, resetting a partially indexed DB to a - known-empty surface first (the chunk worker replays without an - anti-join). Caller holds the write transaction. - """ + full backfill; returns the high-water id. Without ``force`` and with + high_water already set, only repairs a missing progress key. Caller + holds the write transaction.""" existing_hw = _meta_row(conn, "fts_rebuild_high_water") if existing_hw is not None and not force: - if _meta_row(conn, "fts_rebuild_progress") is None: - # high_water without progress: fts_rebuild_step treats missing - # progress as "done by another process" and optimize would - # no-op then stamp. Re-seed progress so the chunk loop runs. - if not self._fts_index_known_empty(conn): - self._reset_fts_index_to_empty(conn) - self.set_meta("fts_rebuild_progress", "0", cursor=conn) + self._reseed_missing_progress(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) return int(hw) def _repair_optimize_bookkeeping(self) -> None: - """Heal interrupted demote/backfill bookkeeping before optimize runs. - - 1. Empty external-content index with messages and no markers (demote - crash window / settle without backfill): seed a full backfill. - 2. high_water without progress: seed progress (resetting a partially - populated index first so the anti-join-free replay cannot - duplicate rows). - - Must not invent markers on a still-legacy inline DB: optimize would - then skip demote and attempt v23 INSERTs against the inline table forever. - """ + """Heal interrupted demote/backfill bookkeeping before optimize runs: + orphan high_water-without-progress gets progress re-seeded; an empty + external index with messages and no markers (demote crash window / + premature stamp) gets a full backfill claim. Never invents markers on a + still-legacy inline DB — optimize would then skip demote and INSERT + against the inline table forever.""" def _do(conn): if _meta_row(conn, "fts_rebuild_high_water") is not None: - # Repair orphan high_water-without-progress only. Never - # invent a fresh claim on a healthy complete index. - if _meta_row(conn, "fts_rebuild_progress") is None: - if not self._fts_index_known_empty(conn): - self._reset_fts_index_to_empty(conn) - self.set_meta("fts_rebuild_progress", "0", cursor=conn) + self._reseed_missing_progress(conn) return - - # No markers. On a still-legacy DB demote owns marker creation. if self._db_has_legacy_inline_fts(conn): - return - - # Non-legacy empty external index (demote crash window / premature - # stamp): seed a full backfill claim. + return # demote owns marker creation if self._fts_external_index_empty_with_messages(conn): _delete_meta(conn, "fts_storage_version") self._seed_fts_rebuild_markers(conn, force=True) @@ -539,25 +460,20 @@ class SessionSearchMixin: if not self._fts_enabled or self.read_only: return False with self._lock: - if self._db_has_legacy_inline_fts(self._conn): + conn = self._conn + if self._db_has_legacy_inline_fts(conn): return True - # Interrupted optimize: legacy vtables already demoted, but the - # transition is unfinished until markers clear and trash is gone. - if _meta_row(self._conn, "fts_rebuild_high_water") is not None: - return True - # CJK-bigram index work — only offerable when THIS process can - # tokenize: a pending backfill (markers set at creation on a - # populated DB) or a stale index awaiting a from-scratch rebuild. + if _meta_row(conn, "fts_rebuild_high_water") is not None: + return True # interrupted optimize: demoted but unfinished + # CJK work is only offerable when THIS process can tokenize. if self._fts_cjk_loaded and ( - _meta_row(self._conn, "fts_cjk_rebuild_high_water") is not None - or _meta_row(self._conn, FTS_CJK_STALE_KEY) is not None + _meta_row(conn, "fts_cjk_rebuild_high_water") is not None + or _meta_row(conn, FTS_CJK_STALE_KEY) is not None ): return True - if self._has_fts_trash(self._conn): + if self._has_fts_trash(conn): return True - # Crash window: empty external index, messages present, no markers, - # no trash. Re-run seeds markers and backfills. - return self._fts_external_index_empty_with_messages(self._conn) + return self._fts_external_index_empty_with_messages(conn) def _demote_legacy_fts_to_trash(self) -> int: """Demote the legacy inline FTS vtables and stage their shadow tables @@ -566,7 +482,7 @@ class SessionSearchMixin: Markers are written in the same BEGIN IMMEDIATE as the demote, BEFORE the empty v23 schema is created (``executescript`` implicitly COMMITs - and cannot run inside that transaction). This closes the crash window + and cannot run inside that transaction), closing the crash window where trash + empty v23 tables exist with no backfill claim. """ def _stage(conn): @@ -594,41 +510,29 @@ class SessionSearchMixin: ] for sh in shadows: conn.execute(f"ALTER TABLE {sh} RENAME TO fts_v22_trash_{sh}") - # Claim the backfill BEFORE empty v23 tables exist so a crash - # before schema ensure resumes instead of stamping an empty index. hw = self._seed_fts_rebuild_markers(conn, force=True) _delete_meta(conn, "fts_optimize_available") return hw hw = int(self._execute_write(_stage)) - - # Outside the write transaction: ``_ensure_fts_schema`` uses - # executescript(), which implicitly commits. Markers are durable. - self._ensure_v23_fts_tables( - "failed to create v23 messages_fts during optimize-storage demote" - ) + # Outside the write transaction (executescript commits); markers are durable. + self._ensure_v23_fts_tables("failed to create v23 messages_fts during optimize-storage demote") return hw def _ensure_v23_fts_tables(self, failure_message: str) -> None: - """Ensure the v23 external-content base + trigram tables under the - lock (IF NOT EXISTS, cheap); raise *failure_message* without the base - table, since the backfill loop would otherwise retry "no such table" - forever.""" + """Ensure the v23 base + trigram tables under the lock (IF NOT EXISTS, + cheap); raise *failure_message* without the base table, since the + backfill loop would otherwise retry "no such table" forever.""" with self._lock: base_ok = self._ensure_fts_schema(self._conn, "messages_fts", FTS_SQL) - trigram_ok = self._ensure_fts_schema( - self._conn, "messages_fts_trigram", FTS_TRIGRAM_SQL - ) + trigram_ok = self._ensure_fts_schema(self._conn, "messages_fts_trigram", FTS_TRIGRAM_SQL) self._trigram_available = bool(trigram_ok) if not base_ok: raise sqlite3.OperationalError(failure_message) self._conn.commit() def optimize_fts_storage( - self, - *, - progress_cb: Optional[Callable[[Dict[str, Any]], None]] = None, - vacuum: bool = True, + self, *, progress_cb: Optional[Callable[[Dict[str, Any]], None]] = None, vacuum: bool = True, ) -> Dict[str, Any]: """Migrate a legacy v22 inline-FTS DB to the v23 external-content schema, foreground and to completion; re-running resumes an @@ -640,28 +544,20 @@ class SessionSearchMixin: if self.read_only: return {"ok": False, "reason": "read_only"} - # Heal empty-index / orphan-marker bookkeeping BEFORE deciding whether - # to demote again, so the phases below actually run. + # Heal bookkeeping BEFORE deciding whether to demote again. self._repair_optimize_bookkeeping() - - # Demote only when still on the legacy shape; a prior demote - # (markers/trash present) skips straight to backfill + teardown. with self._lock: legacy = self._db_has_legacy_inline_fts(self._conn) pending = self.get_meta("fts_rebuild_high_water") is not None if legacy and not pending: self._demote_legacy_fts_to_trash() elif pending and not legacy: - # Resume mid-demote: markers exist, empty v23 tables may still be - # missing if the process died between the staged demote commit and - # schema ensure. - self._ensure_v23_fts_tables( - "failed to re-create v23 messages_fts on optimize-storage resume" - ) + # Resume mid-demote: the process may have died between the staged + # demote commit and schema ensure. + self._ensure_v23_fts_tables("failed to re-create v23 messages_fts on optimize-storage resume") - # A stale CJK index can only be recovered from scratch; reset it so - # the cjk backfill phase rebuilds it. Then ensure table + markers - # exist (a v23 DB gaining the cjk index for the first time). + # A stale CJK index can only be recovered from scratch; then ensure + # table + markers exist (a v23 DB gaining the cjk index for the first time). self._fts_cjk_reset_if_stale() if self._fts_cjk_loaded: with self._lock: @@ -671,42 +567,28 @@ class SessionSearchMixin: def _emit(phase: str) -> None: if progress_cb is None: return - st = self.fts_rebuild_status() - if st is None: - st = self.fts_cjk_rebuild_status() + st = self.fts_rebuild_status() or self.fts_cjk_rebuild_status() progress_cb({ - "phase": phase, - "percent": st["percent"] if st else 100, - "indexed": st["indexed"] if st else 0, - "total": st["total"] if st else 0, + "phase": phase, "percent": st["percent"] if st else 100, + "indexed": st["indexed"] if st else 0, "total": st["total"] if st else 0, }) - def _pause(chunk_seconds: float) -> None: - """Inter-chunk throttle — the single place the duty cycle is - enforced. Without it back-to-back BEGIN IMMEDIATE chunks starve a - live gateway/CLI sharing the DB out of its lock retries.""" - time.sleep(max( - self._FTS_REBUILD_MIN_PAUSE, - chunk_seconds * self._FTS_REBUILD_DUTY_FACTOR, - )) - def _drive(phase: str, step) -> None: - """Run *step* to completion, emitting progress and throttling - between chunks so a live gateway sharing the DB stays responsive.""" + """Run *step* to completion; the inter-chunk sleep is the single + place the duty cycle is enforced — back-to-back BEGIN IMMEDIATE + chunks starve a live gateway/CLI out of its lock retries.""" while True: _t0 = time.monotonic() if not step(): break _emit(phase) - _pause(time.monotonic() - _t0) + time.sleep(max(self._FTS_REBUILD_MIN_PAUSE, (time.monotonic() - _t0) * self._FTS_REBUILD_DUTY_FACTOR)) - # Phase 1: base backfill. Phase 1b: CJK-bigram backfill (its own - # marker pair; no-op without the tokenizer or pending work). + # Phase 1: base backfill; 1b: CJK-bigram backfill (own marker pair). _emit("backfill") _drive("backfill", self.fts_rebuild_step) _emit("backfill") _drive("backfill", self.fts_cjk_rebuild_step) - # Phase 2: tear down the demoted legacy shadow tables in chunks. _emit("teardown") _drive("teardown", self._fts_teardown_trash_step) @@ -719,13 +601,9 @@ class SessionSearchMixin: still_trash = self._has_fts_trash(self._conn) empty_index = self._fts_external_index_empty_with_messages(self._conn) if still_pending or still_trash or empty_index: - reason = ( - "backfill_incomplete" if still_pending or empty_index - else "teardown_incomplete" - ) + reason = "backfill_incomplete" if still_pending or empty_index else "teardown_incomplete" logger.warning( - "FTS storage optimization did not settle (%s): " - "pending=%s trash=%s empty_index=%s", + "FTS storage optimization did not settle (%s): pending=%s trash=%s empty_index=%s", reason, still_pending, still_trash, empty_index, ) return {"ok": False, "reason": reason, "vacuumed": None} @@ -739,31 +617,24 @@ class SessionSearchMixin: self._conn.execute("VACUUM") vacuum_ok = True except sqlite3.OperationalError as exc: - # Usually no free disk for VACUUM's temp copy; the optimization - # still succeeded, space is reclaimed by a later VACUUM. + # Usually no free disk for VACUUM's temp copy; a later VACUUM reclaims. logger.warning("VACUUM after FTS optimize failed: %s", exc) vacuum_ok = False - # Best-effort WAL fold-back. REFUSED (SQLITE_BUSY) while another - # connection holds a WAL read-mark, so callers must NOT size the - # result by stat()ing the file — use :meth:`logical_size_bytes`. - # PASSIVE, not TRUNCATE: a TRUNCATE reset from this transient CLI - # process would race a live gateway writer and tear B-tree pages. + # Best-effort WAL fold-back, REFUSED (SQLITE_BUSY) while another + # connection holds a read-mark — so callers must size the result + # via logical_size_bytes, not stat(). PASSIVE, never TRUNCATE: a + # TRUNCATE reset from a transient CLI would race a live writer. try: with self._lock: self._conn.execute("PRAGMA wal_checkpoint(PASSIVE)") except Exception as exc: - logger.debug( - "WAL checkpoint (PASSIVE) after optimize VACUUM failed: %s", - exc, - ) + logger.debug("WAL checkpoint (PASSIVE) after optimize VACUUM failed: %s", exc) - # Phase 4: stamp the FTS layout (the source of truth for "optimized"), - # clear the "available" flag, and advance schema_version if a DB - # opened only by pre-decoupling code left it behind. + # Phase 4: stamp the FTS layout (source of truth for "optimized"), clear + # the "available" flag, advance schema_version if pre-decoupling code + # left it behind. Re-checked inside the write transaction so a + # concurrent writer cannot race a stamp past incomplete work. def _settle(conn): - # Re-check inside the write transaction so a concurrent writer - # cannot race a stamp past incomplete work. Returns a refusal - # reason (nothing stamped) or None. if _meta_row(conn, "fts_rebuild_high_water") is not None: return "backfill_incomplete" if self._has_fts_trash(conn): @@ -779,24 +650,16 @@ class SessionSearchMixin: return None refusal = self._execute_write(_settle) if refusal is not None: - # A concurrent process changed state since the pre-vacuum check; - # report instead of crashing the CLI — a re-run can still settle. - logger.warning( - "FTS storage optimization settle refused (%s)", refusal - ) + logger.warning("FTS storage optimization settle refused (%s)", refusal) return {"ok": False, "reason": refusal, "vacuumed": vacuum_ok} _emit("done") - logger.info( - "FTS storage optimization complete (layout v%d).", FTS_STORAGE_VERSION - ) + logger.info("FTS storage optimization complete (layout v%d).", FTS_STORAGE_VERSION) return {"ok": True, "vacuumed": vacuum_ok} + # ── Read views ───────────────────────────────────────────────────────── + def get_anchored_view( - self, - session_id: str, - around_message_id: int, - window: int = 5, - bookend: int = 3, + self, session_id: str, around_message_id: int, window: int = 5, bookend: int = 3, keep_roles: Optional[Tuple[str, ...]] = ("user", "assistant"), ) -> Dict[str, Any]: """Anchored window (``get_messages_around``) plus session bookends. @@ -812,34 +675,16 @@ class SessionSearchMixin: resolution in one call. Empty slices + zero counts when the anchor isn't in the session. ``keep_roles=None`` disables role filtering. """ - if bookend < 0: - bookend = 0 - - primitive = self.get_messages_around( - session_id, around_message_id, window=window - ) + bookend = max(bookend, 0) + primitive = self.get_messages_around(session_id, around_message_id, window=window) window_rows = primitive["window"] if not window_rows: - return { - "window": [], - "messages_before": 0, - "messages_after": 0, - "bookend_start": [], - "bookend_end": [], - } + return {"window": [], "messages_before": 0, "messages_after": 0, "bookend_start": [], "bookend_end": []} - # Apply role filter to the window, but never drop the anchor itself. + filtered_window = window_rows if keep_roles is not None: keep_set = set(keep_roles) - filtered_window = [ - m for m in window_rows - if m.get("id") == around_message_id or m.get("role") in keep_set - ] - else: - filtered_window = window_rows - - window_min_id = window_rows[0]["id"] - window_max_id = window_rows[-1]["id"] + filtered_window = [m for m in window_rows if m.get("id") == around_message_id or m.get("role") in keep_set] bookend_start_rows: List[Any] = [] bookend_end_rows: List[Any] = [] @@ -849,7 +694,6 @@ class SessionSearchMixin: if keep_roles is not None: role_clause = f" AND role IN ({','.join('?' for _ in keep_roles)})" role_params = list(keep_roles) - with self._read_ctx() as conn: def _bookend(op: str, boundary_id: int, order: str): return conn.execute( @@ -859,10 +703,9 @@ class SessionSearchMixin: f"ORDER BY id {order} LIMIT ?", (session_id, boundary_id, *role_params, bookend), ).fetchall() - - bookend_start_rows = _bookend("<", window_min_id, "ASC") + bookend_start_rows = _bookend("<", window_rows[0]["id"], "ASC") # End rows come back DESC for the LIMIT cap; flip to ASC. - bookend_end_rows = list(reversed(_bookend(">", window_max_id, "DESC"))) + bookend_end_rows = list(reversed(_bookend(">", window_rows[-1]["id"], "DESC"))) def _hydrate(row) -> Dict[str, Any]: msg = dict(row) @@ -872,9 +715,7 @@ class SessionSearchMixin: try: msg["tool_calls"] = json.loads(msg["tool_calls"]) except (json.JSONDecodeError, TypeError): - logger.warning( - "Failed to deserialize tool_calls in get_anchored_view, falling back to []" - ) + logger.warning("Failed to deserialize tool_calls in get_anchored_view, falling back to []") msg["tool_calls"] = [] if msg.get("display_metadata") is not None: msg["display_metadata"] = self._decode_display_metadata(msg["display_metadata"]) @@ -889,26 +730,22 @@ class SessionSearchMixin: } def list_recent_user_messages( - self, - session_id: str, - limit: int = 20, - include_inactive: bool = False, + self, session_id: str, limit: int = 20, include_inactive: bool = False, ) -> List[Dict[str, Any]]: """The *limit* most-recent real user turns, newest first, as ``{id, timestamp, preview}`` (preview = first 80 chars, whitespace collapsed). Used by /rewind and ``/undo [N]``. - Bookkeeping timeline rows (``display_kind`` set) are excluded: they - are durable ``role='user'`` rows but no client counts them as user - turns, and including them made ``/undo`` soft-delete from a marker - instead of the last real turn. Only active messages by default. + Bookkeeping timeline rows (``display_kind`` set) are excluded: no + client counts them as user turns, and including them made ``/undo`` + soft-delete from a marker instead of the last real turn. Legacy + standalone compaction handoffs are role='user' rows with NO + display_kind — invisible to SQL — so fetch with headroom and drop them + in the decode loop; otherwise ``/undo N`` pairs an in-memory count that + excludes handoffs with a DB pick that includes them. """ active_clause = "" if include_inactive else " AND active = 1" display_clause = " AND (display_kind IS NULL OR display_kind = '')" - # Legacy standalone compaction handoffs are role='user' rows with NO - # display_kind — SQL can't see them, so fetch with headroom and drop - # them in the decode loop; otherwise /undo N pairs an in-memory count - # that excludes handoffs with a DB pick that includes them. fetch_limit = int(limit) * 2 + 5 with self._lock: rows = self._conn.execute( @@ -928,27 +765,19 @@ class SessionSearchMixin: decoded = self._decode_content(row["content"]) if ContextCompressor._is_context_summary_content(decoded): continue # compaction handoff — never a user-originated turn - if isinstance(decoded, list): - # Multimodal — flatten text parts. - text_parts = [ - p.get("text", "") for p in decoded - if isinstance(p, dict) and p.get("type") == "text" - ] - preview = " ".join(t for t in text_parts if t).strip() - if not preview: - preview = "[multimodal content]" - elif isinstance(decoded, str): - # A /skill turn embeds the whole skill body; show what the user - # typed instead of the skill's opening prose. + if isinstance(decoded, str): + # A /skill turn embeds the whole skill body; show what was typed. preview = describe_skill_invocation(decoded) or decoded else: - preview = "" - preview = " ".join(preview.split()) # collapse whitespace + preview = _flatten_text(decoded) + preview = " ".join(preview.split()) if len(preview) > 80: preview = preview[:77] + "..." result.append({"id": row["id"], "timestamp": row["timestamp"], "preview": preview}) return result + # ── Query analysis ───────────────────────────────────────────────────── + @staticmethod def _sanitize_fts5_query(query: str) -> str: """Sanitize user input for FTS5 MATCH (raw special characters raise @@ -959,8 +788,8 @@ class SessionSearchMixin: # Cap before any regex processing so adversarial input stays bounded. query = query[:MAX_FTS5_QUERY_CHARS] - # Step 1: protect balanced quoted phrases via numbered placeholders. - # Linear scan, not regex, so pathological quote runs cannot backtrack. + # 1. Protect balanced quoted phrases via numbered placeholders. Linear + # scan, not regex, so pathological quote runs cannot backtrack. _quoted_parts: list = [] pieces: list[str] = [] i = 0 @@ -972,43 +801,31 @@ class SessionSearchMixin: continue end = query.find('"', i + 1) if end == -1: - # Unmatched quote: replace with whitespace. - pieces.append(" ") + pieces.append(" ") # unmatched quote -> whitespace i += 1 continue _quoted_parts.append(query[i:end + 1]) pieces.append(f"\x00Q{len(_quoted_parts) - 1}\x00") i = end + 1 - sanitized = "".join(pieces) - # Step 2: strip remaining FTS5-special characters (see - # _FTS5_SPECIAL_CHARS); e.g. an unquoted ``TODO: fix`` parses as - # ``column:term`` and raises "no such column". + # 2. Strip FTS5-special characters (an unquoted ``TODO: fix`` parses as + # ``column:term``). ``%`` is only spared for the CJK LIKE fallback. sanitized = _FTS5_SPECIAL_RE.sub(" ", sanitized) - - # Step 2b: ``%`` is only spared for the CJK LIKE fallback; a non-CJK - # query never reaches it, so ``50%`` would hit MATCH raw and raise. if "%" in sanitized and not SessionSearchMixin._contains_cjk(sanitized): sanitized = sanitized.replace("%", " ") - - # Step 3: collapse repeated * and drop leading * (prefix needs a char). + # 3. Collapse repeated * and drop leading * (prefix needs a char). sanitized = re.sub(r"\*+", "*", sanitized) sanitized = re.sub(r"(^|\s)\*", r"\1", sanitized) - - # Step 4: drop dangling boolean operators at start/end (syntax errors). + # 4. Drop dangling boolean operators at start/end (syntax errors). sanitized = re.sub(r"(?i)^(AND|OR|NOT)\b\s*", "", sanitized.strip()) sanitized = re.sub(r"(?i)\s+(AND|OR|NOT)\s*$", "", sanitized.strip()) - - # Step 5: quote dotted/hyphenated/underscored terms in ONE pass (the - # tokenizer splits on them; sequential passes double-quote - # ``my-app.config``). + # 5. Quote dotted/hyphenated/underscored terms in ONE pass (sequential + # passes double-quote ``my-app.config``). sanitized = re.sub(r"\b(\w+(?:[._-]\w+)+)\b", r'"\1"', sanitized) - - # Step 6: restore preserved quoted phrases. + # 6. Restore preserved quoted phrases. for i, quoted in enumerate(_quoted_parts): sanitized = sanitized.replace(f"\x00Q{i}\x00", quoted) - return sanitized.strip() @staticmethod @@ -1023,12 +840,10 @@ class SessionSearchMixin: @staticmethod def _contains_cjk(text: str) -> bool: - """Check if text contains CJK (Chinese, Japanese, Korean) characters.""" return any(SessionSearchMixin._is_cjk_codepoint(ord(ch)) for ch in text) @classmethod def _count_cjk(cls, text: str) -> int: - """Count CJK characters in text.""" return sum(1 for ch in text if cls._is_cjk_codepoint(ord(ch))) @classmethod @@ -1051,9 +866,7 @@ class SessionSearchMixin: """True when every non-operator token is >=3 chars: a shorter token produces no trigrams, and with FTS5's implicit AND one such token makes the whole MATCH return nothing.""" - tokens = [ - t for t in query.strip('"').strip().split() if t.upper() not in _FTS_OPERATORS - ] + tokens = [t for t in query.strip('"').strip().split() if t.upper() not in _FTS_OPERATORS] return bool(tokens) and all(len(t) >= 3 for t in tokens) @classmethod @@ -1067,60 +880,48 @@ class SessionSearchMixin: if t.upper() not in _FTS_OPERATORS and cls._contains_cjk(t) ) - def _run_trigram_search( - self, - raw_query: str, - *, - table: str = "messages_fts_trigram", - order_by_sql: str, - include_inactive: bool, - source_filter: List[str] = None, - exclude_sources: List[str] = None, - role_filter: List[str] = None, - limit: int = 20, - offset: int = 0, - ) -> Optional[List[Dict[str, Any]]]: - """Search a substring-capable index (``messages_fts_trigram`` or - ``messages_fts_cjk``): trigram matches substrings regardless of word - boundaries (CJK phrases unicode61 splits, Latin runs it fuses onto - adjacent CJK like ``修改youer服务端``); cjk-bigram splits Latin runs - off CJK for an exact ranked match. Returns ``None`` when the query - cannot execute (e.g. tokenizer unavailable) so the caller can fall - back.""" - tri_sql, tri_params = self._fts_match_sql( - table, _quote_fts_tokens(raw_query), order_by_sql, - include_inactive=include_inactive, source_filter=source_filter, - exclude_sources=exclude_sources, role_filter=role_filter, - limit=limit, offset=offset, + def _trigram_route_ok(self, raw_query: str) -> bool: + """Per-token CJK length gate for the trigram index: ``广西 OR 桂林 OR + 漓江`` has 6 CJK chars total but 2 per token, so trigram returns 0.""" + return ( + self._count_cjk(raw_query) >= 3 + and not self._has_short_cjk_token(raw_query) + and self._trigram_available ) - with self._read_ctx() as conn: - try: - tri_cursor = conn.execute(tri_sql, tri_params) - except sqlite3.OperationalError: - # Query failed at runtime — let the caller fall back. - return None - return [dict(row) for row in tri_cursor.fetchall()] + + 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" + if not self._contains_cjk(sanitized): + return "fts5" + raw = sanitized.strip('"').strip() + if self._fts_cjk_available and not self._has_lone_cjk_run(raw): + return "fts_cjk" + if self._trigram_route_ok(raw): + return "trigram" + return "like_scan" + except Exception: + return "unknown" + + # ── Query builders / runners ─────────────────────────────────────────── @staticmethod def _fts_match_sql( - table: str, - match_query: str, - order_by_sql: str, - *, - include_inactive: bool, - source_filter: Optional[List[str]], - exclude_sources: Optional[List[str]], - role_filter: Optional[List[str]], - limit: int, - offset: int, + table: str, match_query: str, order_by_sql: str, *, include_inactive: bool, + source_filter: Optional[List[str]], exclude_sources: Optional[List[str]], + role_filter: Optional[List[str]], limit: int, offset: int, ) -> Tuple[str, list]: """MATCH query + params against one FTS5 index joined to messages/sessions.""" where = [f"{table} MATCH ?"] params: list = [match_query] _search_filter_clauses( - where, params, include_inactive=include_inactive, - source_filter=source_filter, exclude_sources=exclude_sources, - role_filter=role_filter, + where, params, include_inactive=include_inactive, source_filter=source_filter, + exclude_sources=exclude_sources, role_filter=role_filter, ) params.extend([limit, offset]) sql = f""" @@ -1136,159 +937,36 @@ class SessionSearchMixin: """ return sql, params - def search_messages( - self, - query: str, - source_filter: List[str] = None, - exclude_sources: List[str] = None, - role_filter: List[str] = None, - limit: int = 20, - offset: int = 0, - sort: str = None, - include_inactive: bool = False, - fields: Optional[Collection[str]] = None, - ) -> List[Dict[str, Any]]: - """Instrumented wrapper around :meth:`_search_messages_impl`: logs one - line per slow search with the routing path taken so latency stays - attributable per query shape. Threshold HERMES_SEARCH_SLOW_MS - (default 1000; 0 logs every call).""" - started = time.time() - rows = None + def _match_rows( + self, table: str, match_query: str, order_by_sql: str, *, fail_open: Optional[str] = None, + operational_debug: Optional[str] = None, **kwargs, + ) -> Optional[List[Dict[str, Any]]]: + """Run one MATCH against *table*; ``None`` when the query cannot + execute (tokenizer unavailable / syntax) so the caller falls back. + + *fail_open* names the index for the substring-capable routes: a + corruption-class ``DatabaseError`` there detaches the derived indexes + (``_enter_fts_fail_open``) and answers from canonical rows — a live + search never performs the unbounded rebuild. Non-FTS corruption, or + any ``DatabaseError`` without *fail_open*, propagates. + """ + sql, params = self._fts_match_sql(table, match_query, order_by_sql, **kwargs) try: - rows = self._search_messages_impl( - query, - source_filter=source_filter, - exclude_sources=exclude_sources, - role_filter=role_filter, - limit=limit, - offset=offset, - sort=sort, - include_inactive=include_inactive, - fields=fields, + return [dict(row) for row in self._read_all(sql, params)] + except sqlite3.OperationalError: + if operational_debug: + logger.debug(operational_debug, exc_info=True) + return None + except sqlite3.DatabaseError as exc: + if fail_open is None or not self._enter_fts_fail_open(exc): + raise + logger.warning( + "%s FTS search hit a corruption error (%s); detached FTS and falling back to canonical LIKE.", + fail_open, exc, ) - return rows - finally: - try: - threshold = float(os.getenv("HERMES_SEARCH_SLOW_MS", "1000")) - except (TypeError, ValueError): - threshold = 1000.0 - elapsed_ms = (time.time() - started) * 1000.0 - if elapsed_ms >= threshold: - logger.info( - "slow session search: path=%s elapsed=%.0fms rows=%s query=%r", - self._describe_search_path(query), - elapsed_ms, - len(rows) if rows is not None else "err", - query[:200], - ) + return None - 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" - if not self._contains_cjk(sanitized): - return "fts5" - raw = sanitized.strip('"').strip() - if self._fts_cjk_available and not self._has_lone_cjk_run(raw): - return "fts_cjk" - if ( - self._count_cjk(raw) >= 3 - and not self._has_short_cjk_token(raw) - and self._trigram_available - ): - return "trigram" - return "like_scan" - 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 (FTS5's implicit conjunction) and - ``NOT`` negates the following term rather than being discarded.""" - 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: - 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"%{_escape_like(term)}%"] * 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})"] - _search_filter_clauses( - where, params, include_inactive=include_inactive, - source_filter=source_filter, exclude_sources=exclude_sources, - role_filter=role_filter, - ) - order = ( - "ASC" - if isinstance(sort, str) and sort.strip().lower() == "oldest" - else "DESC" - ) - return self._like_rows( - where, [snippet_term, *params, limit, offset], - order_by=f"ORDER BY m.timestamp {order}, m.id {order}", limit_sql="LIMIT ? OFFSET ?", - ) - - def _like_rows( - self, where: List[str], params: list, *, order_by: str, limit_sql: str - ) -> List[Dict[str, Any]]: + def _like_rows(self, where: List[str], params: list, *, order_by: str, limit_sql: str) -> List[Dict[str, Any]]: """Canonical-table LIKE scan; ``params[0]`` is the snippet anchor term.""" sql = f""" SELECT m.id, m.session_id, m.role, @@ -1302,14 +980,66 @@ class SessionSearchMixin: """ return [dict(row) for row in self._read_all(sql, params)] + @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 (FTS5's implicit conjunction) and + ``NOT`` negates the following term rather than being discarded.""" + groups: List[List[Tuple[str, bool]]] = [[]] + negate_next = False + for raw_token in _LIKE_TOKEN_RE.findall(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: + clauses.append(f"NOT {_LIKE_COALESCED_COLUMN_SQL}" if negated else _LIKE_COALESCED_COLUMN_SQL) + params.extend(_like_params(term)) + 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, *, limit: int, offset: int, sort: Optional[str], **filters, + ) -> 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})"] + _search_filter_clauses(where, params, **filters) + order = "ASC" if isinstance(sort, str) and sort.strip().lower() == "oldest" else "DESC" + return self._like_rows( + where, [snippet_term, *params, limit, offset], + order_by=f"ORDER BY m.timestamp {order}, m.id {order}", limit_sql="LIMIT ? OFFSET ?", + ) + 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: - stale = self._read_one( - "SELECT 1 FROM state_meta WHERE key = ? LIMIT 1", (FTS_STALE_KEY,) - ) + stale = self._read_one("SELECT 1 FROM state_meta WHERE key = ? LIMIT 1", (FTS_STALE_KEY,)) except sqlite3.Error: return if stale is not None: @@ -1319,16 +1049,12 @@ class SessionSearchMixin: self._fts_cjk_available = False def _finalize_search_matches( - self, - matches: List[Dict[str, Any]], - result_fields: Optional[Collection[str]] = None, + self, matches: List[Dict[str, Any]], result_fields: Optional[Collection[str]] = None, ) -> List[Dict[str, Any]]: """Attach neighboring messages (1 before + after, only when the projection consumes ``context``) and trim full content. Each context query takes its own read transaction, never a lock across N queries.""" - context_matches = ( - matches if result_fields is None or "context" in result_fields else () - ) + 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: @@ -1365,53 +1091,56 @@ class SessionSearchMixin: )""", (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 + match["context"] = [ + {"role": row["role"], "content": _flatten_text(self._decode_content(row["content"]))[:200]} + for row in ctx_cursor.fetchall() + ] except Exception: match["context"] = [] - # No search route selects full content (snippet + metadata only); the - # pop is a guard for any future route that does. + # No route selects full content; the pop guards any future one that does. 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 - ] - + matches = [{field: match[field] for field in result_fields if field in match} for match in matches] return matches + # ── search_messages ──────────────────────────────────────────────────── + + def search_messages( + self, query: str, source_filter: List[str] = None, exclude_sources: List[str] = None, + role_filter: List[str] = None, limit: int = 20, offset: int = 0, sort: str = None, + include_inactive: bool = False, fields: Optional[Collection[str]] = None, + ) -> List[Dict[str, Any]]: + """Instrumented wrapper around :meth:`_search_messages_impl`: logs one + line per slow search with the routing path taken so latency stays + attributable per query shape. Threshold HERMES_SEARCH_SLOW_MS + (default 1000; 0 logs every call).""" + started = time.time() + rows = None + try: + rows = self._search_messages_impl( + query, source_filter=source_filter, exclude_sources=exclude_sources, role_filter=role_filter, + limit=limit, offset=offset, sort=sort, include_inactive=include_inactive, fields=fields, + ) + return rows + finally: + try: + threshold = float(os.getenv("HERMES_SEARCH_SLOW_MS", "1000")) + except (TypeError, ValueError): + threshold = 1000.0 + elapsed_ms = (time.time() - started) * 1000.0 + if elapsed_ms >= threshold: + logger.info( + "slow session search: path=%s elapsed=%.0fms rows=%s query=%r", + self._describe_search_path(query), elapsed_ms, + len(rows) if rows is not None else "err", query[:200], + ) + def _search_messages_impl( - self, - query: str, - source_filter: List[str] = None, - exclude_sources: List[str] = None, - role_filter: List[str] = None, - limit: int = 20, - offset: int = 0, - sort: str = None, - include_inactive: bool = False, - fields: Optional[Collection[str]] = None, + self, query: str, source_filter: List[str] = None, exclude_sources: List[str] = None, + role_filter: List[str] = None, limit: int = 20, offset: int = 0, sort: str = None, + include_inactive: bool = False, fields: Optional[Collection[str]] = None, ) -> List[Dict[str, Any]]: """FTS5 search across session messages (keywords, ``"phrases"``, AND/OR/NOT, ``prefix*``). @@ -1428,10 +1157,8 @@ class SessionSearchMixin: searches every row. """ result_fields = self._search_message_fields(fields) - if not query or not query.strip(): return [] - query = self._sanitize_fts5_query(query) if not query: return [] @@ -1442,153 +1169,38 @@ class SessionSearchMixin: ) self._refresh_fts_stale_state() if self._fts_stale: - matches = self._search_messages_like_fallback( - query, limit=limit, offset=offset, sort=sort, **filters - ) + matches = self._search_messages_like_fallback(query, limit=limit, offset=offset, sort=sort, **filters) return self._finalize_search_matches(matches, result_fields=result_fields) if not self._fts_enabled: return [] - # Normalise sort; anything unknown falls back to rank-only so callers - # can pass through user input. - if isinstance(sort, str): - sort_norm = sort.strip().lower() - if sort_norm not in ("newest", "oldest"): - sort_norm = None - else: - sort_norm = None - - if sort_norm == "newest": - order_by_sql = "ORDER BY m.timestamp DESC, rank" - elif sort_norm == "oldest": - order_by_sql = "ORDER BY m.timestamp ASC, rank" - else: - order_by_sql = "ORDER BY rank" - - # CJK queries bypass the unicode61 table, whose tokenizer splits CJK - # into single characters ("大别山项目" -> "大 AND 别 AND ...": false - # positives, missed phrases). 3+ CJK chars -> trigram; shorter -> - # LIKE (trigram needs 9 UTF-8 bytes = 3 CJK chars). - matches: List[Dict[str, Any]] = [] + order_by_sql = _FTS_ORDER_BY.get(sort.strip().lower() if isinstance(sort, str) else None, "ORDER BY rank") + route = dict(order_by_sql=order_by_sql, limit=limit, offset=offset, **filters) + # Tool rows are excluded from the trigram/cjk indexes (see FTS_TRIGRAM_SQL). + wants_tool_rows = bool(role_filter) and "tool" in role_filter is_cjk = self._contains_cjk(query) if is_cjk: - raw_query = query.strip('"').strip() - _trigram_succeeded = False - # Tool rows are excluded from the trigram/cjk indexes (see - # FTS_TRIGRAM_SQL), so a role='tool' CJK query must use LIKE. - _wants_tool_rows = bool(role_filter) and "tool" in role_filter - - # CJK-bigram route: serves every CJK shape the legacy code split - # between trigram and LIKE full scans, except role='tool' queries - # and LONE 1-char CJK runs (the index stores bigrams for runs >=2, - # so a single-char term only matches isolated chars — LIKE is broader). - if ( - self._fts_cjk_available - and not _wants_tool_rows - and not self._has_lone_cjk_run(raw_query) - ): - cjk_sql, cjk_params = self._fts_match_sql( - "messages_fts_cjk", _quote_fts_tokens(raw_query), order_by_sql, - limit=limit, offset=offset, **filters, - ) - try: - matches = [dict(row) for row in self._read_all(cjk_sql, cjk_params)] - _trigram_succeeded = True - except sqlite3.OperationalError: - # Tokenizer missing / query syntax — trigram + LIKE still answer. - logger.debug( - "messages_fts_cjk query failed; falling back to " - "trigram/LIKE", exc_info=True, - ) - except sqlite3.DatabaseError as exc: - # A live search never performs the unbounded full rebuild: - # detach the derived indexes and answer from canonical rows. - # Non-FTS corruption is not safe to reinterpret here. - if not self._enter_fts_fail_open(exc): - raise - logger.warning( - "CJK-bigram FTS search hit a corruption error (%s); " - "detached FTS and falling back to canonical LIKE.", - exc, - ) - - # Per-token CJK length check: trigram needs >=3 CJK chars per - # token. "广西 OR 桂林 OR 漓江" has 6 CJK chars total but 2 per - # token — trigram returns 0, so such queries take LIKE. - if ( - not _trigram_succeeded - and self._count_cjk(raw_query) >= 3 - and not self._has_short_cjk_token(raw_query) - and self._trigram_available - and not _wants_tool_rows - ): - tri_sql, tri_params = self._fts_match_sql( - "messages_fts_trigram", _quote_fts_tokens(raw_query), order_by_sql, - limit=limit, offset=offset, **filters, - ) - try: - matches = [dict(row) for row in self._read_all(tri_sql, tri_params)] - _trigram_succeeded = True - except sqlite3.OperationalError: - # Trigram query failed at runtime — fall through to LIKE. - pass - except sqlite3.DatabaseError as exc: - # Same bounded recovery as the CJK/main paths; a non-FTS - # storage error stays fatal rather than hidden as a miss. - if not self._enter_fts_fail_open(exc): - raise - logger.warning( - "Trigram FTS search hit a corruption error (%s); " - "detached FTS and falling back to canonical LIKE.", - exc, - ) - if not _trigram_succeeded: - # LIKE substring fallback; one clause per non-operator token so - # "广西 OR 桂林 OR 漓江" matches each term independently. - non_op_tokens = [ - t for t in raw_query.split() if t.upper() not in _FTS_OPERATORS - ] or [raw_query] - like_params: list = [] - for tok in non_op_tokens: - like_params += [f"%{_escape_like(tok)}%"] * 3 - like_where = [ - f"({' OR '.join([_LIKE_ANY_COLUMN_SQL] * len(non_op_tokens))})" - ] - _search_filter_clauses(like_where, like_params, **filters) - # instr() for snippet uses first search token - matches = self._like_rows( - like_where, [non_op_tokens[0], *like_params, limit, offset], - order_by="ORDER BY m.timestamp DESC", limit_sql="LIMIT ? OFFSET ?", - ) + matches = self._search_cjk(query, wants_tool_rows, route) else: - sql, params = self._fts_match_sql( - "messages_fts", query, order_by_sql, limit=limit, offset=offset, **filters - ) + sql, params = self._fts_match_sql("messages_fts", query, **route) try: matches = [dict(row) for row in self._read_all(sql, params)] except sqlite3.OperationalError: - # FTS5 query syntax error despite sanitization — return empty - return [] + return [] # FTS5 syntax error despite sanitization except sqlite3.DatabaseError as exc: - # Corruption parent class (OperationalError is caught above). - # Live search must stay bounded: detach the derived indexes and - # answer from canonical rows; repair paths own the rebuild. + # Corruption parent class: detach the derived indexes and answer + # from canonical rows; repair paths own the rebuild. if not self._enter_fts_fail_open(exc): raise - matches = self._search_messages_like_fallback( - query, limit=limit, offset=offset, sort=sort, **filters - ) + matches = self._search_messages_like_fallback(query, limit=limit, offset=offset, sort=sort, **filters) # Deferred-rebuild supplement: while the backfill is pending the FTS # indexes miss the (progress, high_water] gap; top up with a bounded - # LIKE scan over that range so old messages never silently vanish - # mid-rebuild. The cost decays to zero as the backfill advances. - rebuild_status = self.fts_rebuild_status() - if rebuild_status is not None and len(matches) < limit: + # LIKE scan so old messages never vanish mid-rebuild. Cost decays to + # zero as the backfill advances. + if self.fts_rebuild_status() is not None and len(matches) < limit: try: - gap_matches = self._search_unindexed_gap( - query, limit - len(matches), **filters - ) + gap_matches = self._search_unindexed_gap(query, limit - len(matches), **filters) seen_ids = {m["id"] for m in matches} matches.extend(m for m in gap_matches if m["id"] not in seen_ids) except sqlite3.OperationalError as exc: @@ -1601,36 +1213,52 @@ class SessionSearchMixin: # (needs >=3-char tokens). Gated on a miss so successful searches keep # their ranking; trade-off: "cat" may then match "concatenate". # Skipped for role='tool' (both indexes exclude tool rows). - if ( - not matches - and not is_cjk - and not (bool(role_filter) and "tool" in role_filter) - ): - _fb_query = query.strip('"').strip() - fb_kwargs = dict(order_by_sql=order_by_sql, limit=limit, offset=offset, **filters) + if not matches and not is_cjk and not wants_tool_rows: + fb_query = _quote_fts_tokens(query.strip('"').strip()) if self._fts_cjk_available: - matches = self._run_trigram_search( - _fb_query, table="messages_fts_cjk", **fb_kwargs - ) or matches - if ( - not matches - and self._trigram_available - and self._trigram_eligible_tokens(query) - ): - matches = self._run_trigram_search(_fb_query, **fb_kwargs) or matches - + matches = self._match_rows("messages_fts_cjk", fb_query, **route) or matches + if not matches and self._trigram_available and self._trigram_eligible_tokens(query): + matches = self._match_rows("messages_fts_trigram", fb_query, **route) or matches return self._finalize_search_matches(matches, result_fields=result_fields) - def _search_unindexed_gap( - self, - fts_query: str, - limit: int, - *, - include_inactive: bool = False, - source_filter: Optional[List[str]] = None, - exclude_sources: Optional[List[str]] = None, - role_filter: Optional[List[str]] = None, - ) -> List[Dict[str, Any]]: + def _search_cjk(self, query: str, wants_tool_rows: bool, route: Dict[str, Any]) -> List[Dict[str, Any]]: + """CJK routing: the unicode61 table splits CJK into single characters + ("大别山项目" -> "大 AND 别 AND ...": false positives, missed phrases). + + cjk-bigram serves every shape the legacy code split between trigram + and LIKE, except role='tool' queries and LONE 1-char CJK runs (the + index stores bigrams for runs >=2, so a single-char term only matches + isolated chars — LIKE is broader). Then trigram (>=3 CJK chars per + token), then a LIKE substring scan with one clause per non-operator + token so "广西 OR 桂林 OR 漓江" matches each term independently. + """ + raw_query = query.strip('"').strip() + match_query = _quote_fts_tokens(raw_query) + if self._fts_cjk_available and not wants_tool_rows and not self._has_lone_cjk_run(raw_query): + matches = self._match_rows( + "messages_fts_cjk", match_query, fail_open="CJK-bigram", + operational_debug="messages_fts_cjk query failed; falling back to trigram/LIKE", **route, + ) + if matches is not None: + return matches + if self._trigram_route_ok(raw_query) and not wants_tool_rows: + matches = self._match_rows("messages_fts_trigram", match_query, fail_open="Trigram", **route) + if matches is not None: + return matches + non_op_tokens = [t for t in raw_query.split() if t.upper() not in _FTS_OPERATORS] or [raw_query] + like_params: list = [] + for tok in non_op_tokens: + like_params += _like_params(tok) + like_where = [f"({' OR '.join([_LIKE_ANY_COLUMN_SQL] * len(non_op_tokens))})"] + filters = {k: route[k] for k in ("include_inactive", "source_filter", "exclude_sources", "role_filter")} + _search_filter_clauses(like_where, like_params, **filters) + # instr() for the snippet uses the first search token. + return self._like_rows( + like_where, [non_op_tokens[0], *like_params, route["limit"], route["offset"]], + order_by="ORDER BY m.timestamp DESC", limit_sql="LIMIT ? OFFSET ?", + ) + + def _search_unindexed_gap(self, fts_query: str, limit: int, **filters) -> List[Dict[str, Any]]: """LIKE-scan ids in (fts_rebuild_progress, fts_rebuild_high_water] — the rows the deferred rebuild hasn't indexed yet. The FTS query is degraded to AND-joined substring terms (quoted phrases kept whole): @@ -1638,39 +1266,25 @@ class SessionSearchMixin: status = self.fts_rebuild_status() if status is None or limit <= 0: return [] - progress, high_water = status["indexed"], status["total"] - - terms: List[str] = [] - for raw_tok in re.findall(r'"[^"]+"|\S+', fts_query): - tok = raw_tok.strip('"').strip("*").strip() - if tok and tok.upper() not in {"AND", "OR", "NOT", "NEAR"}: - terms.append(tok) + terms = [ + tok for tok in (t.strip('"').strip("*").strip() for t in _LIKE_TOKEN_RE.findall(fts_query)) + if tok and tok.upper() not in _LIKE_SKIP_TOKENS + ] if not terms: return [] - where = ["m.id > ? AND m.id <= ?"] - params: list = [progress, high_water] + params: list = [status["indexed"], status["total"]] for term in terms: where.append(_LIKE_ANY_COLUMN_SQL) - params += [f"%{_escape_like(term)}%"] * 3 - _search_filter_clauses( - where, params, include_inactive=include_inactive, - source_filter=source_filter, exclude_sources=exclude_sources, - role_filter=role_filter, - ) + params += _like_params(term) + _search_filter_clauses(where, params, **filters) return self._like_rows( - where, [terms[0], *params, limit], - order_by="ORDER BY m.timestamp DESC", limit_sql="LIMIT ?", + where, [terms[0], *params, limit], order_by="ORDER BY m.timestamp DESC", limit_sql="LIMIT ?", ) def search_sessions_by_id( - self, - query: str, - limit: int = 20, - include_archived: bool = True, - source: str = None, - sources: List[str] = None, - exclude_sources: List[str] = None, + self, query: str, limit: int = 20, include_archived: bool = True, source: str = None, + sources: List[str] = None, exclude_sources: List[str] = None, ) -> List[Dict[str, Any]]: """Search surfaced sessions by exact/prefix/substring session id (paste an id from logs and jump to it). Also matches ``_lineage_root_id`` so @@ -1678,19 +1292,12 @@ class SessionSearchMixin: needle = (query or "").strip().lower() if not needle or limit <= 0: return [] - # list_sessions_rich pushes the id LIKE filter (own id + forward - # compression chain) into SQL; over-fetch so the in-Python - # exact/prefix/substring ranking has candidates, then truncate. + # compression chain) into SQL; over-fetch so the in-Python ranking has + # candidates, then truncate. candidates = self.list_sessions_rich( - source=source, - sources=sources, - exclude_sources=exclude_sources, - limit=max(limit * 4, limit), - offset=0, - include_archived=include_archived, - order_by_last_active=True, - id_query=needle, + source=source, sources=sources, exclude_sources=exclude_sources, limit=max(limit * 4, limit), + offset=0, include_archived=include_archived, order_by_last_active=True, id_query=needle, ) def score(row: Dict[str, Any]) -> int: @@ -1702,20 +1309,19 @@ class SessionSearchMixin: return 1 return 2 - ranked = sorted( - enumerate(candidates), - key=lambda item: (score(item[1]), item[0]), - ) + ranked = sorted(enumerate(candidates), key=lambda item: (score(item[1]), item[0])) return [row for _, row in ranked[:limit]] + # ── FTS maintenance commands ─────────────────────────────────────────── + def _fts_table_exists(self, name: str) -> bool: - """True if an FTS5 virtual table is queryable in this DB.""" + """True if an FTS5 virtual table is queryable in this DB ("no such + table" and "vtable constructor failed" — missing tokenizer / + mid-teardown — both count as not queryable).""" try: self._conn.execute(f"SELECT 1 FROM {name} LIMIT 0") return True except sqlite3.DatabaseError: - # "no such table", or "vtable constructor failed" (missing - # tokenizer / mid-teardown) — either way not queryable. return False def optimize_fts(self) -> int: @@ -1723,8 +1329,7 @@ class SessionSearchMixin: Pure maintenance: changes neither results nor ``snippet()`` output, only layout and speed; complementary to VACUUM, which then returns - the freed pages. Skips absent tables, so it is safe unconditionally. - Returns the number of indexes optimized. + the freed pages. Skips absent tables. Returns the number optimized. """ optimized = 0 with self._lock: @@ -1732,14 +1337,10 @@ class SessionSearchMixin: if not self._fts_table_exists(tbl): continue try: - self._conn.execute( - f"INSERT INTO {tbl}({tbl}) VALUES('optimize')" - ) + self._conn.execute(f"INSERT INTO {tbl}({tbl}) VALUES('optimize')") optimized += 1 except sqlite3.OperationalError as exc: - logger.warning( - "FTS optimize failed for %s: %s", tbl, exc - ) + logger.warning("FTS optimize failed for %s: %s", tbl, exc) return optimized def rebuild_fts(self) -> int: @@ -1767,21 +1368,15 @@ class SessionSearchMixin: if not self._fts_table_exists(tbl): continue try: - self._conn.execute( - f"INSERT INTO {tbl}({tbl}) VALUES('rebuild')" - ) + self._conn.execute(f"INSERT INTO {tbl}({tbl}) VALUES('rebuild')") self._conn.commit() rebuilt += 1 except sqlite3.OperationalError as exc: self._conn.rollback() - logger.warning( - "FTS rebuild failed for %s: %s", tbl, exc - ) + logger.warning("FTS rebuild failed for %s: %s", tbl, exc) return rebuilt - def _merge_fts_incrementally( - self, *, max_pages: int, max_commands: Optional[int] = None - ) -> int: + def _merge_fts_incrementally(self, *, max_pages: int, max_commands: Optional[int] = None) -> int: """Run bounded FTS5 ``'merge'`` commands against each present index. A positive merge rank stops after ~that many output pages, so each @@ -1818,19 +1413,11 @@ class SessionSearchMixin: for tbl in self._FTS_TABLES: if not self._fts_table_exists(tbl): continue - # One-time (per instance) usermerge floor; metadata-only write, - # persisted so future connections inherit it. if not self._fts_usermerge_floor_applied: - self._conn.execute( - f"INSERT INTO {tbl}({tbl}, rank) " - "VALUES('usermerge', 2)" - ) + self._conn.execute(f"INSERT INTO {tbl}({tbl}, rank) VALUES('usermerge', 2)") for _ in range(max_commands): before = self._conn.total_changes - self._conn.execute( - f"INSERT INTO {tbl}({tbl}, rank) VALUES('merge', ?)", - (max_pages,), - ) + self._conn.execute(f"INSERT INTO {tbl}({tbl}, rank) VALUES('merge', ?)", (max_pages,)) executed += 1 if self._conn.total_changes - before < 2: break