refactor(hermes_state_registry): alias back-compat names directly, compact docstrings
This commit is contained in:
+48
-77
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user