fix(state): lock guard rides SQLite's own descriptors
OFD locks on the connection's fds die with the connection: no private descriptors to track, retire, or exclude from holder scans.
This commit is contained in:
+8
-21
@@ -522,7 +522,7 @@ class SessionDB(
|
||||
self._retired_capture_lock = threading.Lock()
|
||||
self._retire_connection: Optional[Callable[[Any], None]] = None
|
||||
self._connection_pinned = False # one unmatched C reference taken at most once per handle
|
||||
self._wal_lock_guard_held = False # hermes_state_lockguard.hold() taken by _open_writer
|
||||
self._wal_lock_guard: dict = {} # hermes_state_lockguard.hold() record, see _open_writer
|
||||
self._db_corrupt, self._db_corrupt_reason = False, "" # sticky quarantine (StateDbCorruptError)
|
||||
self._fts_usermerge_floor_applied = False # one-shot usermerge-floor write guard
|
||||
self._fts_enabled = self._fts_stale = self._trigram_available = False
|
||||
@@ -602,11 +602,10 @@ class SessionDB(
|
||||
# FTS optimization is OPT-IN (`hermes db optimize`); no background worker races session lifecycle.
|
||||
self._ensure_db_file_generation()
|
||||
if self._wal_active:
|
||||
# Independent copies of the two POSIX locks that keep a sibling's close from unlinking
|
||||
# this WAL generation: any in-process open()/close() of state.db or -shm cancels SQLite's
|
||||
# own (howtocorrupt §2.2), these survive it. Released in close().
|
||||
_lockguard.hold(self.db_path)
|
||||
self._wal_lock_guard_held = True
|
||||
# OFD copies of the two POSIX locks that keep a sibling's close from unlinking this WAL
|
||||
# generation: any in-process open()/close() of state.db or -shm cancels SQLite's own
|
||||
# (howtocorrupt §2.2); these survive it. Lifted in close().
|
||||
self._wal_lock_guard = _lockguard.hold(self.db_path)
|
||||
|
||||
def _open_read_only(self) -> None:
|
||||
"""Read-only attach for cross-profile aggregation: no schema init, NO write
|
||||
@@ -1111,11 +1110,8 @@ class SessionDB(
|
||||
return False
|
||||
if sys.platform.startswith("linux"):
|
||||
watched = _watched_sqlite_sidecar_paths(self.db_path)
|
||||
guard_fds = _lockguard.owned_fds()
|
||||
try:
|
||||
for target, fd_path in _proc_fd_targets(os.getpid()):
|
||||
if int(fd_path.rsplit("/", 1)[1]) in guard_fds:
|
||||
continue # the lock guard's own descriptor (see hermes_state_lockguard)
|
||||
canonical = _canonical_sqlite_path(target)
|
||||
if (" (deleted)" in target and canonical in watched
|
||||
and _fd_is_truly_unlinked(fd_path, watched[canonical])):
|
||||
@@ -1334,8 +1330,8 @@ class SessionDB(
|
||||
"""
|
||||
if self._quarantine_reason() is not None:
|
||||
return
|
||||
if self._wal_lock_guard_held:
|
||||
_lockguard.refresh(self.db_path) # -shm minted after open, or path re-pointed
|
||||
if self._wal_lock_guard:
|
||||
_lockguard.hold(self.db_path, self._wal_lock_guard) # a -shm minted after open
|
||||
try:
|
||||
with self._lock:
|
||||
result = self._conn.execute("PRAGMA wal_checkpoint(PASSIVE)").fetchone()
|
||||
@@ -1414,10 +1410,7 @@ class SessionDB(
|
||||
self._conn.execute("PRAGMA wal_checkpoint(PASSIVE)")
|
||||
except Exception as exc:
|
||||
logger.debug("WAL checkpoint (PASSIVE) at close failed: %s", exc)
|
||||
# Release the guard first: SQLite's close-time reset then sees only real holders
|
||||
# (a sibling process's own intact locks still refuse the unlink; a true last close
|
||||
# ends the generation, so a later state.db replace never pairs with a stale WAL).
|
||||
self._release_wal_lock_guard()
|
||||
_lockguard.release(self._wal_lock_guard) # before the close: see release()
|
||||
if retire_without_close:
|
||||
self._pin_connection(self._conn)
|
||||
self._conn = None
|
||||
@@ -1427,12 +1420,6 @@ class SessionDB(
|
||||
# Only a clean close ends the generation; retain the recorded
|
||||
# identity when retiring an unsafe handle.
|
||||
self._db_sidecar_identity = {}
|
||||
_lockguard.retire_idle(self.db_path)
|
||||
|
||||
def _release_wal_lock_guard(self) -> None:
|
||||
if self._wal_lock_guard_held:
|
||||
self._wal_lock_guard_held = False
|
||||
_lockguard.release(self.db_path)
|
||||
|
||||
def __del__(self) -> None:
|
||||
"""Safety net: close() if the caller forgot. Attribute access stays
|
||||
|
||||
@@ -287,7 +287,7 @@ def _iter_darwin_fd_targets():
|
||||
yield pid, fd, target, identity
|
||||
|
||||
|
||||
def _iter_darwin_sidecar_holders(db_path, *, own_pid: int, skip_fds=frozenset()) -> List[Tuple[int, str]]:
|
||||
def _iter_darwin_sidecar_holders(db_path) -> List[Tuple[int, str]]:
|
||||
"""The macOS leg of :func:`iter_deleted_sqlite_sidecar_holders`: libproc enumeration matched
|
||||
against the watched sidecar paths, judged by identity.
|
||||
|
||||
@@ -297,9 +297,7 @@ def _iter_darwin_sidecar_holders(db_path, *, own_pid: int, skip_fds=frozenset())
|
||||
base = os.path.realpath(os.path.abspath(os.fspath(db_path)))
|
||||
watched = {os.path.normcase(path): path for path in (base + "-wal", base + "-shm")}
|
||||
holders: List[Tuple[int, str]] = []
|
||||
for pid, fd, target, identity in _iter_darwin_fd_targets():
|
||||
if pid == own_pid and fd in skip_fds:
|
||||
continue # our lock guard's descriptor (hermes_state_lockguard), not a SQLite connection
|
||||
for pid, _fd, target, identity in _iter_darwin_fd_targets():
|
||||
literal = watched.get(os.path.normcase(target))
|
||||
if literal is not None and _identity_is_truly_unlinked(identity, literal):
|
||||
holders.append((pid, target))
|
||||
@@ -322,15 +320,11 @@ def iter_deleted_sqlite_sidecar_holders(db_path) -> List[Tuple[int, str]]:
|
||||
return []
|
||||
holders: List[Tuple[int, str]] = []
|
||||
try:
|
||||
from hermes_state_lockguard import owned_fds
|
||||
own_pid, guard_fds = os.getpid(), owned_fds()
|
||||
if sys.platform == "darwin":
|
||||
holders = _iter_darwin_sidecar_holders(db_path, own_pid=own_pid, skip_fds=guard_fds)
|
||||
holders = _iter_darwin_sidecar_holders(db_path)
|
||||
elif sys.platform.startswith("linux"):
|
||||
watched = _watched_sqlite_sidecar_paths(db_path)
|
||||
for pid, target, fd_path in _iter_proc_fd_targets():
|
||||
if pid == own_pid and int(fd_path.rsplit("/", 1)[1]) in guard_fds:
|
||||
continue # our lock guard's descriptor, not a connection on a dead generation
|
||||
canonical = _canonical_sqlite_path(target)
|
||||
if (" (deleted)" in target and canonical in watched
|
||||
and _fd_is_truly_unlinked(fd_path, watched[canonical])):
|
||||
|
||||
+79
-153
@@ -1,23 +1,19 @@
|
||||
"""Hold a state.db writer's WAL-mode file locks on descriptors SQLite does not own.
|
||||
"""Hold a state.db writer's WAL-mode file locks in a form a stray ``close()`` cannot cancel.
|
||||
|
||||
SQLite protects a live WAL generation with two POSIX advisory locks: a SHARED lock on the main
|
||||
file's lock range and a shared lock on the DMS byte of ``state.db-shm``. A sibling process may
|
||||
checkpoint and unlink ``-wal``/``-shm`` at its close only after taking both EXCLUSIVE. POSIX locks
|
||||
are per process, so any ``open()``/``close()`` of those two files inside the holder — a raw probe,
|
||||
a plugin, a stray ``head -c`` in-process — cancels both (sqlite.org/howtocorrupt.html §2.2) and
|
||||
the next foreign close strands the holder on a deleted generation (``DeletedWalGenerationError``).
|
||||
a plugin, a tool reading ``~/.hermes`` — cancels both (sqlite.org/howtocorrupt.html §2.2) and the
|
||||
next foreign close strands the holder on a deleted generation (``DeletedWalGenerationError``).
|
||||
|
||||
This module re-holds the same two ranges as *open file description* locks (``F_OFD_SETLK``) on
|
||||
private descriptors that are never closed. OFD locks belong to the description, not the process:
|
||||
a stray ``close()`` elsewhere cannot cancel them, and releasing them with ``F_UNLCK`` never
|
||||
disturbs SQLite's own locks. Both lock types conflict with a foreign EXCLUSIVE, so the sibling's
|
||||
close-time unlink is refused for as long as a writer handle is open here. The same refusal applies
|
||||
to THIS process's close: the last writer no longer deletes the sidecars, which is the
|
||||
``SQLITE_DBCONFIG_NO_CKPT_ON_CLOSE`` behaviour on runtimes whose ``sqlite3`` cannot arm it.
|
||||
|
||||
One guard per database path per process, refcounted across writer handles; descriptors are
|
||||
retired (never closed) when the path is re-pointed at a new inode. No-op on Windows and on
|
||||
runtimes without OFD locks.
|
||||
This module adds the same two ranges as *open file description* locks (``F_OFD_SETLK``) on the
|
||||
descriptors SQLite itself holds. OFD locks belong to the description, not the process: a stray
|
||||
``close()`` elsewhere cannot cancel them, they die with the connection's own descriptor (nothing
|
||||
extra to track or retire), and they conflict with a foreign EXCLUSIVE exactly like SQLite's own,
|
||||
so the sibling's close-time unlink is refused while a guarded handle is open. The guard is
|
||||
lifted before the handle's own close so a true last close still ends the generation normally.
|
||||
No-op on Windows and on runtimes without OFD locks.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -26,9 +22,7 @@ import logging
|
||||
import os
|
||||
import struct
|
||||
import sys
|
||||
import threading
|
||||
from pathlib import Path
|
||||
from typing import Dict, List, Optional
|
||||
from typing import Dict, Optional, Tuple
|
||||
|
||||
logger = logging.getLogger("hermes_state")
|
||||
|
||||
@@ -54,6 +48,13 @@ except ImportError: # Windows
|
||||
# struct flock differs per libc: glibc/musl put type+whence first, Darwin/BSD last.
|
||||
_FLOCK_FORMAT = "@qqihh" if sys.platform == "darwin" or "bsd" in sys.platform else "@hhqqi"
|
||||
|
||||
Identity = Tuple[int, int]
|
||||
Held = Dict[int, Identity] # fd -> (st_dev, st_ino) it referenced when locked
|
||||
|
||||
|
||||
def supported() -> bool:
|
||||
return _F_OFD_SETLK is not None
|
||||
|
||||
|
||||
def _flock(lock_type: int, start: int, length: int) -> bytes:
|
||||
if _FLOCK_FORMAT == "@qqihh":
|
||||
@@ -71,150 +72,75 @@ def _ofd_lock(fd: int, lock_type: int, start: int, length: int) -> bool:
|
||||
return True
|
||||
|
||||
|
||||
class _PathGuard:
|
||||
__slots__ = ("main_fd", "main_ident", "shm_fd", "shm_ident", "refs", "main_locked", "shm_locked")
|
||||
|
||||
def __init__(self) -> None:
|
||||
self.main_fd = self.shm_fd = -1
|
||||
self.main_ident = self.shm_ident = None
|
||||
self.refs = 0
|
||||
self.main_locked = self.shm_locked = False
|
||||
|
||||
|
||||
_LOCK = threading.Lock()
|
||||
_GUARDS: Dict[str, _PathGuard] = {}
|
||||
_RETIRED_FDS: List[int] = [] # descriptors for re-pointed paths; closing one would cancel SQLite's locks
|
||||
|
||||
|
||||
def supported() -> bool:
|
||||
return _F_OFD_SETLK is not None
|
||||
|
||||
|
||||
def _bind_fd(path: str, fd: int, ident) -> tuple:
|
||||
"""Return ``(fd, ident)`` for *path*, reusing *fd* while it still names the path's inode."""
|
||||
def _identity(path: str) -> Optional[Identity]:
|
||||
try:
|
||||
st = os.stat(path)
|
||||
except OSError:
|
||||
return fd, ident
|
||||
current = (st.st_dev, st.st_ino)
|
||||
if fd >= 0 and ident == current:
|
||||
return fd, ident
|
||||
if fd >= 0:
|
||||
_RETIRED_FDS.append(fd)
|
||||
return None
|
||||
return (st.st_dev, st.st_ino)
|
||||
|
||||
|
||||
def _own_fds_for(identities: Dict[Identity, Tuple[int, int]]):
|
||||
"""Yield ``(fd, identity, (start, length))`` for every descriptor of this process on one of
|
||||
*identities* (SQLite's own connection descriptors; the cached header-probe fd too, harmless)."""
|
||||
for fd_dir in ("/proc/self/fd", "/dev/fd"):
|
||||
try:
|
||||
names = os.listdir(fd_dir)
|
||||
except OSError:
|
||||
continue
|
||||
for name in names:
|
||||
if not name.isdigit():
|
||||
continue
|
||||
fd = int(name)
|
||||
try:
|
||||
st = os.fstat(fd)
|
||||
except OSError:
|
||||
continue
|
||||
ident = (st.st_dev, st.st_ino)
|
||||
rng = identities.get(ident)
|
||||
if rng is not None:
|
||||
yield fd, ident, rng
|
||||
return
|
||||
|
||||
|
||||
def hold(db_path, held: Optional[Held] = None) -> Held:
|
||||
"""Lock the guard ranges on every descriptor this process has open on ``state.db`` and its
|
||||
``-shm``; returns the record :func:`release` needs (pass it back to extend an existing one:
|
||||
a ``-shm`` minted after open, a reopened connection). Safe to repeat."""
|
||||
held = {} if held is None else held
|
||||
if not supported():
|
||||
return held
|
||||
base = os.fspath(db_path)
|
||||
wanted: Dict[Identity, Tuple[int, int]] = {}
|
||||
for path, rng in ((base, (_SHARED_FIRST, _SHARED_SIZE)), (base + "-shm", (_SHM_DMS_BYTE, 1))):
|
||||
ident = _identity(path)
|
||||
if ident is not None:
|
||||
wanted[ident] = rng
|
||||
try:
|
||||
fd = os.open(path, os.O_RDONLY | getattr(os, "O_CLOEXEC", 0))
|
||||
for fd, ident, (start, length) in _own_fds_for(wanted):
|
||||
if held.get(fd) == ident:
|
||||
continue
|
||||
if _ofd_lock(fd, _F_RDLCK, start, length):
|
||||
held[fd] = ident
|
||||
except OSError:
|
||||
return -1, None
|
||||
return fd, current
|
||||
logger.debug("WAL lock guard unavailable for %s", base, exc_info=True)
|
||||
return held
|
||||
|
||||
|
||||
def _apply_locked(guard: _PathGuard, db_path: str) -> None:
|
||||
guard.main_fd, guard.main_ident = _bind_fd(db_path, guard.main_fd, guard.main_ident)
|
||||
if guard.main_fd >= 0:
|
||||
guard.main_locked = _ofd_lock(guard.main_fd, _F_RDLCK, _SHARED_FIRST, _SHARED_SIZE)
|
||||
guard.shm_fd, guard.shm_ident = _bind_fd(db_path + "-shm", guard.shm_fd, guard.shm_ident)
|
||||
if guard.shm_fd >= 0:
|
||||
guard.shm_locked = _ofd_lock(guard.shm_fd, _F_RDLCK, _SHM_DMS_BYTE, 1)
|
||||
|
||||
|
||||
def hold(db_path: Path) -> None:
|
||||
"""Take (or add a reference to) the guard for *db_path*. Call once per writer handle after
|
||||
its connection is open in WAL mode; pair with :func:`release`."""
|
||||
def release(held: Held) -> None:
|
||||
"""Unlock what :func:`hold` locked, on descriptors that still reference the same inode (a
|
||||
number recycled onto another file is left alone). Call BEFORE the handle's own close so
|
||||
SQLite's close-time reset sees only real holders: a sibling process's intact locks still
|
||||
refuse the unlink, and a true last close ends the generation, so a later ``state.db``
|
||||
replace never pairs with a stale WAL."""
|
||||
if not supported():
|
||||
return
|
||||
key = os.fspath(db_path)
|
||||
with _LOCK:
|
||||
guard = _GUARDS.setdefault(key, _PathGuard())
|
||||
guard.refs += 1
|
||||
for fd, ident in list(held.items()):
|
||||
try:
|
||||
_apply_locked(guard, key)
|
||||
st = os.fstat(fd)
|
||||
if (st.st_dev, st.st_ino) == ident:
|
||||
_ofd_lock(fd, _F_UNLCK, _SHARED_FIRST, _SHARED_SIZE)
|
||||
_ofd_lock(fd, _F_UNLCK, _SHM_DMS_BYTE, 1)
|
||||
except OSError:
|
||||
logger.debug("WAL lock guard unavailable for %s", key, exc_info=True)
|
||||
|
||||
|
||||
def refresh(db_path: Path) -> None:
|
||||
"""Re-arm a held guard: a ``-shm`` that did not exist at :func:`hold` time, or a path re-pointed
|
||||
at a new inode since. Cheap when everything is in place (one ``stat`` per file)."""
|
||||
if not supported():
|
||||
return
|
||||
key = os.fspath(db_path)
|
||||
with _LOCK:
|
||||
guard = _GUARDS.get(key)
|
||||
if guard is None or guard.refs <= 0:
|
||||
return
|
||||
try:
|
||||
_apply_locked(guard, key)
|
||||
except OSError:
|
||||
logger.debug("WAL lock guard refresh failed for %s", key, exc_info=True)
|
||||
|
||||
|
||||
def release(db_path: Path) -> None:
|
||||
"""Drop one reference; the last one unlocks both ranges. Call BEFORE closing the handle's own
|
||||
connection so SQLite's close-time reset sees only real holders (a sibling process's intact
|
||||
locks still refuse the unlink; a true last close ends the generation normally, so a later
|
||||
``state.db`` replace never pairs with a stale WAL). Descriptors are closed by
|
||||
:func:`retire_idle` once no connection to the path remains."""
|
||||
if not supported():
|
||||
return
|
||||
key = os.fspath(db_path)
|
||||
with _LOCK:
|
||||
guard = _GUARDS.get(key)
|
||||
if guard is None or guard.refs <= 0:
|
||||
return
|
||||
guard.refs -= 1
|
||||
if guard.refs:
|
||||
return
|
||||
for fd, start, length in ((guard.main_fd, _SHARED_FIRST, _SHARED_SIZE), (guard.shm_fd, _SHM_DMS_BYTE, 1)):
|
||||
if fd >= 0:
|
||||
try:
|
||||
_ofd_lock(fd, _F_UNLCK, start, length)
|
||||
except OSError:
|
||||
logger.debug("WAL lock guard unlock failed for %s", key, exc_info=True)
|
||||
guard.main_locked = guard.shm_locked = False
|
||||
|
||||
|
||||
def retire_idle(db_path: Path) -> None:
|
||||
"""Close the guard descriptors once no tracked SQLite connection to *db_path* remains in this
|
||||
process. Closing then cancels nothing, and a lingering fd on the path would make another
|
||||
process's holder scan (``hermes doctor`` repair, snapshot restore) count this one as live.
|
||||
While any connection is still open the descriptors stay put: closing would cancel its locks."""
|
||||
if not supported():
|
||||
return
|
||||
key = os.fspath(db_path)
|
||||
with _LOCK:
|
||||
guard = _GUARDS.get(key)
|
||||
if guard is None or guard.refs or _path_has_live_connection(key):
|
||||
return
|
||||
del _GUARDS[key]
|
||||
for fd in (guard.main_fd, guard.shm_fd):
|
||||
if fd >= 0:
|
||||
try:
|
||||
os.close(fd)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
|
||||
def _path_has_live_connection(key: str) -> bool:
|
||||
try:
|
||||
from hermes_cli.sqlite_safe_read import has_live_connection
|
||||
except ImportError:
|
||||
return True # cannot prove quiescence: keep the descriptors
|
||||
return has_live_connection(key)
|
||||
|
||||
|
||||
def held(db_path: Path) -> bool:
|
||||
"""Both ranges currently guarded for *db_path* (diagnostics and tests)."""
|
||||
with _LOCK:
|
||||
guard = _GUARDS.get(os.fspath(db_path))
|
||||
return bool(guard and guard.refs and guard.main_locked and guard.shm_locked)
|
||||
|
||||
|
||||
def owned_fds() -> frozenset:
|
||||
"""Every descriptor this module keeps open (active and retired). Deleted-sidecar holder scans
|
||||
must skip these: a retired guard fd on an unlinked ``-shm`` is not a SQLite connection reading
|
||||
a dead generation, and reporting it would refuse every later open in this process."""
|
||||
with _LOCK:
|
||||
fds = set(_RETIRED_FDS)
|
||||
for guard in _GUARDS.values():
|
||||
fds.update(fd for fd in (guard.main_fd, guard.shm_fd) if fd >= 0)
|
||||
return frozenset(fds)
|
||||
pass
|
||||
held.clear()
|
||||
|
||||
Reference in New Issue
Block a user