From 7d48a84acf55bb23f5c02f56812152564b03197d Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 19:39:50 -0700 Subject: [PATCH] refactor(state): reflow prose docstrings to the 108-col budget (word-preserving) --- hermes_state_common.py | 26 +++---- hermes_state_compression.py | 64 ++++++++--------- hermes_state_gateway.py | 66 ++++++++--------- hermes_state_maintenance.py | 19 +++-- hermes_state_messages.py | 138 ++++++++++++++++-------------------- 5 files changed, 141 insertions(+), 172 deletions(-) diff --git a/hermes_state_common.py b/hermes_state_common.py index 4ea56c4f59..0f31877d8f 100644 --- a/hermes_state_common.py +++ b/hermes_state_common.py @@ -196,10 +196,9 @@ def is_automatic_end_reason(reason) -> bool: def _legacy_reset_child_sql(alias: str, reasons_sql: str) -> str: - """Pre-marker reset-continuation heuristic: child rides its parent's exact non-empty - routing key and the parent ended at a reset boundary. Shared by ``_RESET_CHILD_SQL`` - and ``reopen_session()``'s marker-stamping UPDATE so the two cannot drift; - ``reasons_sql`` is a literal or placeholder list.""" + """Pre-marker reset-continuation heuristic: child rides its parent's exact non-empty routing key and the + parent ended at a reset boundary. Shared by ``_RESET_CHILD_SQL`` and ``reopen_session()``'s + marker-stamping UPDATE so the two cannot drift; ``reasons_sql`` is a literal or placeholder list.""" return ( f"EXISTS (SELECT 1 FROM sessions p" f" WHERE p.id = {alias}.parent_session_id" @@ -799,9 +798,8 @@ def is_advisory_lock_contention(exc: BaseException) -> bool: def _proc_start_ticks(pid: int): - """Kernel start time of *pid* (field 22 of ``/proc//stat``), which with the PID - uniquely identifies a process; None off Linux or on any failure — callers must - treat None as unknowable and FAIL CLOSED.""" + """Kernel start time of *pid* (field 22 of ``/proc//stat``), which with the PID uniquely identifies + a process; None off Linux or on any failure — callers must treat None as unknowable and FAIL CLOSED.""" try: with open(f"/proc/{pid}/stat", "rb") as fh: stat = fh.read() @@ -964,8 +962,7 @@ def _describe_lock_holder(record) -> str: def _acquire_msvcrt_lock(lock_path, handle, timeout): - """Windows counterpart of ``_acquire_db_flock`` (no orphan break); same - True / False / None contract.""" + """Windows counterpart of ``_acquire_db_flock`` (no orphan break); same True / False / None contract.""" import msvcrt deadline = time.monotonic() + timeout @@ -991,12 +988,11 @@ def _acquire_msvcrt_lock(lock_path, handle, timeout): def fts_rebuild_admission(db_path, *, timeout_seconds=None): """Serialize full structural FTS rebuilds on *db_path* across processes. - Yields True when this process holds the authority, False when the bounded acquire - timed out or the lock file could not be opened. On False the caller must NOT - rebuild (fail closed); the stale breadcrumb guarantees a retry. ``db_path`` None - (in-memory DB) yields True. Opportunistic in-process retries pass - ``timeout_seconds=0`` so a live holder never stalls a long-lived writer; the orphan - break still applies.""" + Yields True when this process holds the authority, False when the bounded acquire timed out or the lock + file could not be opened. On False the caller must NOT rebuild (fail closed); the stale breadcrumb + guarantees a retry. ``db_path`` None (in-memory DB) yields True. Opportunistic in-process retries pass + ``timeout_seconds=0`` so a live holder never stalls a long-lived writer; the orphan break still applies. + """ if db_path is None: yield True return diff --git a/hermes_state_compression.py b/hermes_state_compression.py index 6f06ee35a0..7578fc73d8 100644 --- a/hermes_state_compression.py +++ b/hermes_state_compression.py @@ -132,10 +132,9 @@ class SessionCompressionMixin: def _publish_child_session_row(self, conn, parent, *, parent_session_id, child_session_id, source, model, model_config, system_prompt, cwd, profile_name) -> None: - """INSERT the compression child's ``sessions`` row copied from *parent*. Same contract - as _insert_session_row's compression-fork backfill: the child stays on the parent's - profile and keeps gateway routing/origin columns; no owner on either side -> this - store's profile.""" + """INSERT the compression child's ``sessions`` row copied from *parent*. Same contract as + _insert_session_row's compression-fork backfill: the child stays on the parent's profile and keeps + gateway routing/origin columns; no owner on either side -> this store's profile.""" system_prompt_hash = self._store_system_prompt(conn, system_prompt) conn.execute( """INSERT INTO sessions ( @@ -161,19 +160,18 @@ class SessionCompressionMixin: compression_lock_holder: str = None, require_compression_lease: bool = True, require_lease_refresh: bool = False, lease_ttl_seconds: float = 300.0, watermark: Optional[int] = None, watermark_ceiling: Optional[int] = None) -> None: - """Atomically close a parent and publish its durable compression child: closure, - child row, and handoff commit in one transaction, so readers see the live parent - or a complete child, never an ended parent with a missing/empty child. + """Atomically close a parent and publish its durable compression child: closure, child row, and + handoff commit in one transaction, so readers see the live parent or a complete child, never an + ended parent with a missing/empty child. - *watermark* (parent's ``get_active_message_watermark`` at compression start): - parent rows with ``id > watermark`` — appends landed during the slow summary — are - column-cloned into the child AFTER the handoff. *watermark_ceiling* bounds the - clone: the rotation path flushes its OWN transcript to the parent just before - publishing and those rows are already in the handoff, so only - ``(watermark, watermark_ceiling]`` is foreign tail (``None`` = unbounded). - *require_lease_refresh* + *compression_lock_holder* refreshes the lease on the same - ``conn`` before the expiry check (no TOCTOU window), so a refresher that died on - transient DB errors gets one last chance.""" + *watermark* (parent's ``get_active_message_watermark`` at compression start): parent rows with ``id + > watermark`` — appends landed during the slow summary — are column-cloned into the child AFTER the + handoff. *watermark_ceiling* bounds the clone: the rotation path flushes its OWN transcript to the + parent just before publishing and those rows are already in the handoff, so only ``(watermark, + watermark_ceiling]`` is foreign tail (``None`` = unbounded). *require_lease_refresh* + + *compression_lock_holder* refreshes the lease on the same ``conn`` before the expiry check (no + TOCTOU window), so a refresher that died on transient DB errors gets one last chance. + """ from hermes_state import CompressionSessionBusyError def _do(conn): @@ -248,9 +246,8 @@ class SessionCompressionMixin: def record_compression_failure_cooldown( self, session_id: str, cooldown_until: float, error: Optional[str] = None) -> None: - """Persist the active compression-failure cooldown. Merge-max with any longer live - deadline so a later shorter write can't reopen the thrash window; error always - takes the latest diagnostic.""" + """Persist the active compression-failure cooldown. Merge-max with any longer live deadline so a + later shorter write can't reopen the thrash window; error always takes the latest diagnostic.""" if not session_id: return self._write_sql_logged( @@ -384,10 +381,9 @@ class SessionCompressionMixin: return False def try_acquire_compression_lock(self, session_id: str, holder: str, ttl_seconds: float = 300.0) -> bool: - """Try to atomically acquire the compression lock for ``session_id``. ``False``: - another holder owns a live lock and the caller MUST NOT compress (its rotation - would split the lineage). Expired locks and structured holders whose local ``pid=`` - is dead are reclaimed transparently.""" + """Try to atomically acquire the compression lock for ``session_id``. ``False``: another holder owns + a live lock and the caller MUST NOT compress (its rotation would split the lineage). Expired + locks and structured holders whose local ``pid=`` is dead are reclaimed transparently.""" from hermes_state import _compression_lock_holder_process_is_dead if not session_id: return False @@ -480,10 +476,9 @@ class SessionCompressionMixin: wait_seconds: float = 1800.0, poll_interval_seconds: float = 1.0, on_wait=None, wait_notice_interval_seconds: float = 15.0, should_abort=None, acquire_patience_s: float = 0.5, ) -> bool: - """Wait for a cross-process turn lease without holding a SQLite lock. ``on_wait(elapsed)`` - is best-effort: called when the first attempt fails and about every - ``wait_notice_interval_seconds`` after. ``should_abort()`` True (e.g. ``/stop``) returns - False at once.""" + """Wait for a cross-process turn lease without holding a SQLite lock. ``on_wait(elapsed)`` is + best-effort: called when the first attempt fails and about every ``wait_notice_interval_seconds`` + after. ``should_abort()`` True (e.g. ``/stop``) returns False at once.""" from hermes_state import classify_persistence_error deadline = time.monotonic() + max(0.0, float(wait_seconds)) wait_started = None @@ -588,14 +583,13 @@ class SessionCompressionMixin: ) or 0 def get_compression_chain(self, session_id: str) -> List[str]: - """Walk the compression-continuation chain forward: root-first through the tip - (``[session_id]`` when no continuation); ``get_compression_tip`` is the last element. - A continuation is a child of a session with ``end_reason='compression'``. The old - ``child.started_at >= parent.ended_at`` test was too brittle (gateway + compression - races insert the real continuation before ``ended_at`` is written, while a stale - websocket later creates a sibling that passes it). Instead exclude - branch/delegate/tool children and prefer children that continue the chain or are - still live over stale closed siblings such as ``ws_orphan_reap``.""" + """Walk the compression-continuation chain forward: root-first through the tip (``[session_id]`` + when no continuation); ``get_compression_tip`` is the last element. A continuation is a child of + a session with ``end_reason='compression'``. The old ``child.started_at >= parent.ended_at`` test + was too brittle (gateway + compression races insert the real continuation before ``ended_at`` is + written, while a stale websocket later creates a sibling that passes it). Instead exclude + branch/delegate/tool children and prefer children that continue the chain or are still live over + stale closed siblings such as ``ws_orphan_reap``.""" current = session_id chain = [current] if current else [] seen = set(chain) diff --git a/hermes_state_gateway.py b/hermes_state_gateway.py index 883e4d6130..1a50fd7aa4 100644 --- a/hermes_state_gateway.py +++ b/hermes_state_gateway.py @@ -200,15 +200,14 @@ class SessionGatewayMixin: origin_json: str = None, include_compression_ancestors: bool = False) -> None: """Persist the gateway routing peer for an existing session row. - ``display_name`` / ``origin_json``: ``None`` leaves the stored value untouched - (consumers read routing data from state.db, not sessions.json). - ``include_compression_ancestors`` keeps a compression lineage on one routing - peer when an explicit resume moves its tip to another lane; per-turn refreshes - update only the supplied row. Self-healing: a missing target row (deferred - ``create_session`` write, or crash between routing publication and row creation) - is INSERTed with full identity rather than no-opped, so a gateway row is never - first-created by the identity-less lazy writer (``update_token_counts``) and left - unroutable forever.""" + ``display_name`` / ``origin_json``: ``None`` leaves the stored value untouched (consumers read + routing data from state.db, not sessions.json). ``include_compression_ancestors`` keeps a + compression lineage on one routing peer when an explicit resume moves its tip to another lane; + per-turn refreshes update only the supplied row. Self-healing: a missing target row (deferred + ``create_session`` write, or crash between routing publication and row creation) is INSERTed with + full identity rather than no-opped, so a gateway row is never first-created by the identity-less + lazy writer (``update_token_counts``) and left unroutable forever. + """ if not session_id or not session_key: return identity = (session_key, source, user_id, chat_id, chat_type, thread_id, display_name, origin_json) @@ -287,11 +286,10 @@ class SessionGatewayMixin: return {r["session_key"]: r["entry_json"] for r in rows} def list_never_active_keyed_sessions(self, *, older_than_days: float) -> List[Dict[str, Any]]: - """Keyed, still-open rows with no evidence of a single turn (no messages, tokens, - tool/API calls, activity, or title): leaked fixtures or chats routed but never - answered. Safe to drop — the gateway mints a fresh session on the next message. - Needs its own selector because ``bulk prune``/``archive`` are pinned to - ``ended_at IS NOT NULL``. ``pinned``/``archived`` = explicit keep intent.""" + """Keyed, still-open rows with no evidence of a single turn (no messages, tokens, tool/API calls, + activity, or title): leaked fixtures or chats routed but never answered. Safe to drop — the gateway + mints a fresh session on the next message. Needs its own selector because ``bulk prune``/``archive`` + are pinned to ``ended_at IS NOT NULL``. ``pinned``/``archived`` = explicit keep intent.""" cutoff = time.time() - (float(older_than_days) * 86400.0) rows = self._read_all( """ @@ -383,19 +381,18 @@ class SessionGatewayMixin: ) -> Optional[Dict[str, Any]]: """Find the latest recoverable gateway session for a routing peer. - The durable ``session_key`` on the row rebuilds a missing/pruned ``sessions.json`` - mapping. Rows ended only by the old ``agent_close`` bug or a mistaken TUI - ``ws_orphan_reap`` are recoverable; explicit boundaries (/new, /resume switches, - compression splits) are not. Ranked by ``last_activity_at`` (fallback - ``started_at`` — alone it resurrected days-old zombies); rows with messages win, - but an empty keyed row still beats ``None`` (which mints a new id while the - transcript may live under a compression child). Reset fence: a candidate is - rejected when a peer boundary row ended *after* its last activity, or the - has-messages ranking could reach behind a /new. The exact-key fallback requires - the complete peer tuple (never cross chats/threads/users) plus a profile fence: a - Telegram DM's tuple is identical for every bot, so a sibling profile's legacy row - would otherwise be adopted. Ours = profile_name is the owner or NULL; stores - outside the profile tree derive no owner and stay unfenced.""" + The durable ``session_key`` on the row rebuilds a missing/pruned ``sessions.json`` mapping. Rows + ended only by the old ``agent_close`` bug or a mistaken TUI ``ws_orphan_reap`` are recoverable; + explicit boundaries (/new, /resume switches, compression splits) are not. Ranked by + ``last_activity_at`` (fallback ``started_at`` — alone it resurrected days-old zombies); rows with + messages win, but an empty keyed row still beats ``None`` (which mints a new id while the transcript + may live under a compression child). Reset fence: a candidate is rejected when a peer boundary row + ended *after* its last activity, or the has-messages ranking could reach behind a /new. The + exact-key fallback requires the complete peer tuple (never cross chats/threads/users) plus a profile + fence: a Telegram DM's tuple is identical for every bot, so a sibling profile's legacy row would + otherwise be adopted. Ours = profile_name is the owner or NULL; stores outside the profile tree + derive no owner and stay unfenced. + """ if not session_key: return None with self._read_ctx() as conn: @@ -411,14 +408,13 @@ class SessionGatewayMixin: return self._session_row_dict(row) if row else None def find_orphaned_gateway_sessions(self, *, max_gap_s: Optional[float] = None) -> List[Dict[str, Any]]: - """Report message-bearing rows that lost their routing identity (messages, no - ``session_key``). Adoptable only when exactly one keyed predecessor can be named: - ``lineage`` (``parent_session_id`` is a keyed row of the same source; no time - window) or ``contiguity`` (exactly one keyed same-source row with compatible - ``user_id`` fell quiet within *max_gap_s* of the orphan's start and is older than - its last activity). Ambiguity is reported ``adoptable=False`` with a reason, never - guessed — mis-adopting splices one person's conversation into another's chat. - Branch/delegate/tool rows are excluded: unkeyed by design, not damage.""" + """Report message-bearing rows that lost their routing identity (messages, no ``session_key``). + Adoptable only when exactly one keyed predecessor can be named: ``lineage`` (``parent_session_id`` + is a keyed row of the same source; no time window) or ``contiguity`` (exactly one keyed same-source + row with compatible ``user_id`` fell quiet within *max_gap_s* of the orphan's start and is older + than its last activity). Ambiguity is reported ``adoptable=False`` with a reason, never guessed — + mis-adopting splices one person's conversation into another's chat. Branch/delegate/tool rows are + excluded: unkeyed by design, not damage.""" gap = self._ORPHAN_ADOPTION_MAX_GAP_S if max_gap_s is None else float(max_gap_s) records: List[Dict[str, Any]] = [] with self._read_ctx() as conn: diff --git a/hermes_state_maintenance.py b/hermes_state_maintenance.py index e53510f293..22499b2c5c 100644 --- a/hermes_state_maintenance.py +++ b/hermes_state_maintenance.py @@ -354,17 +354,16 @@ class SessionMaintenanceMixin: sessions_dir: Optional[Path] = None, min_vacuum_interval_days: int = 30, min_vacuum_freelist_ratio: float = AUTO_VACUUM_MIN_FREELIST_RATIO, ) -> Dict[str, Any]: - """Idempotent startup auto-maintenance (never raises): prune inactive sessions, reap stale - open state-owned rows, optional VACUUM. Runs at most once per ``min_interval_hours``; - VACUUM has its own ``min_vacuum_interval_days`` throttle and also requires - ``freelist_count / page_count`` > ``min_vacuum_freelist_ratio`` so a small prune on a - dense multi-GB database never triggers a full rewrite. Stale-open reconciliation: - cron/kanban/subagent/one-shot CLI rows never set ``ended_at`` when their process dies and - prune only deletes ended rows, so after pruning, open rows from + """Idempotent startup auto-maintenance (never raises): prune inactive sessions, reap stale open + state-owned rows, optional VACUUM. Runs at most once per ``min_interval_hours``; VACUUM has its own + ``min_vacuum_interval_days`` throttle and also requires ``freelist_count / page_count`` > + ``min_vacuum_freelist_ratio`` so a small prune on a dense multi-GB database never triggers a full + rewrite. Stale-open reconciliation: cron/kanban/subagent/one-shot CLI rows never set ``ended_at`` + when their process dies and prune only deletes ended rows, so after pruning, open rows from :attr:`_AUTO_PRUNE_STALE_OPEN_SOURCES` older than ``retention_days`` are closed - (``startup_orphan_reap``); they stay resumable and age from their close. Returns - ``{"skipped", "pruned", "closed", "vacuumed"}`` plus ``"freelist_ratio"`` when a VACUUM - was considered and ``"error"`` on failure.""" + (``startup_orphan_reap``); they stay resumable and age from their close. Returns ``{"skipped", + "pruned", "closed", "vacuumed"}`` plus ``"freelist_ratio"`` when a VACUUM was considered and + ``"error"`` on failure.""" from hermes_state import _release_auto_maintenance_lock, _try_acquire_auto_maintenance_lock result: Dict[str, Any] = {"skipped": False, "pruned": 0, "closed": 0, "vacuumed": False} maintenance_lock = _try_acquire_auto_maintenance_lock(self.db_path) diff --git a/hermes_state_messages.py b/hermes_state_messages.py index c240742916..65a640d683 100644 --- a/hermes_state_messages.py +++ b/hermes_state_messages.py @@ -111,10 +111,9 @@ class SessionMessagesMixin: """Message append/replace/rewind, reactions, resume conversations, replay dedupe.""" def _bump_conversation_generation(self, conn, session_id: str, end_reason: str) -> None: - """Advance this peer's conversation generation past a boundary, inside the txn that - writes it. Only ``_RESET_END_REASONS`` count (``compression`` continues one - conversation). Never derived from session rows (deletes/prunes could re-emit a - retired affinity identity); it only ever increments.""" + """Advance this peer's conversation generation past a boundary, inside the txn that writes it. Only + ``_RESET_END_REASONS`` count (``compression`` continues one conversation). Never derived from + session rows (deletes/prunes could re-emit a retired affinity identity); it only ever increments.""" if end_reason not in _RESET_END_REASONS: return row = conn.execute("SELECT source, session_key FROM sessions WHERE id = ?", (session_id,)).fetchone() @@ -127,11 +126,10 @@ class SessionMessagesMixin: @classmethod def _encode_content(cls, content: Any) -> Any: - """Serialize list/dict content (multimodal parts) as a sentinel-prefixed JSON string - (sqlite3 binds only str/bytes/int/float/None). Lone UTF-16 surrogates (web-scraped - tool results reach the canonical history unsanitized) are scrubbed here: left raw, - sqlite3 raises UnicodeEncodeError, the flush is abandoned and the session silently - stops persisting. Paired with :meth:`_decode_content`.""" + """Serialize list/dict content (multimodal parts) as a sentinel-prefixed JSON string (sqlite3 binds + only str/bytes/int/float/None). Lone UTF-16 surrogates (web-scraped tool results reach the canonical + history unsanitized) are scrubbed here: left raw, sqlite3 raises UnicodeEncodeError, the flush is + abandoned and the session silently stops persisting. Paired with :meth:`_decode_content`.""" if isinstance(content, str): return _sanitize_surrogates(content) if content is None or isinstance(content, (bytes, int, float)): @@ -174,9 +172,8 @@ class SessionMessagesMixin: @staticmethod def _decode_display_metadata(raw: Any) -> Optional[Dict[str, Any]]: - """Decode a ``display_metadata`` column into a dict (never the raw TEXT — the desktop - does ``'task_count' in meta``). Pre-guard rows are double-encoded: a second string - layer is unwrapped.""" + """Decode a ``display_metadata`` column into a dict (never the raw TEXT — the desktop does + ``'task_count' in meta``). Pre-guard rows are double-encoded: a second string layer is unwrapped.""" if raw is None: return None try: @@ -193,10 +190,9 @@ class SessionMessagesMixin: @staticmethod def _reasoning_json_text(value: Any) -> Optional[str]: - """Serialize a structured reasoning field for its TEXT column. Strings are stored - as-is: round-tripping callers (get_messages -> replace_messages) hand back the raw - TEXT; re-dumping would double-encode it and reasoning-replay consumers - (``isinstance(..., list)``) would drop it.""" + """Serialize a structured reasoning field for its TEXT column. Strings are stored as-is: + round-tripping callers (get_messages -> replace_messages) hand back the raw TEXT; re-dumping + would double-encode it and reasoning-replay consumers (``isinstance(..., list)``) would drop it.""" return None if not value else (value if isinstance(value, str) else json.dumps(value)) def _check_transcript_write_guards( @@ -205,12 +201,11 @@ class SessionMessagesMixin: reject_active_turn_lease: bool = False, reject_active_compression_lock: bool = False, allow_closed_compression_parent: bool = False, ) -> None: - """Transcript-write admission checks, run INSIDE the write txn (shared by every - writer). Ordinary appends do NOT check compression_locks: the lock only stops two - COMPRESSIONS colliding and archive_and_compact() commits against a watermark, so - concurrent appends are safe (blocking them killed turns during slow summaries). - Destructive user mutations opt in via ``reject_active_*`` so a compressor that - captured its watermark cannot resurrect the removed turn.""" + """Transcript-write admission checks, run INSIDE the write txn (shared by every writer). Ordinary + appends do NOT check compression_locks: the lock only stops two COMPRESSIONS colliding and + archive_and_compact() commits against a watermark, so concurrent appends are safe (blocking them + killed turns during slow summaries). Destructive user mutations opt in via ``reject_active_*`` so a + compressor that captured its watermark cannot resurrect the removed turn.""" from hermes_state import CompressionSessionClosedError, SessionCompressionInProgressError, SessionTurnLeaseLostError if reject_active_compression_lock: active_lock = conn.execute(_COMPRESSION_LOCK_ROW_SQL, (session_id,)).fetchone() @@ -297,10 +292,9 @@ class SessionMessagesMixin: display_metadata: Optional[Dict[str, Any]] = None, compression_lock_holder: Optional[str] = None, turn_lease_holder: Optional[str] = None, turn_lease_ttl_seconds: float = 300.0, ) -> int: - """Append one message; returns the row id and bumps the session counters. - ``platform_message_id``: the platform's own id (recall-style flows). ``api_content``: - byte-fidelity sidecar — the exact string sent to the API when it differed from - ``content`` — stored as sent except lone surrogates.""" + """Append one message; returns the row id and bumps the session counters. ``platform_message_id``: + the platform's own id (recall-style flows). ``api_content``: byte-fidelity sidecar — the exact + string sent to the API when it differed from ``content`` — stored as sent except lone surrogates.""" msg = dict(locals()) # every keyword above is a message-dict field of the same name # Encode outside the write txn (display metadata first: log-order parity). msg["display_metadata"] = self._encode_display_metadata(display_metadata) @@ -325,10 +319,9 @@ class SessionMessagesMixin: turn_lease_holder: Optional[str] = None, chunk_rows: Optional[int] = None, turn_lease_ttl_seconds: float = 300.0, ) -> int: - """Append *messages* (``_insert_message_rows`` dict shape) in ONE write txn: all rows - land or none, guards run once. ``chunk_rows`` bounds txn size for LARGE copies - (branch seeds; FTS triggers run per row) by committing in chunks. Returns the - inserted row count.""" + """Append *messages* (``_insert_message_rows`` dict shape) in ONE write txn: all rows land or none, + guards run once. ``chunk_rows`` bounds txn size for LARGE copies (branch seeds; FTS triggers run per + row) by committing in chunks. Returns the inserted row count.""" if not messages: return 0 if chunk_rows is not None and len(messages) > chunk_rows: @@ -502,11 +495,10 @@ class SessionMessagesMixin: self, session_id: str, messages: List[Dict[str, Any]], active_only: bool = False, archive_dropped: bool = False, reject_active_turn_lease: bool = False, ) -> None: - """Atomically replace a session's stored messages (/retry, /undo, /compress). - DESTRUCTIVE by default (rows DELETEd, leave FTS). ``active_only`` spares - soft-archived rows (needed alongside in-place compaction). ``archive_dropped`` - SOFT-archives the live rows rewind-style instead of deleting — what rewind/edit/ - regenerate must use, since a DELETE leaves nothing to recover. ``reject_active_ + """Atomically replace a session's stored messages (/retry, /undo, /compress). DESTRUCTIVE by default + (rows DELETEd, leave FTS). ``active_only`` spares soft-archived rows (needed alongside in-place + compaction). ``archive_dropped`` SOFT-archives the live rows rewind-style instead of deleting — what + rewind/edit/ regenerate must use, since a DELETE leaves nothing to recover. ``reject_active_ turn_lease`` runs the lease check in-txn for user rewrites that don't own it.""" from hermes_state import CompressionSessionClosedError active_clause = " AND active = 1" if active_only else "" @@ -534,9 +526,8 @@ class SessionMessagesMixin: "SELECT 1 FROM messages WHERE session_id = ? AND active = 0 LIMIT 1", (session_id,)) is not None def get_active_message_watermark(self, session_id: str) -> int: - """MAX(id) of the session's active rows — captured at compression START; every active - row above it arrived concurrently and must survive compaction verbatim. 0 for an - empty/unknown session.""" + """MAX(id) of the session's active rows — captured at compression START; every active row above it + arrived concurrently and must survive compaction verbatim. 0 for an empty/unknown session.""" if not session_id: return 0 row = self._read_one("SELECT COALESCE(MAX(id), 0) FROM messages WHERE session_id = ? AND active = 1", (session_id,)) @@ -861,13 +852,12 @@ class SessionMessagesMixin: self, rows, *, session_id: str, include_ancestors: bool, repair_alternation: bool, include_row_ids: bool = False, include_summary_markers: bool = False, ) -> List[Dict[str, Any]]: - """Decode fetched message rows (ordered by id, pre-filtered) into OpenAI format. - Every dict is stamped ``_DB_PERSISTED_MARKER_KEY`` at the source (born durable) so an - identity-losing handoff never re-appends the whole transcript on flush. ``_row_id`` - is opt-in (gateway reactions). Reasoning fields are restored on assistant rows only. - Key order of each dict is stable. ``api_content`` is returned VERBATIM (no - sanitize/strip): the replay path substitutes it to keep the provider prompt cache - byte-stable.""" + """Decode fetched message rows (ordered by id, pre-filtered) into OpenAI format. Every dict is + stamped ``_DB_PERSISTED_MARKER_KEY`` at the source (born durable) so an identity-losing handoff + never re-appends the whole transcript on flush. ``_row_id`` is opt-in (gateway reactions). + Reasoning fields are restored on assistant rows only. Key order of each dict is stable. + ``api_content`` is returned VERBATIM (no sanitize/strip): the replay path substitutes it to keep + the provider prompt cache byte-stable.""" from hermes_state import _strip_background_review_harness, _strip_stale_tool_call_markers messages = [] exact_user_clones: Dict[Tuple[Any, str], Dict[str, Any]] = {} @@ -1006,9 +996,8 @@ class SessionMessagesMixin: for message in lineage if message.get("_row_id") in ancestor_ids] def get_conversation_root(self, session_id: str) -> str: - """ROOT id of *session_id*'s lineage — the stable conversation id across compression - segments and delegate subagents (Nous Portal usage tagging). Unchanged when there - is no recorded parent.""" + """ROOT id of *session_id*'s lineage — the stable conversation id across compression segments and + delegate subagents (Nous Portal usage tagging). Unchanged when there is no recorded parent.""" chain = self._session_lineage_root_to_tip(session_id) return chain[0] if chain and chain[0] else session_id @@ -1093,16 +1082,14 @@ class SessionMessagesMixin: self, session_id: str, target_message_id: int, *, preserve_compaction_handoff: bool = False, expected_active_ids: Optional[List[int]] = None, expected_target_content: Any = None, ) -> Dict[str, Any]: - """Soft-delete (``active=0``) every message with id >= *target_message_id*, the target - included (the caller pre-fills it as the next prompt). Returns ``{"rewound_count", - "target_message", "new_head_id"}``, plus ``replacement_message_id`` with - ``preserve_compaction_handoff`` (archives a composite summary carrier, inserts its - hidden handoff scaffold as the new head). ``ValueError`` when the target is missing - or not a ``user`` row. ``expected_active_ids`` / ``expected_target_content`` pin the - active set and the canonical live payload in-txn before any mutation - (presentation-only metadata changes don't invalidate a rewind). A live turn lease - refuses the rewind; expired/dead holders are reclaimed. ``rewind_count`` always - increments.""" + """Soft-delete (``active=0``) every message with id >= *target_message_id*, the target included (the + caller pre-fills it as the next prompt). Returns ``{"rewound_count", "target_message", + "new_head_id"}``, plus ``replacement_message_id`` with ``preserve_compaction_handoff`` (archives a + composite summary carrier, inserts its hidden handoff scaffold as the new head). ``ValueError`` when + the target is missing or not a ``user`` row. ``expected_active_ids`` / ``expected_target_content`` + pin the active set and the canonical live payload in-txn before any mutation (presentation-only + metadata changes don't invalidate a rewind). A live turn lease refuses the rewind; expired/dead + holders are reclaimed. ``rewind_count`` always increments.""" def _do(conn): self._check_transcript_write_guards( conn, session_id, None, reject_active_turn_lease=True, reject_active_compression_lock=True) @@ -1177,24 +1164,22 @@ class SessionMessagesMixin: return parent_id in markers if parent_id else any(m is not None for m in markers) def is_explicit_fork_child(self, session_id: str) -> bool: - """Public read-only view of :meth:`_is_explicit_fork_child_row`; a missing row is not - a fork (``agent/prompt_cache_scope.py`` keeps a declared conversation key from - crossing the fork boundary).""" + """Public read-only view of :meth:`_is_explicit_fork_child_row`; a missing row is not a fork + (``agent/prompt_cache_scope.py`` keeps a declared conversation key from crossing the fork boundary).""" session = self.get_session(session_id) return bool(session and self._is_explicit_fork_child_row(session)) def latest_conversation_boundary(self, session_key: str, source: str) -> Optional[int]: - """How many conversation boundaries (``_RESET_END_REASONS`` ends) this routing peer - has crossed, or ``None`` when never reset. The peer is ``(session_key, source)`` — - the identity recovery uses — never the key alone (an API caller may legally reuse a - Telegram row's key). Read from ``conversation_generations`` (advanced inside each - boundary's txn), not an aggregate over session rows: deletes/prunes would let an - aggregate re-emit a retired pair. Rows are never garbage-collected, by design - (dropping one would re-issue generation 1 — the ABA this counter prevents). - Wall-clock-free, so a backwards NTP correction cannot reorder it. DBs upgraded - mid-conversation start at no generation and take their first from the next boundary - written (a pre-upgrade reset shares its predecessor's scope once — costs a warm - prompt-cache bucket, never crosses an identity).""" + """How many conversation boundaries (``_RESET_END_REASONS`` ends) this routing peer has crossed, or + ``None`` when never reset. The peer is ``(session_key, source)`` — the identity recovery uses — + never the key alone (an API caller may legally reuse a Telegram row's key). Read from + ``conversation_generations`` (advanced inside each boundary's txn), not an aggregate over session + rows: deletes/prunes would let an aggregate re-emit a retired pair. Rows are never + garbage-collected, by design (dropping one would re-issue generation 1 — the ABA this counter + prevents). Wall-clock-free, so a backwards NTP correction cannot reorder it. DBs upgraded + mid-conversation start at no generation and take their first from the next boundary written (a + pre-upgrade reset shares its predecessor's scope once — costs a warm prompt-cache bucket, never + crosses an identity).""" if not session_key or not source: return None row = self._read_one( @@ -1211,11 +1196,10 @@ class SessionMessagesMixin: self._execute_write(_do) def purge_stale_tool_call_markers(self, *, dry_run: bool = False, backup: bool = True) -> Dict[str, Any]: - """Permanently clear bare tool-call marker content (e.g. "[memory]") left by pre-fix - sessions (``_rows_to_conversation`` already repairs it in memory; this stops the - re-scan). Only ``content`` is touched. ``backup``: ``VACUUM INTO`` snapshot first - (none when nothing changes). Returns ``{"dry_run", "rows_affected", "row_ids", - "backup_path"}``.""" + """Permanently clear bare tool-call marker content (e.g. "[memory]") left by pre-fix sessions + (``_rows_to_conversation`` already repairs it in memory; this stops the re-scan). Only ``content`` + is touched. ``backup``: ``VACUUM INTO`` snapshot first (none when nothing changes). Returns + ``{"dry_run", "rows_affected", "row_ids", "backup_path"}``.""" from hermes_state import _STALE_TOOL_CALL_MARKER_RE def _find_affected(conn) -> List[int]: