fix(sessions): capture a lost WAL generation durably before shutdown settles
After DeletedWalGenerationError the retired frames exist only in an unlinked -wal inode that this process keeps open. #106315 stops them from being checkpointed under wrong page numbers, but the canonical remediation ("stop the gateway, dashboard and cron writers, then reopen") lets the kernel drop the inode with them: committed transactions whose disposition is unknown were destroyed by the prescribed recovery path itself. Capture the exact retired generation next to the database at the first halt, or at close() if that sees the loss first. The WAL is located by the (st_dev, st_ino) recorded at open among this process's own descriptors, never by pathname, so a sibling's deleted WAL or a newer sidecar minted at the same path cannot be mistaken for it; it is read with pread and no descriptor is closed, moved or truncated. The artifact holds the WAL, the -shm when still ours, the main image (or its header past a size cap) and a manifest with identities, digests and the sidecar generation found at the path at capture time. Nothing is merged: whether the frames belong on top of the file now at the path stays an operator decision. close() refuses to settle without the capture: it raises and leaves the handle open. Regressions: capture at halt and at close, refusal to settle on capture failure, inode-not-pathname selection with a second deleted WAL under the same name, refusal to guess without a recorded identity, and a subprocess control where rows committed only in the retired WAL are recovered from the capture alone after the writer process has exited. Refs #105670
This commit is contained in:
+51
-3
@@ -49,7 +49,7 @@ from hermes_state_dbfile import (
|
||||
_canonical_sqlite_path, _connect_tracked_db, _read_sqlite_application_id, _stat_sqlite_sidecar_identity,
|
||||
_watched_sqlite_sidecar_paths, has_invalid_sqlite_header_preopen, is_zeroed_state_db, quarantine_cross_process_lock,
|
||||
quarantine_invalid_state_db,
|
||||
refuse_deleted_wal_generation,
|
||||
RetiredGenerationCaptureError, capture_retired_wal_generation, refuse_deleted_wal_generation,
|
||||
)
|
||||
from hermes_state_messages import SessionMessagesMixin
|
||||
from hermes_state_wal import _WAL_INCOMPAT_MARKERS, apply_database_pragmas, apply_wal_with_fallback
|
||||
@@ -452,6 +452,9 @@ class SessionDB(
|
||||
self._db_file_application_id: int = 0
|
||||
self._db_sidecar_identity: Dict[str, tuple] = {}
|
||||
self._db_replaced = self._db_wal_generation_lost = False
|
||||
# Durable capture of a lost WAL generation (see _capture_retired_generation): once per handle.
|
||||
self._retired_generation_capture: Optional[Path] = None
|
||||
self._retired_capture_lock = threading.Lock()
|
||||
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
|
||||
@@ -1007,9 +1010,37 @@ class SessionDB(
|
||||
if self._db_wal_generation_lost or self._wal_generation_was_lost():
|
||||
self._db_wal_generation_lost = True
|
||||
self._disable_close_time_checkpoint()
|
||||
try:
|
||||
self._capture_retired_generation("halt")
|
||||
except RetiredGenerationCaptureError as exc:
|
||||
logger.error(
|
||||
"Could not capture the retired WAL generation of %s at halt: %s. close() retries "
|
||||
"the capture and refuses to settle without it.", self.db_path, exc,
|
||||
)
|
||||
logger.error(_DELETED_WAL_GENERATION_MSG)
|
||||
raise DeletedWalGenerationError(_DELETED_WAL_GENERATION_MSG)
|
||||
|
||||
def _capture_retired_generation(self, trigger: str) -> Path:
|
||||
"""Durably capture the lost WAL generation this handle still holds open, once per handle.
|
||||
|
||||
The quarantine keeps the retired frames from being checkpointed under wrong page numbers,
|
||||
but they live only in an unlinked inode that dies with this process's last descriptor, and
|
||||
the canonical DeletedWalGenerationError remediation is to stop the writers. Capturing at the
|
||||
first halt (or at close(), whichever sees the loss first) makes "preserve" outlive the
|
||||
process. Raises RetiredGenerationCaptureError; nothing is mutated on failure."""
|
||||
with self._retired_capture_lock:
|
||||
if self._retired_generation_capture is not None:
|
||||
return self._retired_generation_capture
|
||||
artifact = capture_retired_wal_generation(
|
||||
self.db_path, sidecar_identity=dict(self._db_sidecar_identity or {}), trigger=trigger,
|
||||
)
|
||||
self._retired_generation_capture = artifact
|
||||
logger.warning(
|
||||
"Captured the retired WAL generation of %s at %s to %s; read its manifest.json before deciding "
|
||||
"whether those frames belong on top of the file now at the path.", self.db_path, trigger, artifact,
|
||||
)
|
||||
return artifact
|
||||
|
||||
def _raise_if_db_replaced(self) -> None:
|
||||
"""Sticky-flag fast path (no log spam on every write), then the live probe."""
|
||||
if self._db_replaced:
|
||||
@@ -1168,7 +1199,24 @@ class SessionDB(
|
||||
pass
|
||||
with self._lock:
|
||||
if self._conn:
|
||||
quarantine_reason = self._quarantine_reason()
|
||||
generation_lost = not self.read_only and (
|
||||
self._db_wal_generation_lost
|
||||
or (bool(self._db_sidecar_identity) and self._wal_generation_was_lost())
|
||||
)
|
||||
if generation_lost:
|
||||
# Loss is settled here, not at exit: the unlinked WAL inode dies with this process's
|
||||
# last descriptor. Capture the exact retired generation before the handle goes and
|
||||
# refuse to settle without it (the capture raises; the handle stays open).
|
||||
self._db_wal_generation_lost = True
|
||||
self._disable_close_time_checkpoint()
|
||||
artifact = self._capture_retired_generation("close")
|
||||
logger.warning(
|
||||
"Skipping the close-time WAL checkpoint for %s: this handle's WAL/SHM generation "
|
||||
"was deleted or replaced; the retired generation is captured at %s. Stop the other "
|
||||
"writers before reopening and inspect the capture before deciding its disposition.",
|
||||
self.db_path, artifact,
|
||||
)
|
||||
quarantine_reason = None if generation_lost else self._quarantine_reason()
|
||||
if quarantine_reason is not None:
|
||||
logger.warning(
|
||||
"Skipping the close-time WAL checkpoint for %s: this "
|
||||
@@ -1176,7 +1224,7 @@ class SessionDB(
|
||||
"before restarting, then run `hermes sessions recover --source %s --inspect-only`.",
|
||||
self.db_path, quarantine_reason, self.db_path,
|
||||
)
|
||||
elif not self.read_only: # PASSIVE, not TRUNCATE (see docstring)
|
||||
elif not self.read_only and not generation_lost: # PASSIVE, not TRUNCATE (see docstring)
|
||||
try:
|
||||
# Every cron run_agent opens+closes a transient SessionDB, so a TRUNCATE here fires
|
||||
# a full WAL reset many times/hour, racing the gateway's long-lived writer on large
|
||||
|
||||
+218
-3
@@ -10,6 +10,7 @@ call time, so tests that monkeypatch ``hermes_state.<name>`` keep intercepting.
|
||||
from __future__ import annotations
|
||||
|
||||
import contextlib
|
||||
import hashlib
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
@@ -42,13 +43,14 @@ _RETIRED_HEADER_PROBE_FDS: "list[int]" = [] # intentionally never closed
|
||||
_FTS_TABLE_NAMES = ("messages_fts", "messages_fts_trigram", "messages_fts_cjk")
|
||||
|
||||
|
||||
def _pread_db_header(db_path: Path, length: int) -> "Optional[bytes]":
|
||||
"""Lock-safe raw header read of a possibly-live SQLite database: POSIX preads from a cached,
|
||||
def _pread_db_range(db_path: Path, offset: int, length: int) -> "Optional[bytes]":
|
||||
"""Lock-safe raw read of a possibly-live SQLite database: POSIX preads from a cached,
|
||||
never-closed fd (rebound when the path names a new inode); Windows reads plainly, since
|
||||
advisory-lock cancellation is a POSIX-only hazard."""
|
||||
from hermes_state import _IS_WINDOWS
|
||||
if _IS_WINDOWS:
|
||||
with contextlib.suppress(OSError), db_path.open("rb") as handle:
|
||||
handle.seek(offset)
|
||||
return handle.read(length)
|
||||
return None
|
||||
key = str(db_path)
|
||||
@@ -74,10 +76,15 @@ def _pread_db_header(db_path: Path, length: int) -> "Optional[bytes]":
|
||||
return None
|
||||
cached = _HEADER_PROBE_FDS[key] = (fd, fst.st_dev, fst.st_ino)
|
||||
with contextlib.suppress(OSError):
|
||||
return os.pread(cached[0], length, 0)
|
||||
return os.pread(cached[0], length, offset)
|
||||
return None
|
||||
|
||||
|
||||
def _pread_db_header(db_path: Path, length: int) -> "Optional[bytes]":
|
||||
"""Lock-safe raw header read of a possibly-live SQLite database (see :func:`_pread_db_range`)."""
|
||||
return _pread_db_range(db_path, 0, length)
|
||||
|
||||
|
||||
def _read_sqlite_application_id(db_path: Path) -> "Optional[int]":
|
||||
"""application_id from the SQLite header, via the lock-safe :func:`_pread_db_header`."""
|
||||
from hermes_state_errors import _STATE_DB_APPLICATION_ID_OFFSET
|
||||
@@ -149,6 +156,214 @@ def refuse_deleted_wal_generation(db_path) -> None:
|
||||
raise DeletedWalGenerationError(_DELETED_WAL_GENERATION_MSG)
|
||||
|
||||
|
||||
# ── Retired WAL generation capture ──────────────────────────────────────────────────────────────
|
||||
#
|
||||
# When a writer's -wal/-shm generation is deleted or replaced underneath it, the frames committed only in
|
||||
# that WAL survive exactly as long as this process keeps the unlinked inode open. The quarantine (#105670)
|
||||
# stops them from being checkpointed under wrong page numbers, but the canonical remediation -- "stop the
|
||||
# writers, then reopen" -- lets the kernel drop the inode with them. So the exact generation is captured
|
||||
# durably first: located by the recorded (st_dev, st_ino) among this process's OWN descriptors (never by
|
||||
# pathname, so a sibling's deleted WAL or a newer sidecar minted at the same path cannot be mistaken for
|
||||
# it), read with pread (no descriptor is closed, moved or truncated) and written with fsync.
|
||||
|
||||
RETIRED_GENERATION_DIR_SUFFIX = ".retired-wal-"
|
||||
RETIRED_GENERATION_MANIFEST = "manifest.json"
|
||||
RETIRED_GENERATION_MANIFEST_VERSION = 1
|
||||
# Up to this size the main image is copied whole, so the artifact is a self-contained state.db + -wal
|
||||
# pair that `hermes sessions recover --source <dir>/state.db` can open. Above it only the 100-byte
|
||||
# header is kept (the manifest says so): unlike the WAL inode, the main file survives process exit at
|
||||
# its path, and a multi-GB copy inside a shutdown path is a worse failure than a header-only artifact.
|
||||
RETIRED_GENERATION_MAIN_IMAGE_MAX_BYTES = 512 * 1024 * 1024
|
||||
_CAPTURE_CHUNK_BYTES = 8 * 1024 * 1024
|
||||
_SQLITE_HEADER_BYTES = 100
|
||||
|
||||
|
||||
class RetiredGenerationCaptureError(RuntimeError):
|
||||
"""The retired WAL generation could not be captured durably; nothing was mutated."""
|
||||
|
||||
|
||||
def _fsync_path(path: Path) -> None:
|
||||
"""fsync a file or directory we own. Never used on the live database: opening and closing a
|
||||
descriptor on a file SQLite has locked would cancel this process's POSIX advisory locks."""
|
||||
if os.name == "nt":
|
||||
return # directories cannot be opened; file writes fsync their own handle
|
||||
fd = os.open(path, os.O_RDONLY)
|
||||
try:
|
||||
os.fsync(fd)
|
||||
finally:
|
||||
os.close(fd)
|
||||
|
||||
|
||||
def _own_descriptor_for_identity(identity) -> "Optional[int]":
|
||||
"""This process's open descriptor for the inode ``(st_dev, st_ino)``, or None.
|
||||
|
||||
Exact-generation qualified: pathnames are never consulted. The descriptor stays owned by
|
||||
SQLite; callers only pread from it."""
|
||||
if not identity:
|
||||
return None
|
||||
wanted = tuple(identity)
|
||||
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
|
||||
try:
|
||||
st = os.fstat(int(name))
|
||||
except OSError:
|
||||
continue
|
||||
if (st.st_dev, st.st_ino) == wanted:
|
||||
return int(name)
|
||||
return None # the table was readable: a definitive miss
|
||||
return None
|
||||
|
||||
|
||||
def _copy_descriptor(fd: int, dest: Path, *, size: int) -> Dict[str, Any]:
|
||||
"""pread ``size`` bytes of ``fd`` from offset 0 into ``dest`` (temp file, fsync, rename)."""
|
||||
digest = hashlib.sha256()
|
||||
part = dest.with_name(dest.name + ".part")
|
||||
offset = 0
|
||||
with open(part, "wb") as out:
|
||||
while offset < size:
|
||||
chunk = os.pread(fd, min(_CAPTURE_CHUNK_BYTES, size - offset), offset)
|
||||
if not chunk:
|
||||
break
|
||||
out.write(chunk)
|
||||
digest.update(chunk)
|
||||
offset += len(chunk)
|
||||
out.flush()
|
||||
os.fsync(out.fileno())
|
||||
os.replace(part, dest)
|
||||
return {"file": dest.name, "bytes": offset, "sha256": digest.hexdigest()}
|
||||
|
||||
|
||||
def _copy_main_image(db_path: Path, dest: Path, *, size: int) -> Dict[str, Any]:
|
||||
"""Copy the live main file through the lock-safe cached descriptor (see ``_pread_db_range``)."""
|
||||
digest = hashlib.sha256()
|
||||
part = dest.with_name(dest.name + ".part")
|
||||
offset = 0
|
||||
with open(part, "wb") as out:
|
||||
while offset < size:
|
||||
chunk = _pread_db_range(db_path, offset, min(_CAPTURE_CHUNK_BYTES, size - offset))
|
||||
if not chunk:
|
||||
break
|
||||
out.write(chunk)
|
||||
digest.update(chunk)
|
||||
offset += len(chunk)
|
||||
out.flush()
|
||||
os.fsync(out.fileno())
|
||||
os.replace(part, dest)
|
||||
return {"file": dest.name, "bytes": offset, "sha256": digest.hexdigest()}
|
||||
|
||||
|
||||
def _parse_sqlite_header(header: bytes) -> Dict[str, Any]:
|
||||
if len(header) < _SQLITE_HEADER_BYTES or header[:16] != b"SQLite format 3\x00":
|
||||
return {"valid": False}
|
||||
raw_page_size = struct.unpack(">H", header[16:18])[0]
|
||||
fields = {"change_counter": 24, "page_count": 28, "user_version": 60, "application_id": 68,
|
||||
"version_valid_for": 92}
|
||||
parsed = {name: struct.unpack(">I", header[off:off + 4])[0] for name, off in fields.items()}
|
||||
return {"valid": True, "page_size": 65536 if raw_page_size == 1 else raw_page_size, **parsed}
|
||||
|
||||
|
||||
def _write_json_durably(path: Path, payload: Dict[str, Any]) -> None:
|
||||
part = path.with_name(path.name + ".part")
|
||||
with open(part, "w", encoding="utf-8") as out:
|
||||
json.dump(payload, out, indent=2, sort_keys=True)
|
||||
out.write("\n")
|
||||
out.flush()
|
||||
os.fsync(out.fileno())
|
||||
os.replace(part, path)
|
||||
|
||||
|
||||
def capture_retired_wal_generation(
|
||||
db_path, *, sidecar_identity: Dict[str, tuple], trigger: str,
|
||||
main_image_max_bytes: int = RETIRED_GENERATION_MAIN_IMAGE_MAX_BYTES,
|
||||
) -> Path:
|
||||
"""Durably capture the lost WAL generation this process still holds open; return the artifact dir.
|
||||
|
||||
The artifact sits next to the database as ``<name>.retired-wal-<utc ts>-<pid>/`` and holds the
|
||||
retired ``<name>-wal`` (and ``-shm`` when still ours), the main image (or its header past the size
|
||||
cap) and ``manifest.json`` with identities, sizes, digests and the sidecar generation found at the
|
||||
path at capture time. Nothing is merged: whether those frames belong on top of the main file is
|
||||
the operator's decision. Raises :class:`RetiredGenerationCaptureError` when the exact generation
|
||||
cannot be located or written; no descriptor is ever closed, moved or truncated.
|
||||
"""
|
||||
db_path = Path(db_path)
|
||||
wal_identity = tuple(sidecar_identity.get("-wal") or ())
|
||||
if not wal_identity:
|
||||
raise RetiredGenerationCaptureError(
|
||||
f"no recorded WAL generation identity for {db_path}; refusing to guess it by pathname")
|
||||
wal_fd = _own_descriptor_for_identity(wal_identity)
|
||||
if wal_fd is None:
|
||||
raise RetiredGenerationCaptureError(
|
||||
f"this process no longer holds the retired WAL inode {wal_identity} of {db_path}")
|
||||
try:
|
||||
wal_size = os.fstat(wal_fd).st_size
|
||||
except OSError as exc:
|
||||
raise RetiredGenerationCaptureError(f"cannot stat the retired WAL of {db_path}: {exc}") from exc
|
||||
|
||||
stem = f"{db_path.name}{RETIRED_GENERATION_DIR_SUFFIX}{time.strftime('%Y%m%d-%H%M%S', time.gmtime())}-{os.getpid()}"
|
||||
final = db_path.with_name(stem)
|
||||
n = 0
|
||||
while final.exists() or final.with_name(final.name + ".partial").exists():
|
||||
n += 1
|
||||
final = db_path.with_name(f"{stem}-{n}")
|
||||
staging = final.with_name(final.name + ".partial")
|
||||
try:
|
||||
staging.mkdir(parents=True, exist_ok=False)
|
||||
manifest: Dict[str, Any] = {
|
||||
"version": RETIRED_GENERATION_MANIFEST_VERSION,
|
||||
"database": str(db_path),
|
||||
"trigger": trigger,
|
||||
"pid": os.getpid(),
|
||||
"captured_at": time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime()),
|
||||
"python": sys.version.split()[0],
|
||||
"sqlite": sqlite3.sqlite_version,
|
||||
"wal": {"identity": list(wal_identity),
|
||||
**_copy_descriptor(wal_fd, staging / (db_path.name + "-wal"), size=wal_size)},
|
||||
"shm": None,
|
||||
# The generation living at the path when we captured: a newer writer's, if one was minted.
|
||||
"path_generation_at_capture": {
|
||||
suffix: list(ident) for suffix, ident in _stat_sqlite_sidecar_identity(db_path).items()},
|
||||
"note": ("Frames in the captured WAL were committed by the retired generation. Whether they "
|
||||
"belong on top of the main file now at the path is an operator decision; inspect "
|
||||
"the copied image with `hermes sessions recover --inspect-only` first."),
|
||||
}
|
||||
shm_identity = tuple(sidecar_identity.get("-shm") or ())
|
||||
shm_fd = _own_descriptor_for_identity(shm_identity) if shm_identity else None
|
||||
if shm_fd is not None:
|
||||
with contextlib.suppress(OSError): # SQLite rebuilds the index; the WAL is what matters
|
||||
manifest["shm"] = {"identity": list(shm_identity), **_copy_descriptor(
|
||||
shm_fd, staging / (db_path.name + "-shm"), size=os.fstat(shm_fd).st_size)}
|
||||
header = _pread_db_range(db_path, 0, _SQLITE_HEADER_BYTES)
|
||||
if header is None:
|
||||
raise RetiredGenerationCaptureError(f"cannot read the main image header of {db_path}")
|
||||
main_size = os.stat(db_path).st_size
|
||||
main: Dict[str, Any] = {"identity": list(_stat_db_file_identity(db_path) or ()) or None,
|
||||
"size": main_size, "header": _parse_sqlite_header(header)}
|
||||
if main_size <= main_image_max_bytes:
|
||||
main.update(mode="copied", **_copy_main_image(db_path, staging / db_path.name, size=main_size))
|
||||
else:
|
||||
header_file = staging / (db_path.name + ".header")
|
||||
header_file.write_bytes(header)
|
||||
_fsync_path(header_file)
|
||||
main.update(mode="header_only", file=header_file.name, bytes=len(header))
|
||||
manifest["main"] = main
|
||||
_write_json_durably(staging / RETIRED_GENERATION_MANIFEST, manifest)
|
||||
_fsync_path(staging)
|
||||
os.replace(staging, final)
|
||||
_fsync_path(final.parent)
|
||||
except RetiredGenerationCaptureError:
|
||||
raise
|
||||
except OSError as exc:
|
||||
raise RetiredGenerationCaptureError(
|
||||
f"could not write the retired WAL generation of {db_path} under {staging}: {exc}") from exc
|
||||
return final
|
||||
|
||||
|
||||
def _connect_tracked_db(path, tracking_path=None, **kwargs):
|
||||
"""``sqlite3.connect`` that registers the open fd so byte-level probes of a live file are
|
||||
refused (an ``open()``/``close()`` would cancel every POSIX lock, even a running VACUUM's
|
||||
|
||||
@@ -0,0 +1,434 @@
|
||||
"""Durable capture of a lost WAL generation (#105670 follow-up).
|
||||
|
||||
After ``DeletedWalGenerationError`` the retired frames survive only in an unlinked inode this process
|
||||
keeps open, and the prescribed remediation ("stop the writers, then reopen") lets the kernel drop it.
|
||||
These tests lose the sidecars for real and check that the exact generation is captured next to the
|
||||
database at the first halt (or at ``close()``), located by inode rather than by pathname, and that a
|
||||
transaction committed only in the retired WAL is recoverable from the capture alone, including after
|
||||
the writer process has exited.
|
||||
"""
|
||||
|
||||
import contextlib
|
||||
import hashlib
|
||||
import json
|
||||
import os
|
||||
import queue
|
||||
import shutil
|
||||
import sqlite3
|
||||
import subprocess
|
||||
import sys
|
||||
import textwrap
|
||||
import threading
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
import hermes_state
|
||||
import hermes_state_wal
|
||||
from hermes_state import DeletedWalGenerationError, SessionDB
|
||||
from hermes_state_dbfile import (
|
||||
RETIRED_GENERATION_MANIFEST, RetiredGenerationCaptureError, capture_retired_wal_generation,
|
||||
)
|
||||
|
||||
FD_DIRECTORY = "/proc/self/fd" if sys.platform.startswith("linux") else "/dev/fd"
|
||||
not_windows = pytest.mark.skipif(sys.platform == "win32", reason="a held sidecar cannot be unlinked on Windows")
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def force_wal(monkeypatch):
|
||||
"""Pin WAL so this host's vulnerable SQLite still matches production topology."""
|
||||
monkeypatch.setattr(hermes_state_wal, "is_sqlite_wal_reset_vulnerable", lambda version_info=None: False)
|
||||
monkeypatch.setattr(hermes_state_wal, "resolve_journal_mode", lambda: "wal")
|
||||
|
||||
|
||||
def _make_db(path: Path, session_id: str, content: str) -> SessionDB:
|
||||
db = SessionDB(db_path=path)
|
||||
db.create_session(session_id, "cli")
|
||||
db.append_message(session_id, role="user", content=content)
|
||||
return db
|
||||
|
||||
|
||||
def _require_wal(db: SessionDB) -> Path:
|
||||
if not db._wal_active:
|
||||
db.close()
|
||||
pytest.skip("WAL not active on this filesystem")
|
||||
wal = Path(str(db.db_path) + "-wal")
|
||||
if not wal.exists():
|
||||
db.close()
|
||||
pytest.skip("WAL sidecar missing after first write")
|
||||
return wal
|
||||
|
||||
|
||||
def _lose_sidecars(db_path: Path, *, rename: bool) -> None:
|
||||
"""Take the -wal/-shm generation away from the writer, as the field incident did."""
|
||||
for suffix in ("-wal", "-shm"):
|
||||
side = Path(str(db_path) + suffix)
|
||||
if not side.exists():
|
||||
continue
|
||||
if rename:
|
||||
side.rename(db_path.with_name("retired" + side.name))
|
||||
else:
|
||||
side.unlink()
|
||||
|
||||
|
||||
def _descriptor_for(identity: tuple) -> int:
|
||||
for name in os.listdir(FD_DIRECTORY):
|
||||
try:
|
||||
fd = int(name)
|
||||
st = os.fstat(fd)
|
||||
except (ValueError, OSError):
|
||||
continue
|
||||
if (st.st_dev, st.st_ino) == identity:
|
||||
return fd
|
||||
pytest.fail("the owning SQLite connection discarded its WAL inode")
|
||||
|
||||
|
||||
def _descriptor_contents(fd: int) -> bytes:
|
||||
return os.pread(fd, os.fstat(fd).st_size, 0)
|
||||
|
||||
|
||||
def _wal_only_sentinel(db: SessionDB, session_id: str) -> str:
|
||||
"""Checkpoint history, then commit one row that lives only in the WAL."""
|
||||
db._conn.execute("PRAGMA wal_autocheckpoint=0")
|
||||
db._try_wal_checkpoint()
|
||||
sentinel = "committed only in the retired WAL " + "x" * 3000
|
||||
db.append_message(session_id, role="assistant", content=sentinel)
|
||||
return sentinel
|
||||
|
||||
|
||||
def _manifest(artifact: Path) -> dict:
|
||||
return json.loads((artifact / RETIRED_GENERATION_MANIFEST).read_text(encoding="utf-8"))
|
||||
|
||||
|
||||
def _recover(artifact: Path, name: str, workdir: Path) -> sqlite3.Connection:
|
||||
"""Open the captured image + WAL as a fresh database, nothing else from the original path."""
|
||||
workdir.mkdir()
|
||||
shutil.copyfile(artifact / name, workdir / name)
|
||||
shutil.copyfile(artifact / (name + "-wal"), workdir / (name + "-wal"))
|
||||
return sqlite3.connect(str(workdir / name))
|
||||
|
||||
|
||||
def _count(conn: sqlite3.Connection, content_like: str) -> int:
|
||||
return conn.execute("SELECT COUNT(*) FROM messages WHERE content LIKE ?", (content_like,)).fetchone()[0]
|
||||
|
||||
|
||||
def _integrity_ok(conn: sqlite3.Connection) -> bool:
|
||||
try:
|
||||
return [r[0] for r in conn.execute("PRAGMA integrity_check").fetchall()] == ["ok"]
|
||||
except sqlite3.DatabaseError:
|
||||
return False
|
||||
|
||||
|
||||
# ── Capture at the first halt ───────────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
@not_windows
|
||||
@pytest.mark.parametrize("rename", [False, True], ids=["unlink", "rename"])
|
||||
def test_halt_captures_the_exact_retired_generation(tmp_path, force_wal, rename):
|
||||
path = tmp_path / "state.db"
|
||||
db = _make_db(path, "gw-0", "seed")
|
||||
_require_wal(db)
|
||||
sentinel = _wal_only_sentinel(db, "gw-0")
|
||||
identity = db._db_sidecar_identity["-wal"]
|
||||
_lose_sidecars(path, rename=rename)
|
||||
original = _descriptor_contents(_descriptor_for(identity))
|
||||
|
||||
with pytest.raises(DeletedWalGenerationError):
|
||||
db.append_message("gw-0", role="user", content="after the loss")
|
||||
|
||||
artifact = db._retired_generation_capture
|
||||
assert artifact is not None and artifact.is_dir() and artifact.parent == tmp_path
|
||||
manifest = _manifest(artifact)
|
||||
assert manifest["trigger"] == "halt"
|
||||
assert tuple(manifest["wal"]["identity"]) == identity
|
||||
captured = (artifact / "state.db-wal").read_bytes()
|
||||
assert captured == original and captured
|
||||
assert manifest["wal"]["sha256"] == hashlib.sha256(captured).hexdigest()
|
||||
assert manifest["main"]["mode"] == "copied" and manifest["main"]["header"]["valid"]
|
||||
assert (artifact / "state.db").stat().st_size == manifest["main"]["bytes"]
|
||||
|
||||
recovered = _recover(artifact, "state.db", tmp_path / "recovered")
|
||||
try:
|
||||
assert _count(recovered, sentinel) == 1
|
||||
assert _integrity_ok(recovered)
|
||||
finally:
|
||||
recovered.close()
|
||||
|
||||
db.close()
|
||||
assert db._conn is None
|
||||
assert db._retired_generation_capture == artifact # captured once, not again at close
|
||||
|
||||
|
||||
@not_windows
|
||||
def test_close_captures_when_the_loss_is_first_seen_at_close(tmp_path, force_wal):
|
||||
path = tmp_path / "state.db"
|
||||
db = _make_db(path, "gw-0", "seed")
|
||||
_require_wal(db)
|
||||
sentinel = _wal_only_sentinel(db, "gw-0")
|
||||
_lose_sidecars(path, rename=False)
|
||||
|
||||
db.close() # no write in between: close() itself must notice the loss and capture
|
||||
|
||||
artifact = db._retired_generation_capture
|
||||
assert artifact is not None and _manifest(artifact)["trigger"] == "close"
|
||||
assert db._conn is None
|
||||
recovered = _recover(artifact, "state.db", tmp_path / "recovered")
|
||||
try:
|
||||
assert _count(recovered, sentinel) == 1
|
||||
finally:
|
||||
recovered.close()
|
||||
|
||||
|
||||
@not_windows
|
||||
def test_close_refuses_to_settle_without_a_capture(tmp_path, force_wal, monkeypatch):
|
||||
path = tmp_path / "state.db"
|
||||
db = _make_db(path, "gw-0", "seed")
|
||||
_require_wal(db)
|
||||
sentinel = _wal_only_sentinel(db, "gw-0")
|
||||
_lose_sidecars(path, rename=False)
|
||||
|
||||
def refuse(*args, **kwargs):
|
||||
raise RetiredGenerationCaptureError("no space left on device")
|
||||
|
||||
monkeypatch.setattr(hermes_state, "capture_retired_wal_generation", refuse)
|
||||
with pytest.raises(RetiredGenerationCaptureError, match="no space left"):
|
||||
db.close()
|
||||
assert db._conn is not None, "shutdown must not settle while the retired generation is uncaptured"
|
||||
assert db._retired_generation_capture is None
|
||||
|
||||
monkeypatch.undo()
|
||||
db.close()
|
||||
artifact = db._retired_generation_capture
|
||||
assert artifact is not None and db._conn is None
|
||||
recovered = _recover(artifact, "state.db", tmp_path / "recovered")
|
||||
try:
|
||||
assert _count(recovered, sentinel) == 1
|
||||
finally:
|
||||
recovered.close()
|
||||
|
||||
|
||||
@not_windows
|
||||
def test_capture_selects_the_recorded_inode_not_the_pathname(tmp_path, force_wal):
|
||||
"""A second deleted WAL under the same pathname belongs to another owner: it must be neither
|
||||
captured as ours nor touched."""
|
||||
path = tmp_path / "state.db"
|
||||
db = _make_db(path, "gw-0", "seed")
|
||||
wal = _require_wal(db)
|
||||
sentinel = _wal_only_sentinel(db, "gw-0")
|
||||
identity = db._db_sidecar_identity["-wal"]
|
||||
_lose_sidecars(path, rename=False)
|
||||
original = _descriptor_contents(_descriptor_for(identity))
|
||||
|
||||
other = sqlite3.connect(str(tmp_path / "other.db"))
|
||||
try:
|
||||
other.execute("PRAGMA journal_mode=WAL")
|
||||
other.execute("CREATE TABLE independent (content TEXT)")
|
||||
other.execute("INSERT INTO independent VALUES ('another owner committed this')")
|
||||
other.commit()
|
||||
Path(str(tmp_path / "other.db") + "-wal").rename(wal) # same pathname as our lost WAL
|
||||
other_identity = (wal.stat().st_dev, wal.stat().st_ino)
|
||||
wal.unlink()
|
||||
assert other_identity != identity
|
||||
other_before = _descriptor_contents(_descriptor_for(other_identity))
|
||||
|
||||
with pytest.raises(DeletedWalGenerationError):
|
||||
db.append_message("gw-0", role="user", content="after the loss")
|
||||
|
||||
artifact = db._retired_generation_capture
|
||||
assert (artifact / "state.db-wal").read_bytes() == original
|
||||
assert tuple(_manifest(artifact)["wal"]["identity"]) == identity
|
||||
assert _descriptor_contents(_descriptor_for(other_identity)) == other_before
|
||||
recovered = _recover(artifact, "state.db", tmp_path / "recovered")
|
||||
try:
|
||||
assert _count(recovered, sentinel) == 1
|
||||
finally:
|
||||
recovered.close()
|
||||
db.close()
|
||||
finally:
|
||||
other.close()
|
||||
|
||||
|
||||
def test_capture_refuses_to_guess_by_pathname(tmp_path, force_wal):
|
||||
path = tmp_path / "state.db"
|
||||
db = _make_db(path, "gw-0", "seed")
|
||||
try:
|
||||
with pytest.raises(RetiredGenerationCaptureError, match="pathname"):
|
||||
capture_retired_wal_generation(path, sidecar_identity={}, trigger="test")
|
||||
with pytest.raises(RetiredGenerationCaptureError, match="no longer holds"):
|
||||
capture_retired_wal_generation(path, sidecar_identity={"-wal": (1, 1)}, trigger="test")
|
||||
assert not list(tmp_path.glob("state.db.retired-wal-*"))
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
# ── Recoverable after the writer process has exited ──────────────────────────────────────────────────
|
||||
#
|
||||
# The gateway writer A runs in its OWN process, seeds + checkpoints history, leaves rows only in its WAL
|
||||
# and then serves stdin commands. The parent takes A's sidecars away, drives one refused write (halt +
|
||||
# capture), mints and checkpoints a newer generation through the path from THIS process, then has A
|
||||
# close and exit normally. The retired rows must be recoverable from the capture alone afterwards.
|
||||
|
||||
_GATEWAY_CHILD = textwrap.dedent(
|
||||
"""
|
||||
import gc, json, os, sys
|
||||
from pathlib import Path
|
||||
repo, hermes_home, db_path = sys.argv[1], sys.argv[2], sys.argv[3]
|
||||
sys.path.insert(0, repo)
|
||||
os.environ["HERMES_HOME"] = hermes_home
|
||||
import hermes_state_wal
|
||||
if hermes_state_wal.is_sqlite_wal_reset_vulnerable():
|
||||
hermes_state_wal.is_sqlite_wal_reset_vulnerable = lambda version_info=None: False
|
||||
hermes_state_wal.resolve_journal_mode = lambda: "wal"
|
||||
from hermes_state import DeletedWalGenerationError, SessionDB
|
||||
|
||||
def emit(**e):
|
||||
sys.stdout.write(json.dumps(e) + "\\n"); sys.stdout.flush()
|
||||
|
||||
db = SessionDB(db_path=Path(db_path))
|
||||
if not db._wal_active:
|
||||
emit(event="skip"); sys.exit(3)
|
||||
for sid in ("gw-0", "gw-1", "gw-2", "gw-3"):
|
||||
db.create_session(sid, "cli")
|
||||
db.append_message(sid, role="user", content="seed")
|
||||
db._try_wal_checkpoint() # history checkpointed into state.db, like a `sessions optimize` pass
|
||||
db._conn.execute("PRAGMA wal_autocheckpoint=0")
|
||||
for sid in ("gw-0", "gw-1", "gw-2", "gw-3"):
|
||||
db.append_message(sid, role="assistant", content="uncheckpointed " + "x" * 3000)
|
||||
emit(event="ready")
|
||||
|
||||
for line in sys.stdin:
|
||||
cmd = line.strip()
|
||||
if cmd == "write":
|
||||
try:
|
||||
db.append_message("gw-0", role="user", content="post-loss turn")
|
||||
emit(event="write", refused=False, artifact=None)
|
||||
except DeletedWalGenerationError:
|
||||
capture = db._retired_generation_capture
|
||||
emit(event="write", refused=True, artifact=None if capture is None else str(capture))
|
||||
elif cmd == "close":
|
||||
try:
|
||||
db.close()
|
||||
emit(event="closed", error=None)
|
||||
except Exception as exc:
|
||||
emit(event="closed", error=repr(exc))
|
||||
del db
|
||||
gc.collect()
|
||||
elif cmd == "quit":
|
||||
break
|
||||
"""
|
||||
)
|
||||
|
||||
|
||||
def _write_second_generation(db_path: Path, n_rows: int) -> int:
|
||||
"""From a process that does NOT hold this db open, mint a fresh WAL generation through the path,
|
||||
write ``n_rows`` messages and checkpoint them into the main file."""
|
||||
conn = sqlite3.connect(str(db_path), timeout=5.0, isolation_level=None)
|
||||
try:
|
||||
conn.execute("PRAGMA journal_mode")
|
||||
sessions = [r[0] for r in conn.execute("SELECT id FROM sessions ORDER BY id").fetchall()]
|
||||
conn.execute("BEGIN IMMEDIATE")
|
||||
for i in range(n_rows):
|
||||
conn.execute(
|
||||
"INSERT INTO messages (session_id, role, content, timestamp) VALUES (?, ?, ?, ?)",
|
||||
(sessions[i % len(sessions)], "assistant", "gen2 " + "y" * 2000 + f" #{i}", 1.0 + i),
|
||||
)
|
||||
conn.execute("COMMIT")
|
||||
conn.execute("PRAGMA wal_checkpoint(TRUNCATE)")
|
||||
return conn.execute("SELECT COUNT(*) FROM messages").fetchone()[0]
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def _assert_retired_rows_recoverable_after_exit(tmp_path, *, rename):
|
||||
repo_root = os.path.dirname(os.path.abspath(hermes_state.__file__))
|
||||
hermes_home = tmp_path / "home"
|
||||
hermes_home.mkdir()
|
||||
path = tmp_path / "state.db"
|
||||
stderr_path = tmp_path / "writer-stderr.log"
|
||||
with stderr_path.open("w", encoding="utf-8") as stderr:
|
||||
proc = subprocess.Popen(
|
||||
[sys.executable, "-c", _GATEWAY_CHILD, repo_root, str(hermes_home), str(path)],
|
||||
stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=stderr,
|
||||
text=True, encoding="utf-8", bufsize=1,
|
||||
env={**os.environ, "HERMES_STATE_DB_GUARD_BYPASS": "1"},
|
||||
)
|
||||
events: "queue.Queue[str | None]" = queue.Queue()
|
||||
|
||||
def read_events():
|
||||
try:
|
||||
for line in proc.stdout:
|
||||
events.put(line)
|
||||
finally:
|
||||
events.put(None)
|
||||
|
||||
threading.Thread(target=read_events, daemon=True).start()
|
||||
|
||||
def next_event(name):
|
||||
try:
|
||||
line = events.get(timeout=20)
|
||||
except queue.Empty:
|
||||
pytest.fail(f"writer timed out waiting for {name!r}\n" + stderr_path.read_text(encoding="utf-8"))
|
||||
assert line is not None, (
|
||||
f"writer exited early (rc={proc.poll()}) waiting for {name!r}\n"
|
||||
+ stderr_path.read_text(encoding="utf-8"))
|
||||
event = json.loads(line)
|
||||
if name == "ready" and event.get("event") == "skip":
|
||||
pytest.skip("WAL not active on this filesystem")
|
||||
assert event.get("event") == name, event
|
||||
return event
|
||||
|
||||
def send(command):
|
||||
proc.stdin.write(command + "\n")
|
||||
proc.stdin.flush()
|
||||
|
||||
try:
|
||||
next_event("ready")
|
||||
_lose_sidecars(path, rename=rename)
|
||||
|
||||
send("write")
|
||||
refused = next_event("write")
|
||||
assert refused["refused"] is True
|
||||
artifact = Path(refused["artifact"])
|
||||
assert artifact.is_dir() and _manifest(artifact)["trigger"] == "halt"
|
||||
|
||||
expected = _write_second_generation(path, n_rows=400)
|
||||
|
||||
send("close")
|
||||
assert next_event("closed")["error"] is None
|
||||
send("quit")
|
||||
assert proc.wait(timeout=20) == 0, stderr_path.read_text(encoding="utf-8")
|
||||
finally:
|
||||
with contextlib.suppress(BrokenPipeError):
|
||||
proc.stdin.close()
|
||||
try:
|
||||
proc.wait(timeout=5)
|
||||
except subprocess.TimeoutExpired:
|
||||
proc.kill()
|
||||
proc.wait(timeout=5)
|
||||
proc.stdout.close()
|
||||
|
||||
# The writer is gone and its unlinked WAL inode with it. Recovery must come from the capture alone.
|
||||
recovered = _recover(artifact, "state.db", tmp_path / "recovered")
|
||||
try:
|
||||
assert _count(recovered, "uncheckpointed %") == 4, "retired WAL-only rows lost across process exit"
|
||||
assert _integrity_ok(recovered), "captured image + WAL do not form a consistent database"
|
||||
finally:
|
||||
recovered.close()
|
||||
|
||||
if hasattr(sqlite3.Connection, "setconfig"):
|
||||
# With SQLITE_DBCONFIG_NO_CKPT_ON_CLOSE the closed handle also left the newer generation alone.
|
||||
live = sqlite3.connect(str(path))
|
||||
try:
|
||||
assert _integrity_ok(live) and live.execute("SELECT COUNT(*) FROM messages").fetchone()[0] == expected
|
||||
finally:
|
||||
live.close()
|
||||
|
||||
|
||||
@pytest.mark.linux_only
|
||||
def test_retired_rows_recoverable_after_process_exit(tmp_path):
|
||||
_assert_retired_rows_recoverable_after_exit(tmp_path, rename=False)
|
||||
|
||||
|
||||
@pytest.mark.macos_only
|
||||
def test_retired_rows_recoverable_after_process_exit_with_renamed_sidecars(tmp_path):
|
||||
_assert_retired_rows_recoverable_after_exit(tmp_path, rename=True)
|
||||
Reference in New Issue
Block a user