refactor(hermes_state_schema): shared trigger-missing predicate and marker SQL constant, compact docstrings

This commit is contained in:
Teknium
2026-09-02 18:19:11 -07:00
parent 58accf5782
commit e9dfc808ee

View File

@@ -1,9 +1,7 @@
"""Schema creation, column reconciliation, and FTS DDL management for SessionDB.
Plain mixin consumed by ``hermes_state.SessionDB``: no ``__init__``, no state
of its own; methods use host attributes established by ``SessionDB.__init__``.
Must never import hermes_state (cycle) — shared constants live in
hermes_state_common.
Plain mixin for ``hermes_state.SessionDB`` (no ``__init__``/state of its own).
Must never import hermes_state (cycle); shared constants live in hermes_state_common.
"""
import datetime
@@ -26,27 +24,24 @@ from hermes_state_common import (
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.
# Pre-split logger identity so log filtering/capture is unchanged.
logger = logging.getLogger("hermes_state")
_FTS_HOLDER_ESCALATE_ATTEMPTS = 3
_FTS_HOLDER_ESCALATE_SECONDS = 60.0
# In-process retry cadence for a deferred stale-FTS rebuild
# (``retry_deferred_fts_recovery``): startup paid the full admission wait once;
# later retries are non-blocking probes. Each failed retry doubles the spacing
# up to the cap, so a permanent holder costs one warning per hour, not per minute.
# In-process retry cadence for a deferred stale-FTS rebuild (``retry_deferred_fts_recovery``):
# startup paid the full admission wait once; later retries are non-blocking probes whose
# spacing doubles up to the cap (a permanent holder costs one warning per hour).
_FTS_STALE_RETRY_SECONDS = 60.0
_FTS_STALE_RETRY_MAX_SECONDS = 3600.0
# schema_read_probe_statements() cache — deriving it parses SCHEMA_SQL in an
# in-memory SQLite database, so do it once per process.
# schema_read_probe_statements() cache (parses SCHEMA_SQL in an in-memory DB; once per process).
_READ_PROBE_STATEMENTS: Optional[tuple] = None
# Trigram triggers come ONLY from FTS_TRIGRAM_SQL / LEGACY_FTS_TRIGRAM_SQL, whose
# CREATE VIRTUAL TABLE needs the trigram tokenizer (SQLite >= 3.34); without it
# _ensure_fts_schema soft-fails that DDL and "all six present" is unsatisfiable.
# Split so a trigger's absence is measured only against the DDL that can create
# it. Exhaustive and disjoint by construction (test_fts_trigger_subsets_match_the_ddl).
# Trigram triggers come ONLY from FTS_TRIGRAM_SQL / LEGACY_FTS_TRIGRAM_SQL, whose CREATE
# VIRTUAL TABLE needs the trigram tokenizer (SQLite >= 3.34); without it _ensure_fts_schema
# soft-fails that DDL and "all six present" is unsatisfiable. Split so a trigger's absence is
# measured only against the DDL that can create it (test_fts_trigger_subsets_match_the_ddl).
_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)
@@ -80,8 +75,8 @@ _SESSION_MODEL_USAGE_HEAL_DDL = """CREATE TABLE session_model_usage (
last_seen REAL,
PRIMARY KEY (session_id, model, billing_provider, billing_base_url, billing_mode, task)
)"""
# Same table as emitted by the v22 migration: column lines at 35 spaces, closing
# paren at 31 (statement text is pinned by the SQL trace harness).
# Same table as emitted by the v22 migration: column lines at 35 spaces, closing paren at 31
# (statement text is pinned by the SQL trace harness).
_SESSION_MODEL_USAGE_V22_DDL = "\n".join(
[_SESSION_MODEL_USAGE_HEAL_DDL.splitlines()[0]]
+ [" " * 35 + ln.strip() for ln in _SESSION_MODEL_USAGE_HEAL_DDL.splitlines()[1:-1]]
@@ -121,6 +116,11 @@ _TITLE_UNIQUE_INDEX_SQL = (
"CREATE UNIQUE INDEX IF NOT EXISTS idx_sessions_title_unique "
"ON sessions(title) WHERE title IS NOT NULL"
)
_STALE_KEY_UPSERT_SQL = (
"INSERT INTO state_meta (key, value) VALUES (?, '1') "
"ON CONFLICT(key) DO UPDATE SET value = excluded.value"
)
_CLEAR_REBUILD_MARKERS_SQL = "DELETE FROM state_meta WHERE key IN ('fts_rebuild_high_water', 'fts_rebuild_progress')"
def _q(ident: str) -> str:
@@ -129,16 +129,13 @@ def _q(ident: str) -> str:
def schema_read_probe_statements() -> tuple:
"""SELECT statements that fail iff a live store is behind SCHEMA_SQL.
Read-only opens skip ``_reconcile_columns()`` by design (no DDL against
another profile's live DB), so healing callers (``_open_session_db_at_path``
in the web server) run these probes after a read-only open: a missing
table/column raises at prepare time. Derived from SCHEMA_SQL so a column
added there is covered automatically (a hand-maintained list went stale
within days); ``LIMIT 0`` so zero rows are read. Column references are
table-qualified: an unqualified double-quoted identifier that fails to
resolve silently degrades to a string literal (SQLite misfeature), which
would make the probe pass on exactly the stale store it exists to catch."""
"""SELECT statements that fail iff a live store is behind SCHEMA_SQL. Read-only opens
skip ``_reconcile_columns()`` (no DDL against another profile's live DB), so healing
callers run these after a read-only open: a missing table/column raises at prepare
time. Derived from SCHEMA_SQL (a hand-maintained list went stale within days);
``LIMIT 0`` reads zero rows. Column references are table-qualified: an unqualified
double-quoted identifier that fails to resolve silently degrades to a string literal
(SQLite misfeature), making the probe pass on exactly the stale store it should catch."""
global _READ_PROBE_STATEMENTS
if _READ_PROBE_STATEMENTS is None:
tables = SessionSchemaMixin._parse_schema_columns(SCHEMA_SQL)
@@ -154,11 +151,10 @@ class SessionSchemaMixin:
def _dedupe_legacy_system_prompts(self, cursor: sqlite3.Cursor) -> None:
"""Move inline prompt snapshots into the shared content-addressed table.
Contention-safe: any ``OperationalError`` mid-loop returns instead of
raising. Partial migration is safe — the legacy ``system_prompt``
column stays a read fallback and the next schema init resumes.
Propagating the error aborted schema init, left the version below
25, and re-entered this migration on every open (gateway crash loop)."""
Contention-safe: any ``OperationalError`` mid-loop returns instead of raising —
partial migration is safe (the legacy ``system_prompt`` column stays a read
fallback and the next schema init resumes), whereas propagating left the version
below 25 and re-entered this migration on every open (gateway crash loop)."""
try:
rows = cursor.execute(
"SELECT id, system_prompt FROM sessions WHERE system_prompt IS NOT NULL"
@@ -202,8 +198,7 @@ class SessionSchemaMixin:
@staticmethod
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)."""
"""Count how many of *names* exist as triggers (_FTS_BASE_TRIGGERS / _FTS_TRIGRAM_TRIGGERS for one half)."""
if not names:
return 0 # "name IN ()" is a SQLite syntax error
placeholders = ",".join("?" for _ in names)
@@ -214,6 +209,10 @@ class SessionSchemaMixin:
).fetchone()
return int(row[0])
@staticmethod
def _fts_triggers_missing(cursor: sqlite3.Cursor, names: Sequence[str]) -> bool:
return SessionSchemaMixin._fts_trigger_count(cursor, names) < len(names)
@staticmethod
def _fts_update_trigger_needs_narrowing(sql: Optional[str]) -> bool:
"""True when trigger SQL is a broad AFTER UPDATE (missing ``OF``)."""
@@ -223,14 +222,12 @@ class SessionSchemaMixin:
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.
``CREATE TRIGGER IF NOT EXISTS`` never replaces an existing broad
trigger, so it would keep firing on every messages row touch. Drop
still-broad UPDATE triggers and re-apply the current DDL. No FTS
rebuild: correctness was already gated by WHEN clauses; OF only skips
unnecessary trigger evaluation. Returns the number dropped."""
# CJK is v23-only. Decide the layout before selecting destructive
# candidates so the legacy branch never drops a trigger it won't recreate.
"""Replace broad AFTER UPDATE FTS triggers with AFTER UPDATE OF variants (``CREATE
TRIGGER IF NOT EXISTS`` never replaces an existing broad trigger, so it kept
firing on every messages row touch). No FTS rebuild: correctness was already
gated by WHEN clauses; OF only skips trigger evaluation. Returns the number dropped."""
# 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")
if not legacy_layout:
@@ -245,19 +242,16 @@ class SessionSchemaMixin:
if not to_drop:
return 0
for name in to_drop:
# Names come from the literal allowlist above — interpolation-safe.
cursor.execute(f"DROP TRIGGER IF EXISTS {name}")
cursor.execute(f"DROP TRIGGER IF EXISTS {name}") # names from the literal allowlist above
# Re-apply current DDL (legacy vs v23 chosen as _init_schema does) so
# CREATE TRIGGER installs the OF variants.
# Re-apply current DDL (legacy vs v23 chosen as _init_schema does) so CREATE
# TRIGGER installs the OF variants.
base_sql, trigram_sql = _FTS_DDL[legacy_layout]
self._ensure_fts_schema(cursor, "messages_fts", base_sql)
self._ensure_fts_schema(cursor, "messages_fts_trigram", trigram_sql)
# Only recreate the CJK trigger this migration actually dropped (it is
# a candidate only on the v23 layout). ``_ensure_fts_cjk_schema``
# soft-fails OperationalError by clearing availability (never raises),
# so afterwards require a narrowed CJK UPDATE trigger or durable
# quarantine (stale breadcrumb + unavailable).
# Only recreate the CJK trigger this migration actually dropped. ``_ensure_fts_cjk_schema``
# soft-fails OperationalError by clearing 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)
@@ -286,9 +280,9 @@ class SessionSchemaMixin:
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:
"""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."""
"""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."""
self._fts_cjk_available = False
try:
self.set_meta(FTS_CJK_STALE_KEY, "1", cursor=cursor)
@@ -301,36 +295,31 @@ class SessionSchemaMixin:
@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).
'rebuild' indexes EVERY row, so the deferred-backfill markers are
"""v23+ external-content tables: 'rebuild' repopulates the inverted index from the
content source. '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')")
cursor.execute(
"DELETE FROM state_meta WHERE key IN ('fts_rebuild_high_water', 'fts_rebuild_progress')"
)
cursor.execute(_CLEAR_REBUILD_MARKERS_SQL)
@staticmethod
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."""
"""Rebuild the LEGACY inline (pre-v23) FTS indexes from messages: no external-content
'rebuild' source, so DELETE + reinsert the concatenated content the legacy
triggers produced. 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, {_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."""
"""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
@@ -338,8 +327,8 @@ class SessionSchemaMixin:
decode_exc = exc
except sqlite3.OperationalError as exc:
if self._is_fts5_unavailable_error(exc):
# A missing trigram tokenizer only affects trigram search;
# only a missing FTS5 module disables FTS entirely.
# A missing trigram tokenizer only affects trigram search; only a missing
# FTS5 module disables FTS entirely.
if self._is_trigram_unavailable_error(exc):
self._warn_trigram_unavailable(exc)
else:
@@ -361,12 +350,11 @@ class SessionSchemaMixin:
# ── 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."""
"""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 = {}
try:
@@ -428,11 +416,10 @@ class SessionSchemaMixin:
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."""
"""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 and self._defer_stale_fts_for_holders(cursor, foreign_holders):
return False
@@ -447,15 +434,13 @@ class SessionSchemaMixin:
return self._recover_stale_fts_locked(cursor, legacy=legacy)
def retry_deferred_fts_recovery(self) -> bool:
"""Retry a deferred stale-FTS rebuild on this open SessionDB (gateway
housekeeping tick). ``_recover_stale_fts`` fails closed at open when
holders or the rebuild lock are busy, leaving ``_fts_stale`` set and
search on LIKE; live write/search paths must never start a full
rebuild, and a gateway opens state.db once for days, so "next open"
never comes. Bounded backoff (``_FTS_STALE_RETRY_SECONDS`` doubling to
the max), non-blocking admission (``timeout=0``), no new thread.
Returns True only when the index was rebuilt and sync triggers
restored. Never raises."""
"""Retry a deferred stale-FTS rebuild on this open SessionDB (gateway housekeeping
tick). ``_recover_stale_fts`` fails closed at open when holders or the rebuild
lock are busy, leaving ``_fts_stale`` set and search on LIKE; live write/search
paths must never start a full rebuild, and a gateway opens state.db once for
days, so "next open" never comes. Bounded backoff (``_FTS_STALE_RETRY_SECONDS``
doubling to the max), non-blocking admission (``timeout=0``), no new thread.
True only when the index was rebuilt and sync triggers restored. Never raises."""
if not self._fts_stale or self.read_only or self._conn is None:
return False
now = time.monotonic()
@@ -476,8 +461,8 @@ class SessionSchemaMixin:
legacy = self._db_has_legacy_inline_fts(cursor)
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.
# CJK was detached alongside the base indexes; its own ensure path
# decides when it comes back online.
self._ensure_fts_cjk_schema(cursor)
self._fts_stale_retry_interval = 0.0
try:
@@ -494,14 +479,14 @@ 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.
One write transaction closes the dangerous gap: no canonical writer
can slip between the full rebuild and trigger restoration."""
"""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):
# A corrupt vtable may fail even a LIMIT 0 probe; it must still be
# included in the drop-and-recreate below.
# A corrupt vtable may fail even a LIMIT 0 probe; it must still be included
# in the drop-and-recreate below.
trigram_status = True
include_trigram = trigram_status is True
@@ -540,7 +525,7 @@ class SessionSchemaMixin:
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 += _CLEAR_REBUILD_MARKERS_SQL + ";"
recovery_sql = (
"BEGIN IMMEDIATE;"
@@ -557,8 +542,8 @@ class SessionSchemaMixin:
self._conn.rollback()
except sqlite3.Error:
pass
# Stale indexes must stay detached even on SQLite builds whose DDL
# transaction behavior differs.
# Stale indexes must stay detached even on SQLite builds whose DDL transaction
# behavior differs.
self._drop_all_fts_triggers(cursor)
self._conn.commit()
logger.error(
@@ -580,13 +565,11 @@ class SessionSchemaMixin:
@staticmethod
def _parse_schema_columns(schema_sql: str) -> Dict[str, Dict[str, str]]:
"""Expected columns per table, parsed from SCHEMA_SQL.
Executes the DDL in an in-memory SQLite database and reads PRAGMA
table_info, so SQLite handles every syntax edge case (no regex).
The result is memoized on disk keyed by a hash of the DDL (~85ms per
startup otherwise; a pure function of the DDL text). Only the
reference-side parse is cached — diffing the LIVE database still runs
every startup. A corrupt or stale cache degrades to recomputation."""
"""Expected columns per table, parsed from SCHEMA_SQL by executing the DDL in an
in-memory SQLite database and reading PRAGMA table_info (no regex). Memoized on
disk keyed by a hash of the DDL (~85ms per startup otherwise); only the
reference-side parse is cached — diffing the LIVE database still runs every
startup. A corrupt or stale cache degrades to recomputation."""
cache_path = None
schema_hash = hashlib.sha256(schema_sql.encode("utf-8")).hexdigest()
try:
@@ -642,9 +625,9 @@ class SessionSchemaMixin:
return table_columns
def _reconcile_columns(self, cursor: sqlite3.Cursor) -> None:
"""ADD every SCHEMA_SQL column missing from the live tables.
Beets/sqlite-utils pattern: SCHEMA_SQL is the single source of truth;
column additions are declarative and need no version-gated migration."""
"""ADD every SCHEMA_SQL column missing from the live tables (beets/sqlite-utils
pattern: SCHEMA_SQL is the single source of truth; column additions are
declarative and need no version-gated migration)."""
expected = self._parse_schema_columns(SCHEMA_SQL)
for table_name, declared_cols in expected.items():
try:
@@ -665,10 +648,9 @@ class SessionSchemaMixin:
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.
# 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.
raise
# Anything else permanently strands the store behind SCHEMA_SQL — be loud.
logger.warning(
@@ -678,8 +660,8 @@ class SessionSchemaMixin:
@staticmethod
def _live_pk_columns(cursor: sqlite3.Cursor, table: str) -> Optional[List[str]]:
"""PRIMARY KEY column names of *table* in key order; None when the
table is missing or has no columns (SCHEMA_SQL creates it correctly)."""
"""PRIMARY KEY column names of *table* in key order; None when the table is
missing or has no columns (SCHEMA_SQL creates it correctly)."""
try:
rows = cursor.execute(f'PRAGMA table_info("{table}")').fetchall()
except sqlite3.OperationalError:
@@ -691,8 +673,8 @@ class SessionSchemaMixin:
@staticmethod
def _rebuild_table(cursor: sqlite3.Cursor, table: str, legacy_name: str, ddl: str, copy_sql: str, indexes=()) -> None:
"""RENAME *table* to *legacy_name*, CREATE it fresh from *ddl*, copy rows
back with *copy_sql*, DROP the legacy copy, recreate *indexes*."""
"""RENAME *table* to *legacy_name*, CREATE it fresh from *ddl*, copy rows back with
*copy_sql*, DROP the legacy copy, recreate *indexes*."""
cursor.execute(f"ALTER TABLE {table} RENAME TO {legacy_name}")
cursor.execute(ddl)
cursor.execute(copy_sql)
@@ -701,13 +683,12 @@ class SessionSchemaMixin:
cursor.execute(sql)
def _heal_gateway_routing_pk(self, cursor: sqlite3.Cursor) -> None:
"""Rebuild ``gateway_routing`` when its PRIMARY KEY predates scoping.
Early builds used ``session_key TEXT PRIMARY KEY``; the reconciler ADDs
``scope`` but SQLite cannot ALTER a PK, so the composite key never
lands and every routing write fails (ON CONFLICT mismatch / UNIQUE
violation across scopes) with per-save warning spam. Rebuild once,
preserving rows; on a cross-scope session_key collision the newest
row wins (INSERT OR REPLACE in updated_at order)."""
"""Rebuild ``gateway_routing`` when its PRIMARY KEY predates scoping. Early builds
used ``session_key TEXT PRIMARY KEY``; the reconciler ADDs ``scope`` but SQLite
cannot ALTER a PK, so the composite key never lands and every routing write
fails (ON CONFLICT mismatch / UNIQUE violation across scopes). Rebuild once,
preserving rows; on a cross-scope session_key collision the newest row wins
(INSERT OR REPLACE in updated_at order)."""
pk_cols = self._live_pk_columns(cursor, "gateway_routing")
if pk_cols is None or pk_cols == ["scope", "session_key"]:
return
@@ -732,21 +713,16 @@ class SessionSchemaMixin:
)
def _heal_session_model_usage_pk(self, cursor: sqlite3.Cursor) -> None:
"""Rebuild ``session_model_usage`` when its PRIMARY KEY lacks ``task``.
Installs already at v22+ when ``task`` landed carry the 5-column PK;
the reconciler ADDs ``task`` as a bare nullable but SQLite cannot
ALTER a PK, and the version-gated v22 rebuild is unreachable there.
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. OR IGNORE:
COALESCE(task, '') on legacy NULL rows can collide with a genuine
"""Rebuild ``session_model_usage`` when its PRIMARY KEY lacks ``task``. Installs
already at v22+ when ``task`` landed carry the 5-column PK; the reconciler ADDs
``task`` as a bare nullable but SQLite cannot ALTER a PK, and the version-gated
v22 rebuild is unreachable there — every ``_record_model_usage()`` upsert then
fails (ON CONFLICT mismatch), 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 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 none open). OR
IGNORE: COALESCE(task, '') on legacy NULL rows can collide with a genuine
''-task row — keep the first rather than fail."""
pk_cols = self._live_pk_columns(cursor, "session_model_usage")
if pk_cols is None or "task" in pk_cols:
@@ -788,28 +764,25 @@ class SessionSchemaMixin:
def _init_schema(self):
"""Create tables and FTS if missing, reconcile columns, run data migrations.
SCHEMA_SQL is the single source of truth: column additions are
declarative via _reconcile_columns(), so reordered migrations can
never skip a column. schema_version remains for data migrations
(row transforms) that cannot be expressed declaratively."""
# Startup-watchdog progress lease: on multi-GB state.db files the
# reconciliation + data migrations are I/O-bound (near-zero CPU), which
# the watchdog's CPU fallback would misread as a parked deadlock. Single
# lease (clamped to _MAX_LEASE_S=900) is deliberate: a wedged init delays
# supervisor respawn by up to the lease; per-chunk renewal isn't worth it.
SCHEMA_SQL is the single source of truth: column additions are declarative via
_reconcile_columns(), so reordered migrations can never skip a column.
schema_version remains for data migrations (row transforms) only."""
# Startup-watchdog lease: on multi-GB state.db files reconciliation + migrations
# are I/O-bound (near-zero CPU), which the watchdog's CPU fallback would misread as
# a parked deadlock. Single lease (clamped to _MAX_LEASE_S=900) is deliberate.
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
# table-shape repairs ADD COLUMN cannot express (PK rebuilds).
# Idempotent, self-healing column reconciliation, then the two table-shape
# repairs ADD COLUMN cannot express (PK rebuilds).
self._reconcile_columns(cursor)
self._heal_gateway_routing_pk(cursor)
self._heal_session_model_usage_pk(cursor)
# Indexes referencing reconciler-added columns must be created AFTER
# _reconcile_columns — in SCHEMA_SQL the initial executescript would
# fail on legacy DBs (WHERE references a not-yet-existing column).
# _reconcile_columns — in SCHEMA_SQL the initial executescript would fail on
# legacy DBs (WHERE references a not-yet-existing column).
try:
cursor.execute(
"CREATE INDEX IF NOT EXISTS idx_messages_platform_msg_id "
@@ -820,11 +793,10 @@ class SessionSchemaMixin:
logger.debug("idx_messages_platform_msg_id create skipped: %s", exc)
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
# it wrote NULL and ``WHERE active = 1`` loaders hid whole histories.
# Unconditional because a ``current_version < 12`` gate never re-ran
# for already-v12+ databases.
# Heal NULL ``active`` rows on every startup: older reconciler builds added
# ``active`` without its NOT NULL DEFAULT 1, so INSERTs omitting it wrote NULL and
# ``WHERE active = 1`` loaders hid whole histories. 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")
except sqlite3.OperationalError:
@@ -836,14 +808,13 @@ class SessionSchemaMixin:
(FTS_STALE_KEY,),
).fetchone() is not None
if self._fts_stale:
# A prior process detached FTS after corruption; keep every FTS
# writer detached until a full rebuild succeeds.
# A prior process detached FTS after corruption; keep every FTS writer
# detached until a full rebuild succeeds.
self._drop_all_fts_triggers(cursor)
if not fts5_available:
# Existing FTS triggers would still fire on messages writes even
# though this runtime cannot read their targets. Drop only the
# triggers so persistence continues; a future FTS5 runtime's
# _ensure_fts_schema() recreates them.
# Existing FTS triggers would still fire on messages writes even though this
# runtime cannot read their targets. Drop only the triggers so persistence
# continues; a future FTS5 runtime's _ensure_fts_schema() recreates them.
self._drop_fts_triggers(cursor)
row = cursor.execute("SELECT version FROM schema_version LIMIT 1").fetchone()
@@ -865,16 +836,16 @@ class SessionSchemaMixin:
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."""
"""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.
# 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:
@@ -891,9 +862,8 @@ class SessionSchemaMixin:
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.
# 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("
@@ -918,42 +888,40 @@ class SessionSchemaMixin:
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.
# 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.
# v20: seed one session_model_usage row per historical session from the
# sessions aggregates. INSERT OR IGNORE: a row newer code already wrote wins.
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 25 GB), so
# an existing install only gets a flag; `hermes sessions optimize-storage`
# performs it in the foreground. 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.
# 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 25 GB), so an existing install
# only gets a flag; `hermes sessions optimize-storage` performs it in the
# foreground. The FTS layout is tracked by the independent `fts_storage_version`
# marker, so schema_version still advances 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.
# v25: de-duplicate system prompt snapshots into the shared content-addressed
# table; the old column stays a read fallback for partially migrated rows.
self._dedupe_legacy_system_prompts(cursor)
# Stamp the FTS layout version (fresh/optimized DBs); 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.
# Stamp the FTS layout version (fresh/optimized DBs); 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)
@@ -965,16 +933,16 @@ class SessionSchemaMixin:
):
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: claiming current would lie.
# 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: 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=''."""
"""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') "
@@ -1003,9 +971,9 @@ class SessionSchemaMixin:
logger.debug("v22 session_model_usage rebuild skipped: %s", exc)
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."""
"""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_UNIQUE_INDEX_SQL)
except sqlite3.IntegrityError:
@@ -1031,18 +999,17 @@ class SessionSchemaMixin:
pass # Index already exists
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."""
"""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.
# 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
@@ -1051,17 +1018,15 @@ class SessionSchemaMixin:
else:
base_sql, trigram_sql = _FTS_DDL[legacy_fts]
rebuild = self._rebuild_legacy_fts_indexes if legacy_fts else 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)
)
# 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_triggers_missing(cursor, _FTS_BASE_TRIGGERS)
trigram_triggers_missing = self._fts_triggers_missing(cursor, _FTS_TRIGRAM_TRIGGERS)
self._fts_enabled = self._ensure_fts_schema(cursor, "messages_fts", base_sql)
if self._fts_enabled:
# Trigram is optional relative to the main table; without it
# CJK search falls back to LIKE.
# 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):
@@ -1076,18 +1041,16 @@ class SessionSchemaMixin:
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.
Reached when the sync triggers were missing and the DDL just recreated
them: the index has a gap of unknown extent. Two processes opening the
same DB after an update commonly hit this simultaneously (the
interleaving that structurally corrupted state.db in production), so
this admits through ``fts_rebuild_admission`` and FAILS CLOSED. On
deferral the just-repaired triggers are dropped again and the stale
breadcrumb persisted — triggers must never be live over an unrebuilt
gap (``_enter_fts_fail_open``'s ordering contract); the winner's
rebuild, ``retry_deferred_fts_recovery`` or ``_recover_stale_fts`` at
next startup restores index and triggers."""
"""Run a full trigger-repair FTS rebuild under cross-process admission. Reached when
the sync triggers were missing and the DDL just recreated them: the index has a
gap of unknown extent. Two processes opening the same DB after an update commonly
hit this simultaneously (the interleaving that structurally corrupted state.db in
production), so this admits through ``fts_rebuild_admission`` and FAILS CLOSED.
On deferral the just-repaired triggers are dropped again and the stale breadcrumb
persisted — triggers must never be live over an unrebuilt gap
(``_enter_fts_fail_open``'s ordering contract); the winner's rebuild,
``retry_deferred_fts_recovery`` or ``_recover_stale_fts`` at next startup
restores index and triggers."""
with fts_rebuild_admission(self.db_path) as admitted:
if admitted:
rebuild_fn()
@@ -1097,11 +1060,7 @@ class SessionSchemaMixin:
"rebuild authority for this state.db; detaching FTS sync "
"until the stale-index recovery path rebuilds it."
)
cursor.execute(
"INSERT INTO state_meta (key, value) VALUES (?, '1') "
"ON CONFLICT(key) DO UPDATE SET value = excluded.value",
(FTS_STALE_KEY,),
)
cursor.execute(_STALE_KEY_UPSERT_SQL, (FTS_STALE_KEY,))
self._drop_all_fts_triggers(cursor)
self._fts_stale = True
self._fts_enabled = False
@@ -1109,8 +1068,8 @@ class SessionSchemaMixin:
self._fts_cjk_available = False
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."""
"""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"
if not sessions_file.exists():
return