refactor(state): hand-compact remaining long docstrings (all invariants kept)
This commit is contained in:
+11
-19
@@ -861,20 +861,15 @@ def _lock_holder_provably_dead(record) -> bool:
|
||||
|
||||
|
||||
def _acquire_db_flock(lock_path, handle, timeout_seconds, poll_seconds, description):
|
||||
"""Bounded POSIX flock acquire with orphaned-holder staleness break.
|
||||
|
||||
Returns ``(acquired, handle)``; *handle* may have been re-opened and the caller
|
||||
closes whichever comes back. *acquired* is True, False (a holder kept the lock past
|
||||
the deadline), or None (non-contention ``OSError``, already logged; callers treat it
|
||||
as not acquired without the held-by-another-process warning).
|
||||
|
||||
``flock`` belongs to the open file DESCRIPTION, which ``fork()`` duplicates: a holder
|
||||
that forks then dies leaves the lock held forever by a child that never releases.
|
||||
When the process that ACQUIRED is provably dead yet the flock is held, the file is
|
||||
unlinked and retaken on a fresh inode; the orphan's flock stays on the old inode
|
||||
blocking nobody. Every successful acquire verifies its inode still names
|
||||
*lock_path*, so a racer that locked a dead inode retries instead of running
|
||||
alongside the breaker. Indeterminate liveness defers."""
|
||||
"""Bounded POSIX flock acquire with orphaned-holder staleness break. Returns ``(acquired, handle)``;
|
||||
*handle* may have been re-opened and the caller closes whichever comes back. *acquired* is True, False
|
||||
(a holder kept the lock past the deadline), or None (non-contention ``OSError``, already logged; callers
|
||||
treat it as not acquired without the held-by-another-process warning). ``flock`` belongs to the open
|
||||
file DESCRIPTION, which ``fork()`` duplicates: a holder that forks then dies leaves the lock held forever
|
||||
by a child that never releases. When the process that ACQUIRED is provably dead yet the flock is held,
|
||||
the file is unlinked and retaken on a fresh inode; the orphan's flock stays on the old inode blocking
|
||||
nobody. Every successful acquire verifies its inode still names *lock_path*, so a racer that locked a
|
||||
dead inode retries instead of running alongside the breaker. Indeterminate liveness defers."""
|
||||
import fcntl
|
||||
|
||||
deadline = time.monotonic() + timeout_seconds
|
||||
@@ -969,13 +964,10 @@ def _acquire_msvcrt_lock(lock_path, handle, timeout):
|
||||
|
||||
@contextlib.contextmanager
|
||||
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
|
||||
"""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.
|
||||
"""
|
||||
``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
|
||||
|
||||
@@ -177,16 +177,13 @@ class SessionCompressionMixin:
|
||||
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.
|
||||
|
||||
*watermark* (parent's ``get_active_message_watermark`` at compression start): parent rows with ``id
|
||||
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.
|
||||
"""
|
||||
TOCTOU window), so a refresher that died on transient DB errors gets one last chance."""
|
||||
from hermes_state import CompressionSessionBusyError
|
||||
def _do(conn):
|
||||
if require_lease_refresh and compression_lock_holder:
|
||||
@@ -364,14 +361,12 @@ class SessionCompressionMixin:
|
||||
self._write_session_column("compression_recovery_deadline", session_id, normalized or None)
|
||||
|
||||
def refresh_compression_lock(self, session_id: str, holder: str, ttl_seconds: float = 300.0) -> bool:
|
||||
"""Extend the compression lock lease if ``holder`` still owns it.
|
||||
|
||||
Ownership is decided by ``holder`` alone, deliberately NOT ``expires_at``: a live
|
||||
owner whose refresher stalled past its TTL must be able to revive its still-
|
||||
unclaimed row, otherwise it keeps compressing with no lease — the window in which
|
||||
a competing path can fork the lineage. It cannot resurrect a lock someone else
|
||||
took: SQLite serialises writes, so the reclaim (DELETE-expired + INSERT OR IGNORE)
|
||||
never interleaves with this UPDATE."""
|
||||
"""Extend the compression lock lease if ``holder`` still owns it. Ownership is decided by ``holder``
|
||||
alone, deliberately NOT ``expires_at``: a live owner whose refresher stalled past its TTL must be
|
||||
able to revive its still-unclaimed row, otherwise it keeps compressing with no lease — the window in
|
||||
which a competing path can fork the lineage. It cannot resurrect a lock someone else took: SQLite
|
||||
serialises writes, so the reclaim (DELETE-expired + INSERT OR IGNORE) never interleaves with this
|
||||
UPDATE."""
|
||||
if not session_id or not holder:
|
||||
return False
|
||||
expires_at = time.time() + ttl_seconds
|
||||
|
||||
+23
-32
@@ -198,16 +198,14 @@ class SessionGatewayMixin:
|
||||
self, session_id: str, *, source: str, user_id: str = None, session_key: str = None,
|
||||
chat_id: str = None, chat_type: str = None, thread_id: str = None, display_name: str = None,
|
||||
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.
|
||||
"""
|
||||
"""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."""
|
||||
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)
|
||||
@@ -375,9 +373,7 @@ class SessionGatewayMixin:
|
||||
self, *, source: str, user_id: Optional[str] = None, session_key: Optional[str] = None,
|
||||
chat_id: Optional[str] = None, chat_type: Optional[str] = None, thread_id: Optional[str] = None,
|
||||
) -> 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
|
||||
"""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
|
||||
@@ -387,8 +383,7 @@ class SessionGatewayMixin:
|
||||
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.
|
||||
"""
|
||||
derive no owner and stay unfenced."""
|
||||
if not session_key:
|
||||
return None
|
||||
with self._read_ctx() as conn:
|
||||
@@ -646,14 +641,12 @@ class SessionGatewayMixin:
|
||||
|
||||
def fail_handoff(
|
||||
self, session_id: str, error: str, *, only_states: Optional[Tuple[str, ...]] = None) -> bool:
|
||||
"""Mark a handoff failed and record the reason; True when a row transitioned.
|
||||
|
||||
``only_states`` makes the write a compare-and-swap on ``handoff_state``. Waiters
|
||||
that give up (CLI 60s poll, Desktop bounded poll) MUST pass ``only_states=("pending",)``:
|
||||
once the watcher has claimed the row (``running``) it owns the terminal state, and an
|
||||
unconditional waiter-side fail races the dispatch — the gateway later overwrites
|
||||
``failed`` → ``completed`` after the user was told the gateway is down (split-brain:
|
||||
the handoff delivered and ``switch_session`` re-pointed the session). The watcher
|
||||
"""Mark a handoff failed and record the reason; True when a row transitioned. ``only_states`` makes
|
||||
the write a compare-and-swap on ``handoff_state``. Waiters that give up (CLI 60s poll, Desktop
|
||||
bounded poll) MUST pass ``only_states=("pending",)``: once the watcher has claimed the row
|
||||
(``running``) it owns the terminal state, and an unconditional waiter-side fail races the dispatch —
|
||||
the gateway later overwrites ``failed`` → ``completed`` after the user was told the gateway is down
|
||||
(split-brain: the handoff delivered and ``switch_session`` re-pointed the session). The watcher
|
||||
fails its OWN claimed row unconditionally."""
|
||||
states = tuple(only_states) if only_states else ()
|
||||
sql = _HANDOFF_FAIL_SQL + "id = ?" + (
|
||||
@@ -661,15 +654,13 @@ class SessionGatewayMixin:
|
||||
return self._write_rowcount(sql, (error[:500], session_id, *states)) > 0
|
||||
|
||||
def reclaim_stale_running_handoffs(self, error: str) -> List[str]:
|
||||
"""Fail every handoff stuck in ``running``; returns the ids reclaimed.
|
||||
|
||||
Only the gateway watcher sets ``running``, for one in-process dispatch — so a
|
||||
``running`` row at watcher startup belongs to a PREVIOUS gateway that died
|
||||
mid-dispatch. It is poisonous: ``request_handoff`` only accepts NULL/``completed``/
|
||||
``failed``, so the session could never hand off again, with no error surfaced.
|
||||
Failing rather than re-queueing is deliberate: the dead gateway may already have
|
||||
switched the session key and dispatched the synthetic turn, so a blind retry risks
|
||||
double delivery; a clean terminal state the user can retry from is right."""
|
||||
"""Fail every handoff stuck in ``running``; returns the ids reclaimed. Only the gateway watcher sets
|
||||
``running``, for one in-process dispatch — so a ``running`` row at watcher startup belongs to a
|
||||
PREVIOUS gateway that died mid-dispatch. It is poisonous: ``request_handoff`` only accepts
|
||||
NULL/``completed``/``failed``, so the session could never hand off again, with no error surfaced.
|
||||
Failing rather than re-queueing is deliberate: the dead gateway may already have switched the
|
||||
session key and dispatched the synthetic turn, so a blind retry risks double delivery; a clean
|
||||
terminal state the user can retry from is right."""
|
||||
def _do(conn):
|
||||
cur = conn.execute("SELECT id FROM sessions WHERE handoff_state = 'running'")
|
||||
ids = [r[0] for r in cur.fetchall()]
|
||||
|
||||
+14
-17
@@ -122,23 +122,20 @@ class SessionMaintenanceMixin:
|
||||
heartbeat_staleness_seconds: Optional[float] = None,
|
||||
heartbeat_ownership_grace_seconds: Optional[float] = None, respect_gateway_heartbeats: bool = True,
|
||||
) -> List[str]:
|
||||
"""Close session rows orphaned by a dead gateway process (its in-process disconnect grace
|
||||
timer died with it, leaving ``ended_at IS NULL`` forever). Rows of ``sources`` whose
|
||||
``started_at`` AND canonical last activity are both older than ``max_idle_seconds`` get
|
||||
``end_reason='startup_orphan_reap'`` (the ``started_at`` predicate protects fresh
|
||||
compression/branch children whose copied activity is old). Only pass sources whose
|
||||
lifecycle the caller owns — never messaging platforms like ``telegram`` (ending those
|
||||
triggers a routing loop). ``exclude_ids`` spares rows this process still holds.
|
||||
Non-destructive: messages kept, row resumable, first-reason-wins.
|
||||
|
||||
With ``respect_gateway_heartbeats`` a row is reaped only when no live backend (heartbeat
|
||||
within ``heartbeat_staleness_seconds``, default ``2 * max_idle_seconds``) could own it: B
|
||||
owns S if ``B.started_at <= S.started_at + grace`` (default = staleness) — grace covers a
|
||||
migrating backend whose sessions predate its first heartbeat, bounded so a PID-reuse
|
||||
respawn cannot protect rows forever. Disable the gate only for state.db-owned sources.
|
||||
SELECT, live-lease validation and UPDATE run in one ``BEGIN IMMEDIATE`` transaction;
|
||||
active leases/locks spare the row, expired guards are removed so their owner is fenced.
|
||||
"""
|
||||
"""Close session rows orphaned by a dead gateway process (its in-process disconnect grace timer died
|
||||
with it, leaving ``ended_at IS NULL`` forever). Rows of ``sources`` whose ``started_at`` AND
|
||||
canonical last activity are both older than ``max_idle_seconds`` get
|
||||
``end_reason='startup_orphan_reap'`` (the ``started_at`` predicate protects fresh compression/branch
|
||||
children whose copied activity is old). Only pass sources whose lifecycle the caller owns — never
|
||||
messaging platforms like ``telegram`` (ending those triggers a routing loop). ``exclude_ids`` spares
|
||||
rows this process still holds. Non-destructive: messages kept, row resumable, first-reason-wins.
|
||||
With ``respect_gateway_heartbeats`` a row is reaped only when no live backend (heartbeat within
|
||||
``heartbeat_staleness_seconds``, default ``2 * max_idle_seconds``) could own it: B owns S if
|
||||
``B.started_at <= S.started_at + grace`` (default = staleness) — grace covers a migrating backend
|
||||
whose sessions predate its first heartbeat, bounded so a PID-reuse respawn cannot protect rows
|
||||
forever. Disable the gate only for state.db-owned sources. SELECT, live-lease validation and UPDATE
|
||||
run in one ``BEGIN IMMEDIATE`` transaction; active leases/locks spare the row, expired guards are
|
||||
removed so their owner is fenced."""
|
||||
srcs = tuple(s for s in sources if s)
|
||||
if max_idle_seconds <= 0 or not srcs:
|
||||
return []
|
||||
|
||||
+11
-14
@@ -542,20 +542,17 @@ class SessionMessagesMixin:
|
||||
model_config_patch: Optional[Dict[str, Any]] = None, watermark: Optional[int] = None,
|
||||
lock_holder: Optional[str] = None, tail_count: int = 0,
|
||||
) -> int:
|
||||
"""Non-destructive in-place compaction under ONE durable session id: soft-archive the
|
||||
active rows (``active=0, compacted=1`` — "summarized away", still searchable) and
|
||||
insert *compacted_messages* as fresh active rows, atomically. Returns the new active
|
||||
count; ``message_count`` becomes the ACTIVE count.
|
||||
|
||||
*watermark* (captured at compression START): rows with ``id > watermark`` arrived
|
||||
during the slow summary and are re-sequenced after the compacted set by a pure-SQL
|
||||
column clone (fresh ids); ``None`` archives everything. *lock_holder*: the commit
|
||||
verifies in-txn that the lease is still held, so a reclaimed lease fails instead of
|
||||
clobbering the winner. *tail_count*: the LAST N compacted rows are the verbatim
|
||||
carried-forward tail; their originals and the watermark clones' originals are
|
||||
superseded duplicates and get rewind-style flags (``active=0, compacted=0``) so
|
||||
search doesn't return each carried message once per compaction.
|
||||
``model_config_patch`` merges in the same txn (``None`` removes a key)."""
|
||||
"""Non-destructive in-place compaction under ONE durable session id: soft-archive the active rows
|
||||
(``active=0, compacted=1`` — "summarized away", still searchable) and insert *compacted_messages* as
|
||||
fresh active rows, atomically. Returns the new active count (``message_count`` becomes the ACTIVE
|
||||
count). *watermark* (captured at compression START): rows with ``id > watermark`` arrived during the
|
||||
slow summary and are re-sequenced after the compacted set by a pure-SQL column clone (fresh ids);
|
||||
``None`` archives everything. *lock_holder*: the commit verifies in-txn that the lease is still held,
|
||||
so a reclaimed lease fails instead of clobbering the winner. *tail_count*: the LAST N compacted rows
|
||||
are the verbatim carried-forward tail; their originals and the watermark clones' originals are
|
||||
superseded duplicates and get rewind-style flags (``active=0, compacted=0``) so search doesn't return
|
||||
each carried message once per compaction. ``model_config_patch`` merges in the same txn (``None``
|
||||
removes a key)."""
|
||||
from hermes_state import SessionCompressionInProgressError
|
||||
def _do(conn):
|
||||
if lock_holder is not None:
|
||||
|
||||
Reference in New Issue
Block a user