diff --git a/hermes_state_registry.py b/hermes_state_registry.py index 900d4b404d..645c225adf 100644 --- a/hermes_state_registry.py +++ b/hermes_state_registry.py @@ -1,27 +1,21 @@ """Process-wide shared SessionDB registry. -A gateway process opens state.db from many call sites; each bare ``SessionDB()`` -mints its own writer connection, lock, close-time WAL checkpoint and token -writer thread, and one connection's close-time checkpoint can race another's -growth (lost/reordered-page-write corruption). This module owns that boundary: -one shared ``SessionDB`` per resolved path per process, refcounted, with -generation-aware retirement when the file is replaced (snapshot restore, -recovery swap). +Each bare ``SessionDB()`` mints its own writer connection, lock, close-time WAL checkpoint +and token writer thread, and one connection's close-time checkpoint can race another's +growth (lost/reordered-page-write corruption). This module owns that boundary: one shared +``SessionDB`` per resolved path per process, refcounted, with generation-aware retirement +when the file is replaced (snapshot restore, recovery swap). Lifecycle rules: - -- ``acquire(path)`` returns the current generation for *path* and bumps its - refcount. Same path ⇒ same instance ⇒ same writer connection. -- ``close()`` on a shared instance is a NO-OP: the registry owns the connection - lifecycle, so one caller can never tear down a writer others still hold. -- ``release(db)`` decrements the generation *db was acquired from* (object- - keyed, so an inode replacement cannot strand a still-owned generation). The - final release of a retired generation tears it down. -- On inode change the old generation is RETIRED (never lent again) but stays - alive until its holders release. If the replacement open fails the registry - keeps NO path entry, so the next acquire retries. -- All teardown happens OUTSIDE the registry lock: a final release's WAL - checkpoint must never stall acquisition for every state.db. +- ``acquire(path)`` returns the current generation for *path* and bumps its refcount. +- ``close()`` on a shared instance is a NO-OP: the registry owns the connection lifecycle. +- ``release(db)`` decrements the generation *db was acquired from* (object-keyed, so an + inode replacement cannot strand a still-owned generation); the final release of a + retired generation tears it down. +- On inode change the old generation is RETIRED (never lent again) but stays alive until + its holders release. If the replacement open fails the registry keeps NO path entry. +- All teardown happens OUTSIDE the registry lock: a final release's WAL checkpoint must + never stall acquisition for every state.db. """ from __future__ import annotations @@ -52,13 +46,13 @@ class _Generation: _lock = threading.Lock() -# path → live generation; retired generations move to _retired (keyed by -# id(db)) until their last holder releases. +# path → live generation; retired generations move to _retired (keyed by id(db)) until +# their last holder releases. _generations: Dict[Path, _Generation] = {} _retired: Dict[int, _Generation] = {} -# Paths whose next generation is being constructed. Construction runs outside -# _lock (schema reconciliation can take seconds), but peers for the SAME file -# must wait or every cold caller opens its own writer before a winner is chosen. +# Paths whose next generation is being constructed. Construction runs outside _lock +# (schema reconciliation can take seconds), but peers for the SAME file must wait or +# every cold caller opens its own writer before a winner is chosen. _opening: Dict[Path, threading.Event] = {} @@ -89,13 +83,11 @@ def _finish_opening(path: Path, opening: threading.Event) -> None: def acquire(db_path: Optional[Path] = None) -> "SessionDB": - """Return the shared SessionDB for *db_path*, incrementing its refcount. - - If the file was replaced (different inode) since the generation opened, - that generation is RETIRED but stays alive for its holders, and a fresh one - is opened in its place. Raises whatever ``SessionDB.__init__`` raises; on - a replacement-open failure the registry holds NO entry for the path. - """ + """Return the shared SessionDB for *db_path*, incrementing its refcount. If the file was + replaced (different inode) since the generation opened, that generation is RETIRED + but stays alive for its holders, and a fresh one is opened in its place. Raises + whatever ``SessionDB.__init__`` raises; on a replacement-open failure the registry + holds NO entry for the path.""" from hermes_state import _default_db_path raw_path = Path(db_path) if db_path is not None else Path(_default_db_path()) @@ -119,12 +111,12 @@ def acquire(db_path: Optional[Path] = None) -> "SessionDB": if opening is None: opening = _opening[path] = threading.Event() break - # Another caller is constructing this path; wait without holding the - # global lock. A failed opener signals too, so a waiter can retry. + # Another caller is constructing this path; wait without holding the global + # lock. A failed opener signals too, so a waiter can retry. opening.wait() - # Open OUTSIDE the lock; the per-path marker prevents redundant writers - # without serialising other files. + # Open OUTSIDE the lock; the per-path marker prevents redundant writers without + # serialising other files. try: db = _open_session_db(path) db._shared_registry_owned = True @@ -150,11 +142,9 @@ def acquire(db_path: Optional[Path] = None) -> "SessionDB": def _retire_generation_locked(path: Path, generation: _Generation) -> None: - """Retire *generation* so it is never lent again (caller holds _lock). - - It stays alive for its holders, tracked in ``_retired`` by ``id(db)`` so - their releases find it even after the path maps to a new generation. - """ + """Retire *generation* so it is never lent again (caller holds _lock). It stays alive + for its holders, tracked in ``_retired`` by ``id(db)`` so their releases find it even + after the path maps to a new generation.""" generation.retired = True if _generations.get(path) is generation: del _generations[path] @@ -162,13 +152,10 @@ def _retire_generation_locked(path: Path, generation: _Generation) -> None: def release(db: "SessionDB") -> bool: - """Decrement the refcount of a shared SessionDB. - - Returns ``True`` if *db* was shared; ``False`` if it is not registry-managed - (caller owns close()). The final release tears the generation down OUTSIDE - the registry lock. Lookup is object-keyed, so holders of an old generation - release into its retired record, not into whatever the path currently names. - """ + """Decrement the refcount of a shared SessionDB. ``True`` if *db* was shared; ``False`` + if it is not registry-managed (caller owns close()). The final release tears the + generation down OUTSIDE the registry lock. Lookup is object-keyed, so holders of an + old generation release into its retired record, not into whatever the path names.""" if db is None: return False key = id(db) @@ -198,18 +185,16 @@ def release(db: "SessionDB") -> bool: _generations.pop(Path(path), None) except (TypeError, ValueError): pass - # Teardown OUTSIDE the lock: stopping the token writer, WAL checkpoint and - # read-pool drain must not block acquisition for every other state.db. + # Teardown OUTSIDE the lock: stopping the token writer, WAL checkpoint and read-pool + # drain must not block acquisition for every other state.db. if needs_teardown: _teardown(db) return True def close_all() -> int: - """Close every shared SessionDB regardless of refcount; returns the count. - - For gateway shutdown, after all agents and cron jobs finished. Idempotent. - """ + """Close every shared SessionDB regardless of refcount; returns the count. For gateway + shutdown, after all agents and cron jobs finished. Idempotent.""" with _lock: generations = list(_generations.values()) + list(_retired.values()) _generations.clear() @@ -222,12 +207,9 @@ def close_all() -> int: def live_shared_session_dbs() -> List["SessionDB"]: - """Snapshot of every live (non-retired) shared SessionDB (refcounts untouched). - - For in-process maintenance (housekeeping deferred-FTS retry). A concurrent - final release may close an instance, in which case the callee sees - ``_conn is None``. - """ + """Snapshot of every live (non-retired) shared SessionDB (refcounts untouched), for + in-process maintenance. A concurrent final release may close an instance, in which + case the callee sees ``_conn is None``.""" with _lock: return [g.db for g in _generations.values() if not g.retired] @@ -242,26 +224,15 @@ def stats() -> Dict[str, int]: } -# ── Backwards-compatible aliases (hermes_state re-exports them) ── - -def get_shared_session_db(db_path: Optional[Path] = None) -> "SessionDB": - return acquire(db_path) - - -def release_shared_session_db(db: "SessionDB") -> bool: - return release(db) - - -def close_shared_session_dbs() -> int: - return close_all() +# Backwards-compatible aliases (hermes_state re-exports them). +get_shared_session_db = acquire +release_shared_session_db = release +close_shared_session_dbs = close_all def release_or_close(db: "SessionDB") -> None: - """Release a shared instance, or close it when it is not registry-managed. - - Drop-in for a plain ``db.close()``: read-only opens, CLI one-shots and - test fakes fall back to a direct close. - """ + """Release a shared instance, or close it when it is not registry-managed. Drop-in for a + plain ``db.close()``: read-only opens, CLI one-shots and test fakes fall back.""" if not release(db): try: db.close()