diff --git a/hermes_state_common.py b/hermes_state_common.py index 74054a7f65..88318404ca 100644 --- a/hermes_state_common.py +++ b/hermes_state_common.py @@ -33,8 +33,7 @@ _PREVIEW_MAX_CHARS = 60 def escape_like(text: str) -> str: """Escape LIKE wildcards (``%``, ``_``) so derived text matches literally; pair with ``ESCAPE '\\'``. ``_`` is common in branch names, titles and - paths, and a documented substring/prefix match must not silently widen. - """ + paths, and a documented substring/prefix match must not silently widen.""" return text.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") @@ -205,8 +204,7 @@ def _legacy_reset_child_sql(alias: str, reasons_sql: str) -> str: """Pre-marker reset-continuation heuristic: child rides its parent's exact non-empty routing key and the parent ended at a reset boundary. Shared by ``_RESET_CHILD_SQL`` and ``reopen_session()``'s marker-stamping UPDATE so - the two cannot drift; ``reasons_sql`` is a literal or placeholder list. - """ + the two cannot drift; ``reasons_sql`` is a literal or placeholder list.""" return ( f"EXISTS (SELECT 1 FROM sessions p" f" WHERE p.id = {alias}.parent_session_id" @@ -904,7 +902,6 @@ def _reopen_lock(lock_path, handle): def _acquire_db_flock(lock_path, handle, timeout_seconds, poll_seconds, description): """Bounded POSIX flock acquire with orphaned-holder staleness break. - Returns ``(acquired, handle)``; *handle* may have been re-opened and the caller closes whichever comes back. *acquired* is True, False (a holder kept the lock past the deadline), or None (non-contention ``OSError``, @@ -918,8 +915,7 @@ def _acquire_db_flock(lock_path, handle, timeout_seconds, poll_seconds, descript is unlinked and retaken on a fresh inode; the orphan's flock stays on the old inode blocking nobody. Every successful acquire verifies its inode still names *lock_path*, so a racer that locked a dead inode retries - instead of running alongside the breaker. Indeterminate liveness defers. - """ + instead of running alongside the breaker. Indeterminate liveness defers.""" import fcntl deadline = time.monotonic() + timeout_seconds @@ -1023,15 +1019,13 @@ def _acquire_msvcrt_lock(lock_path, handle, timeout): @contextlib.contextmanager def fts_rebuild_admission(db_path, *, timeout_seconds=None): """Serialize full structural FTS rebuilds on *db_path* across processes. - Yields True when this process holds the authority, False when the bounded acquire timed out or the lock file could not be opened. On False the caller must NOT rebuild (fail closed); the stale breadcrumb guarantees a retry. ``db_path`` None (in-memory DB) yields True. Opportunistic in-process retries pass ``timeout_seconds=0`` so a live - holder never stalls a long-lived writer; the orphan break still applies. - """ + holder never stalls a long-lived writer; the orphan break still applies.""" if db_path is None: yield True return diff --git a/hermes_state_portability.py b/hermes_state_portability.py index 7cb754b6a9..125feb0f70 100644 --- a/hermes_state_portability.py +++ b/hermes_state_portability.py @@ -126,15 +126,13 @@ class SessionPortabilityMixin: def list_cron_job_runs(self, job_id: str, limit: int = 20, offset: int = 0) -> List[Dict[str, Any]]: """List the run sessions produced by a single cron job, newest first. - Cron runs are flat sessions with id ``cron_{job_id}_{timestamp}``; they never compress or branch, so this skips ``list_sessions_rich``'s compression-chain CTE / leading-wildcard ``id_query`` path (which seeds from EVERY ``source='cron'`` row). Instead a ``[prefix, prefix_hi)`` index range scan on id filtered to ``source='cron'``, so work scales with the requested window. Returns the - ``list_sessions_rich`` row shape (``preview`` + ``last_active``). - """ + ``list_sessions_rich`` row shape (``preview`` + ``last_active``).""" prefix = f"cron_{job_id}_" # Half-open upper bound: bump the final byte so the range covers exactly # the ids starting with ``prefix``. @@ -249,8 +247,7 @@ class SessionPortabilityMixin: (never deleted) with ``end_reason='adopted_by_profile'`` — deliberately NOT in the recoverable set, so resurrection cannot undo an adoption. Returns the ``import_sessions`` dict plus ``adopted`` and - ``donor_retired`` (True only when EVERY segment's retirement applied). - """ + ``donor_retired`` (True only when EVERY segment's retirement applied).""" payload = donor_db.export_session_lineage(session_id) if not payload: return { @@ -298,14 +295,12 @@ class SessionPortabilityMixin: def _retire_donor_segment(self, donor_db: Any, seg_id: str) -> bool: """Archive one adopted donor segment; False when skipped or failed. - TOCTOU close-out: the divergence guard used EXPORT-TIME counts; re-read both stores right before stamping so donor growth never lands behind a non-recoverable archive (equal-count CONTENT divergence is accepted — bytes stay in the donor either way, only reachability differs). A retirement failure must not fail the adoption (a later resume retries - idempotently), but never claims success it didn't have. - """ + idempotently), but never claims success it didn't have.""" try: donor_now = len(donor_db.get_messages(seg_id)) local_now = len(self.get_messages(seg_id)) @@ -538,7 +533,6 @@ class SessionPortabilityMixin: def import_sessions(self, sessions: List[Dict[str, Any]]) -> Dict[str, Any]: """Import sessions exported by :meth:`export_session` or ``export_all``. - Existing ids are skipped. A child keeps its parent only when the parent exists or is in the same payload; otherwise it is detached so partial imports pass FK validation. Gateway routing, handoff, rewind and other @@ -548,8 +542,7 @@ class SessionPortabilityMixin: Activity contract: export INCLUDES ``last_activity_*`` (durable row fields) but import RESETS them to NULL — resurrecting a stale "working ..." label would fabricate activity the watchdog and listings - act on. Intentional asymmetry, pinned by regression test. - """ + act on. Intentional asymmetry, pinned by regression test.""" if not isinstance(sessions, list): raise ValueError("sessions must be a list") if len(sessions) > self._IMPORT_MAX_SESSIONS: diff --git a/hermes_state_schema.py b/hermes_state_schema.py index 4f20e51a2c..787df66736 100644 --- a/hermes_state_schema.py +++ b/hermes_state_schema.py @@ -128,7 +128,6 @@ 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 @@ -137,8 +136,7 @@ def schema_read_probe_statements() -> tuple: 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. - """ + would make the probe pass on exactly the stale store it exists to catch.""" global _READ_PROBE_STATEMENTS if _READ_PROBE_STATEMENTS is None: tables = SessionSchemaMixin._parse_schema_columns(SCHEMA_SQL) @@ -154,13 +152,11 @@ 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). - """ + 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" @@ -226,13 +222,11 @@ class SessionSchemaMixin: 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. - """ + 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. legacy_layout = self._db_has_legacy_inline_fts(cursor) @@ -336,8 +330,7 @@ class SessionSchemaMixin: 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. - """ + recovery flows. Anything else (malformed schema, corrupt vtable) re-raises.""" try: cursor.execute(f"SELECT * FROM {table_name} LIMIT 0") return True @@ -371,8 +364,7 @@ class SessionSchemaMixin: After ``_FTS_HOLDER_ESCALATE_ATTEMPTS`` deferrals spanning ``_FTS_HOLDER_ESCALATE_SECONDS``, provably inactive orphan Desktop - backends are reaped and the holders re-checked. - """ + backends are reaped and the holders re-checked.""" now = time.time() record = None try: @@ -436,12 +428,10 @@ class SessionSchemaMixin: 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. - """ + 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 @@ -504,10 +494,8 @@ class SessionSchemaMixin: 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. - """ + can slip between the full rebuild and trigger restoration.""" try: trigram_status = self._fts_table_probe(cursor, "messages_fts_trigram") except (sqlite3.DatabaseError, UnicodeDecodeError): @@ -592,14 +580,12 @@ 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. - """ + every startup. A corrupt or stale cache degrades to recomputation.""" cache_path = None schema_hash = hashlib.sha256(schema_sql.encode("utf-8")).hexdigest() try: @@ -656,10 +642,8 @@ class SessionSchemaMixin: 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. - """ + 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: @@ -706,14 +690,12 @@ class SessionSchemaMixin: 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. - """ + row wins.""" pk_cols = self._live_pk_columns(cursor, "gateway_routing") if pk_cols is None or pk_cols == ["scope", "session_key"]: return @@ -743,7 +725,6 @@ 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. @@ -755,8 +736,7 @@ class SessionSchemaMixin: 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. - """ + 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 @@ -802,12 +782,10 @@ 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. - """ + (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 @@ -1051,14 +1029,12 @@ class SessionSchemaMixin: 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. - """ + 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): diff --git a/hermes_state_search.py b/hermes_state_search.py index 0b6faaa78c..ec7d72d176 100644 --- a/hermes_state_search.py +++ b/hermes_state_search.py @@ -89,10 +89,8 @@ def _search_filter_clauses( exclude_sources: Optional[List[str]], role_filter: Optional[List[str]], ) -> None: """Append the visibility/source/role predicates every search route shares. - Live rows (active=1) AND compaction-archived rows (compacted=1) are - discoverable; only rewind/undo rows (active=0, compacted=0) are hidden. - """ + discoverable; only rewind/undo rows (active=0, compacted=0) are hidden.""" if not include_inactive: where.append("(m.active = 1 OR m.compacted = 1)") if source_filter is not None: @@ -129,12 +127,10 @@ class SessionSearchMixin: def _try_incremental_merge_fts(self) -> None: """Run one bounded FTS5 merge pass without failing the completed write. - The canonical write is already committed: no maintenance failure (even the bare SystemError CPython's sqlite3 layer can raise under cross-thread errmsg scrambling) may make the caller replay an - ambiguous, possibly-durable write. - """ + ambiguous, possibly-durable write.""" if not self._fts_enabled: return try: @@ -188,13 +184,11 @@ class SessionSearchMixin: def _fts_rebuild_finish(self) -> None: """Finalize the deferred rebuild: boundary sweep + clear markers. - The sweep is cheap insurance against a write that slipped through the migration-boundary instant (between high_water capture and trigger activation). The trigram half is gated on ``_trigram_available``: without the tokenizer/table an unconditional INSERT raises ``no such - table`` and aborts the whole rebuild (and optimize_fts_storage()). - """ + table`` and aborts the whole rebuild (and optimize_fts_storage()).""" sweeps = [self._BOUNDARY_SWEEP_SQL.format(table="messages_fts", extra="")] if self._trigram_available: sweeps.append(self._BOUNDARY_SWEEP_SQL.format(table="messages_fts_trigram", extra="AND m.role <> 'tool' ")) @@ -292,8 +286,7 @@ class SessionSearchMixin: FTS5 machinery is involved. Integer single-column-key tables drain with a high-water marker so each chunk's scan is bounded (restarting the scan was O(n²)); compound-key tables cannot use a scalar high-water and - keep the chunked ``LIMIT`` delete — they are small by construction. - """ + keep the chunked ``LIMIT`` delete — they are small by construction.""" with self._lock: trash = [ r[0] for r in self._conn.execute( @@ -396,13 +389,11 @@ class SessionSearchMixin: def _reset_fts_index_to_empty(self, conn) -> None: """Truncate the v23 external-content tables via FTS5 ``'delete-all'``. - A plain DELETE is O(rows) on external-content FTS5 and corrupts the index when indexed rows diverged from ``messages`` — exactly the shape this repair handles. The backfill worker replays its id range with no anti-join, so a replay from zero is only safe once the index is known - empty; this is how a partially indexed DB gets there. - """ + empty; this is how a partially indexed DB gets there.""" for tbl in ("messages_fts", "messages_fts_trigram"): try: conn.execute(f"INSERT INTO {tbl}({tbl}) VALUES('delete-all')") @@ -483,8 +474,7 @@ class SessionSearchMixin: Markers are written in the same BEGIN IMMEDIATE as the demote, BEFORE the empty v23 schema is created (``executescript`` implicitly COMMITs and cannot run inside that transaction), closing the crash window - where trash + empty v23 tables exist with no backfill claim. - """ + where trash + empty v23 tables exist with no backfill claim.""" def _stage(conn): self._drop_fts_triggers(conn) conn.execute("DROP VIEW IF EXISTS messages_fts_trigram_src") @@ -737,8 +727,7 @@ class SessionSearchMixin: standalone compaction handoffs are role='user' rows with NO display_kind — invisible to SQL — so fetch with headroom and drop them in the decode loop; otherwise ``/undo N`` pairs an in-memory count that - excludes handoffs with a DB pick that includes them. - """ + excludes handoffs with a DB pick that includes them.""" active_clause = "" if include_inactive else " AND active = 1" display_clause = " AND (display_kind IS NULL OR display_kind = '')" fetch_limit = int(limit) * 2 + 5 @@ -943,8 +932,7 @@ class SessionSearchMixin: corruption-class ``DatabaseError`` there detaches the derived indexes (``_enter_fts_fail_open``) and answers from canonical rows — a live search never performs the unbounded rebuild. Non-FTS corruption, or - any ``DatabaseError`` without *fail_open*, propagates. - """ + any ``DatabaseError`` without *fail_open*, propagates.""" sql, params = self._fts_match_sql(table, match_query, order_by_sql, **kwargs) try: return [dict(row) for row in self._read_all(sql, params)] @@ -1149,8 +1137,7 @@ class SessionSearchMixin: Rewound rows (``active=0, compacted=0``) are excluded by default; compaction-archived rows (``compacted=1``) ARE included so the pre-compaction transcript stays discoverable. ``include_inactive`` - searches every row. - """ + searches every row.""" result_fields = self._search_message_fields(fields) if not query or not query.strip(): return [] @@ -1224,8 +1211,7 @@ class SessionSearchMixin: index stores bigrams for runs >=2, so a single-char term only matches isolated chars — LIKE is broader). Then trigram (>=3 CJK chars per token), then a LIKE substring scan with one clause per non-operator - token so "广西 OR 桂林 OR 漓江" matches each term independently. - """ + token so "广西 OR 桂林 OR 漓江" matches each term independently.""" raw_query = query.strip('"').strip() match_query = _quote_fts_tokens(raw_query) if self._fts_cjk_available and not wants_tool_rows and not self._has_lone_cjk_run(raw_query): @@ -1320,11 +1306,9 @@ class SessionSearchMixin: def optimize_fts(self) -> int: """Merge fragmented FTS5 segments into one per index (``'optimize'``). - Pure maintenance: changes neither results nor ``snippet()`` output, only layout and speed; complementary to VACUUM, which then returns - the freed pages. Skips absent tables. Returns the number optimized. - """ + the freed pages. Skips absent tables. Returns the number optimized.""" optimized = 0 with self._lock: for tbl in self._FTS_TABLES: @@ -1369,7 +1353,6 @@ class SessionSearchMixin: def _merge_fts_incrementally(self, *, max_pages: int, max_commands: Optional[int] = None) -> int: """Run bounded FTS5 ``'merge'`` commands against each present index. - A positive merge rank stops after ~that many output pages, so each command holds the write lock for milliseconds regardless of index size (``'optimize'`` takes 9-18 s per index on a 10 GB DB, exhausting a @@ -1381,8 +1364,7 @@ class SessionSearchMixin: command's own INSERT is 1 change). Each command is its own implicit transaction, so competing processes interleave mid-pass. Missing tables are valid variants (optimize_fts_storage drops + backfills them live) - and are skipped; other SQLite errors propagate. Returns commands executed. - """ + and are skipped; other SQLite errors propagate. Returns commands executed.""" if isinstance(max_pages, bool) or not isinstance(max_pages, int): raise TypeError("max_pages must be an integer") if max_pages <= 0: