diff --git a/hermes_state.py b/hermes_state.py index f2b108f0a0..a93a6d25a1 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -24,9 +24,7 @@ from contextlib import contextmanager from pathlib import Path from agent.message_sanitization import _sanitize_surrogates -# Known-durable message marker, shared with agent.context_compressor. run_agent -# keeps its own copy (cannot import hermes_state: circular), guarded by -# test_marker_constant_in_sync. +# Known-durable message marker (run_agent keeps a copy: circular import; a test pins them in sync). from agent.context_compressor import ( # noqa: F401 (re-exported; tests import it from here) _DB_PERSISTED_MARKER as _DB_PERSISTED_MARKER_KEY, ) @@ -43,34 +41,25 @@ from hermes_state_common import ( # noqa: F401 (re-exported; tests import from stat_db_file_identity as _stat_db_file_identity, ) from hermes_state_errors import ( # noqa: F401 (re-exported; the historical import path) - _DB_CORRUPTION_MARKERS, _DELETED_WAL_GENERATION_MSG, _DISK_FULL_MARKERS, _DISK_IO_ERROR_MARKER, - _MALFORMED_DB_MARKERS, _MALFORMED_SCHEMA_MARKERS, _STATE_DB_APPLICATION_ID_OFFSET, - _STATE_DB_CORRUPT_MSG, _STATE_DB_GENERATION_KEY, _STATE_DB_REPLACED_MSG, - _TRANSIENT_SQLITE_MARKERS, PERSISTENCE_ERROR_CAUSES, CompressionSessionBusyError, - CompressionSessionClosedError, DeletedWalGenerationError, SessionCompressionInProgressError, - SessionTurnLeaseLostError, StateDbCorruptError, StateDbReplacedError, _is_no_more_rows, - classify_persistence_error, is_disk_full_error, is_malformed_db_error, - is_malformed_schema_error, is_transient_sqlite_error, + _DELETED_WAL_GENERATION_MSG, _DISK_IO_ERROR_MARKER, _STATE_DB_APPLICATION_ID_OFFSET, + _STATE_DB_CORRUPT_MSG, _STATE_DB_GENERATION_KEY, _STATE_DB_REPLACED_MSG, PERSISTENCE_ERROR_CAUSES, + CompressionSessionBusyError, CompressionSessionClosedError, DeletedWalGenerationError, + SessionCompressionInProgressError, SessionTurnLeaseLostError, StateDbCorruptError, + StateDbReplacedError, _is_no_more_rows, classify_persistence_error, is_disk_full_error, + is_malformed_db_error, is_malformed_schema_error, is_transient_sqlite_error, ) from hermes_state_guard import ( # noqa: F401 (re-exported; tests patch hermes_state.) - _PYTEST_LAUNCHER_NAMES, _STATE_DB_GUARD_BYPASS_ENV, _TEST_ISOLATION_MARKER_ENV, - _has_pytest_ancestor, _in_test_context, _is_production_state_db, _process_looks_like_pytest, - _real_platform_state_root, _running_under_pytest, _set_last_init_error, get_last_init_error, + _STATE_DB_GUARD_BYPASS_ENV, _in_test_context, _is_production_state_db, + _process_looks_like_pytest, _real_platform_state_root, _running_under_pytest, + _set_last_init_error, get_last_init_error, ) from hermes_state_readpool import ( # noqa: F401 (re-exported; tests import from hermes_state) - _HANDLES_PER_PATH_WARN, _READ_POOL_MAX, _READ_POOL_PROCESS_MAX, _PathReadBudget, - _proc_fd_targets, _process_read_permits, _read_budget_for, + _READ_POOL_MAX, _proc_fd_targets, _process_read_permits, _read_budget_for, ) from hermes_state_sessions import ( # noqa: F401 (re-exported; the historical import path) - SESSION_STATUS_COMPLETE, SESSION_STATUS_EMPTY, SESSION_STATUS_ERROR, SESSION_STATUS_INTERRUPTED, - SessionSessionsMixin, _collect_delegate_child_ids, _cwd_prefix_clause, - _delete_delegate_children, _parse_model_config, _session_filter_where, - classify_session_status, workspace_key, -) -from hermes_state_fts import ( # noqa: F401 (re-exported; the historical import path) - FTS_CJK_TABLE_SQL, FTS_CJK_TRIGGER_SQL, SessionFtsSetupMixin, fts5_cjk_so_path, - load_fts5_cjk_extension, + SessionSessionsMixin, _cwd_prefix_clause, workspace_key, ) +from hermes_state_fts import SessionFtsSetupMixin, load_fts5_cjk_extension # noqa: F401 (re-exported) from hermes_state_portability import SessionPortabilityMixin from hermes_state_telegram import SessionTelegramTopicsMixin from hermes_state_schema import SessionSchemaMixin @@ -123,12 +112,10 @@ MAX_SAFE_EXPORT_MESSAGES = 20_000 def _configured_transcript_limit(key: str, fallback: int) -> int: """``sessions.`` from config.yaml (lazy import: circular at load), else - *fallback*. 0 disables the guard. Not cached: load_config_readonly is - mtime-cached already, and fresh resolution keeps monkeypatching tests working.""" + *fallback*. 0 disables the guard. Not cached (load_config_readonly is).""" try: from hermes_cli.config import load_config_readonly - sessions_cfg = load_config_readonly().get("sessions") or {} - value = sessions_cfg.get(key) + value = (load_config_readonly().get("sessions") or {}).get(key) if value is None: return fallback limit = int(value) @@ -181,31 +168,20 @@ def _system_prompt_hash(system_prompt: str) -> str: def _compression_lock_holder_process_is_dead(holder: str) -> bool: """True only when a ``pid=`` lock holder's local PID is provably gone. - - A process killed mid-compression cannot release its lease and every new - turn would re-attempt compaction until TTL expiry. Reclaim only on kernel - proof; unstructured/same-process holders and any probe doubt stay protected - (PID reuse must never steal a live lease; a wrongly-kept lease self-heals via TTL). - """ + Reclaim on kernel proof only: unstructured/same-process holders (another + thread's live lease) and any probe doubt keep the lease until TTL expiry + (PID reuse must never steal a live lease; a wrongly-kept one self-heals).""" match = _COMPRESSION_LOCK_HOLDER_PID_RE.search(holder or "") - if match is None: - return False - try: - pid = int(match.group(1)) - except (TypeError, ValueError): - return False - # Same-process holder (another thread's live lease): never self-reclaim — - # the lease refresher and release path own it. + pid = int(match.group(1)) if match else 0 if pid <= 0 or pid == os.getpid(): return False if psutil is not None: try: - # Canonical cross-platform liveness answer; recycled PIDs read as alive (conservative). - return not psutil.pid_exists(pid) + return not psutil.pid_exists(pid) # recycled PIDs read as alive (conservative) except Exception: - return False # any doubt → keep the lease until TTL expiry - # psutil-less fallback is POSIX-only: os.kill(pid, 0) is NOT a no-op probe on - # Windows (sig=0 maps to CTRL_C_EVENT and can kill the target's console group). + return False + # psutil-less fallback is POSIX-only: on Windows os.kill(pid, 0) maps sig=0 to + # CTRL_C_EVENT and can kill the target's console group. if os.name == "nt": return False try: @@ -223,10 +199,9 @@ def _scrub_surrogates(value: Any) -> Any: return _sanitize_surrogates(value) if isinstance(value, str) else value -# Billing buckets that aren't a routable provider identity. A session that -# persisted only one of these (never ran /model) falls back to the config -# default rather than restoring a bare bucket. Shared by session_gateway_runtime -# and tui_gateway.server so the two consumers cannot drift. +# Billing buckets that aren't a routable provider identity: a session that +# persisted only one of these (never ran /model) falls back to the config default. +# Shared by session_gateway_runtime and tui_gateway.server so they cannot drift. _BARE_BILLING_PROVIDERS = frozenset({"auto", "custom"}) @@ -234,56 +209,42 @@ T = TypeVar("T") DEFAULT_DB_PATH = get_hermes_home() / "state.db" -# Back off from read-only opens for this long after one fails: long enough -# that an unreadable file isn't retried per query, short enough that transient -# fd pressure doesn't strand the read pool. +# Back off from read-only opens after one fails: not retried per query, but short +# enough that transient fd pressure doesn't strand the read pool. _READ_OPEN_RETRY_SECONDS = 60.0 -# Transient SQLITE_IOERR retry budget for READ-ONLY opens. A WAL writer's -# checkpoint / reset / frame flush can surface "disk I/O error" to a concurrent -# mode=ro reader for a millisecond-wide window (ro cannot do the -shm recovery -# the read needs). NOT attempted on writable opens: a writer owns the -# transition, so an IOERR there is a real storage/fd problem. +# Transient SQLITE_IOERR retry budget for READ-ONLY opens: a WAL writer's +# checkpoint/reset/frame flush surfaces "disk I/O error" to a concurrent mode=ro +# reader for a millisecond-wide window (ro cannot do the -shm recovery). Never +# for writable opens: a writer owns the transition, so an IOERR there is real. _READ_ONLY_IOERR_RETRY_ATTEMPTS = 3 _READ_ONLY_IOERR_RETRY_BACKOFF_S = 0.05 -def _is_transient_read_only_ioerr(exc: sqlite3.OperationalError, *, attempt: int) -> bool: - """Retry a read-only open? See _READ_ONLY_IOERR_RETRY_ATTEMPTS: a - persistent IOERR still exhausts the budget and propagates.""" - return attempt < _READ_ONLY_IOERR_RETRY_ATTEMPTS and _DISK_IO_ERROR_MARKER in str(exc).lower() - -# Import-time snapshot so _default_db_path() can detect a deliberately -# re-pointed DEFAULT_DB_PATH (tests monkeypatch the constant directly). +# Import-time snapshot so _default_db_path() can detect a re-pointed +# DEFAULT_DB_PATH (tests monkeypatch the constant directly). _IMPORT_DEFAULT_DB_PATH = DEFAULT_DB_PATH def _default_db_path() -> Path: - """Default state DB path at CALL time. A re-pointed ``DEFAULT_DB_PATH`` (the - test escape hatch) wins; otherwise ``get_hermes_home()`` is resolved fresh so - a runtime HERMES_HOME redirect works regardless of import order (the frozen - import-time value pointed every default SessionDB() at the real state.db).""" + """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.""" if DEFAULT_DB_PATH != _IMPORT_DEFAULT_DB_PATH: return DEFAULT_DB_PATH return get_hermes_home() / "state.db" # Live-DB guard knobs live HERE (not in hermes_state_guard): the hermetic conftest -# monkeypatches ``hermes_state._STATE_DB_GUARD_BYPASS`` / ``_EXTRA_DENY_ROOTS``. -#: Escape hatch for tests that genuinely need the real DB (conftest sets it for -#: ``@pytest.mark.live_system_guard_bypass``); scripts may set it explicitly. +# monkeypatches ``hermes_state._STATE_DB_GUARD_BYPASS`` (escape hatch for +# ``@pytest.mark.live_system_guard_bypass``) and ``_EXTRA_DENY_ROOTS`` (the +# pre-sandbox root, so custom-HERMES_HOME deployments are covered too). _STATE_DB_GUARD_BYPASS = False - -#: Extra production roots to refuse; conftest injects the pre-sandbox root so -#: custom-HERMES_HOME deployments are covered too. _STATE_DB_GUARD_EXTRA_DENY_ROOTS: Tuple[Path, ...] = () def _production_state_roots() -> List[Path]: - roots: List[Path] = [] - real_root = _real_platform_state_root() - if real_root is not None: - roots.append(real_root) + roots = [r for r in (_real_platform_state_root(),) if r is not None] for extra in _STATE_DB_GUARD_EXTRA_DENY_ROOTS: try: roots.append(Path(extra).expanduser().resolve()) @@ -294,11 +255,8 @@ def _production_state_roots() -> List[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, see :func:`_in_test_context`) - resolves a production DB. No-op outside pytest and for hermetic paths.""" - if _STATE_DB_GUARD_BYPASS or os.environ.get(_STATE_DB_GUARD_BYPASS_ENV): - return - if not _in_test_context(): + 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: resolved = Path(db_path).expanduser().resolve() @@ -318,8 +276,7 @@ def _ensure_test_isolation(db_path: Path) -> None: ) -# Openings of the background-review harness prompts (agent/background_review.py), -# matched case-sensitively against leading user/system content. +# Openings of the background-review harness prompts (agent/background_review.py). _REVIEW_HARNESS_PREFIXES = ( "Review the conversation above and update the skill library", "Review the conversation above and consider saving to memory", @@ -327,8 +284,8 @@ _REVIEW_HARNESS_PREFIXES = ( def _is_background_review_harness_message(msg: Dict[str, Any]) -> bool: - """Persisted background-review harness prompt (older builds wrote the forked - curator's turns into real sessions; replaying them hijacks the session).""" + """Persisted harness prompt (older builds wrote the forked curator's turns + into real sessions; replaying them hijacks the session).""" if not isinstance(msg, dict) or msg.get("role") not in {"user", "system"}: return False content = msg.get("content") @@ -336,8 +293,7 @@ def _is_background_review_harness_message(msg: Dict[str, Any]) -> bool: def _strip_background_review_harness(messages: List[Dict[str, Any]]) -> List[Dict[str, Any]]: - """Drop harness messages and the curator-mode assistant reply that - immediately followed each; everything else passes through in order.""" + """Drop harness messages and the curator-mode assistant reply that immediately followed each.""" if not messages: return messages out: List[Dict[str, Any]] = [] @@ -359,9 +315,8 @@ _STALE_TOOL_CALL_MARKER_RE = re.compile(r"^\[[A-Za-z_][A-Za-z0-9_.-]*\]$") def _is_stale_tool_call_marker_message(msg: Dict[str, Any]) -> bool: - """Assistant tool-call turn whose content is a bare ``[marker]`` — an older - conversation_loop cached a local template's marker and persisted it as the - "final response"; sessions written before the fix still carry these rows.""" + """Assistant tool-call turn whose content is a bare ``[marker]`` (an older + conversation_loop persisted a local template's marker as the final response).""" if not isinstance(msg, dict) or msg.get("role") != "assistant" or not msg.get("tool_calls"): return False content = msg.get("content") @@ -369,24 +324,22 @@ def _is_stale_tool_call_marker_message(msg: Dict[str, Any]) -> bool: def _strip_stale_tool_call_markers(messages: List[Dict[str, Any]]) -> List[Dict[str, Any]]: - """Blank stale ``[marker]`` assistant content: replaying it teaches the model - to keep emitting the marker. Only ``content`` is blanked; tool_call/result - pairing stays intact.""" + """Blank stale ``[marker]`` assistant content (replaying it teaches the model + to keep emitting it); tool_call/result pairing stays intact.""" repaired = 0 for msg in filter(_is_stale_tool_call_marker_message, messages): msg["content"] = "" repaired += 1 if repaired: logger.info( - "Cleared %d stale tool-call marker message(s) while restoring session (#78148)", - repaired, + "Cleared %d stale tool-call marker message(s) while restoring session (#78148)", repaired, ) return messages def format_session_db_unavailable(prefix: str = "Session database not available") -> str: """User-facing "session DB unavailable" message with the captured init cause - (e.g. "locking protocol" from NFS/SMB, with a WAL-docs hint).""" + (plus a WAL-docs hint for NFS/SMB-style locking failures).""" cause = get_last_init_error() if not cause: return f"{prefix}." @@ -401,21 +354,12 @@ def format_session_db_unavailable(prefix: str = "Session database not available" _repair_attempted_paths: set[str] = set() _repair_attempt_lock = threading.Lock() - -# Cross-process schema-surgery lock: ``_repair_attempt_lock`` covers one -# interpreter only, while gateway, Desktop backend, CLI and TUI worker share the -# file and each used to run surgery + VACUUM on top of the winner's. Timeout -# sized for the slowest legitimate holder (VACUUM over a multi-GB DB). +# Cross-process schema-surgery lock timeout (``_repair_attempt_lock`` covers one +# interpreter only); sized for the slowest legitimate holder (VACUUM, multi-GB DB). _REPAIR_LOCK_TIMEOUT_SECONDS = 120.0 _IS_WINDOWS = sys.platform == "win32" -# Repair-loop bounding (hermes_state_repair): unhealable b-tree damage failed -# repair on every start, each pass taking a fresh ~900MB backup (89GB of -# identical copies). A sidecar attempt ledger (fingerprint = size + content -# sample) refuses surgery after _MAX_PERSISTENT_REPAIR_ATTEMPTS, and backups are -# deduped and capped at _MAX_MALFORMED_BACKUPS. - 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.""" @@ -434,23 +378,22 @@ def divert_session_transcript_jsonl(session_id: str, messages) -> "Optional[Path return path -# ── Process-wide shared SessionDB registry ── -# Lives in hermes_state_registry.py; re-exported here for the historical import -# path. Long-lived in-process callers (gateway, tui_gateway, cron, in-process -# tools) share ONE writer connection per resolved path via -# get_shared_session_db(); CLI one-shots, recovery flows and read-only -# cross-profile opens use SessionDB() directly with their own close(). +# Process-wide shared SessionDB registry (hermes_state_registry): long-lived +# in-process callers share ONE writer connection per resolved path via +# get_shared_session_db(); one-shots use SessionDB() with their own close(). from hermes_state_registry import ( # noqa: F401 (re-export) close_shared_session_dbs, get_shared_session_db, release_or_close, ) + class SessionDB( - SessionSessionsMixin, SessionFtsSetupMixin, SessionSearchMixin, SessionSchemaMixin, SessionPortabilityMixin, SessionTelegramTopicsMixin, - SessionCompressionMixin, SessionGatewayMixin, SessionMaintenanceMixin, SessionUsageMixin, - SessionTitlesMixin, SessionMessagesMixin, + SessionSessionsMixin, SessionFtsSetupMixin, SessionSearchMixin, SessionSchemaMixin, + SessionPortabilityMixin, SessionTelegramTopicsMixin, SessionCompressionMixin, + 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); each method opens its own cursor.""" + pattern (many reader threads, one writer via WAL).""" # Only these state-owned producers join automatic stale-open reconciliation; # messaging/UI sources have their own lifecycle owners; unknown sources fail closed. @@ -459,24 +402,19 @@ class SessionDB( ) # ── Write-contention tuning ── - # SQLite's deterministic busy handler convoys under many hermes processes, - # so the SQLite timeout stays short (1s) and retries use random jitter. - # Patience is TIME-based: a sibling legitimately holds the lock for seconds - # (TRUNCATE checkpoint at close, VACUUM after auto-prune, recovery, an older - # process's unbounded FTS optimize); an attempt-counted budget lost that race - # and destroyed the turn as session_persistence_failed on a healthy store. - # Routine writes give up after _WRITE_PATIENCE_S; transcript writes (whose - # failure aborts the user's turn) get _TRANSCRIPT_WRITE_PATIENCE_S. Jitter - # stays small for _WRITE_RETRY_SLOW_AFTER_S, then backs off. + # SQLite's deterministic busy handler convoys under many hermes processes, so + # the SQLite timeout stays short (1s) and retries use random jitter. Patience + # is TIME-based (a sibling legitimately holds the lock for seconds: checkpoint + # at close, VACUUM, recovery, an old process's FTS optimize); attempt-counted + # budgets destroyed turns on a healthy store. Transcript writes (failure + # aborts the user's turn) get the longer budget; observation-only activity + # writes sit on the response-critical path and get a sub-second one. _WRITE_PATIENCE_S = 20.0 _TRANSCRIPT_WRITE_PATIENCE_S = 60.0 - # Observation-only activity heartbeat/label writes sit on the response- - # critical path: sub-second budget; a skipped write retries next window. _ACTIVITY_WRITE_PATIENCE_S = 0.5 - # A live compression lock gets a short budget: compression publishes in a - # couple of seconds, so a brief wait saves most concurrent turns — but the - # lease is a correctness boundary, so a writer still locked out afterwards - # must be refused rather than land 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 rather than landing a stale turn in a wedged compression. _COMPRESSION_BUSY_WAIT_S = 5.0 _WRITE_RETRY_MIN_S = 0.020 # 20ms _WRITE_RETRY_MAX_S = 0.150 # 150ms @@ -485,22 +423,20 @@ class SessionDB( _WRITE_RETRY_SLOW_MAX_S = 1.000 # 1s # PASSIVE WAL checkpoint every N successful writes. _CHECKPOINT_EVERY_N_WRITES = 50 - # Bounded FTS ``'merge'`` (milliseconds 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. usermerge - # is lowered to 2 so levels with >= 2 segments merge (default 4 never converges). + # 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. _FTS_MERGE_EVERY_N_WRITES = 1000 _FTS_MERGE_MAX_PAGES_PER_INDEX = 500 _FTS_MERGE_COMMANDS_PER_PASS = 4 - # Imports cap lower than exports: an import holds one BEGIN IMMEDIATE, so - # bounded batches avoid starving live writers (one dashboard file at a time). + # Imports cap lower than exports: an import holds one BEGIN IMMEDIATE. _IMPORT_MAX_SESSIONS = 500 _IMPORT_MAX_MESSAGES_PER_SESSION = 10_000 _IMPORT_MAX_TOTAL_MESSAGES = 50_000 _IMPORT_MAX_SESSION_BYTES = 5 * 1024 * 1024 _IMPORT_MAX_TOTAL_BYTES = 25 * 1024 * 1024 # Accounting workers retire when idle so a bound-method target can't keep an - # abandoned SessionDB (and its descriptors) alive; a later enqueue restarts one. + # abandoned SessionDB (and its descriptors) alive. _TOKEN_WRITER_IDLE_SECONDS = 30.0 @staticmethod @@ -546,41 +482,31 @@ class SessionDB( self.read_only = read_only self._lock = threading.Lock() # Read-path split (WAL only): reads borrow a read-only connection from a - # BOUNDED pool so they never queue behind writer flushes on self._lock - # (see _read_ctx). The old per-thread scheme pinned one connection (two - # fds) per SessionDB x anyio worker thread for the process lifetime until - # a 256 RLIMIT_NOFILE service hit EMFILE while staying alive, so the - # supervisor's restart-on-exit never fired. + # BOUNDED pool so they never queue behind writer flushes on self._lock (see + # _read_ctx); the old per-thread connections pinned fds for the process + # lifetime and hit EMFILE while staying alive (supervisor never restarted). self._read_pool: "queue.LifoQueue[sqlite3.Connection]" = queue.LifoQueue(maxsize=_READ_POOL_MAX) - # Permits bound PEAK descriptors (the pool bounds only the idle set) and - # are shared per DATABASE PATH (see _PathReadBudget). Acquired - # non-blocking on purpose: a reader without a permit degrades to the - # writer lock — blocking would turn fd exhaustion into a stall. + # Permits bound PEAK descriptors (the pool bounds only the idle set), shared + # per DATABASE PATH (_PathReadBudget); acquired non-blocking so a reader + # without a permit degrades to the writer lock instead of stalling. self._read_budget = _read_budget_for(self.db_path) self._read_budget.register(self) - # Bound to the semaphore itself so every release site is unchanged. self._read_permits = self._read_budget.permits - # Reads that fell back to the writer connection — the only visible - # signal that the ceiling is being reached (diagnostic, not load-bearing). - self._read_permit_exhausted = 0 self._read_conns_lock = threading.Lock() - # Set when close() begins; a reader still in flight then closes its own + # 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 stamp — a TIMESTAMP, not a sticky - # bool: the likeliest trigger is transient EMFILE, and a permanent flag - # would demote every reader (the gateway shares one SessionDB across all - # agents) to the writer lock forever. Expires after _READ_OPEN_RETRY_SECONDS. + # "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. self._read_open_failed_at = 0.0 self._wal_active = False self._write_count = 0 - # File identity of the opened state.db, compared on every write (and - # before FTS fail-open / reopen) so an out-of-band replace cannot limp - # through in-place surgery. Inode catches mv/new-file; application_id - # catches cp onto the same path (same inode, truncate+rewrite). + # File identity of the opened state.db, compared on every write so an + # out-of-band replace cannot limp through in-place surgery. Inode catches + # mv/new-file; application_id catches cp onto the same path. self._db_file_identity: Optional[tuple] = None self._db_file_application_id: int = 0 - self._db_file_generation_token: str = "" self._db_sidecar_identity: Dict[str, tuple] = {} self._db_replaced = self._db_wal_generation_lost = False # Sticky quarantine (see StateDbCorruptError); never cleared. @@ -588,12 +514,12 @@ class SessionDB( self._db_corrupt_reason = "" 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 extension present on the writer connection; + # _fts_cjk_loaded: tokenizer present 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 (queue_token_counts). Distinct from self._lock - # so enqueue/flush bookkeeping never contends with SQLite writes. + # Async token accounting; distinct from self._lock so enqueue/flush + # bookkeeping never contends with SQLite writes. self._token_queue: deque = deque() self._token_queue_cond = threading.Condition(threading.Lock()) self._token_writer_thread: Optional[threading.Thread] = None @@ -609,16 +535,13 @@ class SessionDB( initialization_complete = True return self.db_path.parent.mkdir(parents=True, exist_ok=True) - # Read-only file/sidecar preflight: repair-or-refuse BEFORE the first - # connection, for an actionable message instead of an opaque "attempt - # to write a readonly database" from deep inside _init_schema. - if not read_only: - preflight_db_writability(self.db_path, db_label="state.db") + # 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. + preflight_db_writability(self.db_path, db_label="state.db") # Serialize zero-byte check, quarantine, connect and schema commit so # concurrent openers don't race the absent-path -> schema-commit window. - needs_startup_guard = not read_only and ( - not self.db_path.exists() or is_zeroed_state_db(self.db_path) - ) + needs_startup_guard = not self.db_path.exists() or is_zeroed_state_db(self.db_path) try: self._open_with_optional_startup_guard(needs_startup_guard) except sqlite3.DatabaseError as exc: @@ -631,25 +554,18 @@ class SessionDB( "state.db schema is malformed (%s) — attempting automatic " "repair (a backup copy is made first).", exc, ) - try: - if self._conn is not None: - self._conn.close() - except Exception: - pass - report = repair_state_db_schema(self.db_path) - if not report.get("repaired"): + self._close_connection_quietly(self._conn) + 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, no surprise disk/latency cost on an unattended open. + # auto-started on open (no background worker racing session lifecycle). self._ensure_db_file_generation() self._record_db_file_identity() initialization_complete = True except Exception as exc: - # Surface WHY via /resume and friends; deliberately never cleared on - # success (see _set_last_init_error). Callers keep their - # ``self._session_db = None`` degradation path. + # Surface WHY via /resume and friends (never cleared on success, see + # _set_last_init_error); callers keep their ``_session_db = None`` path. _set_last_init_error(f"{type(exc).__name__}: {exc}") raise finally: @@ -658,14 +574,12 @@ class SessionDB( self._close_connection_quietly(conn) def _open_read_only(self) -> None: - """Read-only attach for cross-profile aggregation: no schema init, NO - write lock (sidebar polling never contends with that profile's backend); - the DB must already exist. FTS flags are probed with SELECTs only, and - the connection is closed on ANY probe failure (malformed schema raises - DatabaseError) so a leaked tracked connection cannot block the forensic - backup the writable heal takes next.""" - open_attempt = 0 - while True: + """Read-only attach for cross-profile aggregation: no schema init, NO write + lock (sidebar polling never contends with that profile's backend); the DB + must exist. FTS flags are probed with SELECTs only, and the connection is + closed on ANY probe failure (malformed schema raises DatabaseError) so a + leaked tracked connection cannot block the forensic backup the writable heal takes next.""" + for attempt in range(_READ_ONLY_IOERR_RETRY_ATTEMPTS + 1): try: self._conn = conn = _connect_tracked_db( f"file:{self.db_path}?mode=ro", tracking_path=self.db_path, uri=True, @@ -682,29 +596,19 @@ class SessionDB( ) except BaseException: self._conn = None - try: - conn.close() - except Exception: - pass + self._close_connection_quietly(conn) raise return except sqlite3.OperationalError as ioerr: - # A WAL checkpoint / reset / frame-flush in flight on the writer - # side can surface SQLITE_IOERR to a concurrent mode=ro reader - # (it cannot perform the -shm recovery the read needs). The - # transition closes in milliseconds; retry a bounded number of - # times before classifying the store as failed. - if not _is_transient_read_only_ioerr(ioerr, attempt=open_attempt): + # Transient SQLITE_IOERR window (see _READ_ONLY_IOERR_RETRY_ATTEMPTS); + # a persistent one exhausts the budget and propagates. + if attempt >= _READ_ONLY_IOERR_RETRY_ATTEMPTS or _DISK_IO_ERROR_MARKER not in str(ioerr).lower(): raise - open_attempt += 1 time.sleep(_READ_ONLY_IOERR_RETRY_BACKOFF_S) 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, do not open the zeroed file (it would fail - opaquely or risk further damage) — raise with the clear message. - """ + """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: @@ -724,25 +628,29 @@ class SessionDB( if qpath is None and self.db_path.exists() and is_zeroed_state_db(self.db_path): raise sqlite3.DatabaseError(msg) - def _connect_and_init(self) -> None: - # Refuse before sqlite3.connect (under the startup lock) so we cannot - # mint a replacement WAL while a live writer still holds a deleted - # sidecar inode. - refuse_deleted_wal_generation(self.db_path) - self._conn = _connect_tracked_db( - str(self.db_path), - check_same_thread=False, - # Short timeout — application-level jittered retry handles - # contention instead of SQLite's internal busy handler (up to 30s). - timeout=1.0, - # None = we manage transactions ourselves (explicit BEGIN IMMEDIATE). - isolation_level=None, + 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.""" + conn = _connect_tracked_db( + str(self.db_path), check_same_thread=False, timeout=1.0, isolation_level=None, ) - self._conn.row_factory = sqlite3.Row - self._wal_active = apply_wal_with_fallback(self._conn, db_label="state.db") == "wal" - apply_database_pragmas(self._conn, db_label="state.db") - self._conn.execute("PRAGMA foreign_keys=ON") - self._fts_cjk_loaded = load_fts5_cjk_extension(self._conn) + try: + conn.row_factory = sqlite3.Row + self._wal_active = apply_wal_with_fallback(conn, db_label="state.db") == "wal" + apply_database_pragmas(conn, db_label="state.db") + conn.execute("PRAGMA foreign_keys=ON") + self._fts_cjk_loaded = load_fts5_cjk_extension(conn) + except BaseException: + self._close_connection_quietly(conn) + raise + return conn + + def _connect_and_init(self) -> None: + # Refuse before sqlite3.connect (under the startup lock) so we cannot mint + # a replacement WAL while a live writer still holds a deleted sidecar inode. + refuse_deleted_wal_generation(self.db_path) + self._conn = self._open_writer_conn() self._init_schema() def _connect_and_init_with_lock_patience(self) -> None: @@ -759,11 +667,7 @@ class SessionDB( err = str(exc).lower() if "locked" not in err and "busy" not in err: raise - try: - if self._conn is not None: - self._conn.close() - except Exception: - pass + self._close_connection_quietly(self._conn) now = time.monotonic() if now >= deadline: raise @@ -790,13 +694,10 @@ class SessionDB( def _get_read_conn(self) -> Optional[sqlite3.Connection]: """Open a fresh read-only connection, or None when unavailable (callers - return it to self._read_pool; this opens, it does not track). - - WAL only: WAL readers never block on the writer, so reads skip - self._lock; under DELETE journal mode (NFS fallback) readers hit - SQLITE_BUSY storms, so the legacy locked path stays. Autocommit reads - see everything committed so far (read-your-writes for flush-then-search). - """ + return it to self._read_pool). WAL only: WAL readers never block on the + writer, so reads skip self._lock; under DELETE journal mode (NFS fallback) + readers hit SQLITE_BUSY storms, so the legacy locked path stays. Autocommit + reads see everything committed so far (read-your-writes for flush-then-search).""" if not self._wal_active or self.read_only: return None with self._read_conns_lock: @@ -804,40 +705,27 @@ class SessionDB( return None if ( self._read_open_failed_at - and time.monotonic() - self._read_open_failed_at - < _READ_OPEN_RETRY_SECONDS + and time.monotonic() - self._read_open_failed_at < _READ_OPEN_RETRY_SECONDS ): return None # Permit BEFORE the open: openers race for permits, not descriptors. if not self._read_budget.acquire(self): - with self._read_conns_lock: - self._read_permit_exhausted += 1 logger.debug( "read pool at capacity (%d) for %s; serving this read from the " - "locked writer connection", - _READ_POOL_MAX, - self.db_path, + "locked writer connection", _READ_POOL_MAX, self.db_path, ) return None conn = None # bound before the try so the handlers can close a half-open one try: + # 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, - # Pooled connections are borrowed by whichever thread reads next - # (sqlite3 otherwise refuses cross-thread use, including close() - # — how the old per-thread connections leaked their fds). - # Exclusive ownership is enforced by pool checkout, not sqlite3. - check_same_thread=False, - timeout=5.0, - isolation_level=None, + f"file:{self.db_path}?mode=ro", tracking_path=self.db_path, uri=True, + check_same_thread=False, timeout=5.0, isolation_level=None, ) conn.row_factory = sqlite3.Row apply_database_pragmas(conn, db_label="state.db") - # The tokenizer registers in the connection's in-memory registry, - # not the file, so mode=ro is fine. - if self._fts_cjk_loaded: + if self._fts_cjk_loaded: # registers in the connection, not the file: ro is fine load_fts5_cjk_extension(conn) except sqlite3.Error: # A half-open connection (open ok, extension load failed) is a live @@ -857,8 +745,7 @@ class SessionDB( def _evict_one_idle_read_conn(self) -> bool: """Close one idle pooled connection (a peer on the same file wants its - permit). Only the idle set is reachable — never pulls a connection out - from under a live reader. Returns whether a permit was released.""" + permit); never pulls a connection out from under a live reader.""" try: conn = self._read_pool.get_nowait() except queue.Empty: @@ -877,14 +764,10 @@ class SessionDB( logger.warning("partially-opened read conn close failed for %s: %s", self.db_path, exc) def _close_read_conn(self, conn) -> None: - """Close a pooled read connection and release its permit. - - A failing close leaks a tracked fd, so it is logged, never swallowed. - The permit is released even then: withholding it would turn one leaked - fd into a permanently narrower read path. Pairs with _get_read_conn(); - over-releasing the BoundedSemaphore raises ValueError rather than - silently widening the ceiling. - """ + """Close a pooled read connection and release its permit. A failing close + leaks a tracked fd (logged, never swallowed) but still releases the permit: + withholding it would permanently narrow the read path. Pairs with + _get_read_conn(); over-releasing the BoundedSemaphore raises ValueError.""" try: conn.close() except Exception as exc: @@ -893,10 +776,8 @@ class SessionDB( 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. The single acquisition seam: a pool hit costs no permit - (the connection already holds one), only _get_read_conn() takes one, so - peak live connections stay bounded however many threads miss at once.""" + """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: @@ -907,10 +788,9 @@ class SessionDB( @contextmanager def _read_ctx(self) -> Iterator[sqlite3.Connection]: """Yield a connection for read-only statements: a pooled read-only - connection with NO lock under WAL (the writer lock was a global choke - point), checked out for the block; otherwise (non-WAL, open failure, - ceiling reached) the writer connection under self._lock — the deliberate - degradation: slower than EMFILE, which the supervisor cannot see.""" + connection with NO lock under WAL; otherwise (non-WAL, open failure, + ceiling reached) the writer connection under self._lock — deliberate + degradation: slower beats EMFILE, which the supervisor cannot see.""" conn = self._checkout_read_conn() if conn is not None: try: @@ -925,9 +805,8 @@ class SessionDB( except queue.Full: pass if not returned: - # close() drained the pool: this connection is surplus. - # queue.Full is unreachable while permits == maxsize, but the - # branch is load-bearing if they ever drift apart (a leak). + # close() drained the pool (or queue.Full: unreachable while + # permits == maxsize, load-bearing if they drift): surplus. self._close_read_conn(conn) return with self._lock: @@ -936,19 +815,18 @@ class SessionDB( yield cast(sqlite3.Connection, self._conn) def _reopen_after_close_locked(self, context: str = "write") -> None: - """Reopen the writer after ``close()`` raced a live caller (a teardown - owner set ``_conn = None`` while a worker still had a transcript flush - to land; the turn's tail was silently dropped). Loud (WARNING) and - bounded (only after an explicit close()); ``__del__`` still releases it. - Caller holds ``self._lock``. A failed reopen names the race in its error.""" + """Reopen the writer after ``close()`` raced a live caller (a teardown owner + set ``_conn = None`` while a worker still had a transcript flush to land). + Loud (WARNING) and bounded (only after an explicit close()). Caller holds + ``self._lock``. No _init_schema: no DDL races with siblings during teardown.""" if self.read_only: raise sqlite3.ProgrammingError( 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 (and its close-time checkpoint) to a damaged file. if self._db_replaced or self._db_file_was_replaced(): self._halt_db_replaced() if self._db_corrupt: @@ -961,57 +839,33 @@ class SessionDB( self._halt_deleted_wal_generation() logger.warning( "state.db connection for %s was closed while a %s was still in " - "flight — reopening (teardown/worker race, #94736)", - self.db_path, - context, + "flight — reopening (teardown/worker race, #94736)", self.db_path, context, ) try: - conn = _connect_tracked_db( - str(self.db_path), check_same_thread=False, timeout=1.0, isolation_level=None, - ) + self._conn = self._open_writer_conn() except Exception as exc: raise sqlite3.OperationalError( f"state.db connection was closed while a {context} was still " f"in flight (a session-teardown path called close() before " f"this worker finished — #94736) and the automatic reopen failed: {exc}" ) from exc - try: - conn.row_factory = sqlite3.Row - self._wal_active = (apply_wal_with_fallback(conn, db_label="state.db") == "wal") - apply_database_pragmas(conn, db_label="state.db") - conn.execute("PRAGMA foreign_keys=ON") - self._fts_cjk_loaded = load_fts5_cjk_extension(conn) - except Exception as exc: - self._close_connection_quietly(conn) - raise sqlite3.OperationalError( - f"state.db reopen after close() succeeded but connection setup failed: {exc}" - ) from exc - # Schema was initialised by the original open; no _init_schema here (no - # DDL races with siblings during teardown). - self._conn = conn def _execute_write( self, fn: Callable[[sqlite3.Connection], T], patience_s: Optional[float] = None, ) -> T: """Run *fn(conn)* inside BEGIN IMMEDIATE with jittered lock retry; commit is handled here (callers must not commit). Returns *fn*'s result. - BEGIN IMMEDIATE takes the WAL write lock up front so contention surfaces - immediately; on locked/busy the Python lock is released, a random jitter - slept, and the WHOLE callback retried (see the class tuning comment for - the two patience budgets and the jitter schedule). *fn* must therefore - stay idempotent under retry. - """ + immediately; on locked/busy the Python lock is released, a jitter slept, + and the WHOLE callback retried — *fn* must stay idempotent under retry.""" 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, not from the start of the write. + # 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: the - # callback has not run, so nothing is replayed (exactly-once safe). Once - # it has started, an IOERR leaves settlement unknown and must propagate — - # this helper owns non-idempotent transcript/counter mutations. + # 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. ioerr_begin_retried = False while True: self._raise_if_db_corrupt() @@ -1040,9 +894,8 @@ 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 and sends the - # operator hunting disk space that was never the problem. + # Transient (see _COMPRESSION_BUSY_WAIT_S): without a wait, a steer + # landing mid-compression aborts the turn. if compression_deadline is None: compression_deadline = min( time.monotonic() + self._COMPRESSION_BUSY_WAIT_S, deadline @@ -1073,9 +926,8 @@ class SessionDB( and not ioerr_begin_retried and self._sleep_before_write_retry(deadline, patience_s) ): - # Retry on the SAME connection. Never close()+reopen to - # "heal": close() cancels this process's POSIX locks on the - # file for every sibling connection (howtocorrupt §2.2). + # Retry on the SAME connection: close()+reopen cancels 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 @@ -1090,10 +942,9 @@ class SessionDB( 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. Never rebuild FTS - # from this live path (minutes of writer lock on a multi-GB DB): - # detach the derived indexes atomically and retry the write. + # 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). if self._enter_fts_fail_open(exc): continue # What survives both checks is structural damage: quarantine. @@ -1155,7 +1006,6 @@ class SessionDB( ).fetchone() if row and row[0]: token = str(row[0]) - self._db_file_generation_token = token pragma_row = self._conn.execute("PRAGMA application_id").fetchone() current = int(pragma_row[0] or 0) if pragma_row else 0 if current == 0: @@ -1207,14 +1057,10 @@ class SessionDB( raise StateDbReplacedError(_STATE_DB_REPLACED_MSG) def _wal_generation_was_lost(self) -> bool: - """True when the WAL/SHM generation this handle opened is gone. - - Recorded generation: pure stat (missing/replaced inode = split); no - /proc walk on healthy writes. Empty identity (fresh DB whose WAL appears - after open, or cleared by a clean close()): probe /proc/self/fd for - deleted sidecars and adopt the current ones once clean. The full - /proc/*/fd walk is reserved for refuse_deleted_wal_generation on open. - """ + """True when the WAL/SHM generation this handle opened is gone. Recorded + generation: pure stat (no /proc walk on healthy writes). Empty identity + (WAL appeared after open, or cleared by a clean close()): probe + /proc/self/fd for deleted sidecars and adopt the current ones once clean.""" recorded = self._db_sidecar_identity or {} base = os.fspath(self.db_path) if recorded: @@ -1295,17 +1141,15 @@ class SessionDB( raise err from exc def _disable_close_time_checkpoint(self) -> None: - """Best-effort SQLITE_DBCONFIG_NO_CKPT_ON_CLOSE (Python 3.12+): skipping - our explicit checkpoint isn't enough, sqlite3's close() still runs the - internal last-connection checkpoint that wrote the incident's 15 pages - under wrong page numbers. See StateDbCorruptError.""" + """Best-effort SQLITE_DBCONFIG_NO_CKPT_ON_CLOSE (Python 3.12+): sqlite3's + close() otherwise runs the internal last-connection checkpoint that wrote + the incident's pages under wrong page numbers (see StateDbCorruptError). + <3.12 has no setconfig; the residual checkpoint only carries + pre-quarantine committed frames, which is tolerable.""" flag = getattr(sqlite3, "SQLITE_DBCONFIG_NO_CKPT_ON_CLOSE", None) conn = self._conn setconfig = getattr(conn, "setconfig", None) - if flag is None or conn is None or setconfig is None: - # <3.12 has no setconfig: the residual close checkpoint is tolerable — it - # can only carry pre-quarantine committed frames; this handle accepts no - # further writes. + if flag is None or setconfig is None: return try: setconfig(flag, True) @@ -1335,12 +1179,10 @@ class SessionDB( return True def _foreign_state_db_holders(self) -> List[Tuple[int, str]]: - """Foreign processes holding this DB or its WAL sidecars. Automatic FTS - repair is structural maintenance and 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. - """ + """Foreign processes holding this DB or its WAL sidecars: automatic FTS + 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. if _IS_WINDOWS: @@ -1353,9 +1195,8 @@ class SessionDB( _canonical_sqlite_path(db_path + "-shm"), } holders: List[Tuple[int, str]] = [] - # Linux: read /proc//fd directly. psutil.open_files() stats the - # literal path, so an unlinked "state.db-wal (deleted)" entry is silently - # dropped and the split-brain holder never seen; readlink keeps the suffix. + # Linux: readlink /proc//fd directly; psutil.open_files() stats the + # literal path and silently drops "state.db-wal (deleted)" entries. if sys.platform.startswith("linux"): try: own_pid = os.getpid() @@ -1368,9 +1209,8 @@ class SessionDB( try: targets = list(_proc_fd_targets(pid)) except OSError: - # Unreadable fd table (other user: root gateway vs user - # desktop). cmdline is world-readable: flag only - # uninspectable holders that look like Hermes. + # Unreadable fd table (other user); cmdline is world-readable: + # flag only uninspectable holders that look like Hermes. cmdline = _read_proc_cmdline(pid) if cmdline is not None and _looks_like_hermes(cmdline): holders.append((pid, f"uninspectable holder: {cmdline[:80]}")) @@ -1379,10 +1219,8 @@ class SessionDB( except Exception as exc: return self._foreign_holder_scan_failed(holders, exc) return holders - # macOS / BSD: psutil.open_files(). macOS does not use the "(deleted)" - # suffix convention, so psutil's filtering is safe here. psutil's - # as_dict() converts AccessDenied to None -> empty iteration; acceptable - # on macOS (the root-gateway/user-desktop topology is Linux-specific). + # macOS / BSD: psutil.open_files() (no "(deleted)" suffix convention there; + # AccessDenied -> None -> empty iteration is acceptable on macOS). try: for process in psutil.process_iter(["pid", "open_files"]): pid = int(process.info["pid"]) @@ -1400,16 +1238,14 @@ class SessionDB( def _foreign_holder_scan_failed(holders: List[Tuple[int, str]], exc: Exception) -> List[Tuple[int, str]]: logger.warning( "Could not prove state.db has no foreign holders; " - "deferring automatic FTS maintenance: %s", - exc, + "deferring automatic FTS maintenance: %s", exc, ) return holders or [(-1, f"open-file scan failed: {exc}")] def _try_wal_checkpoint(self) -> None: """Best-effort PASSIVE WAL checkpoint; never raises. PASSIVE never blocks - writers and leaves the WAL at its high-water mark (bounded by - journal_size_limit); the old TRUNCATE strategy corrupted B-trees on - 65K+ page databases under exclusive-lock I/O pressure.""" + writers; the old TRUNCATE strategy corrupted B-trees on 65K+ page + databases under exclusive-lock I/O pressure.""" if self._db_corrupt: return # quarantined: never checkpoint over a damaged image try: @@ -1421,12 +1257,8 @@ class SessionDB( logger.warning("WAL checkpoint (PASSIVE) failed: %s", exc) def __enter__(self) -> "SessionDB": - """``with SessionDB(path) as db:`` closes the handle on exit. Owners must - release deterministically: a started token writer used to pin the - instance (bound-method target + strong atexit hook) so __del__ never ran - for exactly the handles that leaked descriptors; the writer now retires - when idle and the hook is weak, but "eventually after a GC cycle" is not - a release policy. close() stays idempotent.""" + """``with SessionDB(path) as db:`` closes on exit; owners must release + deterministically ("eventually after a GC cycle" is not a release policy).""" return self def __exit__(self, exc_type, exc, tb) -> bool: @@ -1435,13 +1267,10 @@ class SessionDB( return False def close(self): - """Close the connection: drain queued token deltas (the writer needs the - connection), then a PASSIVE checkpoint on writable handles (NOT - TRUNCATE: per-cron-run connections close many times an hour and a full - WAL reset races the gateway's live writer, tearing B-tree pages). - A registry-shared instance RELEASES one refcount instead, so one - caller's close cannot tear down a connection others still use.""" - if getattr(self, "_shared_registry_owned", False): + """Drain queued token deltas, then a PASSIVE checkpoint on writable handles + (NOT TRUNCATE: a full WAL reset races the gateway's live writer, tearing + B-tree pages). A registry-shared instance RELEASES one refcount instead.""" + if self._shared_registry_owned: from hermes_state_registry import release release(self) return @@ -1449,8 +1278,7 @@ class SessionDB( hook, self._token_atexit_hook = self._token_atexit_hook, None if hook is not None: atexit.unregister(hook) - # Closed flag first (under the lock): an in-flight reader then closes its - # own connection instead of re-populating the drained pool. + # Closed flag first: an in-flight reader then closes its own connection. with self._read_conns_lock: self._read_conns_closed = True while self._evict_one_idle_read_conn(): @@ -1474,15 +1302,13 @@ 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 last close lets SQLite unlink the sidecars — a - # legitimate end of this generation, not a split. Drop it so a - # teardown-race reopen re-adopts what exists instead of halting. + # 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: - """Safety net: close() if the caller forgot (read pool, token writer and - atexit hook too). Attribute access stays guarded: module teardown order - is undefined.""" + """Safety net: close() if the caller forgot. Attribute access stays + guarded: module teardown order is undefined.""" if self.__dict__.get("_conn") is None: return try: @@ -1491,14 +1317,10 @@ class SessionDB( pass # ── Async token accounting (SessionUsageMixin) ── - # update_token_counts() stalls the turn thread for tens-hundreds of ms on a - # cold multi-GB DB after EVERY API call; queue_token_counts() reduces the - # critical path to a deque append, a single-writer thread applies deltas in - # order, coalescing consecutive same-route deltas. Exact readers call - # flush_token_counts() first. Route fields must be equal for two deltas to - # merge (model/billing_* feed COALESCE backfill and the per-model - # attribution key; cost_status/source are last-non-None-wins) so the merged - # UPDATE is byte-for-byte equivalent to applying the deltas sequentially. + # queue_token_counts() reduces the critical path to a deque append; a + # single-writer thread applies deltas in order, coalescing consecutive + # same-route deltas (route fields must be EQUAL to merge so the merged UPDATE + # equals applying them sequentially). Exact readers call flush_token_counts(). _TOKEN_DELTA_SUM_FIELDS = ( "input_tokens", "output_tokens", "cache_read_tokens", "cache_write_tokens", "reasoning_tokens", "api_call_count", @@ -1519,21 +1341,17 @@ class SessionDB( TITLE_SOURCE_USER = "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 (no session-id - # pointer exists); 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) ── - # Prefix distinguishing JSON-encoded structured content (multimodal parts) - # from plain strings; NUL is not legal in normal text, so it cannot collide. + # Prefix marking JSON-encoded structured content; NUL cannot collide with text. _CONTENT_JSON_PREFIX = "\x00json:" - #: Reactions live inside ``display_metadata`` (not a side table) so they - #: survive rewind/compaction row rewrites with the row itself. + #: Reactions live inside ``display_metadata`` so they survive row rewrites. REACTIONS_METADATA_KEY = "reactions" - # Columns every conversation projection decodes (model-fed and display - # views share one SELECT); ``active`` rides along so a display read can - # split compaction-archived rows from the live set 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, " @@ -1593,12 +1411,13 @@ class SessionDB( ``prefix`` (LIKE wildcards escaped) — e.g. ``loop:`` rows.""" if not prefix: return [] - escaped = prefix.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") rows = self._read_all( - "SELECT key, value FROM state_meta WHERE key LIKE ? ESCAPE '\\'", (escaped + "%",), + "SELECT key, value FROM state_meta WHERE key LIKE ? ESCAPE '\\'", + (_escape_like(prefix) + "%",), ) return [(row[0], row[1]) for row in rows] + 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).""" diff --git a/tests/hermes_state/test_session_lifecycle_status.py b/tests/hermes_state/test_session_lifecycle_status.py index 9223dfa351..e56b2e8266 100644 --- a/tests/hermes_state/test_session_lifecycle_status.py +++ b/tests/hermes_state/test_session_lifecycle_status.py @@ -7,12 +7,12 @@ the delete wiring the picker's 'd' key relies on. import pytest -from hermes_state import ( +from hermes_state import SessionDB +from hermes_state_sessions import ( SESSION_STATUS_COMPLETE, SESSION_STATUS_EMPTY, SESSION_STATUS_ERROR, SESSION_STATUS_INTERRUPTED, - SessionDB, classify_session_status, ) diff --git a/tests/test_delegate_cascade_49148.py b/tests/test_delegate_cascade_49148.py index 3369a95aa1..211a3061b1 100644 --- a/tests/test_delegate_cascade_49148.py +++ b/tests/test_delegate_cascade_49148.py @@ -11,7 +11,7 @@ the deletion set, permanently deleting the parent session and its messages. import json import sqlite3 -from hermes_state import _collect_delegate_child_ids, _delete_delegate_children +from hermes_state_sessions import _collect_delegate_child_ids, _delete_delegate_children def _make_conn(): diff --git a/tests/test_session_db_read_conn_pool.py b/tests/test_session_db_read_conn_pool.py index 3afb88f6e1..acfb8cff83 100644 --- a/tests/test_session_db_read_conn_pool.py +++ b/tests/test_session_db_read_conn_pool.py @@ -499,7 +499,8 @@ def test_idle_permits_are_reclaimed_from_a_peer_instance(db): def test_peak_is_bounded_across_many_database_files(tmp_path): """Read connections must be capped for the PROCESS, not just per file.""" import hermes_state - from hermes_state import SessionDB, _READ_POOL_MAX, _READ_POOL_PROCESS_MAX + from hermes_state import SessionDB, _READ_POOL_MAX + from hermes_state_readpool import _READ_POOL_PROCESS_MAX n_files = (_READ_POOL_PROCESS_MAX // _READ_POOL_MAX) + 2 dbs = [] @@ -548,7 +549,8 @@ def test_peak_is_bounded_across_many_database_files(tmp_path): @pytest.mark.requires_wal def test_idle_connections_are_reclaimed_across_database_files(tmp_path): """A quiet profile's idle connections must not starve the busy one.""" - from hermes_state import SessionDB, _READ_POOL_MAX, _READ_POOL_PROCESS_MAX + from hermes_state import SessionDB, _READ_POOL_MAX + from hermes_state_readpool import _READ_POOL_PROCESS_MAX quiet = [] try: @@ -629,7 +631,8 @@ def test_duplicate_handles_on_one_path_are_reported(db, caplog): """Writer connections cannot be capped, so duplicates must be visible.""" import logging - from hermes_state import SessionDB, _HANDLES_PER_PATH_WARN + from hermes_state import SessionDB + from hermes_state_readpool import _HANDLES_PER_PATH_WARN extra = [] try: