diff --git a/hermes_state.py b/hermes_state.py index dc79f35f8d..86589cd6e8 100644 --- a/hermes_state.py +++ b/hermes_state.py @@ -183,7 +183,6 @@ def _scrub_surrogates(value: Any) -> Any: # Shared by session_gateway_runtime and tui_gateway.server so they cannot drift. _BARE_BILLING_PROVIDERS = frozenset({"auto", "custom"}) - T = TypeVar("T") DEFAULT_DB_PATH = get_hermes_home() / "state.db" @@ -199,7 +198,6 @@ _READ_OPEN_RETRY_SECONDS = 60.0 _READ_ONLY_IOERR_RETRY_ATTEMPTS = 3 _READ_ONLY_IOERR_RETRY_BACKOFF_S = 0.05 - # 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 @@ -375,12 +373,11 @@ 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: 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 + # SQLite's deterministic busy handler convoys under many hermes processes: keep its + # timeout short (1s) and retry with random jitter. Patience is TIME-based (a sibling + # legitimately holds the lock for seconds: checkpoint at close, VACUUM, recovery, FTS + # optimize); attempt-counted budgets destroyed turns on a healthy store. Transcript + # 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 = 20.0 _TRANSCRIPT_WRITE_PATIENCE_S = 60.0 @@ -418,7 +415,8 @@ class SessionDB( return None prompt_hash = hashlib.sha256(system_prompt.encode("utf-8")).hexdigest() conn.execute( - "INSERT OR IGNORE INTO system_prompts (hash, prompt) VALUES (?, ?)", (prompt_hash, system_prompt), + "INSERT OR IGNORE INTO system_prompts (hash, prompt) VALUES (?, ?)", + (prompt_hash, system_prompt), ) return prompt_hash @@ -460,14 +458,12 @@ class SessionDB( _ensure_test_isolation(self.db_path) # before any connection/pragma/mkdir 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 connections pinned fds for the process - # lifetime and hit EMFILE while staying alive (supervisor never restarted). + # Read-path split (WAL only): reads borrow from a BOUNDED read-only pool so they + # never queue behind writer flushes on self._lock (see _read_ctx); unbounded + # per-thread connections pinned fds for the process lifetime and hit EMFILE. self._read_pool: "queue.LifoQueue[sqlite3.Connection]" = queue.LifoQueue(maxsize=_READ_POOL_MAX) - # 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. + # Permits bound PEAK descriptors (the pool bounds only the idle set), shared per + # DATABASE PATH; acquired non-blocking so a permitless reader degrades to the writer lock. self._read_budget = _read_budget_for(self.db_path) self._read_budget.register(self) self._read_permits = self._read_budget.permits @@ -481,9 +477,8 @@ class SessionDB( 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 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. + # File identity of the opened state.db, compared on every write so an out-of-band + # replace cannot limp through in-place surgery (inode: mv/new-file; application_id: cp). self._db_file_identity: Optional[tuple] = None self._db_file_application_id: int = 0 self._db_sidecar_identity: Dict[str, tuple] = {} @@ -573,11 +568,7 @@ class SessionDB( 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, - check_same_thread=False, timeout=1.0, isolation_level=None, - ) - conn.row_factory = sqlite3.Row + self._conn = conn = self._connect_read_only(timeout=1.0) try: apply_database_pragmas(conn, db_label="state.db") cursor = conn.cursor() @@ -594,10 +585,22 @@ class SessionDB( except sqlite3.OperationalError as ioerr: # 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(): + 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.""" + 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, + ) + conn.row_factory = sqlite3.Row + 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.""" @@ -663,10 +666,8 @@ class SessionDB( now = time.monotonic() if now >= deadline: raise - time.sleep(min( - random.uniform(self._WRITE_RETRY_SLOW_MIN_S, self._WRITE_RETRY_SLOW_MAX_S), - max(deadline - now, 0.001), - )) + jitter = random.uniform(self._WRITE_RETRY_SLOW_MIN_S, self._WRITE_RETRY_SLOW_MAX_S) + time.sleep(min(jitter, max(deadline - now, 0.001))) # ── Read-path split ── @@ -679,12 +680,9 @@ class SessionDB( if not self._wal_active or self.read_only: return None with self._read_conns_lock: - if self._read_conns_closed: - return None - if ( - self._read_open_failed_at - and time.monotonic() - self._read_open_failed_at < _READ_OPEN_RETRY_SECONDS - ): + failed_at = self._read_open_failed_at + backing_off = failed_at and time.monotonic() - failed_at < _READ_OPEN_RETRY_SECONDS + if self._read_conns_closed or backing_off: return None # Permit BEFORE the open: openers race for permits, not descriptors. if not self._read_budget.acquire(self): @@ -695,13 +693,7 @@ class SessionDB( 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, - check_same_thread=False, timeout=5.0, isolation_level=None, - ) - conn.row_factory = sqlite3.Row + conn = self._connect_read_only(timeout=5.0) apply_database_pragmas(conn, db_label="state.db") if self._fts_cjk_loaded: # registers in the connection, not the file: ro is fine load_fts5_cjk_extension(conn) @@ -860,9 +852,7 @@ class SessionDB( # 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 - ) + compression_deadline = min(time.monotonic() + self._COMPRESSION_BUSY_WAIT_S, deadline) if self._sleep_before_write_retry( compression_deadline, self._COMPRESSION_BUSY_WAIT_S ): @@ -900,8 +890,7 @@ class SessionDB( # 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) + "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() @@ -1077,9 +1066,7 @@ class SessionDB( def _corrupt_error(self, prefix: str = "") -> "StateDbCorruptError": """Build the quarantine error for this handle (message assembled once).""" - return StateDbCorruptError( - f"{prefix}{_STATE_DB_CORRUPT_MSG} (cause: {self._db_corrupt_reason})" - ) + return StateDbCorruptError(f"{prefix}{_STATE_DB_CORRUPT_MSG} (cause: {self._db_corrupt_reason})") def _halt_db_corrupt(self, exc: BaseException) -> None: """Quarantine this handle and raise; never run in-file repair here.""" @@ -1090,9 +1077,7 @@ class SessionDB( "state.db %s reported structural corruption outside the FTS " "indexes (%s); quarantining this handle: no further writes, no " "automatic reopen, no explicit WAL checkpoint at close. Stop the " - "gateway and run `hermes sessions recover --source %s --inspect-only`.", - self.db_path, - exc, + "gateway and run `hermes sessions recover --source %s --inspect-only`.", self.db_path, exc, self.db_path, ) err = self._corrupt_error() @@ -1132,10 +1117,10 @@ class SessionDB( if now >= deadline: return False slow = now - (deadline - patience_s) >= self._WRITE_RETRY_SLOW_AFTER_S - jitter = ( - random.uniform(self._WRITE_RETRY_SLOW_MIN_S, self._WRITE_RETRY_SLOW_MAX_S) if slow - else random.uniform(self._WRITE_RETRY_MIN_S, self._WRITE_RETRY_MAX_S) - ) + jitter = random.uniform(*( + (self._WRITE_RETRY_SLOW_MIN_S, self._WRITE_RETRY_SLOW_MAX_S) if slow + else (self._WRITE_RETRY_MIN_S, self._WRITE_RETRY_MAX_S) + )) time.sleep(min(jitter, max(deadline - now, 0.001))) return True @@ -1158,10 +1143,7 @@ class SessionDB( if sys.platform.startswith("linux"): try: own_pid = os.getpid() - for pid_str in os.listdir("/proc"): - if not pid_str.isdigit(): - continue - pid = int(pid_str) + for pid in (int(p) for p in os.listdir("/proc") if p.isdigit()): if pid == own_pid: continue try: @@ -1193,10 +1175,11 @@ class SessionDB( return holders @staticmethod - def _foreign_holder_scan_failed(holders: List[Tuple[int, str]], exc: Exception) -> List[Tuple[int, str]]: + 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, + "Could not prove state.db has no foreign holders; deferring automatic FTS maintenance: %s", exc, ) return holders or [(-1, f"open-file scan failed: {exc}")] @@ -1248,10 +1231,8 @@ class SessionDB( "Skipping the close-time WAL checkpoint for %s: this " "handle observed structural corruption (%s). Take a " "snapshot of state.db, -wal and -shm before restarting, " - "then run `hermes sessions recover --source %s --inspect-only`.", - self.db_path, - self._db_corrupt_reason, - self.db_path, + "then run `hermes sessions recover --source %s --inspect-only`.", self.db_path, + self._db_corrupt_reason, self.db_path, ) elif not self.read_only: # PASSIVE, not TRUNCATE (see docstring) try: @@ -1274,10 +1255,9 @@ class SessionDB( pass # ── Async token accounting (SessionUsageMixin) ── - # 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(). + # queue_token_counts() is a deque append; a single-writer thread applies deltas in + # order, coalescing consecutive deltas whose route fields are EQUAL (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", @@ -1290,9 +1270,8 @@ class SessionDB( MAX_TITLE_LENGTH = 100 - # Title provenance, lowest to highest authority: auto-titling may only - # replace a strictly lower-authority title, so ``derived`` upgrades to - # ``llm`` exactly once and nothing generated clobbers a user-typed name. + # Title provenance, lowest to highest authority: auto-titling may only replace a + # strictly lower-authority title (``derived`` -> ``llm`` once; never a user-typed name). TITLE_SOURCE_DERIVED = "derived" TITLE_SOURCE_LLM = "llm" TITLE_SOURCE_USER = "user" diff --git a/hermes_state_sessions.py b/hermes_state_sessions.py index 2efc8e2cff..0b19db65a8 100644 --- a/hermes_state_sessions.py +++ b/hermes_state_sessions.py @@ -11,8 +11,7 @@ from pathlib import Path from typing import Any, Callable, Dict, List, Optional, Tuple from agent.session_activity import ( - ActivityProvenance, bound_activity_description, build_activity_snapshot, - normalize_activity_provenance, + ActivityProvenance, bound_activity_description, build_activity_snapshot, normalize_activity_provenance, ) from hermes_state_common import ( _LISTABLE_CHILD_SQL, _PREVIEW_ELIGIBLE_SQL, _PREVIEW_RAW_SELECT, _RECOVERABLE_END_REASONS, @@ -102,25 +101,18 @@ def _session_filter_where( where: List[str] = [] params: List[Any] = [] if exclude_children: - where.append(_LISTABLE_CHILD_SQL) - where.append(f"{_delegate_from_json('s.model_config')} IS NULL") + where += [_LISTABLE_CHILD_SQL, f"{_delegate_from_json('s.model_config')} IS NULL"] include_sources = [source] if source else list(sources or []) - if include_sources: - where.append(f"s.source IN ({_session_ids_placeholders(include_sources)})") - params.extend(include_sources) - if session_key: - where.append("s.session_key = ?") - params.append(session_key) - if exclude_sources: - where.append(f"s.source NOT IN ({_session_ids_placeholders(exclude_sources)})") - params.extend(exclude_sources) - if cwd_prefix: - clause, clause_params = _cwd_prefix_clause(cwd_prefix) - where.append(clause) - params.extend(clause_params) - if min_message_count > 0: - where.append("s.message_count >= ?") - params.append(min_message_count) + for clause, values in ( + (f"s.source IN ({_session_ids_placeholders(include_sources)})", include_sources), + ("s.session_key = ?", [session_key] if session_key else []), + (f"s.source NOT IN ({_session_ids_placeholders(exclude_sources or ())})", exclude_sources or []), + (_cwd_prefix_clause(cwd_prefix) if cwd_prefix else ("", [])), + ("s.message_count >= ?", [min_message_count] if min_message_count > 0 else []), + ): + if values: + where.append(clause) + params.extend(values) if archived_only: where.append("s.archived = 1") elif not include_archived: @@ -141,8 +133,7 @@ def _collect_delegate_child_ids(conn, parent_ids: List[str]) -> List[str]: ph = _session_ids_placeholders(frontier) cursor = conn.execute( f"SELECT id FROM sessions WHERE {df} IN ({ph}) " - f"OR (parent_session_id IN ({ph}) AND {df} IS NOT NULL)", - frontier + frontier, + f"OR (parent_session_id IN ({ph}) AND {df} IS NOT NULL)", frontier + frontier, ) frontier = [row["id"] for row in cursor.fetchall() if row["id"] not in found] found.update(frontier) @@ -171,9 +162,7 @@ SESSION_STATUS_EMPTY = "empty" _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: +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).""" @@ -229,8 +218,6 @@ _INHERIT_PARENT_ROUTING_SQL = ( class SessionSessionsMixin: """Session rows: create/inherit, lifecycle flags, model_config, listing, deletion.""" - _PROFILE_DIR_RE = re.compile(r"^[a-z0-9][a-z0-9_-]{0,63}$") - def _own_profile_name(self) -> Optional[str]: """The profile owning THIS store, from ``db_path`` alone (``/state.db`` → default, ``/profiles//state.db`` → name); path-based because a @@ -242,7 +229,7 @@ class SessionSessionsMixin: parent = Path(self.db_path).resolve().parent if parent == root: return "default" - if parent.parent == root / "profiles" and self._PROFILE_DIR_RE.match(parent.name): + if parent.parent == root / "profiles" 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) @@ -339,17 +326,21 @@ class SessionSessionsMixin: self._execute_write(_do, patience_s=self._TRANSCRIPT_WRITE_PATIENCE_S) def create_session(self, session_id: str, source: str, **kwargs) -> str: - """Create a new session record. Returns the session_id.""" + """Create (upsert) a session record. Returns the session_id.""" self._insert_session_row(session_id, source, **kwargs) return session_id + def ensure_session(self, session_id: str, source: str = "unknown", model: str = None, **kwargs) -> str: + """Ensure a session row exists (upsert). Accepts optional kwargs.""" + self._insert_session_row(session_id, source, model=model, **kwargs) + return session_id + def set_expiry_finalized(self, session_id: str, finalized: bool = True) -> None: """Mirror ``SessionEntry.expiry_finalized`` so it survives a lost sessions.json.""" if not session_id: return self._write_sql( - "UPDATE sessions SET expiry_finalized = ? WHERE id = ?", - (1 if finalized else 0, session_id), + "UPDATE sessions SET expiry_finalized = ? WHERE id = ?", (1 if finalized else 0, session_id), ) def find_session_by_origin( @@ -372,8 +363,7 @@ class SessionSessionsMixin: if thread_id is not None: query += " AND COALESCE(thread_id, '') = ?" params.append(str(thread_id)) - query += " ORDER BY started_at DESC" - rows = [dict(r) for r in self._read_all(query, params)] + rows = [dict(r) for r in self._read_all(query + " ORDER BY started_at DESC", params)] if not rows: return None if user_id: @@ -426,8 +416,7 @@ class SessionSessionsMixin: conn.execute( "UPDATE sessions AS child SET model_config = json_set(" "COALESCE(child.model_config, '{}'), '$._reset_from', child.parent_session_id) " - "WHERE child.parent_session_id = ? " - "AND json_extract(COALESCE(child.model_config, '{}'), " + "WHERE child.parent_session_id = ? AND json_extract(COALESCE(child.model_config, '{}'), " " '$._reset_from') IS NULL " f"AND {_legacy_reset_child_sql('child', _session_ids_placeholders(_RESET_END_REASONS))}", (session_id, *_RESET_END_REASONS), @@ -449,8 +438,7 @@ class SessionSessionsMixin: # generation advances here too — same transaction, only when written. try: return bool(self._execute_write(lambda conn: self._end_and_bump( - conn, - "UPDATE sessions SET ended_at = ?, end_reason = ? WHERE id = ? AND (ended_at IS NULL " + conn, "UPDATE sessions SET ended_at = ?, end_reason = ? WHERE id = ? AND (ended_at IS NULL " f"OR end_reason IN ({_RECOVERABLE_END_REASONS_SQL}))", (now, reason, session_id), session_id, reason, ))) @@ -597,9 +585,7 @@ class SessionSessionsMixin: 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: + 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 @@ -613,8 +599,7 @@ class SessionSessionsMixin: if provider: patch["provider"] = provider self._write_model_config_patch( - session_id, patch, - "UPDATE sessions SET model = ?, model_config = ?, " + session_id, patch, "UPDATE sessions SET model = ?, model_config = ?, " "system_prompt = NULL, system_prompt_hash = NULL WHERE id = ?", lambda merged: (model, merged, session_id), ) @@ -703,13 +688,6 @@ class SessionSessionsMixin: enable the bypass by accident).""" return bool(_parse_model_config((session_meta or {}).get("model_config")).get("yolo_mode")) - def ensure_session( - self, session_id: str, source: str = "unknown", model: str = None, **kwargs, - ) -> str: - """Ensure a session row exists (upsert). Accepts optional kwargs.""" - self._insert_session_row(session_id, source, model=model, **kwargs) - return session_id - def get_session(self, session_id: str) -> Optional[Dict[str, Any]]: """Get a session by ID (drains queued token deltas first so cost readers see exact totals).""" self.flush_token_counts() @@ -954,10 +932,9 @@ class SessionSessionsMixin: self, source: str = None, sources: List[str] = None, exclude_sources: List[str] = None, cwd_prefix: str = None, limit: int = 20, offset: int = 0, include_children: bool = False, min_message_count: int = 0, project_compression_tips: bool = True, - order_by_last_active: bool = False, include_archived: bool = False, - archived_only: bool = False, id_query: str = None, search_query: str = None, - compact_rows: bool = False, include_pinned: bool = False, session_key: str = None, - include_hidden: bool = False, + order_by_last_active: bool = False, include_archived: bool = False, archived_only: bool = False, + 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 @@ -965,19 +942,17 @@ class SessionSessionsMixin: 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, exclude_sources=exclude_sources, cwd_prefix=cwd_prefix, - min_message_count=min_message_count, archived_only=archived_only, - include_archived=include_archived, + exclude_children=not include_children, source=source, sources=sources, session_key=session_key, + exclude_sources=exclude_sources, cwd_prefix=cwd_prefix, min_message_count=min_message_count, + archived_only=archived_only, include_archived=include_archived, ) if not include_hidden: where_clauses.append("s.hidden = 0") where_sql = _where_sql(where_clauses) base_where_params = list(params) # pinned back-fill reuses the WHERE before LIMIT/OFFSET - prompt_select, prompt_join = ("", "") if compact_rows else ( - ", COALESCE(sp.prompt, s.system_prompt) AS _system_prompt_resolved", - "LEFT JOIN system_prompts sp ON sp.hash = s.system_prompt_hash", - ) + prompt_select, prompt_join = ( + "", "", + ) if compact_rows else(", COALESCE(sp.prompt, s.system_prompt) AS _system_prompt_resolved", "LEFT JOIN system_prompts sp ON sp.hash = s.system_prompt_hash",) _sel = self._compact_session_cols() if compact_rows else "s.*" if order_by_last_active: # The CTE walks compression-continuation edges forward from the admitted @@ -1188,8 +1163,7 @@ class SessionSessionsMixin: rows = conn.execute( "SELECT COALESCE(NULLIF(s.source, ''), 'cli') AS source, COUNT(*) AS count " f"FROM sessions s{_where_sql(where_clauses, ' ')} " - "GROUP BY COALESCE(NULLIF(s.source, ''), 'cli') ORDER BY count DESC", - params, + "GROUP BY COALESCE(NULLIF(s.source, ''), 'cli') ORDER BY count DESC", params, ).fetchall() return {str(row["source"]): int(row["count"] or 0) for row in rows} @@ -1374,15 +1348,14 @@ class SessionSessionsMixin: Never raises: {"skipped", "archived", "error"?}.""" result: Dict[str, Any] = {"skipped": False, "archived": 0} try: - last_raw = self.get_meta("last_auto_archive") now = time.time() - if last_raw: - try: - if now - float(last_raw) < min_interval_hours * 3600: - result["skipped"] = True - return result - except (TypeError, ValueError): - pass # corrupt meta; treat as no prior run + try: + last = float(self.get_meta("last_auto_archive") or 0.0) + except (TypeError, ValueError): + last = 0.0 # corrupt meta; treat as no prior run + if last and now - last < min_interval_hours * 3600: + result["skipped"] = True + return result archived = result["archived"] = self.archive_stale_sessions(idle_days, exclude_pinned=exclude_pinned) # Record even a zero-archive run so we don't re-sweep every call. self.set_meta("last_auto_archive", str(now))