diff --git a/hermes_state_schema.py b/hermes_state_schema.py index 0130fd434f..81ad88b4f4 100644 --- a/hermes_state_schema.py +++ b/hermes_state_schema.py @@ -21,21 +21,9 @@ from typing import Dict, List, Optional, Sequence from hermes_constants import get_hermes_home from hermes_startup_watchdog import report_startup_progress from hermes_state_common import ( - DEFERRED_INDEX_SQL, - FTS_CJK_STALE_KEY, - FTS_REBUILD_DEFERRAL_KEY, - FTS_STALE_KEY, - FTS_SQL, - FTS_STORAGE_VERSION, - FTS_TRIGRAM_SQL, - LEGACY_FTS_SQL, - LEGACY_FTS_TRIGRAM_SQL, - SCHEMA_SQL, - SCHEMA_VERSION, - _FTS_CJK_TRIGGERS, - _FTS_TRIGGERS, - _ephemeral_child_sql, - fts_rebuild_admission, + DEFERRED_INDEX_SQL, FTS_CJK_STALE_KEY, FTS_REBUILD_DEFERRAL_KEY, FTS_STALE_KEY, FTS_SQL, + FTS_STORAGE_VERSION, FTS_TRIGRAM_SQL, LEGACY_FTS_SQL, LEGACY_FTS_TRIGRAM_SQL, SCHEMA_SQL, + SCHEMA_VERSION, _FTS_CJK_TRIGGERS, _FTS_TRIGGERS, _ephemeral_child_sql, fts_rebuild_admission, ) # Keep the pre-split logger identity so log filtering/capture is unchanged. @@ -64,6 +52,103 @@ _READ_PROBE_STATEMENTS: Optional[tuple] = None _FTS_TRIGRAM_TRIGGERS = tuple(n for n in _FTS_TRIGGERS if "_trigram_" in n) _FTS_BASE_TRIGGERS = tuple(n for n in _FTS_TRIGGERS if n not in _FTS_TRIGRAM_TRIGGERS) +_LEGACY_INLINE_CONCAT_SQL = ( + "COALESCE(content, '') || ' ' || " + "COALESCE(tool_name, '') || ' ' || " + "COALESCE(tool_calls, '') " +) +_SESSION_MODEL_USAGE_INDEX_SQL = ( + "CREATE INDEX IF NOT EXISTS idx_session_model_usage_session ON session_model_usage(session_id)", + "CREATE INDEX IF NOT EXISTS idx_session_model_usage_model ON session_model_usage(model)", +) +_SESSION_MODEL_USAGE_COLS = """session_id, model, billing_provider, billing_base_url, + billing_mode, task, api_call_count, input_tokens, + output_tokens, cache_read_tokens, cache_write_tokens, + reasoning_tokens, estimated_cost_usd, actual_cost_usd, + cost_status, cost_source, first_seen, last_seen""" +_SESSION_MODEL_USAGE_HEAL_DDL = """CREATE TABLE session_model_usage ( + session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE, + model TEXT NOT NULL, + billing_provider TEXT NOT NULL DEFAULT '', + billing_base_url TEXT NOT NULL DEFAULT '', + billing_mode TEXT NOT NULL DEFAULT '', + task TEXT NOT NULL DEFAULT '', + api_call_count INTEGER NOT NULL DEFAULT 0, + input_tokens INTEGER NOT NULL DEFAULT 0, + output_tokens INTEGER NOT NULL DEFAULT 0, + cache_read_tokens INTEGER NOT NULL DEFAULT 0, + cache_write_tokens INTEGER NOT NULL DEFAULT 0, + reasoning_tokens INTEGER NOT NULL DEFAULT 0, + estimated_cost_usd REAL NOT NULL DEFAULT 0, + actual_cost_usd REAL NOT NULL DEFAULT 0, + cost_status TEXT, + cost_source TEXT, + first_seen REAL, + last_seen REAL, + PRIMARY KEY (session_id, model, billing_provider, billing_base_url, billing_mode, task) +)""" +# Same table, v22-migration indentation (statement text is pinned by the SQL trace harness). +_SESSION_MODEL_USAGE_V22_DDL = """CREATE TABLE session_model_usage ( + session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE, + model TEXT NOT NULL, + billing_provider TEXT NOT NULL DEFAULT '', + billing_base_url TEXT NOT NULL DEFAULT '', + billing_mode TEXT NOT NULL DEFAULT '', + task TEXT NOT NULL DEFAULT '', + api_call_count INTEGER NOT NULL DEFAULT 0, + input_tokens INTEGER NOT NULL DEFAULT 0, + output_tokens INTEGER NOT NULL DEFAULT 0, + cache_read_tokens INTEGER NOT NULL DEFAULT 0, + cache_write_tokens INTEGER NOT NULL DEFAULT 0, + reasoning_tokens INTEGER NOT NULL DEFAULT 0, + estimated_cost_usd REAL NOT NULL DEFAULT 0, + actual_cost_usd REAL NOT NULL DEFAULT 0, + cost_status TEXT, + cost_source TEXT, + first_seen REAL, + last_seen REAL, + PRIMARY KEY (session_id, model, billing_provider, billing_base_url, billing_mode, task) + )""" +# Statement text pinned by the SQL trace harness (whitespace included). +_SESSION_MODEL_USAGE_V20_SEED_SQL = """INSERT OR IGNORE INTO session_model_usage ( + session_id, model, billing_provider, + billing_base_url, billing_mode, + api_call_count, input_tokens, + output_tokens, cache_read_tokens, + cache_write_tokens, reasoning_tokens, + estimated_cost_usd, actual_cost_usd, + cost_status, cost_source, first_seen, last_seen + ) + SELECT id, COALESCE(model, 'unknown'), + COALESCE(billing_provider, ''), + COALESCE(billing_base_url, ''), + COALESCE(billing_mode, ''), + COALESCE(api_call_count, 0), + COALESCE(input_tokens, 0), + COALESCE(output_tokens, 0), + COALESCE(cache_read_tokens, 0), + COALESCE(cache_write_tokens, 0), + COALESCE(reasoning_tokens, 0), + COALESCE(estimated_cost_usd, 0), + COALESCE(actual_cost_usd, 0), + cost_status, cost_source, + started_at, COALESCE(ended_at, started_at) + FROM sessions + WHERE COALESCE(input_tokens, 0) + + COALESCE(output_tokens, 0) + + COALESCE(cache_read_tokens, 0) + + COALESCE(cache_write_tokens, 0) + + COALESCE(reasoning_tokens, 0) > 0""" +_TITLE_UNIQUE_INDEX_SQL = ( + "CREATE UNIQUE INDEX IF NOT EXISTS idx_sessions_title_unique " + "ON sessions(title) WHERE title IS NOT NULL" +) + + +def _q(ident: str) -> str: + """Double-quote an SQL identifier.""" + return '"' + ident.replace('"', '""') + '"' + def schema_read_probe_statements() -> tuple: """SELECT statements that fail iff a live store is behind SCHEMA_SQL. @@ -86,15 +171,7 @@ def schema_read_probe_statements() -> tuple: if _READ_PROBE_STATEMENTS is None: tables = SessionSchemaMixin._parse_schema_columns(SCHEMA_SQL) _READ_PROBE_STATEMENTS = tuple( - 'SELECT {} FROM "{}" LIMIT 0'.format( - ", ".join( - '"{}"."{}"'.format( - table.replace('"', '""'), col.replace('"', '""') - ) - for col in cols - ), - table.replace('"', '""'), - ) + "SELECT {} FROM {} LIMIT 0".format(", ".join(f"{_q(table)}.{_q(col)}" for col in cols), _q(table)) for table, cols in sorted(tables.items()) ) return _READ_PROBE_STATEMENTS @@ -119,7 +196,6 @@ class SessionSchemaMixin: ).fetchall() except sqlite3.OperationalError: return - for session_id, prompt in rows: try: prompt_hash = self._store_system_prompt(cursor, prompt) @@ -158,15 +234,11 @@ class SessionSchemaMixin: pass @staticmethod - def _fts_trigger_count( - cursor: sqlite3.Cursor, - names: Sequence[str] = _FTS_TRIGGERS, - ) -> int: + def _fts_trigger_count(cursor: sqlite3.Cursor, names: Sequence[str] = _FTS_TRIGGERS) -> int: """Count how many of *names* currently exist as triggers (pass _FTS_BASE_TRIGGERS / _FTS_TRIGRAM_TRIGGERS to check one half).""" if not names: - # "name IN ()" is a SQLite syntax error. - return 0 + return 0 # "name IN ()" is a SQLite syntax error placeholders = ",".join("?" for _ in names) row = cursor.execute( f"SELECT COUNT(*) FROM sqlite_master " @@ -177,16 +249,11 @@ class SessionSchemaMixin: @staticmethod def _fts_update_trigger_needs_narrowing(sql: Optional[str]) -> bool: - """True when trigger SQL is missing AFTER UPDATE OF (still broad).""" + """True when trigger SQL is a broad AFTER UPDATE (missing ``OF``).""" if not sql: return False - # Collapse whitespace so multi-line DDL still matches. - compact = " ".join(sql.split()).upper() - # Already narrowed. - if "AFTER UPDATE OF " in compact: - return False - # Broad UPDATE trigger that we still need to replace. - return "AFTER UPDATE ON " in compact + compact = " ".join(sql.split()).upper() # multi-line DDL still matches + return "AFTER UPDATE OF " not in compact and "AFTER UPDATE ON " in compact def _migrate_broad_fts_update_triggers(self, cursor: sqlite3.Cursor) -> int: """Replace broad AFTER UPDATE FTS triggers with AFTER UPDATE OF variants. @@ -200,10 +267,7 @@ class SessionSchemaMixin: # CJK is v23-only. Decide the layout before selecting destructive # candidates so the legacy branch never drops a trigger it won't recreate. legacy_layout = self._db_has_legacy_inline_fts(cursor) - update_names = ( - "messages_fts_update", - "messages_fts_trigram_update", - ) + update_names = ("messages_fts_update", "messages_fts_trigram_update") if not legacy_layout: update_names += ("messages_fts_cjk_update",) placeholders = ", ".join("?" for _ in update_names) @@ -212,43 +276,31 @@ class SessionSchemaMixin: f"WHERE type = 'trigger' AND name IN ({placeholders})", update_names, ).fetchall() - to_drop = [ - name for name, sql in rows - if self._fts_update_trigger_needs_narrowing(sql) - ] + to_drop = [name for name, sql in rows if self._fts_update_trigger_needs_narrowing(sql)] if not to_drop: return 0 - for name in to_drop: - # Names are drawn from the update_names literal allowlist above — - # never user input — so the identifier is interpolation-safe. + # Names come from the literal allowlist above — interpolation-safe. cursor.execute(f"DROP TRIGGER IF EXISTS {name}") - # Re-apply current DDL so CREATE TRIGGER installs the OF variants. - # Choose legacy vs v23 the same way _init_schema does. + # Re-apply current DDL (legacy vs v23 chosen as _init_schema does) so + # CREATE TRIGGER installs the OF variants. if legacy_layout: self._ensure_fts_schema(cursor, "messages_fts", LEGACY_FTS_SQL) - self._ensure_fts_schema( - cursor, "messages_fts_trigram", LEGACY_FTS_TRIGRAM_SQL - ) + self._ensure_fts_schema(cursor, "messages_fts_trigram", LEGACY_FTS_TRIGRAM_SQL) else: self._ensure_fts_schema(cursor, "messages_fts", FTS_SQL) - self._ensure_fts_schema( - cursor, "messages_fts_trigram", FTS_TRIGRAM_SQL - ) + self._ensure_fts_schema(cursor, "messages_fts_trigram", FTS_TRIGRAM_SQL) # Only recreate the CJK trigger this migration actually dropped. # ``_ensure_fts_cjk_schema`` soft-fails OperationalError by clearing - # availability (never raises), so after ensure require a narrowed - # CJK UPDATE trigger or durable quarantine (stale breadcrumb + - # unavailable). + # availability (never raises), so afterwards require a narrowed CJK + # UPDATE trigger or durable quarantine (stale breadcrumb + unavailable). if "messages_fts_cjk_update" in to_drop: try: self._ensure_fts_cjk_schema(cursor) except Exception: self._quarantine_cjk_after_update_of_migration(cursor) - logger.exception( - "CJK FTS re-ensure after UPDATE OF migration failed" - ) + logger.exception("CJK FTS re-ensure after UPDATE OF migration failed") raise if not self._cjk_update_trigger_is_narrowed(cursor): self._quarantine_cjk_after_update_of_migration(cursor) @@ -256,7 +308,6 @@ class SessionSchemaMixin: "CJK FTS UPDATE trigger missing or still broad after " "UPDATE OF migration; marked stale and unavailable" ) - logger.info( "Migrated %d broad FTS UPDATE trigger(s) to AFTER UPDATE OF " "(no rebuild required)", @@ -273,9 +324,7 @@ class SessionSchemaMixin: ).fetchone() return bool(row) and not self._fts_update_trigger_needs_narrowing(row[0]) - def _quarantine_cjk_after_update_of_migration( - self, cursor: sqlite3.Cursor - ) -> None: + def _quarantine_cjk_after_update_of_migration(self, cursor: sqlite3.Cursor) -> None: """Fail closed after dropping the CJK UPDATE trigger mid-migration: clear availability, persist ``fts_cjk_stale``, drop any residual CJK UPDATE trigger so a later open cannot IF-NOT-EXISTS over a gap.""" @@ -283,72 +332,50 @@ class SessionSchemaMixin: try: self.set_meta(FTS_CJK_STALE_KEY, "1", cursor=cursor) except Exception: - logger.debug( - "Could not persist CJK FTS stale breadcrumb", - exc_info=True, - ) + logger.debug("Could not persist CJK FTS stale breadcrumb", exc_info=True) try: cursor.execute("DROP TRIGGER IF EXISTS messages_fts_cjk_update") except Exception: - logger.debug( - "Could not drop residual CJK UPDATE trigger after quarantine", - exc_info=True, - ) - + logger.debug("Could not drop residual CJK UPDATE trigger after quarantine", exc_info=True) @staticmethod - def _rebuild_fts_indexes( - cursor: sqlite3.Cursor, - *, - include_trigram: bool = True, - ) -> None: - # v23+ external-content tables: 'rebuild' repopulates the inverted - # index from the content source (messages / messages_fts_trigram_src). + def _rebuild_fts_indexes(cursor: sqlite3.Cursor, *, include_trigram: bool = True) -> None: + """v23+ external-content tables: 'rebuild' repopulates the inverted + index from the content source (messages / messages_fts_trigram_src). + 'rebuild' indexes EVERY row, so the deferred-backfill markers are + cleared or the worker would re-insert covered rows (duplicates).""" cursor.execute("INSERT INTO messages_fts(messages_fts) VALUES('rebuild')") if include_trigram: - cursor.execute( - "INSERT INTO messages_fts_trigram(messages_fts_trigram) VALUES('rebuild')" - ) - # 'rebuild' indexes EVERY row: clear deferred-backfill markers or the - # worker would re-insert rows already covered (duplicates). + cursor.execute("INSERT INTO messages_fts_trigram(messages_fts_trigram) VALUES('rebuild')") cursor.execute( "DELETE FROM state_meta WHERE key IN " "('fts_rebuild_high_water', 'fts_rebuild_progress')" ) @staticmethod - def _rebuild_legacy_fts_indexes( - cursor: sqlite3.Cursor, - *, - include_trigram: bool = True, - ) -> None: + def _rebuild_legacy_fts_indexes(cursor: sqlite3.Cursor, *, include_trigram: bool = True) -> None: """Rebuild the LEGACY inline (pre-v23) FTS indexes from messages. - Inline tables have no external-content 'rebuild' source, so DELETE + reinsert the concatenated content the legacy triggers produced. - Never touches the v23 shape. - """ + Never touches the v23 shape.""" tables = ("messages_fts", "messages_fts_trigram") if include_trigram else ("messages_fts",) for tbl in tables: cursor.execute(f"DELETE FROM {tbl}") - cursor.execute( - f"INSERT INTO {tbl}(rowid, content) " - "SELECT id, " - "COALESCE(content, '') || ' ' || " - "COALESCE(tool_name, '') || ' ' || " - "COALESCE(tool_calls, '') " - "FROM messages" - ) + cursor.execute(f"INSERT INTO {tbl}(rowid, content) SELECT id, {_LEGACY_INLINE_CONCAT_SQL}FROM messages") def _fts_table_probe(self, cursor: sqlite3.Cursor, table_name: str) -> Optional[bool]: + """True = queryable, False = absent, None = FTS module/tokenizer missing + or content undecodable (index degraded, store accessible). + + Invalid UTF-8 in FTS content surfaces as a bare UnicodeDecodeError on + some builds and as OperationalError("Could not decode to UTF-8 ...") on + others; both are caught so the probe never raises into writable-init / + recovery flows. Anything else (malformed schema, corrupt vtable) re-raises. + """ try: cursor.execute(f"SELECT * FROM {table_name} LIMIT 0") return True except (sqlite3.OperationalError, UnicodeDecodeError) as exc: - # Invalid UTF-8 in FTS content surfaces as a bare UnicodeDecodeError - # (not sqlite3.Error) on some builds and as OperationalError("Could - # not decode to UTF-8 ...") on others; catch both so the probe never - # raises into writable-init/recovery flows. if isinstance(exc, sqlite3.OperationalError): if self._is_fts5_unavailable_error(exc): # A missing trigram tokenizer only affects trigram search; @@ -360,11 +387,8 @@ class SessionSchemaMixin: return None if "no such table" in str(exc).lower(): return False - # Anything else (malformed schema, corrupt vtable) re-raises. if "decode to utf-8" not in str(exc).lower(): raise - # Decode error: index degraded, store accessible; writable init / - # recovery schedules a rebuild or degrades to LIKE. logger.warning( "%s probe encountered invalid UTF-8 in FTS content; " "search may return incomplete results until FTS is rebuilt: %s", @@ -373,90 +397,91 @@ class SessionSchemaMixin: ) return None - def _recover_stale_fts( - self, cursor: sqlite3.Cursor, *, legacy: bool, timeout_seconds=None - ) -> bool: + # ── Stale-FTS recovery ───────────────────────────────────────────────── + + def _defer_stale_fts_for_holders(self, cursor: sqlite3.Cursor, foreign_holders) -> bool: + """Record a deferral diagnostic for the foreign processes holding the + DB and decide whether to defer; True = defer (holders remain). + + After ``_FTS_HOLDER_ESCALATE_ATTEMPTS`` deferrals spanning + ``_FTS_HOLDER_ESCALATE_SECONDS``, provably inactive orphan Desktop + backends are reaped and the holders re-checked. + """ + now = time.time() + record = None + try: + row = cursor.execute( + "SELECT value FROM state_meta WHERE key = ? LIMIT 1", + (FTS_REBUILD_DEFERRAL_KEY,), + ).fetchone() + if row: + parsed = json.loads(row[0]) + if isinstance(parsed, dict): + record = parsed + except (sqlite3.Error, TypeError, ValueError, json.JSONDecodeError): + record = None + try: + first_seen = float((record or {}).get("first_seen", now)) + attempts = int((record or {}).get("attempts", 0)) + 1 + except (TypeError, ValueError): + first_seen = now + attempts = 1 + if first_seen > now or first_seen < 0: + first_seen = now + diagnostic = { + "first_seen": first_seen, + "last_seen": now, + "attempts": attempts, + "holder_pids": sorted({pid for pid, _path in foreign_holders if pid > 0}), + } + cursor.execute( + "INSERT INTO state_meta (key, value) VALUES (?, ?) " + "ON CONFLICT(key) DO UPDATE SET value = excluded.value", + (FTS_REBUILD_DEFERRAL_KEY, json.dumps(diagnostic, sort_keys=True)), + ) + if attempts >= _FTS_HOLDER_ESCALATE_ATTEMPTS and now - first_seen >= _FTS_HOLDER_ESCALATE_SECONDS: + reaped = self._reap_inactive_orphan_desktop_holders( + foreign_holders, min_age_seconds=_FTS_HOLDER_ESCALATE_SECONDS, + ) + if reaped: + logger.error( + "Reaped inactive orphan Desktop backend(s) %s after %d " + "state.db FTS rebuild deferrals; checking holders again.", + reaped, + attempts, + ) + foreign_holders = self._foreign_state_db_holders() + if foreign_holders: + logger.error( + "state.db FTS repair remains blocked after %d deferrals " + "by holder(s) %s. Stop the listed processes, then run " + "`hermes sessions optimize-storage` with the gateway stopped. " + "`hermes doctor` reports this degraded state.", + attempts, + foreign_holders, + ) + if not foreign_holders: + return False + logger.warning( + "Deferred stale state.db FTS rebuild while foreign processes " + "hold the database or WAL sidecars (%s); canonical writes and " + "LIKE search remain available (deferral %d).", + foreign_holders, + attempts, + ) + return True + + def _recover_stale_fts(self, cursor: sqlite3.Cursor, *, legacy: bool, timeout_seconds=None) -> bool: """Atomically rebuild stale base/trigram indexes and resume syncing. *timeout_seconds* bounds the cross-process admission wait; None uses the full startup budget, ``0`` is the non-blocking in-process retry. + Fails closed: foreign holders or a lost admission race leave the + breadcrumb set and defer to a later retry. """ foreign_holders = self._foreign_state_db_holders() - if foreign_holders: - now = time.time() - record = None - try: - row = cursor.execute( - "SELECT value FROM state_meta WHERE key = ? LIMIT 1", - (FTS_REBUILD_DEFERRAL_KEY,), - ).fetchone() - if row: - parsed = json.loads(row[0]) - if isinstance(parsed, dict): - record = parsed - except (sqlite3.Error, TypeError, ValueError, json.JSONDecodeError): - record = None - - try: - first_seen = float((record or {}).get("first_seen", now)) - attempts = int((record or {}).get("attempts", 0)) + 1 - except (TypeError, ValueError): - first_seen = now - attempts = 1 - if first_seen > now or first_seen < 0: - first_seen = now - holder_pids = sorted({pid for pid, _path in foreign_holders if pid > 0}) - diagnostic = { - "first_seen": first_seen, - "last_seen": now, - "attempts": attempts, - "holder_pids": holder_pids, - } - cursor.execute( - "INSERT INTO state_meta (key, value) VALUES (?, ?) " - "ON CONFLICT(key) DO UPDATE SET value = excluded.value", - (FTS_REBUILD_DEFERRAL_KEY, json.dumps(diagnostic, sort_keys=True)), - ) - - escalated = ( - attempts >= _FTS_HOLDER_ESCALATE_ATTEMPTS - and now - first_seen >= _FTS_HOLDER_ESCALATE_SECONDS - ) - if escalated: - reaped = self._reap_inactive_orphan_desktop_holders( - foreign_holders, - min_age_seconds=_FTS_HOLDER_ESCALATE_SECONDS, - ) - if reaped: - logger.error( - "Reaped inactive orphan Desktop backend(s) %s after %d " - "state.db FTS rebuild deferrals; checking holders again.", - reaped, - attempts, - ) - foreign_holders = self._foreign_state_db_holders() - if foreign_holders: - logger.error( - "state.db FTS repair remains blocked after %d deferrals " - "by holder(s) %s. Stop the listed processes, then run " - "`hermes sessions optimize-storage` with the gateway stopped. " - "`hermes doctor` reports this degraded state.", - attempts, - foreign_holders, - ) - - if foreign_holders: - logger.warning( - "Deferred stale state.db FTS rebuild while foreign processes " - "hold the database or WAL sidecars (%s); canonical writes and " - "LIKE search remain available (deferral %d).", - foreign_holders, - attempts, - ) - return False - # Full structural rebuild: admit through the cross-process authority - # (fail closed). Losing the race means another process is doing this - # exact recovery; the breadcrumb stays set and we retry later. + if foreign_holders and self._defer_stale_fts_for_holders(cursor, foreign_holders): + return False with fts_rebuild_admission(self.db_path, timeout_seconds=timeout_seconds) as admitted: if not admitted: logger.warning( @@ -491,8 +516,7 @@ class SessionSchemaMixin: interval = _FTS_STALE_RETRY_SECONDS self._fts_stale_retry_after = now + interval self._fts_stale_retry_interval = min( - max(interval, _FTS_STALE_RETRY_SECONDS, 1.0) * 2.0, - _FTS_STALE_RETRY_MAX_SECONDS, + max(interval, _FTS_STALE_RETRY_SECONDS, 1.0) * 2.0, _FTS_STALE_RETRY_MAX_SECONDS, ) try: with self._lock: @@ -500,9 +524,7 @@ class SessionSchemaMixin: return False cursor = self._conn.cursor() legacy = self._db_has_legacy_inline_fts(cursor) - recovered = self._recover_stale_fts( - cursor, legacy=legacy, timeout_seconds=0.0 - ) + recovered = self._recover_stale_fts(cursor, legacy=legacy, timeout_seconds=0.0) if recovered: # CJK was detached alongside the base indexes; its own # ensure path decides when it comes back online. @@ -521,10 +543,12 @@ class SessionSchemaMixin: ) return False - def _recover_stale_fts_locked( - self, cursor: sqlite3.Cursor, *, legacy: bool - ) -> bool: - """Body of :meth:`_recover_stale_fts`; caller holds rebuild authority.""" + def _recover_stale_fts_locked(self, cursor: sqlite3.Cursor, *, legacy: bool) -> bool: + """Body of :meth:`_recover_stale_fts`; caller holds rebuild authority. + + One write transaction closes the dangerous gap: no canonical writer + can slip between the full rebuild and trigger restoration. + """ try: trigram_status = self._fts_table_probe(cursor, "messages_fts_trigram") except (sqlite3.DatabaseError, UnicodeDecodeError): @@ -533,9 +557,7 @@ class SessionSchemaMixin: trigram_status = True include_trigram = trigram_status is True - drop_sql = "".join( - f"DROP TRIGGER IF EXISTS {trigger};" for trigger in _FTS_TRIGGERS - ) + 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;" @@ -567,21 +589,11 @@ class SessionSchemaMixin: schema_sql = FTS_SQL if include_trigram: schema_sql += FTS_TRIGRAM_SQL - rebuild_sql = schema_sql + ( - "INSERT INTO messages_fts(messages_fts) VALUES('rebuild');" - ) + 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');" - ) + 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 @@ -617,6 +629,8 @@ class SessionSchemaMixin: ) return True + # ── Declarative column reconciliation ────────────────────────────────── + @staticmethod def _parse_schema_columns(schema_sql: str) -> Dict[str, Dict[str, str]]: """Expected columns per table, parsed from SCHEMA_SQL. @@ -642,8 +656,7 @@ class SessionSchemaMixin: ): tables = blob["tables"] if all( - isinstance(cols, dict) - and all(isinstance(v, str) for v in cols.values()) + isinstance(cols, dict) and all(isinstance(v, str) for v in cols.values()) for cols in tables.values() ): return tables @@ -676,9 +689,7 @@ class SessionSchemaMixin: if cache_path is not None: try: cache_path.parent.mkdir(parents=True, exist_ok=True) - fd, tmp = tempfile.mkstemp( - dir=str(cache_path.parent), prefix=".schema_columns." - ) + fd, tmp = tempfile.mkstemp(dir=str(cache_path.parent), prefix=".schema_columns.") with os.fdopen(fd, "w", encoding="utf-8") as fh: json.dump({"schema_hash": schema_hash, "tables": table_columns}, fh) os.replace(tmp, cache_path) @@ -695,41 +706,33 @@ class SessionSchemaMixin: expected = self._parse_schema_columns(SCHEMA_SQL) for table_name, declared_cols in expected.items(): try: - rows = cursor.execute( - f'PRAGMA table_info("{table_name}")' - ).fetchall() + rows = cursor.execute(f'PRAGMA table_info("{table_name}")').fetchall() except sqlite3.OperationalError: continue # Table doesn't exist yet (shouldn't happen after executescript) # PRAGMA table_info rows: (cid, name, type, notnull, dflt_value, pk) live_cols = {row[1] for row in rows} - for col_name, col_type in declared_cols.items(): - if col_name not in live_cols: - safe_name = col_name.replace('"', '""') - try: - cursor.execute( - f'ALTER TABLE "{table_name}" ADD COLUMN "{safe_name}" {col_type}' - ) - except sqlite3.OperationalError as exc: - message = str(exc).lower() - if "duplicate column" in message: - # A sibling process won the ADD race; store is correct. - logger.debug( - "reconcile %s.%s: %s", table_name, col_name, exc, - ) - continue - if "locked" in message or "busy" in message: - # Lock contention: swallowing it left the store - # half-reconciled ("no such column" on every read). - # Re-raise so _connect_and_init_with_lock_patience - # retries the WHOLE init (idempotent) with backoff. - raise - # Anything else is a schema mistake that permanently - # strands the store behind SCHEMA_SQL — be loud. - logger.warning( - "reconcile %s.%s failed; store remains behind " - "SCHEMA_SQL: %s", table_name, col_name, exc, - ) + if col_name in live_cols: + continue + try: + cursor.execute(f'ALTER TABLE "{table_name}" ADD COLUMN {_q(col_name)} {col_type}') + except sqlite3.OperationalError as exc: + message = str(exc).lower() + if "duplicate column" in message: + # A sibling process won the ADD race; store is correct. + logger.debug("reconcile %s.%s: %s", table_name, col_name, exc) + continue + if "locked" in message or "busy" in message: + # Swallowing lock contention left the store half-reconciled + # ("no such column" on every read). Re-raise so + # _connect_and_init_with_lock_patience retries the WHOLE + # init (idempotent) with backoff. + raise + # Anything else permanently strands the store behind SCHEMA_SQL — be loud. + logger.warning( + "reconcile %s.%s failed; store remains behind " + "SCHEMA_SQL: %s", table_name, col_name, exc, + ) @staticmethod def _live_pk_columns(cursor: sqlite3.Cursor, table: str) -> Optional[List[str]]: @@ -757,15 +760,12 @@ class SessionSchemaMixin: pk_cols = self._live_pk_columns(cursor, "gateway_routing") if pk_cols is None or pk_cols == ["scope", "session_key"]: return - logger.info( "gateway_routing has legacy primary key %r; rebuilding with " "composite (scope, session_key) key", pk_cols, ) - cursor.execute( - "ALTER TABLE gateway_routing RENAME TO gateway_routing_legacy_pk" - ) + cursor.execute("ALTER TABLE gateway_routing RENAME TO gateway_routing_legacy_pk") cursor.execute( """CREATE TABLE gateway_routing ( scope TEXT NOT NULL DEFAULT '', @@ -793,50 +793,25 @@ class SessionSchemaMixin: Every ``_record_model_usage()`` upsert then fails (ON CONFLICT mismatch), aborting the write transaction and silently zeroing token and cost accounting. Idempotent; no-op on healthy databases. + + FK-off window: INSERT OR IGNORE does NOT suppress foreign-key + violations, so an orphaned usage row (partial prune while accounting + was broken) would abort the whole rebuild. PRAGMA foreign_keys is a + no-op inside a transaction — fine here, _init_schema runs on an + isolation_level=None connection with no transaction open. """ pk_cols = self._live_pk_columns(cursor, "session_model_usage") if pk_cols is None or "task" in pk_cols: return - logger.info( "session_model_usage has legacy primary key %r (missing task); " "rebuilding with composite 6-column key", sorted(pk_cols), ) - # FK-off window: INSERT OR IGNORE does NOT suppress foreign-key - # violations, so an orphaned usage row (partial prune while accounting - # was broken) would abort the whole rebuild. PRAGMA foreign_keys is a - # no-op inside a transaction — fine here, _init_schema runs on an - # isolation_level=None connection with no transaction open. cursor.execute("PRAGMA foreign_keys=OFF") try: - cursor.execute( - "ALTER TABLE session_model_usage " - "RENAME TO session_model_usage_legacy_pk" - ) - cursor.execute( - """CREATE TABLE session_model_usage ( - session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE, - model TEXT NOT NULL, - billing_provider TEXT NOT NULL DEFAULT '', - billing_base_url TEXT NOT NULL DEFAULT '', - billing_mode TEXT NOT NULL DEFAULT '', - task TEXT NOT NULL DEFAULT '', - api_call_count INTEGER NOT NULL DEFAULT 0, - input_tokens INTEGER NOT NULL DEFAULT 0, - output_tokens INTEGER NOT NULL DEFAULT 0, - cache_read_tokens INTEGER NOT NULL DEFAULT 0, - cache_write_tokens INTEGER NOT NULL DEFAULT 0, - reasoning_tokens INTEGER NOT NULL DEFAULT 0, - estimated_cost_usd REAL NOT NULL DEFAULT 0, - actual_cost_usd REAL NOT NULL DEFAULT 0, - cost_status TEXT, - cost_source TEXT, - first_seen REAL, - last_seen REAL, - PRIMARY KEY (session_id, model, billing_provider, billing_base_url, billing_mode, task) -)""" - ) + cursor.execute("ALTER TABLE session_model_usage RENAME TO session_model_usage_legacy_pk") + cursor.execute(_SESSION_MODEL_USAGE_HEAL_DDL) # OR IGNORE: COALESCE(task, '') on legacy NULL rows can collide # with a genuine ''-task row — keep the first rather than fail. cursor.execute( @@ -859,19 +834,15 @@ class SessionSchemaMixin: FROM session_model_usage_legacy_pk""" ) cursor.execute("DROP TABLE session_model_usage_legacy_pk") - cursor.execute( - "CREATE INDEX IF NOT EXISTS idx_session_model_usage_session " - "ON session_model_usage(session_id)" - ) - cursor.execute( - "CREATE INDEX IF NOT EXISTS idx_session_model_usage_model " - "ON session_model_usage(model)" - ) + for sql in _SESSION_MODEL_USAGE_INDEX_SQL: + cursor.execute(sql) except sqlite3.OperationalError as exc: logger.debug("session_model_usage PK heal skipped: %s", exc) finally: cursor.execute("PRAGMA foreign_keys=ON") + # ── _init_schema ─────────────────────────────────────────────────────── + def _init_schema(self): """Create tables and FTS if missing, reconcile columns, run data migrations. @@ -887,9 +858,7 @@ class SessionSchemaMixin: # _MAX_LEASE_S=900): a genuinely wedged init delays supervisor respawn # by up to the lease; per-chunk renewal isn't worth the complexity. report_startup_progress(600.0, phase="state_db_init_schema") - cursor = self._conn.cursor() - cursor.executescript(SCHEMA_SQL) # Idempotent, self-healing column reconciliation, then the two @@ -909,9 +878,7 @@ class SessionSchemaMixin: ) except sqlite3.OperationalError as exc: logger.debug("idx_messages_platform_msg_id create skipped: %s", exc) - - # Same ordering constraint (idx_messages_session_active on ``active``). - cursor.executescript(DEFERRED_INDEX_SQL) + cursor.executescript(DEFERRED_INDEX_SQL) # same ordering constraint (``active``) # Heal NULL ``active`` rows on every startup: older reconciler builds # added ``active`` without its NOT NULL DEFAULT 1, so INSERTs omitting @@ -919,14 +886,11 @@ class SessionSchemaMixin: # Unconditional because a ``current_version < 12`` gate never re-ran # for already-v12+ databases. try: - cursor.execute( - "UPDATE messages SET active = 1 WHERE active IS NULL" - ) + cursor.execute("UPDATE messages SET active = 1 WHERE active IS NULL") except sqlite3.OperationalError: pass 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,), @@ -944,161 +908,149 @@ class SessionSchemaMixin: row = cursor.execute("SELECT version FROM schema_version LIMIT 1").fetchone() if row is None: - cursor.execute( - "INSERT INTO schema_version (version) VALUES (?)", - (SCHEMA_VERSION,), - ) + cursor.execute("INSERT INTO schema_version (version) VALUES (?)", (SCHEMA_VERSION,)) # Store provenance so fresh vs wiped stores are distinguishable. now_iso = datetime.datetime.now(datetime.timezone.utc).isoformat() instance_id = str(uuid.uuid4()) cursor.executemany( "INSERT OR IGNORE INTO state_meta (key, value) VALUES (?, ?)", - [ - ("store_instance_id", instance_id), - ("store_created_at_utc", now_iso), - ], + [("store_instance_id", instance_id), ("store_created_at_utc", now_iso)], ) - else: - current_version = row[0] - # Renew the lease: the version-gated chain can rewrite whole tables - # on large DBs (same single-lease trade-off as above). - report_startup_progress(600.0, phase="state_db_data_migrations") - # Version-gated chain for data migrations only (row backfills, - # version-specific index changes); column additions never belong here. - if current_version < 10 and SCHEMA_VERSION == 10: - # v10: one-time trigram backfill. Only when v10 itself is the - # target: v11+ drops and rebuilds both FTS tables, so the - # backfill would only burn startup time and WAL space. - if fts5_available: - _fts_trigram_exists = self._fts_table_probe( - cursor, "messages_fts_trigram" - ) - if _fts_trigram_exists is False: - if self._ensure_fts_schema( - cursor, "messages_fts_trigram", FTS_TRIGRAM_SQL - ): - cursor.execute( - "INSERT INTO messages_fts_trigram(rowid, content) " - "SELECT id, content FROM messages WHERE content IS NOT NULL" - ) - else: - fts_migrations_complete = False - elif _fts_trigram_exists is None: - fts_migrations_complete = False - else: - fts_migrations_complete = False - # (v11 inline FTS re-index was superseded by v23 and removed.) - if current_version < 16: - # v16: tag delegate subagent rows so pickers stay clean after - # parent deletes orphan them. The shared predicate excludes - # user-visible reset children. - try: - cursor.execute( - "UPDATE sessions SET model_config = json_set(" - "COALESCE(model_config, '{}'), '$._delegate_from', parent_session_id) " - f"WHERE parent_session_id IS NOT NULL " - "AND json_extract(COALESCE(model_config, '{}'), '$._delegate_from') IS NULL " - f"AND {_ephemeral_child_sql('sessions')}" - ) - cursor.execute( - "UPDATE sessions SET model_config = json_set(" - "COALESCE(model_config, '{}'), '$._delegate_from', '__orphaned__') " - "WHERE parent_session_id IS NULL " - "AND json_extract(COALESCE(model_config, '{}'), '$._delegate_from') IS NULL " - "AND json_extract(COALESCE(model_config, '{}'), '$._branched_from') IS NULL " - "AND title IS NULL " - "AND message_count <= 25 " - "AND EXISTS (SELECT 1 FROM messages m " - " WHERE m.session_id = sessions.id AND m.role = 'tool') " - "AND NOT EXISTS (SELECT 1 FROM sessions ch " - " WHERE ch.parent_session_id = sessions.id)" - ) - except sqlite3.OperationalError: - pass - if current_version < 18: - # v18: backfill gateway metadata from sessions.json. Best-effort: - # missing metadata just means consumers fall back to - # sessions.json until the gateway rewrites those rows. - try: - self._backfill_gateway_metadata_from_sessions_json(cursor) - except Exception as exc: - logger.debug("v18 gateway metadata backfill skipped: %s", exc) - if current_version < 20: - # v20: seed one session_model_usage row per historical session - # from the sessions aggregates. INSERT OR IGNORE: a row newer - # code already wrote wins over the stale aggregate. - try: - cursor.execute( - """INSERT OR IGNORE INTO session_model_usage ( - session_id, model, billing_provider, - billing_base_url, billing_mode, - api_call_count, input_tokens, - output_tokens, cache_read_tokens, - cache_write_tokens, reasoning_tokens, - estimated_cost_usd, actual_cost_usd, - cost_status, cost_source, first_seen, last_seen - ) - SELECT id, COALESCE(model, 'unknown'), - COALESCE(billing_provider, ''), - COALESCE(billing_base_url, ''), - COALESCE(billing_mode, ''), - COALESCE(api_call_count, 0), - COALESCE(input_tokens, 0), - COALESCE(output_tokens, 0), - COALESCE(cache_read_tokens, 0), - COALESCE(cache_write_tokens, 0), - COALESCE(reasoning_tokens, 0), - COALESCE(estimated_cost_usd, 0), - COALESCE(actual_cost_usd, 0), - cost_status, cost_source, - started_at, COALESCE(ended_at, started_at) - FROM sessions - WHERE COALESCE(input_tokens, 0) - + COALESCE(output_tokens, 0) - + COALESCE(cache_read_tokens, 0) - + COALESCE(cache_write_tokens, 0) - + COALESCE(reasoning_tokens, 0) > 0""" - ) - except sqlite3.OperationalError: - pass - if current_version < 22: - # v22: ``task`` joins the session_model_usage PRIMARY KEY ('' = - # main loop; 'vision'/'compression'/... = aux calls). SQLite - # cannot ALTER a PK, so rebuild; existing rows are main-loop - # accounting → task=''. - try: - legacy_pk = cursor.execute( - "SELECT COUNT(*) FROM pragma_table_info('session_model_usage') " - "WHERE name = 'task' AND pk > 0" - ).fetchone()[0] - if not legacy_pk: - cursor.execute("ALTER TABLE session_model_usage RENAME TO session_model_usage_v21") + self._run_data_migrations(cursor, row[0], fts5_available) + + self._ensure_unique_title_index(cursor) + if fts5_available: + self._init_fts(cursor) + self._conn.commit() + + def _run_data_migrations(self, cursor: sqlite3.Cursor, current_version: int, fts5_available: bool) -> None: + """Version-gated chain for DATA migrations only (row backfills, + version-specific index changes); column additions never belong here. + Advances schema_version at the end unless FTS work could not complete.""" + # Renew the lease: the chain can rewrite whole tables on large DBs. + report_startup_progress(600.0, phase="state_db_data_migrations") + fts_migrations_complete = True + if current_version < 10 and SCHEMA_VERSION == 10: + # v10: one-time trigram backfill. Only when v10 itself is the + # target: v11+ drops and rebuilds both FTS tables, so the + # backfill would only burn startup time and WAL space. + if fts5_available: + _fts_trigram_exists = self._fts_table_probe(cursor, "messages_fts_trigram") + if _fts_trigram_exists is False: + if self._ensure_fts_schema(cursor, "messages_fts_trigram", FTS_TRIGRAM_SQL): cursor.execute( - """CREATE TABLE session_model_usage ( - session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE, - model TEXT NOT NULL, - billing_provider TEXT NOT NULL DEFAULT '', - billing_base_url TEXT NOT NULL DEFAULT '', - billing_mode TEXT NOT NULL DEFAULT '', - task TEXT NOT NULL DEFAULT '', - api_call_count INTEGER NOT NULL DEFAULT 0, - input_tokens INTEGER NOT NULL DEFAULT 0, - output_tokens INTEGER NOT NULL DEFAULT 0, - cache_read_tokens INTEGER NOT NULL DEFAULT 0, - cache_write_tokens INTEGER NOT NULL DEFAULT 0, - reasoning_tokens INTEGER NOT NULL DEFAULT 0, - estimated_cost_usd REAL NOT NULL DEFAULT 0, - actual_cost_usd REAL NOT NULL DEFAULT 0, - cost_status TEXT, - cost_source TEXT, - first_seen REAL, - last_seen REAL, - PRIMARY KEY (session_id, model, billing_provider, billing_base_url, billing_mode, task) - )""" + "INSERT INTO messages_fts_trigram(rowid, content) " + "SELECT id, content FROM messages WHERE content IS NOT NULL" ) - cursor.execute( - """INSERT INTO session_model_usage ( + else: + fts_migrations_complete = False + elif _fts_trigram_exists is None: + fts_migrations_complete = False + else: + fts_migrations_complete = False + # (v11 inline FTS re-index was superseded by v23 and removed.) + if current_version < 16: + # v16: tag delegate subagent rows so pickers stay clean after + # parent deletes orphan them. The shared predicate excludes + # user-visible reset children. + try: + cursor.execute( + "UPDATE sessions SET model_config = json_set(" + "COALESCE(model_config, '{}'), '$._delegate_from', parent_session_id) " + f"WHERE parent_session_id IS NOT NULL " + "AND json_extract(COALESCE(model_config, '{}'), '$._delegate_from') IS NULL " + f"AND {_ephemeral_child_sql('sessions')}" + ) + cursor.execute( + "UPDATE sessions SET model_config = json_set(" + "COALESCE(model_config, '{}'), '$._delegate_from', '__orphaned__') " + "WHERE parent_session_id IS NULL " + "AND json_extract(COALESCE(model_config, '{}'), '$._delegate_from') IS NULL " + "AND json_extract(COALESCE(model_config, '{}'), '$._branched_from') IS NULL " + "AND title IS NULL " + "AND message_count <= 25 " + "AND EXISTS (SELECT 1 FROM messages m " + " WHERE m.session_id = sessions.id AND m.role = 'tool') " + "AND NOT EXISTS (SELECT 1 FROM sessions ch " + " WHERE ch.parent_session_id = sessions.id)" + ) + except sqlite3.OperationalError: + pass + if current_version < 18: + # v18: backfill gateway metadata from sessions.json. Best-effort: + # consumers fall back to sessions.json until the gateway rewrites. + try: + self._backfill_gateway_metadata_from_sessions_json(cursor) + except Exception as exc: + logger.debug("v18 gateway metadata backfill skipped: %s", exc) + if current_version < 20: + # v20: seed one session_model_usage row per historical session + # from the sessions aggregates. INSERT OR IGNORE: a row newer + # code already wrote wins over the stale aggregate. + try: + cursor.execute(_SESSION_MODEL_USAGE_V20_SEED_SQL) + except sqlite3.OperationalError: + pass + if current_version < 22: + self._migrate_v22_session_model_usage(cursor) + # v23: FTS storage redesign (external-content tables; inline v11 + # tables were ~75% of state.db on heavy installs). OPT-IN, NOT + # AUTOMATIC: the transition is disk-heavy (~2x transient) and long + # (hours on a 25 GB DB), so an existing install only gets a flag + # advertising it; `hermes sessions optimize-storage` performs it as + # a deliberate foreground operation. DECOUPLED VERSIONING: the FTS + # layout is tracked by the independent `fts_storage_version` + # marker, so schema_version still advances here and future + # migrations land for legacy-FTS users too. + if current_version < 23 and fts5_available and self._db_has_legacy_inline_fts(cursor): + self.set_meta("fts_optimize_available", "1", cursor=cursor) + if current_version < 25: + # v25: de-duplicate system prompt snapshots into the shared + # content-addressed table; the old column stays a read fallback + # for partially migrated or externally written rows. + self._dedupe_legacy_system_prompts(cursor) + + # Stamp the FTS layout version (fresh/optimized DBs) so the main + # version can always advance; a legacy DB keeps its absent/0 marker + # until optimize-storage runs. An INTERRUPTED optimize (rebuild + # markers, trash tables, or an empty external index against + # non-empty messages) is NOT stamped: the marker is the source of + # truth for "fully optimized" and keeps the resume offer alive. + if ( + fts5_available + and not self._db_has_legacy_inline_fts(cursor) + and cursor.execute( + "SELECT 1 FROM state_meta " + "WHERE key = 'fts_rebuild_high_water' LIMIT 1" + ).fetchone() is None + and not self._has_fts_trash(cursor) + and not self._fts_external_index_empty_with_messages(cursor) + ): + self.set_meta("fts_storage_version", str(FTS_STORAGE_VERSION), cursor=cursor) + + # Advance schema_version — deliberately NOT gated on the FTS opt-in + # (that would block every future migration for a user who never + # optimizes). FTS5 unavailable is the one skip: we can't have + # created the current FTS objects, so claiming current would lie. + if current_version < SCHEMA_VERSION and fts_migrations_complete and fts5_available: + cursor.execute("UPDATE schema_version SET version = ?", (SCHEMA_VERSION,)) + + def _migrate_v22_session_model_usage(self, cursor: sqlite3.Cursor) -> None: + """v22: ``task`` joins the session_model_usage PRIMARY KEY ('' = main + loop; 'vision'/'compression'/... = aux calls). SQLite cannot ALTER a + PK, so rebuild; existing rows are main-loop accounting → task=''.""" + try: + legacy_pk = cursor.execute( + "SELECT COUNT(*) FROM pragma_table_info('session_model_usage') " + "WHERE name = 'task' AND pk > 0" + ).fetchone()[0] + if legacy_pk: + return + cursor.execute("ALTER TABLE session_model_usage RENAME TO session_model_usage_v21") + cursor.execute(_SESSION_MODEL_USAGE_V22_DDL) + cursor.execute( + """INSERT INTO session_model_usage ( session_id, model, billing_provider, billing_base_url, billing_mode, task, api_call_count, input_tokens, output_tokens, cache_read_tokens, cache_write_tokens, @@ -1111,84 +1063,20 @@ class SessionSchemaMixin: reasoning_tokens, estimated_cost_usd, actual_cost_usd, cost_status, cost_source, first_seen, last_seen FROM session_model_usage_v21""" - ) - cursor.execute("DROP TABLE session_model_usage_v21") - cursor.execute( - "CREATE INDEX IF NOT EXISTS idx_session_model_usage_session " - "ON session_model_usage(session_id)" - ) - cursor.execute( - "CREATE INDEX IF NOT EXISTS idx_session_model_usage_model " - "ON session_model_usage(model)" - ) - except sqlite3.OperationalError as exc: - logger.debug("v22 session_model_usage rebuild skipped: %s", exc) - # v23: FTS storage redesign (external-content tables; inline v11 - # tables were ~75% of state.db on heavy installs). OPT-IN, NOT - # AUTOMATIC: the transition is disk-heavy (~2x transient) and long - # (hours on a 25 GB DB), so an existing install only gets a flag - # advertising it; `hermes sessions optimize-storage` performs it as - # a deliberate foreground operation. DECOUPLED VERSIONING: the FTS - # layout is tracked by the independent `fts_storage_version` - # marker, so schema_version still advances here and future - # migrations land for legacy-FTS users too. - if ( - current_version < 23 - and fts5_available - and self._db_has_legacy_inline_fts(cursor) - ): - self.set_meta("fts_optimize_available", "1", cursor=cursor) + ) + cursor.execute("DROP TABLE session_model_usage_v21") + for sql in _SESSION_MODEL_USAGE_INDEX_SQL: + cursor.execute(sql) + except sqlite3.OperationalError as exc: + logger.debug("v22 session_model_usage rebuild skipped: %s", exc) - if current_version < 25: - # v25: de-duplicate system prompt snapshots into the shared - # content-addressed table; the old column stays a read fallback - # for partially migrated or externally written rows. - self._dedupe_legacy_system_prompts(cursor) - - # Stamp the FTS layout version (fresh/optimized DBs) so the main - # version can always advance; a legacy DB keeps its absent/0 marker - # until optimize-storage runs. An INTERRUPTED optimize (rebuild - # markers, trash tables, or an empty external index against - # non-empty messages) is NOT stamped: the marker is the source of - # truth for "fully optimized" and keeps the resume offer alive. - if ( - fts5_available - and not self._db_has_legacy_inline_fts(cursor) - and cursor.execute( - "SELECT 1 FROM state_meta " - "WHERE key = 'fts_rebuild_high_water' LIMIT 1" - ).fetchone() is None - and not self._has_fts_trash(cursor) - and not self._fts_external_index_empty_with_messages(cursor) - ): - self.set_meta( - "fts_storage_version", str(FTS_STORAGE_VERSION), cursor=cursor - ) - - # Advance schema_version — deliberately NOT gated on the FTS opt-in - # (that would block every future migration for a user who never - # optimizes). FTS5 unavailable is the one skip: we can't have - # created the current FTS objects, so claiming current would lie. - if ( - current_version < SCHEMA_VERSION - and fts_migrations_complete - and fts5_available - ): - cursor.execute( - "UPDATE schema_version SET version = ?", - (SCHEMA_VERSION,), - ) - - # Unique title index. Older DBs may hold duplicate aliases from before - # the constraint; keep every session, the newest retains the alias. - title_index_sql = ( - "CREATE UNIQUE INDEX IF NOT EXISTS idx_sessions_title_unique " - "ON sessions(title) WHERE title IS NOT NULL" - ) + def _ensure_unique_title_index(self, cursor: sqlite3.Cursor) -> None: + """Unique title index. Older DBs may hold duplicate aliases from before + the constraint; keep every session, the newest retains the alias. The + index must never abort opening the DB, so the repair is guarded too.""" try: - cursor.execute(title_index_sql) + cursor.execute(_TITLE_UNIQUE_INDEX_SQL) except sqlite3.IntegrityError: - # The index must never abort opening the DB; guard the repair too. try: cursor.execute( """UPDATE sessions AS older @@ -1204,79 +1092,61 @@ class SessionSchemaMixin: "Cleared %d duplicate session title(s) while restoring the unique index", cursor.rowcount, ) - cursor.execute(title_index_sql) + cursor.execute(_TITLE_UNIQUE_INDEX_SQL) except sqlite3.Error: - logger.exception( - "Could not repair duplicate session titles; " - "unique title index not created" - ) + logger.exception("Could not repair duplicate session titles; unique title index not created") except sqlite3.OperationalError: pass # Index already exists - if fts5_available: - # Run the FTS DDL even when the vtable exists so CREATE TRIGGER IF - # NOT EXISTS repairs trigger-only degradation from a no-FTS5 runtime. - # OPT-IN v23 boundary: a legacy v22 inline install must keep its - # inline schema + triggers — the v23 external-content DDL would - # create the trigram source VIEW and leave a mixed state — so it - # gets the legacy DDL only; fresh/opted-in DBs get v23. - 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 base indexes and has its - # own stale marker; its ensure path decides when it returns. - self._ensure_fts_cjk_schema(cursor) - else: - self._fts_enabled = False - self._trigram_available = False - self._fts_cjk_available = False + def _init_fts(self, cursor: sqlite3.Cursor) -> None: + """Create/repair the FTS objects on an FTS5-capable runtime. + + The DDL runs even when the vtable exists so CREATE TRIGGER IF NOT + EXISTS repairs trigger-only degradation from a no-FTS5 runtime. + OPT-IN v23 boundary: a legacy v22 inline install must keep its inline + schema + triggers (the v23 external-content DDL would create the + trigram source VIEW and leave a mixed state), so it gets the legacy + DDL only; fresh/opted-in DBs get v23. + """ + 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 base indexes and has its own + # stale marker; its ensure path decides when it returns. + self._ensure_fts_cjk_schema(cursor) else: - base_sql, trigram_sql, rebuild = ( - (LEGACY_FTS_SQL, LEGACY_FTS_TRIGRAM_SQL, self._rebuild_legacy_fts_indexes) - if legacy_fts - else (FTS_SQL, FTS_TRIGRAM_SQL, self._rebuild_fts_indexes) - ) - # Measure BEFORE the DDL below runs, so these describe the - # pre-repair state. Whether the trigram half is even - # creatable is only known AFTER _ensure_fts_schema, which is - # why the two halves are combined at the `if`, not here. - base_triggers_missing = ( - self._fts_trigger_count(cursor, _FTS_BASE_TRIGGERS) - < len(_FTS_BASE_TRIGGERS) - ) - trigram_triggers_missing = ( - self._fts_trigger_count(cursor, _FTS_TRIGRAM_TRIGGERS) - < len(_FTS_TRIGRAM_TRIGGERS) - ) - self._fts_enabled = self._ensure_fts_schema( - cursor, "messages_fts", base_sql - ) - if self._fts_enabled: - # Trigram FTS5 for CJK/substring search is optional - # relative to the main table; if it cannot be created, - # CJK search falls back to LIKE. - trigram_enabled = self._ensure_fts_schema( - cursor, "messages_fts_trigram", trigram_sql - ) - self._trigram_available = trigram_enabled - if base_triggers_missing or ( - trigram_enabled and trigram_triggers_missing - ): - self._run_admitted_startup_rebuild( - cursor, - lambda: rebuild(cursor, include_trigram=trigram_enabled), - ) - if not legacy_fts: - # CJK-bigram index (cjk_unicode61): strictly additive - # and gated on the loadable tokenizer. - self._ensure_fts_cjk_schema(cursor) - - # Replace any pre-existing broad AFTER UPDATE triggers with - # AFTER UPDATE OF variants. IF NOT EXISTS cannot rewrite them. + self._fts_enabled = False + self._trigram_available = False + self._fts_cjk_available = False + else: + base_sql, trigram_sql, rebuild = ( + (LEGACY_FTS_SQL, LEGACY_FTS_TRIGRAM_SQL, self._rebuild_legacy_fts_indexes) + if legacy_fts + else (FTS_SQL, FTS_TRIGRAM_SQL, self._rebuild_fts_indexes) + ) + # Measure BEFORE the DDL below runs (pre-repair state). Whether the + # trigram half is even creatable is only known AFTER + # _ensure_fts_schema, which is why the halves combine at the `if`. + base_triggers_missing = self._fts_trigger_count(cursor, _FTS_BASE_TRIGGERS) < len(_FTS_BASE_TRIGGERS) + trigram_triggers_missing = ( + self._fts_trigger_count(cursor, _FTS_TRIGRAM_TRIGGERS) < len(_FTS_TRIGRAM_TRIGGERS) + ) + self._fts_enabled = self._ensure_fts_schema(cursor, "messages_fts", base_sql) if self._fts_enabled: - self._migrate_broad_fts_update_triggers(cursor) - - self._conn.commit() + # Trigram is optional relative to the main table; without it + # CJK search falls back to LIKE. + trigram_enabled = self._ensure_fts_schema(cursor, "messages_fts_trigram", trigram_sql) + self._trigram_available = trigram_enabled + if base_triggers_missing or (trigram_enabled and trigram_triggers_missing): + self._run_admitted_startup_rebuild( + cursor, lambda: rebuild(cursor, include_trigram=trigram_enabled), + ) + if not legacy_fts: + # CJK-bigram index: strictly additive, gated on the loadable tokenizer. + self._ensure_fts_cjk_schema(cursor) + # IF NOT EXISTS cannot rewrite pre-existing broad AFTER UPDATE triggers. + if self._fts_enabled: + self._migrate_broad_fts_update_triggers(cursor) def _run_admitted_startup_rebuild(self, cursor, rebuild_fn) -> None: """Run a full trigger-repair FTS rebuild under cross-process admission. @@ -1312,9 +1182,7 @@ class SessionSchemaMixin: self._trigram_available = False self._fts_cjk_available = False - def _backfill_gateway_metadata_from_sessions_json( - self, cursor: sqlite3.Cursor - ) -> None: + def _backfill_gateway_metadata_from_sessions_json(self, cursor: sqlite3.Cursor) -> None: """One-time v18 backfill of gateway metadata from sessions.json. Only fills NULL columns — never overwrites data written by newer code.""" sessions_file = get_hermes_home() / "sessions" / "sessions.json" @@ -1331,6 +1199,7 @@ class SessionSchemaMixin: if not session_id: continue origin = entry.get("origin") + origin_dict = origin if isinstance(origin, dict) else None cursor.execute( """UPDATE sessions SET session_key = COALESCE(session_key, ?), @@ -1346,11 +1215,11 @@ class SessionSchemaMixin: WHERE id = ?""", ( entry.get("session_key") or key, - (origin or {}).get("chat_id") if isinstance(origin, dict) else None, + origin_dict.get("chat_id") if origin_dict is not None else None, entry.get("chat_type"), - (origin or {}).get("thread_id") if isinstance(origin, dict) else None, + origin_dict.get("thread_id") if origin_dict is not None else None, entry.get("display_name"), - json.dumps(origin) if isinstance(origin, dict) else None, + json.dumps(origin) if origin_dict is not None else None, 1 if entry.get("expiry_finalized") or entry.get("memory_flushed") else 0, str(session_id), ),