From ff3ebf509f7e43ced54a127812e71feca120bd9d Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 20:18:08 -0700 Subject: [PATCH] refactor(state): dict-dispatch persistence-error classifier, unify WHERE/placeholder builders, compact docstrings across state modules --- hermes_state_errors.py | 140 +++++----- hermes_state_fts.py | 113 ++++---- hermes_state_guard.py | 59 ++--- hermes_state_readpool.py | 70 +++-- hermes_state_sessions.py | 539 ++++++++++++++++----------------------- 5 files changed, 372 insertions(+), 549 deletions(-) diff --git a/hermes_state_errors.py b/hermes_state_errors.py index 806c81b7ea..af735220da 100644 --- a/hermes_state_errors.py +++ b/hermes_state_errors.py @@ -5,14 +5,12 @@ strings as well as live sqlite3 exceptions.""" import errno import sqlite3 -# --------------------------------------------------------------------------- -# Malformed-schema recovery: ``sqlite_master`` itself is inconsistent (typically -# a DUPLICATE ``CREATE VIRTUAL TABLE messages_fts`` row). SQLite parses the -# whole schema while preparing the FIRST statement, so EVERY statement raises — -# including ``PRAGMA journal_mode`` (it trips in apply_wal_with_fallback during -# __init__, before _init_schema) and plain ``DROP TABLE``; only -# ``PRAGMA writable_schema=ON`` + sqlite_master surgery still work. Canonical -# sessions/messages are intact; recovery rebuilds only the FTS layer. +# Malformed schema: ``sqlite_master`` itself is inconsistent (typically a DUPLICATE +# ``CREATE VIRTUAL TABLE messages_fts`` row). SQLite parses the whole schema while +# preparing the FIRST statement, so EVERY statement raises (even ``PRAGMA +# journal_mode`` during __init__); only ``PRAGMA writable_schema=ON`` + +# sqlite_master surgery still work. Canonical rows are intact; recovery rebuilds +# only the FTS layer. _MALFORMED_SCHEMA_MARKERS = ("malformed database schema",) _MALFORMED_DB_MARKERS = (*_MALFORMED_SCHEMA_MARKERS, "database disk image is malformed") @@ -25,51 +23,42 @@ def is_malformed_db_error(exc: BaseException) -> bool: ) -# SQLITE_IOERR as a substring (wrapped strings still classify); shared by the -# read-only open retry and the write-path BEGIN retry. +# SQLITE_IOERR as a substring (wrapped strings still classify). _DISK_IO_ERROR_MARKER = "disk i/o error" -# "Store BUSY, not gone" — HTTP callers map these to 503 instead of 500. -# Corruption deliberately absent: a malformed store must surface, not be -# retried into a timeout. +# "Store BUSY, not gone" — HTTP callers map these to 503 instead of 500. Corruption +# is deliberately absent: a malformed store must surface, not be retried into a timeout. _TRANSIENT_SQLITE_MARKERS = ( _DISK_IO_ERROR_MARKER, "database is locked", "database table is locked", "busy", ) def _is_no_more_rows(exc: sqlite3.Error) -> bool: - """Transient engine error on contended WAL appends; the identical write succeeds - standalone, so it retries like locked/busy. Message-scoped because some builds - raise it as InterfaceError (outside DatabaseError).""" + """Transient engine error on contended WAL appends (retries like locked/busy); + message-scoped because some builds raise it as InterfaceError.""" return "no more rows available" in str(exc).lower() def is_transient_sqlite_error(exc: BaseException) -> bool: - """"Busy right now", not "damaged". One predicate so the read-only open - retry and the HTTP 503-vs-500 split cannot drift apart.""" + """"Busy right now", not "damaged": one predicate so retry and the HTTP + 503-vs-500 split cannot drift apart.""" return isinstance(exc, sqlite3.OperationalError) and any( marker in str(exc).lower() for marker in _TRANSIENT_SQLITE_MARKERS ) def is_malformed_schema_error(exc: BaseException) -> bool: - """Only SQLite's explicit malformed-schema text. A generic "disk image is - malformed" (SQLITE_CORRUPT) may be any B-tree/freelist page and does not - prove canonical rows intact, so runtime repair must fail closed on it.""" + """Only SQLite's explicit malformed-schema text: a generic "disk image is + malformed" may be any B-tree page, so runtime repair must fail closed on it.""" return isinstance(exc, sqlite3.DatabaseError) and any( marker in str(exc).lower() for marker in _MALFORMED_SCHEMA_MARKERS ) -# "Filesystem cannot accept another write" substrings (OSError, sqlite3, and -# wrapped RPC strings all match the same helper). +# "Filesystem cannot accept another write" substrings (OSError, sqlite3, wrapped RPC strings). _DISK_FULL_MARKERS = ( - "no space left on device", - "not enough space", - "database or disk is full", # SQLITE_FULL - "disk full", - "full disk", - "enospc", + "no space left on device", "not enough space", "database or disk is full", # SQLITE_FULL + "disk full", "full disk", "enospc", ) @@ -83,67 +72,21 @@ def is_disk_full_error(exc: BaseException | str | None) -> bool: return any(marker in lowered for marker in _DISK_FULL_MARKERS) -# Every classify_persistence_error bucket; consumers enumerate this tuple so a -# new bucket can never silently desynchronize them. +# Every classify_persistence_error bucket; consumers enumerate this tuple. PERSISTENCE_ERROR_CAUSES = ( "locked", "compression", "compression_closed", "turn_lease", "corrupt", "replaced", "disk", "unknown", ) -# "Database FILE structurally damaged" substrings. NOTE: "database disk image is +# "Database FILE structurally damaged" substrings. "database disk image is # malformed" contains "disk", so this check MUST run before the disk bucket in # classify_persistence_error or B-tree corruption reads as "free some disk space". _DB_CORRUPTION_MARKERS = ( - "malformed", # "database disk image is malformed" (SQLITE_CORRUPT) - "file is not a database", # SQLITE_NOTADB (also connection-level poisoning) - "not a database", - "database corruption", + "malformed", "file is not a database", "not a database", "database corruption", ) -def classify_persistence_error(exc_or_str) -> str: - """Coarse cause bucket (PERSISTENCE_ERROR_CAUSES) so the user's guidance - matches: "locked" = busy, retry; "disk" = full/read-only/permissions; - "compression" = a live lease refused the write; "compression_closed" = adopt - the rotated session id; "turn_lease" = fencing, not storage; "corrupt" = - file damage (repair path, not disk space); "replaced" = stop writing.""" - if exc_or_str is None: - return "unknown" - # Lease refusals contain neither "locked" nor "busy": match by type, then by - # phrase for strings that survived RPC wrapping. - if isinstance(exc_or_str, SessionTurnLeaseLostError): - return "turn_lease" - if isinstance(exc_or_str, CompressionSessionClosedError): - return "compression_closed" - if isinstance(exc_or_str, CompressionSessionBusyError): - return "compression" - if isinstance(exc_or_str, StateDbReplacedError): # incl. DeletedWalGenerationError - return "replaced" - if isinstance(exc_or_str, StateDbCorruptError): - return "corrupt" - text = str(exc_or_str).lower() - if "turn lease" in text: - return "turn_lease" - if "closed by compression" in text: - return "compression_closed" - if "being compressed" in text or "compression lease" in text: - return "compression" - if "was replaced underneath" in text: - return "replaced" - if "deleted state.db-wal" in text or "deleted state.db-shm" in text: - return "replaced" - # Corruption BEFORE the lock/disk buckets: "disk image is malformed" - # contains "disk" and some wrapped strings mention "locked" recovery. - if any(marker in text for marker in _DB_CORRUPTION_MARKERS): - return "corrupt" - if "locked" in text or "busy" in text: - return "locked" - if is_disk_full_error(exc_or_str) or "disk" in text or "readonly" in text or "read-only" in text: - return "disk" - return "unknown" - - class CompressionSessionClosedError(RuntimeError): """A durable write targeted a parent already closed by compression.""" @@ -223,3 +166,44 @@ _STATE_DB_CORRUPT_MSG = ( "--inspect-only` or restore a snapshot. Unwritten transcripts are diverted to " "sessions/.jsonl (and the gateway pending_messages spool)." ) + + +_PERSISTENCE_CAUSE_BY_TYPE = ( + (SessionTurnLeaseLostError, "turn_lease"), + (CompressionSessionClosedError, "compression_closed"), + (CompressionSessionBusyError, "compression"), + (StateDbReplacedError, "replaced"), + (StateDbCorruptError, "corrupt"), +) +_PERSISTENCE_CAUSE_BY_PHRASE = ( + (("turn lease",), "turn_lease"), + (("closed by compression",), "compression_closed"), + (("being compressed", "compression lease"), "compression"), + (("was replaced underneath", "deleted state.db-wal", "deleted state.db-shm"), "replaced"), + (_DB_CORRUPTION_MARKERS, "corrupt"), + (("locked", "busy"), "locked"), +) + + +def classify_persistence_error(exc_or_str) -> str: + """Coarse cause bucket (PERSISTENCE_ERROR_CAUSES) so the user's guidance + matches: "locked" = busy, retry; "disk" = full/read-only/permissions; + "compression" = a live lease refused the write; "compression_closed" = adopt + the rotated session id; "turn_lease" = fencing, not storage; "corrupt" = + file damage (repair path, not disk space); "replaced" = stop writing.""" + if exc_or_str is None: + return "unknown" + # Lease refusals contain neither "locked" nor "busy": match by type first, + # then by phrase for strings that survived RPC wrapping. Order matters: + # StateDbReplacedError covers DeletedWalGenerationError; corruption comes + # BEFORE the lock/disk buckets ("disk image is malformed" contains "disk"). + for exc_type, cause in _PERSISTENCE_CAUSE_BY_TYPE: + if isinstance(exc_or_str, exc_type): + return cause + text = str(exc_or_str).lower() + for markers, cause in _PERSISTENCE_CAUSE_BY_PHRASE: + if any(marker in text for marker in markers): + return cause + if is_disk_full_error(exc_or_str) or any(m in text for m in ("disk", "readonly", "read-only")): + return "disk" + return "unknown" diff --git a/hermes_state_fts.py b/hermes_state_fts.py index 5349f6565d..ce3c977b98 100644 --- a/hermes_state_fts.py +++ b/hermes_state_fts.py @@ -14,23 +14,18 @@ from hermes_state_common import FTS_CJK_STALE_KEY, FTS_STALE_KEY, _FTS_CJK_TRIGG logger = logging.getLogger("hermes_state") # ── CJK-bigram FTS index (replaces the trigram index when available) ──── -# Trigram needs >=3 chars per term, so 1-2 char CJK terms fell through to a -# LIKE table scan (3-6s CPU per query on multi-GB installs). ``cjk_unicode61`` -# (native/fts5_cjk/, loadable) re-emits CJK runs as overlapping bigrams; FTS5 -# phrase semantics then give exact substring matching down to 2 chars. +# Trigram needs >=3 chars per term, so 1-2 char CJK terms fell through to a LIKE +# table scan; ``cjk_unicode61`` (native/fts5_cjk/, loadable) re-emits CJK runs as +# overlapping bigrams. Same v23 discipline as the trigram table: external-content +# over a tool-row-excluding view, triggers gated on a DEDICATED marker pair +# (fts_cjk_rebuild_high_water / _progress). The table exists ONLY when the +# tokenizer loads; a process that cannot load it drops the cjk triggers (writes +# keep working; the index goes stale until the next optimize-storage). # -# Same v23 discipline as the trigram table: external-content over a -# tool-row-excluding view, triggers gated on a DEDICATED marker pair -# (fts_cjk_rebuild_high_water / _progress) so a cjk-only backfill never gates -# the complete messages_fts triggers. The table exists ONLY when the tokenizer -# loads (~/.hermes/lib/libfts5_cjk.so); a process that cannot load it drops the -# cjk triggers (writes keep working; the index goes stale until the next -# optimize-storage on a capable host). -# -# Split DDL: the table/view is safe to ensure any time; triggers are created -# ONLY while the index is complete-or-marker-gated. A stale index must keep its +# Split DDL: the table/view is safe to ensure any time; triggers are created ONLY +# while the index is complete-or-marker-gated. A stale index must keep its # triggers DROPPED — an external-content 'delete' for a rowid the index never -# held is the canonical FTS5 corruption hazard the marker gating prevents. +# held is the canonical FTS5 corruption hazard. FTS_CJK_TABLE_SQL = """ CREATE VIEW IF NOT EXISTS messages_fts_cjk_src AS SELECT id, role, content, tool_name, tool_calls @@ -90,12 +85,11 @@ BEGIN END; """ + def fts5_cjk_so_path() -> Path: """Location of the cjk_unicode61 loadable extension.""" env = os.getenv("HERMES_FTS5_CJK_SO") - if env: - return Path(env).expanduser() - return get_hermes_home() / "lib" / "libfts5_cjk.so" + return Path(env).expanduser() if env else get_hermes_home() / "lib" / "libfts5_cjk.so" def _cjk_fts_config_enabled() -> bool: @@ -104,13 +98,10 @@ def _cjk_fts_config_enabled() -> bool: def load_fts5_cjk_extension(conn: sqlite3.Connection) -> bool: - """Best-effort load of the cjk_unicode61 tokenizer. False (never raises) - when the .so is absent, ``sessions.cjk_fts`` is off, or extension loading - is compiled out — callers then behave as before the cjk index existed.""" - if not _cjk_fts_config_enabled(): - return False + """Best-effort load of the cjk_unicode61 tokenizer; False (never raises) when + the .so is absent, ``sessions.cjk_fts`` is off, or loading is compiled out.""" path = fts5_cjk_so_path() - if not path.exists(): + if not _cjk_fts_config_enabled() or not path.exists(): return False try: conn.enable_load_extension(True) @@ -129,10 +120,7 @@ class SessionFtsSetupMixin: @staticmethod def _is_fts5_unavailable_error(exc: sqlite3.OperationalError) -> bool: - # Builds with FTS5 but without the optional trigram tokenizer raise - # "no such tokenizer: trigram" instead of "no such module"; the loadable - # cjk_unicode61 tokenizer shows the same capability-error shape. Scoped - # to those two tokenizers so unrelated tokenizer errors aren't masked. + """No FTS5 module, or an optional tokenizer missing (same capability-error shape).""" err = str(exc).lower() return ("no such module" in err and "fts5" in err) or SessionFtsSetupMixin._is_trigram_unavailable_error(exc) @@ -141,14 +129,12 @@ class SessionFtsSetupMixin: """Only an optional tokenizer is missing (trigram needs SQLite >= 3.34; cjk_unicode61 is loadable): "this one index can't be served", never "disable FTS".""" err = str(exc).lower() - return ("no such tokenizer: trigram" in err or "no such tokenizer: cjk_unicode61" in err) + return "no such tokenizer: trigram" in err or "no such tokenizer: cjk_unicode61" in err @staticmethod def _db_has_legacy_inline_fts(cursor: sqlite3.Cursor) -> bool: - """messages_fts exists in ANY pre-v23 shape. v23 is external-content over - content/tool_name/tool_calls; every legacy shape (inline single-column - v11..v22, or the v10-era external single-column) lacks tool_name, so - "stored CREATE lacks tool_name" catches both. False when absent (fresh DB).""" + """messages_fts exists in ANY pre-v23 shape: every legacy shape lacks + tool_name, so "stored CREATE lacks tool_name" catches them all. False when absent.""" row = cursor.execute( "SELECT sql FROM sqlite_master WHERE type = 'table' AND name = 'messages_fts'" ).fetchone() @@ -156,7 +142,7 @@ class SessionFtsSetupMixin: def _warn_trigram_unavailable(self, exc: sqlite3.OperationalError) -> None: """Log once that the trigram tokenizer is missing; base FTS5 stays enabled.""" - if getattr(self, "_trigram_unavailable_warned", False): + if getattr(self, "_trigram_unavailable_warned", False): # attr is lazily created here return self._trigram_unavailable_warned = True logger.info( @@ -182,12 +168,11 @@ class SessionFtsSetupMixin: ) def _ensure_fts_cjk_schema(self, cursor) -> None: - """Create / repair / self-heal the CJK-bigram index (see the module - comment). Sets ``_fts_cjk_available``; never raises. Loaded + absent → - create (a populated DB gets the backfill markers and is NOT served until - optimize-storage backfills); loaded + present → ensure triggers, honour - the stale breadcrumb; NOT loaded + live triggers → drop them so INSERTs - don't fail at trigger time and leave the breadcrumb.""" + """Create / repair / self-heal the CJK-bigram index (see the module comment). + Sets ``_fts_cjk_available``; never raises. Loaded + absent → create (a + populated DB gets backfill markers and is NOT served until optimize-storage + backfills); loaded + present → ensure triggers, honour the stale breadcrumb; + NOT loaded + live triggers → drop them (INSERTs must not fail at trigger time).""" try: cjk_present = bool(cursor.execute( "SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = 'messages_fts_cjk'" @@ -200,8 +185,7 @@ class SessionFtsSetupMixin: _FTS_CJK_TRIGGERS, ).fetchall()] if live: - # Breadcrumb FIRST (a crash between the two statements - # is merely conservative), then drop. + # Breadcrumb FIRST (a crash between the two is merely conservative). logger.warning( "messages_fts_cjk triggers present but the " "cjk_unicode61 tokenizer is unavailable (%s) — " @@ -230,12 +214,10 @@ class SessionFtsSetupMixin: try: cursor.executescript(FTS_CJK_TABLE_SQL) if not cjk_present: - # Any old stale breadcrumb refers to a table that no longer exists. + # An old stale breadcrumb refers to a table that no longer exists. cursor.execute("DELETE FROM state_meta WHERE key = ?", (FTS_CJK_STALE_KEY,)) - # Empty DB: index complete by construction (triggers cover everything), - # no markers. Populated DB: the marker pair keeps the id-gated triggers - # correct while old rows await optimize-storage; the index is NOT - # served until that backfill completes. + # Empty DB: complete by construction, no markers. Populated DB: the + # marker pair keeps the id-gated triggers correct until backfill. if cursor.execute("SELECT COUNT(*) FROM messages WHERE role <> 'tool'").fetchone()[0] > 0: hw = cursor.execute("SELECT COALESCE(MAX(id), 0) FROM messages").fetchone()[0] for k, v in ( @@ -247,9 +229,7 @@ class SessionFtsSetupMixin: (k, v), ) if cursor.execute("SELECT 1 FROM state_meta WHERE key = ?", (FTS_CJK_STALE_KEY,)).fetchone(): - # Gap of unknown extent: do NOT reinstall triggers (an - # external-content 'delete' for an unindexed rowid corrupts the - # index); the next optimize-storage rebuilds from scratch. + # Gap of unknown extent: do NOT reinstall triggers (see module comment). self._fts_cjk_available = False return cursor.executescript(FTS_CJK_TRIGGER_SQL) @@ -292,9 +272,8 @@ class SessionFtsSetupMixin: @staticmethod def _is_fts_write_corruption_error(exc: sqlite3.DatabaseError) -> bool: - """Corruption SQLite identifies as FTS-scoped: SQLITE_CORRUPT_VTAB, or - (older builds) an ``fts5:`` message. A bare malformed-image error is - structural and must not trigger live FTS maintenance.""" + """Corruption SQLite identifies as FTS-scoped (SQLITE_CORRUPT_VTAB, or an + ``fts5:`` message on older builds); a bare malformed image is structural.""" error_code = getattr(exc, "sqlite_errorcode", None) if error_code is not None: return error_code == getattr(sqlite3, "SQLITE_CORRUPT_VTAB", 267) @@ -302,10 +281,9 @@ class SessionFtsSetupMixin: return msg.startswith("fts5:") and "corrupt structure" in msg def _enter_fts_fail_open(self, exc: sqlite3.DatabaseError) -> bool: - """Detach corrupt FTS indexes so canonical writes can continue. Stale - breadcrumb + trigger drop commit atomically: once triggers are absent - the index has a gap of unknown extent, so no process may reinstall them - without rebuilding every row.""" + """Detach corrupt FTS indexes so canonical writes can continue. Breadcrumb + + trigger drop commit atomically: once triggers are absent the index has a + gap of unknown extent, so nobody may reinstall them without a full rebuild.""" if not self._fts_enabled or not self._is_fts_write_corruption_error(exc): return False self._raise_if_db_corrupt() @@ -355,21 +333,18 @@ class SessionFtsSetupMixin: return True # ── Chunked FTS rebuild engine (v23 opt-in optimize) ── - # One blocking rebuild held the write lock ~16 minutes on a 25 GB DB, so the - # backfill runs in small chunks, each its own short transaction (resumable - # from fts_rebuild_progress; concurrent runners claim chunks by CAS). - # THROTTLING: a greedy loop owned the lock ~85% of the time and starved - # other processes' writers; 500-row chunks plus a pause of max(MIN_PAUSE, - # chunk cost x DUTY_FACTOR) cap our duty cycle unconditionally (works - # cross-process, unlike any same-process activity stamp). + # One blocking rebuild held the write lock ~16 min on a 25 GB DB, so the + # backfill runs in small chunks (each its own short transaction, resumable from + # fts_rebuild_progress, claimed by CAS). A greedy loop starved other writers: + # a pause of max(MIN_PAUSE, chunk cost x DUTY_FACTOR) caps the duty cycle + # cross-process, unlike any same-process activity stamp. _FTS_REBUILD_CHUNK_ROWS = 500 _FTS_REBUILD_DUTY_FACTOR = 4.0 # sleep >= 4x chunk cost (≤20% duty) _FTS_REBUILD_MIN_PAUSE = 0.2 # seconds — floor between chunks # Demoted v22 FTS shadow tables awaiting teardown: DROP of a multi-GB vtable - # blocks for minutes, so the v23 migration demotes the vtable definitions - # out of sqlite_master and renames the orphaned shadow tables (now plain - # tables) to fts_v22_trash_*; the worker empties them in chunks, then drops. + # blocks for minutes, so the v23 migration renames the orphaned shadow tables + # to fts_v22_trash_*; the worker empties them in chunks, then drops. _FTS_TRASH_PREFIX = "fts_v22_trash_" def _has_fts_trash(self, conn) -> bool: @@ -380,6 +355,6 @@ class SessionFtsSetupMixin: (self._FTS_TRASH_PREFIX.replace("_", "\\_") + "%",), ).fetchone()) - # FTS5 tables merged on optimize; trigram may be disabled and cjk exists only - # with the loadable tokenizer, so each is probed before touching (optimize_fts). + # FTS5 tables merged on optimize; each is probed before touching (trigram may + # be disabled, cjk exists only with the loadable tokenizer). _FTS_TABLES = ("messages_fts", "messages_fts_trigram", "messages_fts_cjk") diff --git a/hermes_state_guard.py b/hermes_state_guard.py index e9506dfc5a..2f4f528992 100644 --- a/hermes_state_guard.py +++ b/hermes_state_guard.py @@ -15,9 +15,7 @@ except ImportError: # pragma: no cover - stripped/scaffold installs only psutil = None # type: ignore[assignment] # Field evidence: pytest fixture rows landed in the production state.db and a -# pytest-spawned child flipped the journal mode under the live WAL writer, -# destroying committed transcripts; any HERMES_HOME escape (fixture ordering, a -# child spawned without it, a shell exporting the real home) fell through silently. +# pytest-spawned child flipped the journal mode under the live WAL writer. #: Env twin of ``_STATE_DB_GUARD_BYPASS`` for child processes (a module global #: cannot cross a process boundary, and ancestry arms the guard there). @@ -26,29 +24,23 @@ _STATE_DB_GUARD_BYPASS_ENV = "HERMES_STATE_DB_GUARD_BYPASS" def _real_platform_state_root() -> Optional[Path]: """The REAL platform-default Hermes root. Avoids ``Path.home()`` / - ``hermes_constants``: tests monkeypatch Path.home to a tempdir while this - module is imported lazily, which would misidentify the hermetic home as - production or miss the real one. ``expanduser`` reads HOME/passwd, which the - conftest never rewrites.""" + ``hermes_constants`` (tests monkeypatch Path.home to a tempdir); ``expanduser`` + reads HOME/passwd, which the conftest never rewrites.""" try: + home = Path(os.path.expanduser("~")) if sys.platform == "win32": base = os.environ.get("LOCALAPPDATA", "").strip() - root = ( - Path(base) / "hermes" - if base - else Path(os.path.expanduser("~")) / "AppData" / "Local" / "hermes" - ) + root = Path(base) / "hermes" if base else home / "AppData" / "Local" / "hermes" else: - root = Path(os.path.expanduser("~")) / ".hermes" + root = home / ".hermes" return root.resolve() except Exception: return None -#: Exported by the hermetic conftest alongside the HERMES_HOME redirect (value: -#: the isolation root). Unlike PYTEST_* (scrubbed by tests that rebuild a child -#: env) it is OURS and inherits by default, so a child carrying it that resolves -#: a production DB is by definition an isolation escape. +#: Exported by the hermetic conftest alongside the HERMES_HOME redirect. Unlike +#: PYTEST_* it is OURS and inherits by default, so a child carrying it that +#: resolves a production DB is by definition an isolation escape. _TEST_ISOLATION_MARKER_ENV = "HERMES_TEST_ISOLATION" @@ -70,18 +62,15 @@ _PYTEST_ANCESTOR: Optional[bool] = None def _process_looks_like_pytest(proc: Any) -> bool: - """True when *proc*'s command line is a pytest invocation (``pytest ...`` or - ``python -m pytest``). Unreadable cmdline => not pytest: guessing the other - way would refuse production opens for unrelated reasons.""" + """True when *proc*'s command line is a pytest invocation. Unreadable cmdline + => not pytest: guessing the other way would refuse production opens.""" try: cmdline = proc.cmdline() or [] except Exception: return False for arg in cmdline: try: - # Split on both separators on every host: os.path.basename is - # POSIX-only under Linux and would leave a Windows-style path - # intact, making the matcher's answer depend on the platform. + # Split on both separators on every host so the answer is platform-independent. name = str(arg).strip('"').strip("'").replace("\\", "/").rsplit("/", 1)[-1].lower() except Exception: continue @@ -91,10 +80,9 @@ def _process_looks_like_pytest(proc: Any) -> bool: def _has_pytest_ancestor() -> bool: - """True when an ancestor process is a pytest run. A child spawned with a - rebuilt env loses PYTEST_* and the HERMES_HOME redirect together — aiming at - production AND disarming the guard in one step; ancestry survives that. - Fails open without psutil / on walk errors (never block real user runs).""" + """True when an ancestor process is a pytest run: a child spawned with a + rebuilt env loses PYTEST_* and the HERMES_HOME redirect together, ancestry + survives that. Fails open without psutil / on walk errors.""" global _PYTEST_ANCESTOR if _PYTEST_ANCESTOR is not None: return _PYTEST_ANCESTOR @@ -109,15 +97,13 @@ def _has_pytest_ancestor() -> bool: def _in_test_context() -> bool: - """Test run by environment or ancestry. Env first (two dict lookups); the - memoised ancestry walk runs at most once per real ``hermes`` invocation.""" + """Test run by environment or ancestry (memoised; env checked first).""" return _running_under_pytest() or _has_pytest_ancestor() def _is_production_state_db(resolved: Path, root: Path) -> bool: - """*resolved* is ``/state.db`` or ``/profiles//state.db``. - Deeper scratch paths (repo worktrees under ~/.hermes/hermes-agent/...) are - deliberately NOT matched so hermetic tests cannot false-positive.""" + """*resolved* is ``/state.db`` or ``/profiles//state.db``; + deeper scratch paths (repo worktrees) are deliberately NOT matched.""" if resolved.parent == root: return True try: @@ -128,16 +114,15 @@ def _is_production_state_db(resolved: Path, root: Path) -> bool: # Last SessionDB() init error, per-process; surfaced by /resume-style slash -# commands so users know WHY. Only SessionDB.__init__ writes it (kanban_db -# failures are reported via their own callers, by design). +# commands so users know WHY. Only SessionDB.__init__ writes it. _last_init_error: Optional[str] = None _last_init_error_lock = threading.Lock() def _set_last_init_error(msg: Optional[str]) -> None: - """Record (or clear with None) the most recent state.db init failure. - __init__ only SETs on failure and never clears on success: a concurrent - successful open would erase the cause another thread's /resume is about to format.""" + """Record (or clear with None) the most recent init failure. __init__ never + clears on success: a concurrent open would erase the cause another thread's + /resume is about to format.""" global _last_init_error with _last_init_error_lock: _last_init_error = msg diff --git a/hermes_state_readpool.py b/hermes_state_readpool.py index b74158fe29..f55087d55d 100644 --- a/hermes_state_readpool.py +++ b/hermes_state_readpool.py @@ -18,33 +18,28 @@ if TYPE_CHECKING: # pragma: no cover # caplog tests pin the "hermes_state" logger name. logger = logging.getLogger("hermes_state") -# Ceiling on read-only connections ALIVE at once against one database FILE -# (idle pooled + checked out, summed over every SessionDB on that file). One -# constant for both the pool maxsize and the permit count: a LifoQueue only caps -# how many are *returned*; with open-on-miss, N readers hitting an empty pool -# all open and peak at N, and EMFILE is a peak-instant condition. So a -# connection holds a permit for its whole lifetime (_get_read_conn -> -# _close_read_conn); once permits are gone reads degrade to the locked writer -# connection — slower, but not a process-wide wedge the supervisor can't see. +# Ceiling on read-only connections ALIVE at once against one database FILE (idle +# pooled + checked out, over every SessionDB on that file). One constant for both +# the pool maxsize and the permit count: a LifoQueue only caps how many are +# *returned*, and EMFILE is a peak-instant condition, so a connection holds a +# permit for its whole lifetime; once permits are gone reads degrade to the +# locked writer connection — slower, but not a wedge the supervisor can't see. _READ_POOL_MAX = 8 -# Ceiling on read-only connections ALIVE in this PROCESS across every state.db -# (a multiplexed gateway opens one per profile, so a per-file cap still scales -# with profile count). Three profiles' worth; past it readers degrade to the -# writer connection for the same reason as _READ_POOL_MAX. +# Ceiling ALIVE in this PROCESS across every state.db (a multiplexed gateway +# opens one per profile); three profiles' worth, then readers degrade likewise. _READ_POOL_PROCESS_MAX = 24 -# Warn past this many SessionDB handles on one file in one process. Diagnostic -# only: writer connections cannot be rationed the way read connections can. +# Warn past this many SessionDB handles on one file in one process (diagnostic: +# writer connections cannot be rationed the way read connections can). _HANDLES_PER_PATH_WARN = 4 # Descriptors kept in reserve for everything that is NOT this module (httpx -# sockets, terminal pipes, log files): SQLite's share is only part of the fd -# table, and the EMFILE it pushes over surfaces elsewhere (terminal_tool). +# sockets, terminal pipes, log files): the EMFILE SQLite pushes over surfaces elsewhere. _FD_HEADROOM_RESERVE = 64 # The fd count is a directory listing; cache it briefly so a read burst isn't a -# syscall per query. Staleness lets through at most the ceiling's worth of opens. +# syscall per query (staleness lets through at most the ceiling's worth of opens). _FD_USAGE_CACHE_SECONDS = 0.25 _process_read_permits = threading.BoundedSemaphore(_READ_POOL_PROCESS_MAX) @@ -70,8 +65,7 @@ def _proc_fd_targets(pid: int) -> Iterator[str]: def _open_fd_count() -> Optional[int]: """Open descriptors in THIS process; None when unmeasurable (Windows: no fd - dir and no RLIMIT_NOFILE, correctly inert — its limit is thousands); -1 when - the probe itself hit EMFILE/ENFILE (that IS the answer: no headroom).""" + dir, correctly inert); -1 when the probe itself hit EMFILE/ENFILE (no headroom).""" for fd_dir in ("/proc/self/fd", "/dev/fd"): try: return len(os.listdir(fd_dir)) @@ -97,10 +91,9 @@ def _fd_soft_limit() -> Optional[int]: def _fd_headroom_ok() -> bool: - """Can the process spare a descriptor for a new read connection? - Fails OPEN when unmeasurable (refusing every read there would be a - self-inflicted convoy); fails CLOSED only on evidence (measured shortfall, - or a probe that couldn't get a descriptor itself).""" + """Can the process spare a descriptor for a new read connection? Fails OPEN + when unmeasurable (refusing every read would be a self-inflicted convoy); + fails CLOSED only on evidence (measured shortfall or a starved probe).""" soft = _fd_soft_limit() if soft is None: return True @@ -127,11 +120,10 @@ def _reclaim_idle_read_conn_anywhere() -> bool: class _PathReadBudget: - """Read-connection permits for ONE database file, shared process-wide: - per-instance semaphores let N SessionDBs on one file peak at N x (1 + MAX) - and walk into EMFILE. An idle pooled connection keeps its permit, so a - permit miss first reclaims an IDLE connection from a peer on the same path - (idle descriptors are transferable, in-use ones are not).""" + """Read-connection permits for ONE database file, shared process-wide + (per-instance semaphores let N SessionDBs peak at N x (1 + MAX)). An idle + pooled connection keeps its permit, so a permit miss first reclaims an IDLE + connection from a peer on the same path.""" def __init__(self) -> None: self.permits = threading.BoundedSemaphore(_READ_POOL_MAX) @@ -148,22 +140,19 @@ class _PathReadBudget: if warn: self._duplicate_handles_warned = True if warn: - # Writer connections cannot be capped (a SessionDB without one cannot - # write); the only bound is not opening redundant handles. Make the - # next duplicate visible before it becomes an incident. + # Writer connections cannot be capped; the only bound is not opening + # redundant handles, so make the duplicate visible before it's an incident. logger.warning( "%d live SessionDB handles on %s in this process; each holds " "its own writer connection (read connections are capped at %d " "for the file). A long-lived process should share one handle per path.", - handles, - db.db_path, - _READ_POOL_MAX, + handles, db.db_path, _READ_POOL_MAX, ) def acquire(self, requester: "SessionDB") -> bool: - """Take a permit for a new read connection, or refuse (caller then reads - via the locked writer connection — slower, never an error). Gates, - broadest first: fd headroom, process-wide ceiling, this file's ceiling.""" + """Take a permit for a new read connection, or refuse (caller degrades to the + locked writer connection). Gates, broadest first: fd headroom, process + ceiling, this file's ceiling.""" if not _fd_headroom_ok(): global _read_open_denied_fd_headroom with _read_budgets_lock: @@ -182,8 +171,7 @@ class _PathReadBudget: _process_read_permits.release() def _acquire_process_permit(self) -> bool: - # Another thread may take a freed permit first; that is a legitimate - # loss, and the caller degrades to the writer lock rather than looping. + # Another thread may take a freed permit first: legitimate loss, no looping. return _process_read_permits.acquire(blocking=False) or ( _reclaim_idle_read_conn_anywhere() and _process_read_permits.acquire(blocking=False) ) @@ -201,8 +189,8 @@ class _PathReadBudget: return any(member._evict_one_idle_read_conn() for member in members) -# canonical db path -> permits for that file. Weak values: the budget lives as -# long as some SessionDB on the path holds it, so tmp_path churn can't grow this. +# canonical db path -> permits for that file. Weak values: the budget lives only +# while some SessionDB on the path holds it, so tmp_path churn can't grow this. _read_budgets: "weakref.WeakValueDictionary[str, _PathReadBudget]" = (weakref.WeakValueDictionary()) _read_budgets_lock = threading.Lock() diff --git a/hermes_state_sessions.py b/hermes_state_sessions.py index f2e2facf68..36c9a7e6da 100644 --- a/hermes_state_sessions.py +++ b/hermes_state_sessions.py @@ -10,7 +10,10 @@ import time from pathlib import Path from typing import Any, Callable, Dict, List, Optional, Tuple -from agent.session_activity import ActivityProvenance +from agent.session_activity import ( + ActivityProvenance, bound_activity_description, build_activity_snapshot, + normalize_activity_provenance, +) from hermes_state_common import ( _LISTABLE_CHILD_SQL, _PREVIEW_ELIGIBLE_SQL, _PREVIEW_RAW_SELECT, _RECOVERABLE_END_REASONS, _RECOVERABLE_END_REASONS_SQL, _RESET_END_REASONS, _legacy_reset_child_sql, _shape_preview, @@ -22,8 +25,8 @@ logger = logging.getLogger("hermes_state") def workspace_key(row: Dict[str, Any]) -> Optional[str]: - """Workspace grouping key: git repo root when known, else cwd, else None. - Branch is deliberately excluded so a checkout doesn't fragment history.""" + """Workspace grouping key: git repo root, else cwd, else None (branch is + deliberately excluded so a checkout doesn't fragment history).""" return (row.get("git_repo_root") or "").strip() or (row.get("cwd") or "").strip() or None @@ -51,10 +54,8 @@ def _parse_model_config(raw: Any) -> Dict[str, Any]: def _cwd_prefix_clause(cwd_prefix: str) -> Tuple[str, List[str]]: prefix = cwd_prefix.rstrip("/\\") or cwd_prefix - # ``_``/``%`` are LIKE wildcards but ordinary path characters (``my_project``): - # unescaped, a prefix also matches sibling directories. The ``=`` arm is an - # exact compare and keeps the raw prefix; the Windows separator backslash - # in the LIKE pattern needs escaping too. + # ``_``/``%`` are LIKE wildcards but ordinary path characters: unescaped, a + # prefix also matches sibling directories. The ``=`` arm keeps the raw prefix. esc = _escape_like(prefix) return ( "(s.cwd = ? OR s.cwd LIKE ? ESCAPE '\\' OR s.cwd LIKE ? ESCAPE '\\')", @@ -64,8 +65,7 @@ def _cwd_prefix_clause(cwd_prefix: str) -> Tuple[str, List[str]]: def _workspace_key_clause(key: str) -> Tuple[str, List[str]]: """WHERE for ``workspace_key(row) == key``: git_repo_root equals ``key``, or - (rows predating per-session git metadata) cwd is at/under ``key``. Used by - ``hermes -c``/``--resume`` to pick the current workspace's MRU, not the global one.""" + (rows predating per-session git metadata) cwd is at/under ``key``.""" prefix = key.rstrip("/\\") or key cwd_clause, cwd_params = _cwd_prefix_clause(prefix) return ( @@ -74,8 +74,8 @@ def _workspace_key_clause(key: str) -> Tuple[str, List[str]]: ) -# First user message of a session, shaped by _shape_preview() in Python. The -# indentation is part of the list_sessions_rich SQL text. +# First user message of a session, shaped by _shape_preview() in Python. +# The indentation is part of the list_sessions_rich SQL text. _PREVIEW_COL_SQL = f"""COALESCE( (SELECT {_PREVIEW_RAW_SELECT} FROM messages m @@ -86,17 +86,25 @@ _PREVIEW_COL_SQL = f"""COALESCE( ) AS _preview_raw""" +def _session_ids_placeholders(ids) -> str: + return ",".join("?" * len(ids)) + + +def _where_sql(clauses: List[str], lead: str = "") -> str: + """``WHERE a AND b`` (with *lead* prefix) or "" when there are no clauses.""" + return f"{lead}WHERE {' AND '.join(clauses)}" if clauses else "" + + def _session_filter_where( *, exclude_children: bool = False, source: str = None, sources: List[str] = None, session_key: str = None, exclude_sources: List[str] = None, cwd_prefix: str = None, min_message_count: int = 0, archived_only: bool = False, include_archived: bool = False, ) -> Tuple[List[str], List[Any]]: - """Shared ``sessions s`` WHERE builder so session counts line up with the - listed rows. ``exclude_children`` hides sub-agent runs and compression - continuations but keeps branch/reset children: ``_LISTABLE_CHILD_SQL`` uses - the stable ``_branched_from`` marker (survives a re-ended parent) OR'd with - the legacy parent-ended-'branched' heuristic for pre-marker rows. Clause - order is part of the SQL text contract.""" + """Shared ``sessions s`` WHERE builder so counts line up with listed rows. + ``exclude_children`` hides sub-agent runs and compression continuations but + keeps branch/reset children (``_LISTABLE_CHILD_SQL``: stable ``_branched_from`` + marker OR the legacy heuristic for pre-marker rows). Clause order is part of + the SQL text contract.""" where: List[str] = [] params: List[Any] = [] if exclude_children: @@ -104,13 +112,13 @@ def _session_filter_where( where.append(f"{_delegate_from_json('s.model_config')} IS NULL") include_sources = [source] if source else list(sources or []) if include_sources: - where.append(f"s.source IN ({','.join('?' for _ in include_sources)})") + where.append(f"s.source IN ({_session_ids_placeholders(include_sources)})") params.extend(include_sources) if session_key: where.append("s.session_key = ?") params.append(session_key) if exclude_sources: - where.append(f"s.source NOT IN ({','.join('?' for _ in exclude_sources)})") + where.append(f"s.source NOT IN ({_session_ids_placeholders(exclude_sources)})") params.extend(exclude_sources) if cwd_prefix: clause, clause_params = _cwd_prefix_clause(cwd_prefix) @@ -127,19 +135,16 @@ def _session_filter_where( def _collect_delegate_child_ids(conn, parent_ids: List[str]) -> List[str]: - """Delegate-subagent ids (``_delegate_from`` marker) to cascade-delete with - *parent_ids*; untagged children keep the orphan-don't-delete contract. - Walks marker chains recursively so an orchestrator's own delegates go too.""" + """Delegate-subagent ids (``_delegate_from`` marker, walked recursively) to + cascade-delete with *parent_ids*; untagged children stay orphaned, not deleted.""" df = _delegate_from_json() seeds = {sid for sid in parent_ids if sid} - # Seed visited with the parents: a marker chain can loop back onto a parent - # (cycle, or a parent that is another parent's delegate child in one batch) - # and it would be collected as its own descendant and cascade-deleted. - # Callers delete parents separately; never return them as children. + # Seed visited with the parents: a marker chain can loop back onto a parent, + # which would then be collected as its own descendant. Never return parents. found: set[str] = set(seeds) frontier = list(seeds) while frontier: - ph = ",".join("?" * len(frontier)) + ph = _session_ids_placeholders(frontier) cursor = conn.execute( f"SELECT id FROM sessions WHERE {df} IN ({ph}) " f"OR (parent_session_id IN ({ph}) AND {df} IS NOT NULL)", @@ -153,7 +158,7 @@ def _collect_delegate_child_ids(conn, parent_ids: List[str]) -> List[str]: def _delete_delegate_children(conn, parent_ids: List[str]) -> List[str]: ids = _collect_delegate_child_ids(conn, parent_ids) if ids: - ph = ",".join("?" * len(ids)) + ph = _session_ids_placeholders(ids) conn.execute(f"DELETE FROM messages WHERE session_id IN ({ph})", ids) # FK safety: orphan any untagged stragglers pointing at a doomed row. conn.execute( @@ -163,8 +168,9 @@ def _delete_delegate_children(conn, parent_ids: List[str]) -> List[str]: return ids + # Lifecycle statuses surfaced by session pickers; classified from the final -# message row ONLY (role, tool_calls, finish_reason) so it stays O(1) per session. +# message row ONLY so it stays O(1) per session. SESSION_STATUS_COMPLETE = "complete" SESSION_STATUS_INTERRUPTED = "interrupted" SESSION_STATUS_ERROR = "error" @@ -177,10 +183,9 @@ _ERROR_FINISH_REASONS = frozenset({"error", "agent_error", "content_filter"}) def classify_session_status( role: Optional[str], has_tool_calls: bool, finish_reason: Optional[str], ) -> str: - """Lifecycle from the final message: error finish → ``error``; assistant - with pending tool_calls (result never landed), or a trailing user/tool row → - ``interrupted``; normal assistant finish or unknown shape → ``complete`` - (benign default; pickers must not alarm on unknown shapes).""" + """Error finish → ``error``; assistant with pending tool_calls or a trailing + user/tool row → ``interrupted``; otherwise ``complete`` (benign default: + pickers must not alarm on unknown shapes).""" if (finish_reason or "").strip().lower() in _ERROR_FINISH_REASONS: return SESSION_STATUS_ERROR r = (role or "").strip().lower() @@ -191,9 +196,8 @@ def classify_session_status( return SESSION_STATUS_COMPLETE -# Parent→child profile_name inheritance fence: keyless rows (CLI / subagent) -# inherit freely; two ``agent::...`` keyed rows must agree on the namespace -# so a default child forked from a sibling profile's row isn't mislabelled. +# Parent→child profile_name inheritance fence: keyless rows inherit freely; two +# ``agent::...`` keyed rows must agree on the namespace. _SAME_KEY_NAMESPACE_SQL = ( "p.session_key IS NULL OR sessions.session_key IS NULL" " OR substr(p.session_key, 1, instr(substr(p.session_key, 7), ':') + 6)" @@ -208,10 +212,9 @@ class SessionSessionsMixin: def _own_profile_name(self) -> Optional[str]: """The profile owning THIS store, from ``db_path`` alone (``/state.db`` - → default, ``/profiles//state.db`` → name). Path-based, not - get_active_profile_name(): a gateway serving a NON-launch profile opens - that profile's store and rows must carry the store's owner. None outside - the profile tree — keep NULL rather than a fabricated owner.""" + → default, ``/profiles//state.db`` → name); path-based because a + gateway serving a NON-launch profile opens that profile's store. None + outside the profile tree — NULL beats a fabricated owner.""" try: from hermes_constants import get_default_hermes_root root = get_default_hermes_root().resolve() @@ -226,14 +229,12 @@ class SessionSessionsMixin: @staticmethod def _inherit_parent_session_metadata(conn, session_id: str) -> None: - """NULL-fill a child's cwd/git/profile from its parent (child creators - didn't propagate them, so lineages dropped out of the project sidebar); - profile_name only within the same ``agent::`` namespace. The second - UPDATE inherits gateway routing columns ONLY for compression forks: a - crash before the gateway re-records the peer would otherwise strand the - child unroutable, while delegate children are spawned under a live - parent and must NOT inherit routing keys (peer recovery could repoint - gateway traffic into a subagent's session).""" + """NULL-fill a child's cwd/git/profile from its parent (profile_name only + within the same ``agent::`` namespace). The second UPDATE inherits + gateway routing columns ONLY for compression forks (a crash before the + gateway re-records the peer would strand the child unroutable); delegate + children must NOT inherit them (peer recovery could repoint gateway + traffic into a subagent's session).""" conn.execute( f"""UPDATE sessions SET cwd = COALESCE(sessions.cwd, @@ -291,13 +292,11 @@ class SessionSessionsMixin: parent_session_id: str = None, cwd: str = None, profile_name: Optional[str] = None, git_repo_root: str = None, origin_json: str = None, display_name: str = None, ) -> None: - """Upsert a session row, COALESCE-filling NULL columns and never - overwriting what an earlier writer set (the gateway creates a bare row - before the agent's create_session carries the real model/prompt; a later - bare source="unknown" cannot clobber it). chat_id/thread_id scope gateway - /resume (IDOR). Children backfill from the parent - (:meth:`_inherit_parent_session_metadata`); a missing profile_name is - stamped with THIS store's own profile (NULL reads as unowned).""" + """Upsert a session row, COALESCE-filling NULL columns and never overwriting + what an earlier writer set (the gateway creates a bare row before + create_session carries the real model/prompt). chat_id/thread_id scope + gateway /resume (IDOR). Children backfill from the parent; a missing + profile_name is stamped with THIS store's own (NULL reads as unowned).""" if not (profile_name or "").strip(): profile_name = self._own_profile_name() def _do(conn): @@ -362,7 +361,7 @@ class SessionSessionsMixin: self._delete_unreferenced_system_prompts(conn) if parent_session_id: self._inherit_parent_session_metadata(conn, session_id) - # Transcript-critical: a failed row creation aborts the turn. Ride out long holds. + # Transcript-critical: a failed row creation aborts the turn. self._execute_write(_do, patience_s=self._TRANSCRIPT_WRITE_PATIENCE_S) def create_session(self, session_id: str, source: str, **kwargs) -> str: @@ -379,16 +378,13 @@ class SessionSessionsMixin: (1 if finalized else 0, session_id), ) - # ── Gateway routing index (replaces sessions.json) ──── - def find_session_by_origin( self, *, platform: str, chat_id: str, thread_id: Optional[str] = None, user_id: Optional[str] = None, ) -> Optional[str]: """Most recent live session_id for source + chat_id (+ thread_id). With - ``user_id``, exact sender matches win; if several distinct users share - the chat and none matches, None rather than contaminating another - participant's session.""" + ``user_id`` exact sender matches win; several distinct users and no match + → None rather than contaminating another participant's session.""" if not platform or chat_id in (None, ""): return None query = """ @@ -418,21 +414,14 @@ class SessionSessionsMixin: return None return str(rows[0]["id"]) - # ── Orphaned gateway-session repair (``hermes sessions repair-routing``) ── - # A write-path failure between routing publication and row creation leaves - # the live transcript in a row without identity columns, invisible to - # recovery (the chat resolves to a days-older keyed row). Widest plausible - # gap between a keyed predecessor going quiet and its unkeyed successor: - # the reported incident was ~60s; 15 minutes stays generous without - # spanning unrelated conversations. + # Orphaned gateway-session repair: widest plausible gap between a keyed + # predecessor going quiet and its unkeyed successor (incident was ~60s; 15 min + # stays generous without spanning unrelated conversations). _ORPHAN_ADOPTION_MAX_GAP_S = 900.0 - # Children with a ``parent_session_id`` that are NOT compression - # continuations (branches, delegate runs, tool sessions). Markers are bound - # to the queried parent id: compression continuations inherit the rotated - # agent's model_config verbatim, so a delegate's continuation carries - # ``_delegate_from=`` and presence-matching - # misclassified real continuations as delegate children. + # Children that are NOT compression continuations (branches, delegates, tool + # sessions). Markers are bound to the queried parent id: continuations inherit + # model_config verbatim, so presence-matching misclassified them as delegates. _NON_CONTINUATION_CHILD_FILTER_SQL = ( " AND COALESCE(json_extract(COALESCE({alias}model_config, '{{}}')," " '$._branched_from'), '') != ?\n" @@ -441,10 +430,9 @@ class SessionSessionsMixin: ) def end_session(self, session_id: str, end_reason: str) -> None: - """Mark a session ended. The first end_reason wins (no-op when already - ended): a compression split must keep ``'compression'`` even if a stale - desynced-CLI end_session() targets it later. reopen_session() first to - deliberately re-end with a new reason.""" + """Mark a session ended; the first end_reason wins (a compression split must + keep ``'compression'`` even if a stale end_session() targets it later). + reopen_session() first to deliberately re-end with a new reason.""" def _do(conn): changed = conn.execute( "UPDATE sessions SET ended_at = ?, end_reason = ? " @@ -457,12 +445,11 @@ class SessionSessionsMixin: self._execute_write(_do) def reopen_session(self, session_id: str) -> None: - """Clear ended_at/end_reason so a session can be resumed. First stamp + """Clear ended_at/end_reason so a session can be resumed; first stamp markerless legacy reset children that depend on the parent's mutable - end_reason (WHERE shared with the listing predicate via - _legacy_reset_child_sql so the two cannot drift).""" + end_reason (WHERE shared with the listing predicate so they cannot drift).""" def _do(conn): - placeholders = ",".join("?" for _ in _RESET_END_REASONS) + placeholders = _session_ids_placeholders(_RESET_END_REASONS) conn.execute( "UPDATE sessions AS child SET model_config = json_set(" "COALESCE(child.model_config, '{}'), '$._reset_from', child.parent_session_id) " @@ -480,11 +467,9 @@ class SessionSessionsMixin: def promote_to_session_reset(self, session_id: str, reason: str = "session_reset") -> bool: """Durably mark an intentional reset boundary on live rows or rows with a - *recoverable* accidental end_reason; explicit boundaries are preserved - (first writer wins). Plain end_session() no-ops on an ended row, so an - ``agent_close`` row would stay recoverable and stale-route recovery would - resurrect the reset session. Keep in sync with - find_latest_gateway_session_for_peer. True when promoted.""" + *recoverable* accidental end_reason (explicit boundaries are preserved): + an ``agent_close`` row left recoverable would be resurrected by + stale-route recovery. Keep in sync with find_latest_gateway_session_for_peer.""" if not session_id: return False now = time.time() @@ -509,13 +494,11 @@ class SessionSessionsMixin: self, session_id: str, cwd: str, git_branch: Optional[str] = None, git_repo_root: Optional[str] = None, replace_git_meta: bool = False, ) -> Optional[int]: - """Persist the authoritative cwd and claim a Git metadata generation. - git fields are written only when non-empty (a probe failure never - clobbers a captured value) except under ``replace_git_meta`` (a - workspace MOVE must overwrite the old repo identity even when the new - cwd has none). Each call bumps ``git_metadata_generation``; async probes - publish via :meth:`publish_session_git_metadata` with that generation so - an older worker cannot overwrite a newer claim (A -> B -> A).""" + """Persist the authoritative cwd and claim a Git metadata generation. git + fields are written only when non-empty (a probe failure never clobbers a + value) except under ``replace_git_meta`` (a workspace MOVE overwrites the + old repo identity). Async probes publish with the returned generation so an + older worker cannot overwrite a newer claim (A -> B -> A).""" if not session_id or not cwd: return None branch = (git_branch or "").strip() @@ -532,12 +515,11 @@ class SessionSessionsMixin: if current_cwd != cwd or replace_git_meta: sets.extend(("git_branch = ?", "git_repo_root = ?")) params.extend((branch or None, repo_root or None)) - elif branch: - sets.append("git_branch = ?") - params.append(branch) - if repo_root and current_cwd == cwd and not replace_git_meta: - sets.append("git_repo_root = ?") - params.append(repo_root) + else: # same cwd: only overwrite with captured (non-empty) values + for col, val in (("git_branch", branch), ("git_repo_root", repo_root)): + if val: + sets.append(f"{col} = ?") + params.append(val) params.append(session_id) conn.execute(f"UPDATE sessions SET {', '.join(sets)} WHERE id = ?", params) row = conn.execute( @@ -551,36 +533,24 @@ class SessionSessionsMixin: git_repo_root: Optional[str] = None, ) -> bool: """Publish async Git enrichment only while its cwd claim is current.""" - if ( - not session_id - or not cwd - or isinstance(generation, bool) - or not isinstance(generation, int) - or generation < 1 - ): + valid_generation = isinstance(generation, int) and not isinstance(generation, bool) and generation >= 1 + if not session_id or not cwd or not valid_generation: return False - branch = (git_branch or "").strip() - repo_root = (git_repo_root or "").strip() - if not branch and not repo_root: + fields = [ + (col, val) for col, val in ( + ("git_branch", (git_branch or "").strip()), ("git_repo_root", (git_repo_root or "").strip()), + ) if val + ] + if not fields: return False - sets: List[str] = [] - params: List[Any] = [] - if branch: - sets.append("git_branch = ?") - params.append(branch) - if repo_root: - sets.append("git_repo_root = ?") - params.append(repo_root) - params.extend((session_id, cwd, generation)) return self._write_rowcount( - f"UPDATE sessions SET {', '.join(sets)} " + f"UPDATE sessions SET {', '.join(f'{col} = ?' for col, _ in fields)} " "WHERE id = ? AND cwd = ? AND git_metadata_generation = ?", - params, + [val for _, val in fields] + [session_id, cwd, generation], ) == 1 def backfill_repo_roots(self, cwd_to_root: Dict[str, str]) -> None: - """Backfill git repo roots for cwds without one (pre-column sessions); - never clobbers a recorded root; empty roots are skipped.""" + """Backfill git repo roots for cwds without one; never clobbers a recorded root.""" pairs = [(root, cwd) for cwd, root in cwd_to_root.items() if root and cwd] if pairs: self._write_sql( @@ -589,22 +559,15 @@ class SessionSessionsMixin: pairs, many=True, ) - # Compression locks (atomic per-session, keyed by session_id, recovered via - # expires_at) live in SessionCompressionMixin; they stop two AIAgents that - # share a session_id from both rotating it into two orphan children. - def touch_session_activity( self, session_id: str, ts: Optional[float] = None, *, description: Optional[str] = None, provenance: Optional[ActivityProvenance] = None, ) -> None: - """Stamp durable mid-turn activity (observation-only; rate-limited by - AIAgent._touch_activity) so surfaces see API/tool/compaction activity - before any message row lands. Never moves ``last_activity_at`` backwards.""" + """Stamp durable mid-turn activity (observation-only; rate-limited by the + caller) so surfaces see activity before any message row lands. Never moves + ``last_activity_at`` backwards.""" if not session_id: return - from agent.session_activity import ( - bound_activity_description, normalize_activity_provenance, - ) when = float(ts if ts is not None else time.time()) desc = bound_activity_description(description) prov = normalize_activity_provenance(provenance).value @@ -617,17 +580,13 @@ class SessionSessionsMixin: ) def clear_session_activity_labels(self, session_id: str) -> None: - """Clear activity labels after a turn (keep ``last_activity_at`` so idle / - watchdog clocks stay continuous; an idle turn must not keep advertising - "compressing"). Runs in the turn's finally: a no-op clear skips the - write transaction, a real one uses the short activity budget.""" + """Clear activity labels after a turn (``last_activity_at`` is kept so idle / + watchdog clocks stay continuous). A no-op clear skips the write transaction.""" if not session_id: return - from agent.session_activity import ActivityProvenance try: row = self._read_one( - "SELECT last_activity_description, last_activity_provenance " - "FROM sessions WHERE id = ?", + "SELECT last_activity_description, last_activity_provenance FROM sessions WHERE id = ?", (session_id,), ) except sqlite3.Error: @@ -646,7 +605,6 @@ class SessionSessionsMixin: row = self.get_session(session_id) if session_id else None if not row: return None - from agent.session_activity import build_activity_snapshot return build_activity_snapshot( last_activity_at=row.get("last_activity_at"), last_activity_description=row.get("last_activity_description"), @@ -675,27 +633,21 @@ class SessionSessionsMixin: self._execute_write(_do) def update_session_tool_names(self, session_id: str, tool_names: Optional[List[str]]) -> None: - """Persist the resolved ``tools[]`` name order so a rebuilt AIAgent - (agent-cache eviction) can't fork the cached tool prefix on a flipped - check_fn verdict. ``None`` clears the pin.""" + """Persist the resolved ``tools[]`` name order so a rebuilt AIAgent can't + fork the cached tool prefix on a flipped check_fn verdict; ``None`` clears.""" payload = json.dumps(list(tool_names)) if tool_names is not None else None self._write_sql("UPDATE sessions SET tool_names = ? WHERE id = ?", (payload, session_id)) def update_session_model( self, session_id: str, model: str, provider: Optional[str] = None ) -> None: - """Set the model after a mid-session /model switch (unconditionally, - unlike update_token_counts' COALESCE), null system_prompt so stale - Model:/Provider: footers rebuild, and replace any confirmed Browser - runtime lock while keeping lineage markers. *provider* is merged into - model_config so resume recombines the model with the provider that - actually serves it, not the config.yaml primary.""" - # This write bypasses the token queue: a still-queued first delta carries - # the pre-switch route and, applied after this UPDATE, would trip the - # first_accounted_route overwrite and resurrect the old model/provider. + """Set the model after a mid-session /model switch (unconditionally), null + system_prompt so stale Model:/Provider: footers rebuild, and drop any + Browser runtime lock (lineage markers survive). *provider* is merged into + model_config so resume recombines the model with the provider that serves it.""" + # Flush first: a still-queued pre-switch delta applied after this UPDATE + # would trip the first_accounted_route overwrite and resurrect the old route. self.flush_token_counts() - # browser_model_lock is deleted via a None patch value (same semantics - # as the old json_remove); lineage markers survive the merge. patch: Dict[str, Any] = {"browser_model_lock": None} if model: patch["model"] = model @@ -710,20 +662,18 @@ class SessionSessionsMixin: ) def _write_model_config_patch( - self, session_id: str, patch: Dict[str, Any], sql: str, - params: Callable[[Optional[str]], tuple], *, clear_prompts: bool = False, + self, session_id: str, patch: Dict[str, Any], + sql: str = "UPDATE sessions SET model_config = ? WHERE id = ?", + params: Optional[Callable[[Optional[str]], tuple]] = None, *, clear_prompts: bool = False, ) -> None: - """Merge ``patch`` into model_config then run ``sql`` with ``params(merged)``. - - One write transaction; no-op when the row doesn't exist. ``clear_prompts`` - additionally garbage-collects unreferenced system_prompts (for writers - that NULL the row's system_prompt_hash). - """ + """Merge ``patch`` into model_config then run ``sql`` with ``params(merged)`` + (default: plain model_config UPDATE) in one write transaction; no-op when + the row doesn't exist. ``clear_prompts`` also GCs unreferenced system_prompts.""" def _do(conn): merged = self._merge_model_config_json(conn, session_id, patch) if merged is _MODEL_CONFIG_ROW_MISSING: return - conn.execute(sql, params(merged)) + conn.execute(sql, params(merged) if params else (merged, session_id)) if clear_prompts: self._delete_unreferenced_system_prompts(conn) self._execute_write(_do) @@ -731,12 +681,10 @@ class SessionSessionsMixin: def _merge_model_config_json( self, conn, session_id: str, patch: Dict[str, Any], *, on_missing: str = "skip", ): - """SELECT + tolerant-parse + merge ``patch`` into model_config — the one - place the merge discipline keeping ``_branched_from``/``_delegate_from`` - alive lives. ``None`` deletes a key. Runs inside the caller's write - transaction. Returns serialized JSON (``None`` when empty, matching - create_session's NULL) or ``_MODEL_CONFIG_ROW_MISSING`` when the row - doesn't exist (``on_missing="raise"`` raises ValueError instead).""" + """SELECT + tolerant-parse + merge ``patch`` into model_config (the one place + that keeps ``_branched_from``/``_delegate_from`` alive); ``None`` deletes a + key. Returns serialized JSON (``None`` when empty, matching create_session's + NULL) or ``_MODEL_CONFIG_ROW_MISSING`` (``on_missing="raise"`` → ValueError).""" row = conn.execute( "SELECT model_config FROM sessions WHERE id = ?", (session_id,), ).fetchone() @@ -754,14 +702,10 @@ class SessionSessionsMixin: def patch_session_model_config(self, session_id: str, patch: Dict[str, Any]) -> None: """Merge ``patch`` into model_config atomically (``None`` removes a key); - no-op when the row or patch is empty. The transcript-coupled path is - archive_and_compact's ``model_config_patch``.""" + no-op when the row or patch is empty.""" if not session_id or not patch: return - self._write_model_config_patch( - session_id, patch, "UPDATE sessions SET model_config = ? WHERE id = ?", - lambda merged: (merged, session_id), - ) + self._write_model_config_patch(session_id, patch) def get_session_model_config_value(self, session_id: str, key: str, default: Any = None) -> Any: """Read one key out of a session's model_config JSON (tolerant parse).""" @@ -793,21 +737,16 @@ class SessionSessionsMixin: ) def set_session_yolo(self, session_id: str, enabled: bool) -> None: - """Persist the per-session YOLO flag into model_config so ``/yolo`` or - ``--yolo`` survives ``hermes --resume``. No-op when the row doesn't exist - yet (creation-time model_config carries the flag for --yolo launches).""" + """Persist the per-session YOLO flag so ``/yolo`` survives ``--resume``; + no-op when the row doesn't exist yet.""" if not session_id: return - self._write_model_config_patch( - session_id, {"yolo_mode": bool(enabled)}, - "UPDATE sessions SET model_config = ? WHERE id = ?", - lambda merged: (merged, session_id), - ) + self._write_model_config_patch(session_id, {"yolo_mode": bool(enabled)}) @staticmethod def session_yolo_enabled(session_meta: Optional[Dict[str, Any]]) -> bool: - """Persisted YOLO flag from a session row (JSON string or parsed dict); - False on any parse failure — resume must never enable the bypass by accident.""" + """Persisted YOLO flag; False on any parse failure (resume must never + enable the bypass by accident).""" return bool(_parse_model_config((session_meta or {}).get("model_config")).get("yolo_mode")) def ensure_session( @@ -829,9 +768,8 @@ class SessionSessionsMixin: return self._session_row_dict(row) if row else None def get_dominant_session_model_route(self, session_id: str) -> Optional[Dict[str, Any]]: - """Main-loop model route that served most API calls. ``sessions`` is a - legacy aggregate mixing route changes; ``session_model_usage`` keeps the - coherent per-call tuple, so status/billing reads prefer it.""" + """Main-loop model route that served most API calls (``session_model_usage`` + keeps the coherent per-call tuple; ``sessions`` mixes route changes).""" self.flush_token_counts() row = self._read_one( """SELECT model, billing_provider, billing_base_url, billing_mode, @@ -863,11 +801,9 @@ class SessionSessionsMixin: return matches[0] if len(matches) == 1 else None def backfill_null_session_profiles(self, profile_name: str) -> int: - """Stamp this store's own profile onto legacy ``profile_name IS NULL`` - rows, which the fail-closed owner ladder cannot route once a Desktop - registers a second connection (pre-ownership sessions became - unresumable). Single-match, not a guess: a store belongs to exactly one - profile. Never overwrites a non-NULL owner; idempotent. Returns rows stamped.""" + """Stamp this store's own profile onto legacy ``profile_name IS NULL`` rows, + which the fail-closed owner ladder cannot route (a store belongs to exactly + one profile). Never overwrites a non-NULL owner. Returns rows stamped.""" stamp = (profile_name or "").strip() if not stamp: return 0 @@ -879,10 +815,9 @@ class SessionSessionsMixin: ) or 0) def _set_lineage_column(self, column: str, session_id: str, value: Any) -> bool: - """Set one ``sessions`` column across a whole compression lineage - (ancestors + descendants joined by end_reason='compression'): Desktop - projects roots forward to their tip, and updating only the displayed tip - would let the untouched root resurrect it on refresh. True if any row changed.""" + """Set one ``sessions`` column across a whole compression lineage: Desktop + projects roots forward to their tip, so updating only the displayed tip + would let the untouched root resurrect it on refresh.""" return self._write_rowcount( f""" WITH RECURSIVE @@ -917,19 +852,16 @@ class SessionSessionsMixin: ) > 0 def set_session_archived(self, session_id: str, archived: bool) -> bool: - """Soft-hide (or unhide) a session and its whole compression lineage; - messages are kept. True when at least one row changed.""" + """Soft-hide (or unhide) a session and its compression lineage; messages are kept.""" return self._set_lineage_column('archived', session_id, 1 if archived else 0) - # Accidental end reasons recovery treats as resumable; the same constant is - # interpolated into the recovery/promotion SQL so literals cannot drift. + # Accidental end reasons recovery treats as resumable (also interpolated into + # the recovery/promotion SQL so literals cannot drift). RECOVERABLE_END_REASONS = _RECOVERABLE_END_REASONS def unarchive_recoverable_session(self, session_id: str) -> bool: """Un-archive a session archived by a recoverable accident (ws_orphan_reap, - agent_close) — used by registry lookups like Bot Mode's canonical chat. - Deliberate archives (no end_reason, or an explicit boundary) are left - alone. True only when a recoverable row was un-archived (whole lineage).""" + agent_close); deliberate archives are left alone. True when un-archived.""" if not session_id: return False try: @@ -938,44 +870,39 @@ class SessionSessionsMixin: return False if not row or not row.get("archived"): return False - # The accidental stamp lives on the live TIP (the registry row keeps - # end_reason='compression'); judge recoverability there. + # The accidental stamp lives on the live TIP; judge recoverability there. tip = row try: tip_id = self.get_compression_tip(session_id) or session_id if tip_id != session_id: tip = self.get_session(tip_id) or row except Exception: - tip_id = session_id + pass if (tip.get("end_reason") or "") not in self.RECOVERABLE_END_REASONS: return False if not self.set_session_archived(session_id, False): return False - # Clear the accidental end stamp, or a LATER deliberate archive (which - # never writes end_reason) would auto-resurrect on the next lookup. + # Clear the accidental end stamp, or a LATER deliberate archive (which never + # writes end_reason) would auto-resurrect on the next lookup. self._write_sql( "UPDATE sessions SET ended_at = NULL, end_reason = NULL WHERE id = ?", (tip["id"],), ) return True def set_session_pinned(self, session_id: str, pinned: bool) -> bool: - """Pin/unpin a session and its compression lineage. Pinned sessions are - exempt from the ``sessions.auto_archive`` sweep; Desktop mirrors its - sidebar pins here so backend sweeps honour them.""" + """Pin/unpin a session and its compression lineage (pins are exempt from the + ``sessions.auto_archive`` sweep).""" return self._set_lineage_column('pinned', session_id, 1 if pinned else 0) def set_session_hidden(self, session_id: str, hidden: bool) -> bool: - """Hide/unhide a session and its compression lineage from the default - list_sessions_rich listing; it stays resumable by the owning surface - (plugins such as kanban manage their own sessions).""" + """Hide/unhide a session and its compression lineage from the default listing; + it stays resumable by the owning surface.""" return self._set_lineage_column('hidden', session_id, 1 if hidden else 0) def set_session_read(self, session_id: str, read: bool = True) -> bool: """Mark read/unread across the compression lineage. ``last_read_at`` is a - watermark, not a flag: unread when activity postdates it, so new - messages flip it back without any write on the message path. NULL = - never tracked = read (shipping the column doesn't badge all history); - 0 = explicitly unread; timestamp = read up to then.""" + watermark: unread when activity postdates it (no write on the message + path). NULL = never tracked = read; 0 = explicitly unread.""" return self._set_lineage_column('last_read_at', session_id, time.time() if read else 0.0) @staticmethod @@ -987,8 +914,7 @@ class SessionSessionsMixin: last_active = session_row.get("last_active") or session_row.get("started_at") return float(last_active or 0) > float(last_read) - # compact_rows excludes only payload-heavy blobs no list consumer renders; - # the projection derives from SCHEMA_SQL so new columns join automatically. + # compact_rows excludes only payload-heavy blobs no list consumer renders. _SESSION_COMPACT_EXCLUDED = frozenset( {"system_prompt", "system_prompt_hash", "git_metadata_generation"} ) @@ -997,10 +923,9 @@ class SessionSessionsMixin: @staticmethod def _chain_search_where(where_sql: str, id_needle: str, search_needle: str) -> Tuple[str, List[Any]]: """Extend ``where_sql`` with the id_query / search_query filters: a row is - admitted when its own id or any id in its forward compression chain - matches (search also matches titles and a punctuation-stripped form so - ``an94`` finds ``AN-94``). Leading-wildcard LIKE can't use an index but - chain membership keeps it bounded — far cheaper than scanning in Python.""" + admitted when its own id or any id in its forward compression chain matches + (search also matches titles and a punctuation-stripped form so ``an94`` + finds ``AN-94``); chain membership keeps the leading-wildcard LIKE bounded.""" params: List[Any] = [] clauses: List[str] = [] def _like_pattern(needle: str) -> str: @@ -1033,11 +958,9 @@ class SessionSessionsMixin: return (f"{where_sql} AND {combined}" if where_sql else f"WHERE {combined}"), params def _project_compression_tips(self, sessions: List[Dict[str, Any]], compact_rows: bool) -> List[Dict[str, Any]]: - """Replace each compression root's surfaced fields with its live tip's - (root ``started_at`` kept for stable ordering); tip rows are fetched in - one batched query. ``_lineage_ids`` carries every id on the chain: a - persisted tile can hold a MIDDLE segment's id, and with only root/tip a - surface cannot prove it names this conversation (one chat open twice).""" + """Replace each compression root's surfaced fields with its live tip's (root + ``started_at`` kept for stable ordering), one batched query. ``_lineage_ids`` + carries every id on the chain: a persisted tile can hold a MIDDLE segment's id.""" tip_ids_by_root: Dict[str, str] = {} chain_by_root: Dict[str, List[str]] = {} for s in sessions: @@ -1089,13 +1012,10 @@ class SessionSessionsMixin: compact_rows: bool = False, include_pinned: bool = False, session_key: str = None, include_hidden: bool = False, ) -> List[Dict[str, Any]]: - """List sessions with preview and ``last_active`` in one query. Subagent - runs / compression continuations are hidden unless ``include_children``; - ``project_compression_tips`` shows each chain as its live tip; - ``order_by_last_active`` sorts by the chain TIP via a recursive CTE (the - only path honouring ``id_query`` / ``search_query``); ``compact_rows`` - omits the system_prompt blob; ``include_pinned`` back-fills pins the page - missed ("always reachable"), still obeying the other filters.""" + """List sessions with preview and ``last_active`` in one query. + ``order_by_last_active`` sorts by the chain TIP via a recursive CTE (the only + path honouring ``id_query`` / ``search_query``); ``include_pinned`` back-fills + pins the page missed, still obeying the other filters.""" self.flush_token_counts() # rows carry token/cost totals where_clauses, params = _session_filter_where( exclude_children=not include_children, source=source, sources=sources, @@ -1105,7 +1025,7 @@ class SessionSessionsMixin: ) if not include_hidden: where_clauses.append("s.hidden = 0") - where_sql = f"WHERE {' AND '.join(where_clauses)}" if where_clauses else "" + where_sql = _where_sql(where_clauses) base_where_params = list(params) # pinned back-fill reuses the WHERE before LIMIT/OFFSET prompt_select = ( "" if compact_rows @@ -1119,12 +1039,10 @@ class SessionSessionsMixin: id_needle = (id_query or "").strip().lower() search_needle = (search_query or "").strip().lower() if order_by_last_active: - # The CTE seeds from rows the outer WHERE admits and walks - # compression-continuation edges forward; MAX over the chain gives - # effective_last_active so ORDER BY + LIMIT happen in SQL. Do NOT + # The CTE walks compression-continuation edges forward from the admitted + # rows; MAX over the chain gives effective_last_active in SQL. Do NOT # require child.started_at >= parent.ended_at: races insert the - # continuation before the parent's ended_at is written, while stale - # websocket siblings could pass the timestamp test and hijack projection. + # continuation before ended_at is written. outer_where, id_params = self._chain_search_where(where_sql, id_needle, search_needle) query = f""" WITH RECURSIVE chain(root_id, cur_id) AS ( @@ -1171,8 +1089,8 @@ class SessionSessionsMixin: """ params.extend([limit, offset]) sessions = [self._list_row(row) for row in self._read_all(query, params)] - # Pinned back-fill runs BEFORE compression projection so a back-filled - # root projects to its tip like any other row. One query, never N+1. + # Pinned back-fill runs BEFORE compression projection so a back-filled root + # projects to its tip like any other row. if include_pinned: seen_ids = {s["id"] for s in sessions} pinned_where = (f"{where_sql} AND s.pinned = 1" if where_sql else "WHERE s.pinned = 1") @@ -1201,14 +1119,13 @@ class SessionSessionsMixin: return sessions def session_lifecycle_statuses(self, session_ids: List[str]) -> Dict[str, str]: - """``{session_id: status}`` from each session's LAST message row (see - :func:`classify_session_status`; ``'empty'`` when no messages). One query: - MAX(id) per session (index seek) joined back for that row — never scans transcripts.""" + """``{session_id: status}`` from each session's LAST message row (``'empty'`` + when none); one query, MAX(id) per session joined back — never scans transcripts.""" ids = [sid for sid in (session_ids or []) if sid] if not ids: return {} statuses: Dict[str, str] = {sid: "empty" for sid in ids} - placeholders = ",".join("?" for _ in ids) + placeholders = _session_ids_placeholders(ids) query = f""" SELECT m.session_id, m.role, m.tool_calls IS NOT NULL AS has_tool_calls, @@ -1230,10 +1147,9 @@ class SessionSessionsMixin: return statuses def assert_export_safe(self, session_id: str, max_messages: Optional[int] = None) -> int: - """Active row count of this segment (compression ancestors excluded), or - raise SessionExportTooLargeError. The LIMITed subquery stops once it - proves the bound is exceeded. ``None`` resolves ``sessions.max_export_messages``; - 0 disables the guard (returns 0 without counting).""" + """Active row count of this segment, or raise SessionExportTooLargeError (the + LIMITed subquery stops once the bound is exceeded). ``None`` resolves + ``sessions.max_export_messages``; 0 disables the guard.""" from hermes_state import SessionExportTooLargeError, resolved_max_export_messages if max_messages is None: max_messages = resolved_max_export_messages() @@ -1252,8 +1168,8 @@ class SessionSessionsMixin: return message_count def _is_explicit_branch_session(self, session_id: str) -> bool: - """Copied user-facing branch (``_branched_from`` marker)? Branches own a - copied transcript; compression continuations need the parent's archived rows.""" + """Copied user-facing branch (``_branched_from``)? Branches own a copied + transcript; compression continuations need the parent's archived rows.""" if not session_id: return False row = self._read_one("SELECT model_config FROM sessions WHERE id = ?", (session_id,)) @@ -1284,9 +1200,8 @@ class SessionSessionsMixin: def search_sessions( self, source: str = None, limit: int = 20, offset: int = 0, workspace_key: str = None, ) -> List[Dict[str, Any]]: - """Sessions MRU-first with a computed ``last_active``; ``workspace_key`` - scopes to one workspace (:func:`workspace_key` semantics) so - ``hermes -c``/``--resume`` picks the current workspace's last session.""" + """Sessions MRU-first with a computed ``last_active``; ``workspace_key`` scopes + to one workspace so ``hermes -c``/``--resume`` picks its last session.""" select_with_last_active = ( "SELECT s.*, COALESCE(sp.prompt, s.system_prompt) AS _system_prompt_resolved, " f"{_sql_session_last_active('s')} AS last_active " @@ -1301,7 +1216,7 @@ class SessionSessionsMixin: ws_clause, ws_params = _workspace_key_clause(workspace_key) where_clauses.append(ws_clause) params.extend(ws_params) - where_sql = f" WHERE {' AND '.join(where_clauses)}" if where_clauses else "" + where_sql = _where_sql(where_clauses, " ") params.extend([limit, offset]) return [self._session_row_dict(row) for row in self._read_all( f"{select_with_last_active}{where_sql} " @@ -1314,22 +1229,21 @@ class SessionSessionsMixin: min_message_count: int = 0, include_archived: bool = False, archived_only: bool = False, exclude_children: bool = False, exclude_sources: List[str] = None, ) -> int: - """Count sessions with the same filters as list_sessions_rich, so a - paired "load more" total matches the listable rows (children or a - cron-excluded page would otherwise inflate it and never settle).""" + """Count sessions with list_sessions_rich's filters so a paired "load more" + total matches the listable rows.""" where_clauses, params = _session_filter_where( exclude_children=exclude_children, source=source, sources=sources, exclude_sources=exclude_sources, cwd_prefix=cwd_prefix, min_message_count=min_message_count, archived_only=archived_only, include_archived=include_archived, ) - where_sql = f" WHERE {' AND '.join(where_clauses)}" if where_clauses else "" - return self._read_one(f"SELECT COUNT(*) FROM sessions s{where_sql}", params)[0] + return self._read_one( + f"SELECT COUNT(*) FROM sessions s{_where_sql(where_clauses, ' ')}", params, + )[0] def session_count_ge(self, n: int = 1) -> bool: - """At least N sessions exist (archived included — "has this install ever - had sessions"). LIMIT short-circuits: 4us vs session_count()'s 543us - index scan on a 20k-session DB.""" + """At least N sessions exist (archived included); LIMIT short-circuits + instead of session_count()'s index scan.""" rows = self._read_all("SELECT 1 FROM sessions LIMIT ?", (n,)) return len(rows) >= n @@ -1337,14 +1251,12 @@ class SessionSessionsMixin: self, *, include_archived: bool = False, archived_only: bool = False, exclude_children: bool = False, ) -> Dict[str, int]: - """``{source: count}`` via one GROUP BY (uses idx_sessions_source unless - ``exclude_children``, whose predicates need a table scan like - list_sessions_rich). ``exclude_children`` mirrors listing visibility.""" + """``{source: count}`` via one GROUP BY; ``exclude_children`` mirrors listing visibility.""" where_clauses, params = _session_filter_where( exclude_children=exclude_children, archived_only=archived_only, include_archived=include_archived, ) - where_sql = f" WHERE {' AND '.join(where_clauses)}" if where_clauses else "" + where_sql = _where_sql(where_clauses, " ") with self._read_ctx() as conn: if self._conn is None: raise RuntimeError("SessionDB connection is closed") @@ -1357,9 +1269,8 @@ class SessionSessionsMixin: return {str(row["source"]): int(row["count"] or 0) for row in rows} def declared_scope_identity(self, session_id: str) -> Tuple[bool, str]: - """(is_fork_child, source) for *session_id* in ONE read — prompt_cache_scope - needs both from the same row. A missing row is (False, ""); DB errors - propagate so the caller fails closed.""" + """(is_fork_child, source) in ONE read (prompt_cache_scope needs both from the + same row). Missing row → (False, ""); DB errors propagate (fail closed).""" session = self.get_session(session_id) if not session: return False, "" @@ -1373,7 +1284,6 @@ class SessionSessionsMixin: return targets = [sessions_dir / f"{session_id}{suffix}" for suffix in (".json", ".jsonl")] try: - # request_dump files use session_id as a prefix component targets.extend(sessions_dir.glob(f"request_dump_{session_id}_*.json")) except OSError: pass @@ -1384,12 +1294,11 @@ class SessionSessionsMixin: pass def get_session_delete_targets(self, session_id: str) -> List[str]: - """Rows :meth:`delete_session` would remove: the session, then its - recursive delegate children (branch/compression children are orphaned, not deleted).""" + """Rows :meth:`delete_session` would remove: the session, then its recursive + delegate children (branch/compression children are orphaned, not deleted).""" with self._read_ctx() as conn: if not conn.execute("SELECT 1 FROM sessions WHERE id = ? LIMIT 1", (session_id,)).fetchone(): return [] - # The borrowed read connection, never self._conn (unlocked writer use). delegate_ids = _collect_delegate_child_ids(conn, [session_id]) return [session_id, *sorted(delegate_ids)] @@ -1397,12 +1306,10 @@ class SessionSessionsMixin: self, session_id: str, sessions_dir: Optional[Path] = None, expected_delete_ids: Optional[List[str]] = None, ) -> bool: - """Delete a session and its messages. Delegate children cascade (they'd - resurface as orphans in pickers); branch/compression children are - orphaned (parent -> NULL). *sessions_dir*: also remove transcript files. - *expected_delete_ids*: proceed only if parent + delegate cascade still - equals that set (export-before-delete fails closed if a new delegate - appeared); the tree is re-walked inside the transaction on purpose (TOCTOU).""" + """Delete a session and its messages; delegate children cascade, + branch/compression children are orphaned. *expected_delete_ids*: proceed + only if parent + delegate cascade still equals that set (re-walked inside + the transaction on purpose: export-before-delete fails closed).""" removed_delegate_ids: List[str] = [] expected_ids = set(expected_delete_ids) if expected_delete_ids is not None else None def _do(conn): @@ -1428,10 +1335,8 @@ class SessionSessionsMixin: return bool(deleted) def delete_session_if_empty(self, session_id: str, sessions_dir: Optional[Path] = None) -> bool: - """Delete *session_id* only if it has no messages, no title and no - children (a parent that spawned work is not "empty"), so start-and-quit - sessions don't pile up in /resume. Check and delete share one - transaction so a concurrently flushed message can't be lost.""" + """Delete *session_id* only if it has no messages, no title and no children; + check and delete share one transaction so a concurrent flush can't be lost.""" def _do(conn): cursor = conn.execute( """ @@ -1457,28 +1362,22 @@ class SessionSessionsMixin: return bool(deleted) def delete_sessions(self, session_ids: List[str], sessions_dir: Optional[Path] = None) -> int: - """Bulk delete (dashboard multi-select) with :meth:`delete_session` - semantics per row, in ONE transaction so a partial failure can't leave - "messages gone, row still there". Unknown ids are skipped (UI selection - can race another tab's delete: succeed-on-the-rest). Returns the number - that actually existed and were deleted.""" - if not session_ids: - return 0 - unique_ids = list({sid for sid in session_ids if isinstance(sid, str) and sid}) + """Bulk delete with :meth:`delete_session` semantics per row, in ONE + transaction. Unknown ids are skipped (UI selection can race another tab's + delete). Returns the number that existed and were deleted.""" + unique_ids = list({sid for sid in session_ids or () if isinstance(sid, str) and sid}) if not unique_ids: return 0 removed_ids: list[str] = [] - removed_delegate_ids: list[str] = [] def _do(conn): - # Filter to IDs that actually exist: return the real deleted count. existing = [row["id"] for row in conn.execute( - f"SELECT id FROM sessions WHERE id IN ({','.join('?' * len(unique_ids))})", + f"SELECT id FROM sessions WHERE id IN ({_session_ids_placeholders(unique_ids)})", unique_ids, ).fetchall()] if not existing: return 0 - existing_placeholders = ",".join("?" * len(existing)) - removed_delegate_ids.extend(_delete_delegate_children(conn, existing)) + existing_placeholders = _session_ids_placeholders(existing) + removed_ids.extend(_delete_delegate_children(conn, existing)) conn.execute( # orphan children whose parent is in the kill list (FK) f"UPDATE sessions SET parent_session_id = NULL " f"WHERE parent_session_id IN ({existing_placeholders})", @@ -1492,31 +1391,26 @@ class SessionSessionsMixin: removed_ids.extend(existing) return len(existing) count = self._execute_write(_do) - for sid in removed_delegate_ids + removed_ids: + for sid in removed_ids: self._remove_session_files(sessions_dir, sid) return count #: Shared by count_empty_sessions / delete_empty_sessions so badge and sweep - #: agree. ``message_count`` counts live rows only — rewind and compaction - #: reset it to 0 while keeping dropped turns as ``active = 0`` (the only - #: recoverable copy) — so NOT EXISTS is the authority; message_count = 0 is - #: a cheap prefilter. + #: agree. ``message_count`` counts live rows only (rewind/compaction keep + #: dropped turns as ``active = 0``), so NOT EXISTS is the authority. _EMPTY_SESSION_WHERE = ( "message_count = 0 AND ended_at IS NOT NULL AND archived = 0 AND NOT EXISTS (" "SELECT 1 FROM messages WHERE messages.session_id = sessions.id)" ) def count_empty_sessions(self) -> int: - """Count of empty, ended, non-archived sessions (:data:`_EMPTY_SESSION_WHERE`). - The ended_at guard matches prune_sessions: a fresh session whose first - message hasn't landed is never sniped out from under the runtime.""" + """Count of empty, ended, non-archived sessions; the ended_at guard means a + fresh session whose first message hasn't landed is never sniped.""" return self._read_one(f"SELECT COUNT(*) FROM sessions WHERE {self._EMPTY_SESSION_WHERE}")[0] def delete_empty_sessions(self, sessions_dir: Optional[Path] = None) -> int: - """Delete every empty, ended, non-archived session (:data:`_EMPTY_SESSION_WHERE`) - in one transaction, orphaning (not cascading) children so branch/subagent - transcripts survive. Transcript files are swept too: the gateway can - leave a stub request_dump_* if it crashed before the first reply.""" + """Delete every empty, ended, non-archived session in one transaction, + orphaning (not cascading) children; transcript files are swept too.""" removed_ids: list[str] = [] def _do(conn): session_ids = {row["id"] for row in conn.execute( @@ -1526,13 +1420,12 @@ class SessionSessionsMixin: return 0 conn.execute( f"UPDATE sessions SET parent_session_id = NULL " - f"WHERE parent_session_id IN ({','.join('?' * len(session_ids))})", + f"WHERE parent_session_id IN ({_session_ids_placeholders(session_ids)})", list(session_ids), ) for sid in session_ids: - # DELETE FROM messages is paranoia — the selector's NOT EXISTS - # probe proved these own no rows — but a row inserted between - # the SELECT and here would otherwise dangle (clean FK state). + # DELETE FROM messages: a row inserted between the SELECT and here + # would otherwise dangle (clean FK state). conn.execute("DELETE FROM messages WHERE session_id = ?", (sid,)) conn.execute("DELETE FROM sessions WHERE id = ?", (sid,)) removed_ids.append(sid) @@ -1546,9 +1439,8 @@ class SessionSessionsMixin: def archive_sessions( self, older_than_days: Optional[float] = None, source: str = None, **filters, ) -> int: - """Bulk soft-hide with prune_sessions' filter surface, via - set_session_archived so each lineage flips as a unit. ``archived`` - defaults to False so repeat runs are idempotent. Returns matches.""" + """Bulk soft-hide with prune_sessions' filter surface, via set_session_archived + so each lineage flips as a unit; idempotent. Returns matches.""" filters.setdefault("archived", False) rows = self.list_prune_candidates(older_than_days=older_than_days, source=source, **filters) for row in rows: @@ -1558,9 +1450,8 @@ class SessionSessionsMixin: def maybe_auto_archive( self, idle_days: float = 3, min_interval_hours: int = 24, exclude_pinned: bool = True, ) -> Dict[str, Any]: - """Idempotent auto-archive of sessions idle for ``idle_days`` (ages on last - activity, non-destructive). ``state_meta['last_auto_archive']`` gates - runs within ``min_interval_hours``; safe to call opportunistically. + """Idempotent, non-destructive auto-archive of sessions idle for ``idle_days``; + ``state_meta['last_auto_archive']`` gates runs within ``min_interval_hours``. Never raises: {"skipped", "archived", "error"?}.""" result: Dict[str, Any] = {"skipped": False, "archived": 0} try: