diff --git a/hermes_state.py b/hermes_state.py index aed5807cba..f0055c3df2 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -103,8 +103,7 @@ _MAX_SAFE_MESSAGES = 20_000 # resume/export guard default def _configured_transcript_limit(key: str, fallback: int = _MAX_SAFE_MESSAGES) -> int: - """``sessions.`` from config.yaml (lazy import: circular at load), else - *fallback*. 0 disables the guard. Not cached (load_config_readonly is).""" + """``sessions.`` from config.yaml (lazy import: circular at load), else *fallback*; 0 disables.""" try: from hermes_cli.config import load_config_readonly value = (load_config_readonly().get("sessions") or {}).get(key) @@ -173,8 +172,7 @@ def _compression_lock_holder_process_is_dead(holder: str) -> bool: def _scrub_surrogates(value: Any) -> Any: - """Replace lone surrogates in text (sqlite3 raises UnicodeEncodeError on them, - aborting the whole write); pass anything else through.""" + """Replace lone surrogates in text (sqlite3 raises UnicodeEncodeError, aborting the whole write).""" return _sanitize_surrogates(value) if isinstance(value, str) else value @@ -199,9 +197,8 @@ _READ_ONLY_IOERR_RETRY_ATTEMPTS, _READ_ONLY_IOERR_RETRY_BACKOFF_S = 3, 0.05 def _default_db_path() -> Path: - """Default state DB path at CALL time: a re-pointed ``DEFAULT_DB_PATH`` wins, - else ``get_hermes_home()`` is resolved fresh so a runtime HERMES_HOME redirect - works regardless of import order.""" + """Default state DB path at CALL time: a re-pointed ``DEFAULT_DB_PATH`` wins, else + ``get_hermes_home()`` is resolved fresh (a runtime HERMES_HOME redirect works regardless of import).""" return DEFAULT_DB_PATH if DEFAULT_DB_PATH != _IMPORT_DEFAULT_DB_PATH else get_hermes_home() / "state.db" @@ -213,8 +210,8 @@ _STATE_DB_GUARD_EXTRA_DENY_ROOTS: Tuple[Path, ...] = () def _ensure_test_isolation(db_path: Path) -> None: - """Raise RuntimeError before any connection/mkdir/pragma/byte probe when a - pytest-context process (env OR ancestry) resolves a production DB.""" + """Raise before any connection/mkdir/pragma/byte probe when a pytest-context process + (env OR ancestry) resolves a production DB.""" if _STATE_DB_GUARD_BYPASS or os.environ.get(_STATE_DB_GUARD_BYPASS_ENV) or not _in_test_context(): return try: @@ -303,8 +300,7 @@ def _strip_stale_tool_call_markers(messages: List[Dict[str, Any]]) -> List[Dict[ def format_session_db_unavailable(prefix: str = "Session database not available") -> str: - """User-facing "session DB unavailable" message with the captured init cause - (plus a WAL-docs hint for NFS/SMB-style locking failures).""" + """User-facing message with the captured init cause (+ WAL-docs hint for NFS/SMB locking failures).""" cause = get_last_init_error() if not cause: return f"{prefix}." @@ -323,8 +319,8 @@ _IS_WINDOWS = sys.platform == "win32" def divert_session_transcript_jsonl(session_id: str, messages) -> "Optional[Path]": - """Append pending messages to HERMES_HOME/sessions/.jsonl (state.db was - replaced under a live process). Returns the path, or None if nothing to write.""" + """Append pending messages to HERMES_HOME/sessions/.jsonl (state.db was replaced under a + live process). Returns the path, or None if nothing to write.""" sid = str(session_id or "").strip() if not sid or not messages: return None @@ -352,11 +348,10 @@ class SessionDB( SessionGatewayMixin, SessionMaintenanceMixin, SessionUsageMixin, SessionTitlesMixin, SessionMessagesMixin, ): - """SQLite-backed session storage with FTS5 search. Thread-safe for the gateway - pattern (many reader threads, one writer via WAL).""" + """SQLite-backed session storage with FTS5 search; many reader threads, one writer (WAL).""" - # Only these state-owned producers join automatic stale-open reconciliation; - # messaging/UI sources have their own lifecycle owners; unknown sources fail closed. + # Only these state-owned producers join automatic stale-open reconciliation; messaging/UI + # sources have their own lifecycle owners; unknown sources fail closed. _AUTO_PRUNE_STALE_OPEN_SOURCES: Tuple[str, ...] = ( "cli", "cron", "kanban", "acp", "api_server", "subagent", "tool", ) @@ -369,24 +364,21 @@ class SessionDB( # writes (failure aborts the turn) get the long budget; observation-only activity # writes sit on the response-critical path and get a sub-second one. _WRITE_PATIENCE_S, _TRANSCRIPT_WRITE_PATIENCE_S, _ACTIVITY_WRITE_PATIENCE_S = 20.0, 60.0, 0.5 - # A live compression lock gets a short wait (compression publishes in seconds), - # but the lease is a correctness boundary: a writer still locked out afterwards - # is refused rather than landing a stale turn in a wedged compression. + # A live compression lock gets a short wait (compression publishes in seconds), but the lease + # is a correctness boundary: a writer still locked out afterwards is refused. _COMPRESSION_BUSY_WAIT_S = 5.0 _WRITE_RETRY_MIN_S, _WRITE_RETRY_MAX_S = 0.020, 0.150 # fast jitter for the first _SLOW_AFTER_S _WRITE_RETRY_SLOW_AFTER_S = 2.0 _WRITE_RETRY_SLOW_MIN_S, _WRITE_RETRY_SLOW_MAX_S = 0.250, 1.000 # PASSIVE WAL checkpoint every N successful writes. _CHECKPOINT_EVERY_N_WRITES = 50 - # Bounded FTS ``'merge'`` (ms of lock each) instead of ``'optimize'`` (9-18s per - # index on a 10GB DB — longer than a writer's patience); up to - # _FTS_MERGE_COMMANDS_PER_PASS per index, stopping on no-progress. + # Bounded FTS ``'merge'`` (ms of lock each) instead of ``'optimize'`` (9-18s per index on a 10GB + # DB, longer than a writer's patience); up to _COMMANDS_PER_PASS per index, stopping on no-progress. _FTS_MERGE_EVERY_N_WRITES, _FTS_MERGE_MAX_PAGES_PER_INDEX, _FTS_MERGE_COMMANDS_PER_PASS = 1000, 500, 4 # Imports cap lower than exports: an import holds one BEGIN IMMEDIATE. _IMPORT_MAX_SESSIONS, _IMPORT_MAX_MESSAGES_PER_SESSION, _IMPORT_MAX_TOTAL_MESSAGES = 500, 10_000, 50_000 _IMPORT_MAX_SESSION_BYTES, _IMPORT_MAX_TOTAL_BYTES = 5 * 1024 * 1024, 25 * 1024 * 1024 - # Accounting workers retire when idle so a bound-method target can't keep an - # abandoned SessionDB (and its descriptors) alive. + # Accounting workers retire when idle so a bound-method target can't keep an abandoned SessionDB alive. _TOKEN_WRITER_IDLE_SECONDS = 30.0 @staticmethod @@ -448,12 +440,11 @@ class SessionDB( self._read_budget.register(self) self._read_permits = self._read_budget.permits self._read_conns_lock = threading.Lock() - # Set when close() begins; an in-flight reader then closes its own - # connection instead of re-populating a pool nobody will drain again. + # Set when close() begins; an in-flight reader then closes its own connection + # instead of re-populating a pool nobody will drain again. self._read_conns_closed = False - # "read-only opens are failing" backoff: a TIMESTAMP, not a sticky bool — - # the likeliest trigger is transient EMFILE, and a permanent flag would - # demote every reader to the writer lock forever. + # Read-open failure backoff is a TIMESTAMP, not a sticky bool: the likeliest trigger + # is transient EMFILE, and a permanent flag would demote every reader forever. self._read_open_failed_at = 0.0 self._wal_active, self._write_count = False, 0 # File identity of the opened state.db, compared on every write so an out-of-band @@ -465,12 +456,11 @@ class SessionDB( self._db_corrupt, self._db_corrupt_reason = False, "" # sticky quarantine (StateDbCorruptError) self._fts_usermerge_floor_applied = False # one-shot usermerge-floor write guard self._fts_enabled = self._fts_stale = self._trigram_available = False - # _fts_cjk_loaded: tokenizer present on the writer connection; - # _fts_cjk_available: messages_fts_cjk is queryable AND not marked stale. + # _fts_cjk_loaded: tokenizer on the writer connection; _fts_cjk_available: messages_fts_cjk + # is queryable AND not marked stale. self._fts_cjk_loaded = self._fts_cjk_available = self._fts_unavailable_warned = False self._conn = None - # Async token accounting; distinct from self._lock so enqueue/flush - # bookkeeping never contends with SQLite writes. + # Async token accounting; distinct from self._lock so enqueue/flush never contends with writes. self._token_queue: deque = deque() self._token_queue_cond = threading.Condition(threading.Lock()) self._token_writer_thread: Optional[threading.Thread] = None @@ -487,8 +477,7 @@ class SessionDB( self._record_db_file_identity() initialization_complete = True except Exception as exc: - # Surface WHY via /resume and friends (never cleared on success, see - # _set_last_init_error); callers keep their ``_session_db = None`` path. + # Surface WHY via /resume and friends; callers keep their ``_session_db = None`` path. _set_last_init_error(f"{type(exc).__name__}: {exc}") raise finally: @@ -497,16 +486,15 @@ class SessionDB( self._close_connection_quietly(conn) def _open_writer(self) -> None: - """Writable open: preflight, quarantine/zero-byte guard, connect + schema, - one in-place schema repair on a malformed sqlite_master, generation stamp.""" + """Writable open: preflight, zero-byte quarantine, connect + schema (one in-place repair of a + malformed sqlite_master), generation stamp.""" self.db_path.parent.mkdir(parents=True, exist_ok=True) - # Read-only file/sidecar preflight BEFORE the first connection: an - # actionable message instead of an opaque "attempt to write a readonly - # database" from deep inside _init_schema. + # Read-only file/sidecar preflight BEFORE the first connection: an actionable message + # instead of an opaque "attempt to write a readonly database" from inside _init_schema. preflight_db_writability(self.db_path, db_label="state.db") try: - # Serialize zero-byte check, quarantine, connect and schema commit so - # concurrent openers don't race the absent-path -> schema-commit window. + # Serialize zero-byte check, quarantine, connect and schema commit so concurrent + # openers don't race the absent-path -> schema-commit window. if not self.db_path.exists() or is_zeroed_state_db(self.db_path): with quarantine_cross_process_lock(self.db_path) as lock_acquired: if not lock_acquired: @@ -520,9 +508,8 @@ class SessionDB( self._handle_quarantine_if_zeroed(already_locked=False) self._connect_and_init_with_lock_patience() except sqlite3.DatabaseError as exc: - # Malformed schema fails on the very first statement (before - # _init_schema), so it can't be caught at the FTS-rebuild layer: - # repair sqlite_master in place (backup first) and reopen once. + # A malformed schema fails on the very first statement (before _init_schema), so the + # FTS-rebuild layer never sees it: repair sqlite_master in place (backup first), reopen once. if not is_malformed_schema_error(exc) or not _claim_repair_attempt(self.db_path): raise logger.error( @@ -533,8 +520,7 @@ class SessionDB( if not repair_state_db_schema(self.db_path).get("repaired"): raise self._connect_and_init_with_lock_patience() - # The v23 FTS optimization is OPT-IN (`hermes db optimize`), never - # auto-started on open (no background worker racing session lifecycle). + # FTS optimization is OPT-IN (`hermes db optimize`); no background worker races session lifecycle. self._ensure_db_file_generation() def _open_read_only(self) -> None: @@ -560,17 +546,15 @@ class SessionDB( raise return except sqlite3.OperationalError as ioerr: - # Transient SQLITE_IOERR window (see _READ_ONLY_IOERR_RETRY_ATTEMPTS); - # a persistent one exhausts the budget and propagates. + # Transient SQLITE_IOERR window (see _READ_ONLY_IOERR_RETRY_ATTEMPTS). transient = _DISK_IO_ERROR_MARKER in str(ioerr).lower() if attempt >= _READ_ONLY_IOERR_RETRY_ATTEMPTS or not transient: raise time.sleep(_READ_ONLY_IOERR_RETRY_BACKOFF_S) def _connect_read_only(self, timeout: float) -> sqlite3.Connection: - """``mode=ro`` tracked connection with Row factory. check_same_thread=False: pooled - connections are borrowed by whichever thread reads next; exclusive ownership is - enforced by pool checkout.""" + """``mode=ro`` tracked connection with Row factory. check_same_thread=False: pooled connections + are borrowed by whichever thread reads next; exclusive ownership is enforced by pool checkout.""" conn = _connect_tracked_db( f"file:{self.db_path}?mode=ro", tracking_path=self.db_path, uri=True, check_same_thread=False, timeout=timeout, isolation_level=None, @@ -579,8 +563,8 @@ class SessionDB( return conn def _handle_quarantine_if_zeroed(self, already_locked: bool = False) -> None: - """Quarantine a zero-byte/headerless state.db so a fresh one can open; if - quarantine failed, raise the clear message instead of opening the zeroed file.""" + """Quarantine a zero-byte/headerless state.db so a fresh one can open; if quarantine failed, + raise the clear message instead of opening the zeroed file.""" if not (self.db_path.exists() and is_zeroed_state_db(self.db_path)): return try: @@ -601,9 +585,9 @@ class SessionDB( raise sqlite3.DatabaseError(msg) def _open_writer_conn(self) -> sqlite3.Connection: - """Connect + WAL/pragma/tokenizer setup for a writer connection (no schema init). - Short timeout: application-level jittered retry handles contention, not - SQLite's busy handler; isolation_level=None: explicit BEGIN IMMEDIATE.""" + """Connect + WAL/pragma/tokenizer setup for a writer connection (no schema init). Short timeout: + jittered application-level retry handles contention, not SQLite's busy handler; + isolation_level=None: explicit BEGIN IMMEDIATE.""" conn = _connect_tracked_db( str(self.db_path), check_same_thread=False, timeout=1.0, isolation_level=None, ) @@ -675,10 +659,9 @@ class SessionDB( if self._fts_cjk_loaded: # registers in the connection, not the file: ro is fine load_fts5_cjk_extension(conn) except BaseException as exc: - # A half-open connection (open ok, extension load failed) is a live - # tracked descriptor — the leak shape this pool exists to fix; a - # stranded permit would permanently shrink the read path by one slot. - # (Not _close_read_conn: callers release their own permit.) + # A half-open connection (open ok, extension load failed) is a live tracked descriptor, + # the leak shape this pool exists to fix; a stranded permit would shrink the read + # path by one slot forever. (Not _close_read_conn: callers release their own permit.) if conn is not None: self._close_conn_logged(conn, "partially-opened read conn") self._read_budget.release() @@ -691,8 +674,7 @@ class SessionDB( return conn def _evict_one_idle_read_conn(self) -> bool: - """Close one idle pooled connection (a peer on the same file wants its - permit); never pulls a connection out from under a live reader.""" + """Close one idle pooled connection (a peer on the same file wants its permit); never a live one.""" try: conn = self._read_pool.get_nowait() except queue.Empty: @@ -701,17 +683,16 @@ class SessionDB( return True def _close_read_conn(self, conn) -> None: - """Close a pooled read connection and release its permit even when the close - fails (withholding it would permanently narrow the read path). Pairs with - _get_read_conn(); over-releasing the BoundedSemaphore raises ValueError.""" + """Close a pooled read connection and release its permit even when the close fails (withholding + it would narrow the read path forever). Over-releasing the BoundedSemaphore raises ValueError.""" try: self._close_conn_logged(conn, "read-conn") finally: self._read_budget.release() def _checkout_read_conn(self) -> Optional[sqlite3.Connection]: - """Borrow a read connection, opening on a miss; None when the read path is - unavailable. A pool hit costs no permit (the connection already holds one).""" + """Borrow a read connection, opening on a miss; None when the read path is unavailable. + A pool hit costs no permit (the connection already holds one).""" if not self._wal_active or self.read_only: return None try: @@ -758,9 +739,8 @@ class SessionDB( f"SessionDB for {self.db_path} was closed (read-only handle); " f"cannot serve a {context} after close()" ) - # A reopen resolves the PATH again: a replaced file would be written through - # stale WAL/shm assumptions; a quarantined handle must never hand a fresh - # connection (and its close-time checkpoint) to a damaged file. + # A reopen resolves the PATH again: a replaced file would be written through stale WAL/shm + # assumptions; a quarantined handle must never hand a fresh connection to a damaged file. if self._db_corrupt and not (self._db_replaced or self._db_file_was_replaced()): raise self._corrupt_error( f"state.db connection for {self.db_path} is quarantined after " @@ -792,11 +772,9 @@ class SessionDB( if patience_s is None: patience_s = self._WRITE_PATIENCE_S deadline = time.monotonic() + patience_s - # Set on the first compression-busy collision: the short wait is measured from then. - compression_deadline: Optional[float] = None - # One retry for SQLITE_IOERR raised by BEGIN IMMEDIATE itself (callback not - # run: nothing replayed). Once it has started, an IOERR leaves settlement - # unknown and must propagate — this helper owns non-idempotent mutations. + compression_deadline: Optional[float] = None # set on the first compression-busy collision + # One retry for SQLITE_IOERR raised by BEGIN IMMEDIATE itself (callback not run: nothing + # replayed). Once fn has started, an IOERR leaves settlement unknown and must propagate. ioerr_begin_retried = False while True: self._raise_if_db_corrupt() @@ -825,8 +803,7 @@ class SessionDB( self._try_incremental_merge_fts() return result except SessionCompressionInProgressError: - # Transient (see _COMPRESSION_BUSY_WAIT_S): without a wait, a steer - # landing mid-compression aborts the turn. + # Transient (see _COMPRESSION_BUSY_WAIT_S): a steer landing mid-compression must not abort. if compression_deadline is None: compression_deadline = min(time.monotonic() + self._COMPRESSION_BUSY_WAIT_S, deadline) if self._sleep_before_write_retry( @@ -855,24 +832,23 @@ class SessionDB( and not ioerr_begin_retried and self._sleep_before_write_retry(deadline, patience_s) ): - # Retry on the SAME connection: close()+reopen cancels this - # process's POSIX locks for every sibling (howtocorrupt §2.2). + # Retry on the SAME connection: close()+reopen would cancel this process's POSIX + # locks for every sibling (howtocorrupt §2.2). ioerr_begin_retried = True continue raise # non-lock error, callback already ran, or patience exhausted except sqlite3.DatabaseError as exc: if _is_no_more_rows(exc) and self._sleep_before_write_retry(deadline, patience_s): continue - # An out-of-band replace surfaces as this same corruption class; - # in-file repair on a NEW generation amplifies the damage. + # An out-of-band replace surfaces as this same corruption class; in-file repair + # on a NEW generation amplifies the damage. if ( "not a database" in str(exc).lower() or is_malformed_db_error(exc) or self._is_fts_write_corruption_error(exc) ): self._raise_if_db_replaced() - # Corrupt FTS shadow tables fail every write via the sync triggers - # while canonical rows are intact: detach the derived indexes - # atomically and retry (never rebuild from the live write path). + # Corrupt FTS shadow tables fail every write via the sync triggers while canonical rows + # are intact: detach the derived indexes atomically and retry (never rebuild here). if self._enter_fts_fail_open(exc): continue # What survives both checks is structural damage: quarantine. @@ -880,8 +856,7 @@ class SessionDB( self._halt_db_corrupt(exc) raise except sqlite3.Error as exc: - # Builds raising 'no more rows' as InterfaceError (sibling of - # DatabaseError); anything else propagates untouched. + # Some builds raise 'no more rows' as InterfaceError (sibling of DatabaseError). if _is_no_more_rows(exc) and self._sleep_before_write_retry(deadline, patience_s): continue raise @@ -915,9 +890,8 @@ class SessionDB( return conn.execute(sql, params).fetchall() def _ensure_db_file_generation(self) -> None: - """Mint a once-per-file generation stamp (state_meta + application_id). - First opener wins (INSERT OR IGNORE); application_id is written only while - 0 so racers converge. PASSIVE checkpoint only — never TRUNCATE.""" + """Mint a once-per-file generation stamp (state_meta + application_id). First opener wins (INSERT + OR IGNORE); application_id is written only while 0 so racers converge. PASSIVE checkpoint only.""" if self.read_only or self._conn is None: return token = uuid.uuid4().hex @@ -968,8 +942,7 @@ class SessionDB( recorded_app = int(self._db_file_application_id or 0) if not recorded_app: return False - # Header 0 = WAL not yet checkpointed, not a replace; any real replacement - # (a copied Hermes DB minted its own id) is nonzero. + # Header 0 = WAL not yet checkpointed, not a replace; a real replacement is nonzero. disk_app = _read_sqlite_application_id(self.db_path) return bool(disk_app and disk_app != recorded_app) @@ -1023,8 +996,8 @@ class SessionDB( @classmethod def _is_structural_corruption_error(cls, exc: BaseException) -> bool: - """Bare SQLITE_CORRUPT/NOTADB with no FTS provenance: canonical B-tree / - schema / freelist damage, never repairable from the live write path.""" + """Bare SQLITE_CORRUPT/NOTADB with no FTS provenance: canonical B-tree/schema/freelist damage, + never repairable from the live write path.""" return ( isinstance(exc, sqlite3.DatabaseError) and not isinstance(exc, StateDbCorruptError) @@ -1078,9 +1051,8 @@ class SessionDB( raise self._corrupt_error() def _sleep_before_write_retry(self, deadline: float, patience_s: float) -> bool: - """Sleep one jitter interval if the budget allows; True = retry, False = - deadline passed. Small jitter for the first _WRITE_RETRY_SLOW_AFTER_S, - then backs off; never overshoots the deadline by a full slow-jitter.""" + """Sleep one jitter interval if the budget allows; True = retry, False = deadline passed. Small + jitter for the first _WRITE_RETRY_SLOW_AFTER_S, then slow; never overshoots the deadline.""" now = time.monotonic() if now >= deadline: return False @@ -1097,8 +1069,8 @@ class SessionDB( repair must not run while another process is attached (a sidecar reset under it splits the WAL inodes). A scan failure is reported as an unknown holder — skipping optional maintenance beats assuming quiescence.""" - # Split-brain needs POSIX unlink semantics (Windows refuses to replace - # open sidecars); psutil.open_files() there can block for minutes. + # Split-brain needs POSIX unlink semantics (Windows refuses to replace open sidecars); + # psutil.open_files() there can block for minutes. if _IS_WINDOWS: return [] if psutil is None: @@ -1109,24 +1081,22 @@ class SessionDB( own_pid = os.getpid() try: if sys.platform.startswith("linux"): - # readlink /proc//fd directly; psutil.open_files() stats the - # literal path and silently drops "state.db-wal (deleted)" entries. + # readlink /proc//fd directly; psutil.open_files() stats the literal path + # and silently drops "state.db-wal (deleted)" entries. for pid in (int(p) for p in os.listdir("/proc") if p.isdigit()): if pid == own_pid: continue try: targets = list(_proc_fd_targets(pid)) except OSError: - # Unreadable fd table (other user); cmdline is world-readable: - # flag only uninspectable holders that look like Hermes. + # Unreadable fd table (other user); flag only Hermes-looking holders via cmdline. cmdline = _read_proc_cmdline(pid) if cmdline is not None and _looks_like_hermes(cmdline): holders.append((pid, f"uninspectable holder: {cmdline[:80]}")) continue holders.extend((pid, t) for t in targets if _canonical_sqlite_path(t) in watched) else: - # macOS / BSD: psutil.open_files() (no "(deleted)" suffix convention there; - # AccessDenied -> None -> empty iteration is acceptable on macOS). + # macOS / BSD: psutil.open_files() (no "(deleted)" convention; AccessDenied -> empty). for process in psutil.process_iter(["pid", "open_files"]): pid = int(process.info["pid"]) if pid == own_pid: @@ -1198,8 +1168,8 @@ class SessionDB( logger.debug("WAL checkpoint (PASSIVE) at close failed: %s", exc) conn, self._conn = self._conn, None self._close_connection_quietly(conn) - # A clean close lets SQLite unlink the sidecars (a legitimate end of - # the generation, not a split): a teardown-race reopen must re-adopt. + # A clean close lets SQLite unlink the sidecars (a legitimate end of the + # generation, not a split): a teardown-race reopen must re-adopt. self._db_sidecar_identity = {} def __del__(self) -> None: @@ -1232,8 +1202,8 @@ class SessionDB( TITLE_SOURCE_DERIVED, TITLE_SOURCE_LLM, TITLE_SOURCE_USER = "derived", "llm", "user" _TITLE_SOURCE_RANK = {TITLE_SOURCE_DERIVED: 0, TITLE_SOURCE_LLM: 1, TITLE_SOURCE_USER: 2} - # Bot Mode's canonical chat is resolved by exact-title lookup: the title IS the - # identity, so _set_session_title refuses renames of a hidden row holding it. + # Bot Mode's canonical chat is resolved by exact-title lookup: the title IS the identity, + # so _set_session_title refuses renames of a hidden row holding it. CANONICAL_BOT_CHAT_TITLE = "Bot Chat" # ── Message storage constants (SessionMessagesMixin) ── @@ -1241,8 +1211,8 @@ class SessionDB( _CONTENT_JSON_PREFIX = "\x00json:" #: Reactions live inside ``display_metadata`` so they survive row rewrites. REACTIONS_METADATA_KEY = "reactions" - # Columns every conversation projection decodes; ``active`` rides along so a - # display read can split compaction-archived rows without a second query. + # Columns every conversation projection decodes; ``active`` rides along so a display read + # can split compaction-archived rows without a second query. _CONVERSATION_ROW_COLUMNS = ( "id, role, content, tool_call_id, tool_calls, tool_name, effect_disposition, " "finish_reason, reasoning, reasoning_content, reasoning_details, " @@ -1253,15 +1223,15 @@ class SessionDB( # ── Meta key/value (scheduler bookkeeping) ── def get_meta(self, key: str) -> Optional[str]: - """Read state_meta[key] on self._lock (not _read_ctx): fts_rebuild_step reads - progress before its write transaction and a WAL reader would not see it.""" + """Read state_meta[key] on self._lock (not _read_ctx): fts_rebuild_step reads progress before its + write transaction and a WAL reader would not see it.""" with self._lock: row = self._conn.execute("SELECT value FROM state_meta WHERE key = ?", (key,)).fetchone() return None if row is None else row[0] def set_meta(self, key: str, value: str, *, cursor: Optional[sqlite3.Cursor] = None) -> None: - """Upsert state_meta[key]; with ``cursor`` the write is inline (the caller - already holds a transaction — nesting BEGIN IMMEDIATE would deadlock).""" + """Upsert state_meta[key]; with ``cursor`` the write is inline (the caller already holds a + transaction — nesting BEGIN IMMEDIATE would deadlock).""" sql = ( "INSERT INTO state_meta (key, value) VALUES (?, ?) " "ON CONFLICT(key) DO UPDATE SET value = excluded.value" @@ -1272,8 +1242,8 @@ class SessionDB( self._write_sql(sql, (key, value)) def retag_kanban_worker_sessions(self, workspaces_root: str) -> int: - """Retag legacy kanban worker rows from ``cli`` to ``kanban`` by cwd under the - board's workspaces root; gated once per root via state_meta. Returns rows retagged.""" + """Retag legacy kanban worker rows from ``cli`` to ``kanban`` by cwd under the board's workspaces + root; gated once per root via state_meta. Returns rows retagged.""" prefix = str(workspaces_root).rstrip("/\\") if not prefix: return 0 @@ -1304,8 +1274,8 @@ class SessionDB( class AsyncSessionDB: - """Async door onto SessionDB: each call is offloaded via asyncio.to_thread so a - blocking SQLite call never freezes the event loop (no method returns a live cursor).""" + """Async door onto SessionDB: every call runs via asyncio.to_thread so a blocking SQLite call + never freezes the event loop (no method returns a live cursor).""" def __init__(self, db: "SessionDB") -> None: self._db = db diff --git a/hermes_state_sessions.py b/hermes_state_sessions.py index 71902f3dd4..174733b887 100644 --- a/hermes_state_sessions.py +++ b/hermes_state_sessions.py @@ -25,8 +25,8 @@ logger = logging.getLogger("hermes_state") def workspace_key(row: Dict[str, Any]) -> Optional[str]: - """Workspace grouping key: git repo root, 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 excluded: a checkout must not + fragment history).""" return (row.get("git_repo_root") or "").strip() or (row.get("cwd") or "").strip() or None @@ -61,8 +61,8 @@ 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``.""" + """WHERE for ``workspace_key(row) == key``: git_repo_root equals ``key``, or (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 ( @@ -93,11 +93,9 @@ def _session_filter_where( 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 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.""" + """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``). Clause order is part of the SQL text contract.""" where: List[str] = [] params: List[Any] = [] if exclude_children: @@ -121,8 +119,8 @@ def _session_filter_where( def _collect_delegate_child_ids(conn, parent_ids: List[str]) -> List[str]: - """Delegate-subagent ids (``_delegate_from`` marker, walked recursively) to - cascade-delete with *parent_ids*; untagged children stay orphaned, not deleted.""" + """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, @@ -163,9 +161,8 @@ _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: - """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).""" + """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() @@ -229,17 +226,17 @@ class SessionSessionsMixin: """Session rows: create/inherit, lifecycle flags, model_config, listing, deletion.""" 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 because a - gateway serving a NON-launch profile opens that profile's store. None - outside the profile tree — NULL beats a fabricated owner.""" + """The profile owning THIS store, from ``db_path`` alone (``/state.db`` → default, + ``/profiles//state.db`` → name): 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() parent = Path(self.db_path).resolve().parent if parent == root: return "default" - if parent.parent == root / "profiles" and re.fullmatch(r"[a-z0-9][a-z0-9_-]{0,63}", parent.name): + is_profile_dir = parent.parent == root / "profiles" + if is_profile_dir and re.fullmatch(r"[a-z0-9][a-z0-9_-]{0,63}", parent.name): return parent.name except Exception: logger.debug("own-profile derivation failed", exc_info=True) @@ -247,12 +244,10 @@ class SessionSessionsMixin: @staticmethod def _inherit_parent_session_metadata(conn, session_id: str) -> None: - """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).""" + """NULL-fill a child's cwd/git/profile from its parent (profile_name only within the same + ``agent::`` namespace). Gateway routing columns are inherited ONLY by 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 traffic into a subagent's session).""" conn.execute(_INHERIT_PARENT_META_SQL, (session_id,)) conn.execute(_INHERIT_PARENT_ROUTING_SQL, (session_id,)) @@ -263,11 +258,10 @@ 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 - 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).""" + """Upsert a session row, 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): @@ -348,9 +342,8 @@ class SessionSessionsMixin: 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; several distinct users and no match - → None rather than contaminating another participant's session.""" + """Most recent live session_id for source + chat_id (+ thread_id). With ``user_id`` exact sender + matches win; several distinct users and no match → None (never another participant's session).""" if not platform or chat_id in (None, ""): return None query = """ @@ -377,14 +370,13 @@ class SessionSessionsMixin: return None return str(rows[0]["id"]) - # 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). + # Orphaned gateway-session repair: widest plausible gap between a keyed predecessor going + # quiet and its unkeyed successor (incident was ~60s; 15 min without spanning conversations). _ORPHAN_ADOPTION_MAX_GAP_S = 900.0 - # 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. + # 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" @@ -393,9 +385,8 @@ class SessionSessionsMixin: ) def end_session(self, session_id: str, end_reason: str) -> None: - """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.""" + """Mark a session ended; the first end_reason wins (a compression split must keep + ``'compression'`` even if a stale end_session() lands later); reopen_session() to re-end.""" self._execute_write(lambda conn: self._end_and_bump( conn, "UPDATE sessions SET ended_at = ?, end_reason = ? WHERE id = ? AND ended_at IS NULL", (time.time(), end_reason, session_id), session_id, end_reason, @@ -410,9 +401,9 @@ class SessionSessionsMixin: return changed def reopen_session(self, session_id: str) -> None: - """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 so they cannot drift).""" + """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 + so they cannot drift).""" def _do(conn): conn.execute( "UPDATE sessions AS child SET model_config = json_set(" @@ -428,10 +419,9 @@ class SessionSessionsMixin: self._execute_write(_do) 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): - an ``agent_close`` row left recoverable would be resurrected by - stale-route recovery. Keep in sync with find_latest_gateway_session_for_peer.""" + """Durably mark an intentional reset boundary on live rows or rows with a *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() @@ -450,11 +440,10 @@ 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 - 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).""" + """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() @@ -514,9 +503,8 @@ class SessionSessionsMixin: 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 the - caller) so surfaces see 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 when = float(ts if ts is not None else time.time()) @@ -532,8 +520,8 @@ class SessionSessionsMixin: ) def clear_session_activity_labels(self, session_id: str) -> None: - """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.""" + """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 try: @@ -571,18 +559,17 @@ 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 can't - fork the cached tool prefix on a flipped check_fn verdict; ``None`` clears.""" + """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), 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. + """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 model and provider.""" + # 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() patch: Dict[str, Any] = {"browser_model_lock": None} if model: @@ -600,9 +587,8 @@ class SessionSessionsMixin: sql: str = "UPDATE sessions SET model_config = ? WHERE id = ?", params: Optional[Callable[[Optional[str]], tuple]] = None, ) -> None: - """Merge ``patch`` into model_config then run ``sql`` with ``params(merged)`` - in one write transaction; no-op when the row doesn't exist. A custom ``sql`` - (the prompt-nulling variants) also GCs unreferenced system_prompts.""" + """Merge ``patch`` into model_config then run ``sql`` with ``params(merged)`` in one write + transaction; no-op when the row doesn't exist. Custom ``sql`` (prompt-nulling) also GCs prompts.""" def _do(conn): merged = self._merge_model_config_json(conn, session_id, patch) if merged is _MODEL_CONFIG_ROW_MISSING: @@ -615,10 +601,9 @@ 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 - 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).""" + """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) or ``_MODEL_CONFIG_ROW_MISSING`` (``on_missing="raise"`` → ValueError).""" row = conn.execute("SELECT model_config FROM sessions WHERE id = ?", (session_id,)).fetchone() if row is None: if on_missing == "raise": @@ -649,8 +634,8 @@ class SessionSessionsMixin: model_options: Optional[Dict[str, Any]] = None, route_source: Optional[str] = None, confirmed: bool = False, ) -> None: - """Persist a Browser / API-client runtime lock into model_config (lineage - markers survive); null system_prompt so cached footers cannot lie.""" + """Persist a Browser / API-client runtime lock into model_config (lineage markers survive); null + system_prompt so cached footers cannot lie.""" lock = { "provider": provider or "", "model": model or "", "model_options": model_options or {}, "route_source": route_source or "", "confirmed": bool(confirmed), "updated_at": time.time(), @@ -667,16 +652,14 @@ class SessionSessionsMixin: ) def set_session_yolo(self, session_id: str, enabled: bool) -> None: - """Persist the per-session YOLO flag so ``/yolo`` survives ``--resume``; - no-op when the row doesn't exist yet.""" + """Persist the per-session YOLO flag so ``/yolo`` survives ``--resume``; no-op without a row.""" if not session_id: return 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; 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).""" return bool(_parse_model_config((session_meta or {}).get("model_config")).get("yolo_mode")) def get_session(self, session_id: str) -> Optional[Dict[str, Any]]: @@ -690,8 +673,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 (``session_model_usage`` - keeps the coherent per-call tuple; ``sessions`` mixes route changes).""" + """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, @@ -722,9 +705,8 @@ class SessionSessionsMixin: return matches[0]["id"] 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 (a store belongs to exactly - one profile). Never overwrites a non-NULL owner. Returns rows stamped.""" + """Stamp this store's own profile onto legacy ``profile_name IS NULL`` rows, which the fail-closed + owner ladder cannot route. Never overwrites a non-NULL owner. Returns rows stamped.""" stamp = (profile_name or "").strip() if not stamp: return 0 @@ -736,9 +718,8 @@ 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: Desktop - projects roots forward to their tip, so updating only the displayed tip - would let the untouched root resurrect it on refresh.""" + """Set one ``sessions`` column across a whole compression lineage: Desktop projects roots + forward to their tip, so updating only the tip would let the root resurrect it on refresh.""" return self._write_rowcount( f""" WITH RECURSIVE @@ -781,8 +762,8 @@ class SessionSessionsMixin: 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); deliberate archives are left alone. True when un-archived.""" + """Un-archive a session archived by a recoverable accident (ws_orphan_reap, agent_close); + deliberate archives are left alone. True when un-archived.""" if not session_id: return False try: @@ -811,19 +792,16 @@ class SessionSessionsMixin: return True def set_session_pinned(self, session_id: str, pinned: bool) -> bool: - """Pin/unpin a session and its compression lineage (pins are exempt from the - ``sessions.auto_archive`` sweep).""" + """Pin/unpin a session and its compression lineage (pins are exempt from the auto_archive sweep).""" return self._set_lineage_column("pinned", session_id, int(pinned)) def set_session_hidden(self, session_id: str, hidden: bool) -> bool: - """Hide/unhide a session and its compression lineage from the default listing; - it stays resumable by the owning surface.""" + """Hide/unhide a session and its compression lineage from the default listing; still resumable.""" return self._set_lineage_column("hidden", session_id, int(hidden)) 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: unread when activity postdates it (no write on the message - path). NULL = never tracked = read; 0 = explicitly unread.""" + """Mark read/unread across the compression lineage. ``last_read_at`` is a watermark: unread when + activity postdates it (no write on the message path). NULL = never tracked = read; 0 = unread.""" return self._set_lineage_column("last_read_at", session_id, time.time() if read else 0.0) @staticmethod @@ -843,10 +821,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``); chain membership keeps the leading-wildcard LIKE bounded.""" + """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``); chain membership bounds the LIKE.""" params: List[Any] = [] clauses: List[str] = [] def like(needle: str) -> str: @@ -878,9 +855,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), one batched query. ``_lineage_ids`` - carries every id on the chain: a persisted tile can hold a MIDDLE segment's id.""" + """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 chain id (a tile may + hold a MIDDLE segment's id).""" chain_by_root: Dict[str, List[str]] = {} # only roots whose tip differs from themselves for s in sessions: if s.get("end_reason") == "compression": @@ -927,10 +904,9 @@ class SessionSessionsMixin: id_query: str = None, search_query: str = None, 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. - ``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.""" + """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, session_key=session_key, @@ -947,7 +923,9 @@ class SessionSessionsMixin: + ("" if compact_rows else ", COALESCE(sp.prompt, s.system_prompt) AS _system_prompt_resolved") + f",\n {_PREVIEW_COL_SQL},\n " ) - prompt_join = "" if compact_rows else "LEFT JOIN system_prompts sp ON sp.hash = s.system_prompt_hash" + prompt_join = ( + "" if compact_rows else "LEFT JOIN system_prompts sp ON sp.hash = s.system_prompt_hash" + ) from_sessions = f"FROM sessions s\n {prompt_join}" if order_by_last_active: # The CTE walks compression-continuation edges forward from the admitted @@ -1024,8 +1002,8 @@ 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 (``'empty'`` - when none); one query, MAX(id) per session joined back — 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 {} @@ -1050,9 +1028,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, or raise SessionExportTooLargeError (the - LIMITed subquery stops once the bound is exceeded). ``None`` resolves - ``sessions.max_export_messages``; 0 disables the guard.""" + """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() @@ -1070,8 +1048,8 @@ class SessionSessionsMixin: return message_count def _is_explicit_branch_session(self, session_id: str) -> bool: - """Copied user-facing branch (``_branched_from``)? 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,)) @@ -1096,8 +1074,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 so ``hermes -c``/``--resume`` picks its last session.""" + """Sessions MRU-first with a computed ``last_active``; ``workspace_key`` scopes to one workspace + so ``hermes -c``/``--resume`` picks its last session.""" where_clauses = [] params: list = [] if source: @@ -1121,8 +1099,7 @@ 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 list_sessions_rich's filters so a paired "load more" - total matches the listable rows.""" + """Count sessions with list_sessions_rich's filters so a paired "load more" total matches.""" 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, @@ -1131,8 +1108,7 @@ class SessionSessionsMixin: 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); LIMIT short-circuits - instead of session_count()'s index scan.""" + """At least N sessions exist (archived included); LIMIT short-circuits session_count()'s scan.""" return len(self._read_all("SELECT 1 FROM sessions LIMIT ?", (n,))) >= n def session_count_by_source( @@ -1155,8 +1131,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) in ONE read (prompt_cache_scope needs both from the - same row). Missing row → (False, ""); DB errors propagate (fail 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, "" @@ -1164,8 +1140,8 @@ class SessionSessionsMixin: @staticmethod def _remove_session_files(sessions_dir: Optional[Path], session_id: str) -> None: - """Remove ``.json``/``.jsonl`` and gateway ``request_dump__*.json``; - OSError is swallowed so a filesystem hiccup never blocks a DB operation.""" + """Remove ``.json``/``.jsonl`` and gateway ``request_dump__*.json``; OSError is swallowed + so a filesystem hiccup never blocks a DB operation.""" if sessions_dir is None: return targets = [sessions_dir / f"{session_id}{suffix}" for suffix in (".json", ".jsonl")] @@ -1180,8 +1156,8 @@ 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 [] @@ -1192,10 +1168,9 @@ 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, - 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).""" + """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_ids: List[str] = [] expected_ids = set(expected_delete_ids) if expected_delete_ids is not None else None def _do(conn): @@ -1220,8 +1195,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; - check and delete share one transaction so a concurrent flush 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( """ @@ -1247,9 +1222,8 @@ class SessionSessionsMixin: return deleted def delete_sessions(self, session_ids: List[str], sessions_dir: Optional[Path] = None) -> int: - """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.""" + """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 deleted.""" unique_ids = list({sid for sid in session_ids or () if isinstance(sid, str) and sid}) if not unique_ids: return 0 @@ -1276,22 +1250,20 @@ class SessionSessionsMixin: 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/compaction keep - #: dropped turns as ``active = 0``), so NOT EXISTS is the authority. + # Shared by count_empty_sessions / delete_empty_sessions so badge and sweep agree. message_count + # counts live rows only (rewind/compaction keep dropped turns as active = 0): NOT EXISTS is 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; the ended_at guard means a - fresh session whose first message hasn't landed is never sniped.""" + """Count of empty, ended, non-archived sessions; ended_at guards a fresh session's first message.""" 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 in one transaction, - orphaning (not cascading) children; transcript files are swept too.""" + """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( @@ -1319,8 +1291,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; 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: @@ -1330,9 +1302,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, 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"?}.""" + """Idempotent, non-destructive auto-archive of sessions idle for ``idle_days``; state_meta + ``last_auto_archive`` gates runs within ``min_interval_hours``. Never raises.""" result: Dict[str, Any] = {"skipped": False, "archived": 0} try: now = time.time()