refactor(state): compact narrative docstrings and comments in hermes_state and hermes_state_sessions (keep every WHY)

This commit is contained in:
Teknium
2026-09-03 00:04:03 -07:00
parent cfb268a6c9
commit d944616a6d
2 changed files with 210 additions and 269 deletions

View File

@@ -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.<key>`` from config.yaml (lazy import: circular at load), else
*fallback*. 0 disables the guard. Not cached (load_config_readonly is)."""
"""``sessions.<key>`` 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/<id>.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/<id>.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/<pid>/fd directly; psutil.open_files() stats the
# literal path and silently drops "state.db-wal (deleted)" entries.
# readlink /proc/<pid>/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

View File

@@ -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 (``<root>/state.db``
→ default, ``<root>/profiles/<name>/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 (``<root>/state.db`` → default,
``<root>/profiles/<name>/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:<ns>:`` 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:<ns>:`` 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 ``<id>.json``/``.jsonl`` and gateway ``request_dump_<id>_*.json``;
OSError is swallowed so a filesystem hiccup never blocks a DB operation."""
"""Remove ``<id>.json``/``.jsonl`` and gateway ``request_dump_<id>_*.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()