Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions agent/turn_explainers.py
Original file line number Diff line number Diff line change
Expand Up @@ -119,6 +119,21 @@
"sessions/<session_id>.jsonl and, on the gateway, "
"pending_messages/pending-*.json."
),
"deleted_wal": (
"the turn was stopped because a live Hermes process held a retired "
"state.db-wal generation after its pathname was deleted or "
"replaced. Stop the gateway, dashboard, and cron writers; "
"do not overwrite the current state.db or delete its sidecars. "
"Check the logs for whether Hermes captured the retired generation, "
"then read the adjacent state.db.retired-wal-*/manifest.json. If "
"manifest.main.mode is `copied`, inspect that artifact with `hermes "
"sessions recover --source <state.db.retired-wal-*/state.db> "
"--inspect-only` before deciding whether its committed frames belong "
"on the current database. A `header_only` artifact is forensic and "
"does not contain a copied state.db to inspect. Unwritten messages "
"were diverted to sessions/<session_id>.jsonl and, on the gateway, "
"pending_messages/pending-*.json."
),
"corrupt": (
"the turn was stopped because the state database "
"reported structural corruption (the transcript would "
Expand Down
1 change: 1 addition & 0 deletions contributors/emails/gaoanze@meituan.com
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
gaoanze888
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
nikkoxgonzales
10 changes: 10 additions & 0 deletions gateway/run_notifications.py
Original file line number Diff line number Diff line change
Expand Up @@ -823,6 +823,16 @@ async def _send_session_db_warning_notifications(self) -> None:
error = getattr(self, "_session_db_init_error", None)
if not error:
return
# Re-check the live store before warning: a startup `database is locked` routinely clears while
# the adapters are still connecting, and a borrowed store handle comes back once its owner
# releases it. The cache's opener clears ``_session_db_init_error`` on recovery, so a stale
# startup failure must not be broadcast as current (#108031).
if getattr(self, "_session_db_handle_cache", None) is not None:
self._open_session_db_for_active_scope()
error = self._session_db_init_error
if not error:
logger.info("state.db recovered before the home-channel warning went out; not broadcasting")
return
from hermes_constants import get_default_hermes_root
from hermes_state import _default_db_path, classify_persistence_error, format_session_db_unavailable
if classify_persistence_error(error) == "corrupt":
Expand Down
56 changes: 46 additions & 10 deletions hermes_state.py
Original file line number Diff line number Diff line change
Expand Up @@ -46,14 +46,16 @@
from hermes_state_schema import SessionSchemaMixin
import hermes_state_holders as _state_holders
from hermes_state_dbfile import (
_canonical_sqlite_path, _connect_tracked_db, _prepare_connection_retirement,
_canonical_sqlite_path, _connect_tracked_db, _fd_is_truly_unlinked, _prepare_connection_retirement,
_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,
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
from hermes_state_wal import (
_WAL_INCOMPAT_MARKERS, _on_disk_journal_mode, apply_database_pragmas, apply_wal_with_fallback,
)
from hermes_state_repair import _claim_repair_attempt, preflight_db_writability, repair_state_db_schema
from hermes_state_titles import SessionTitlesMixin
from hermes_state_usage import SessionUsageMixin
Expand Down Expand Up @@ -615,7 +617,12 @@ def _open_writer_conn(self) -> sqlite3.Connection:
)
try:
conn.row_factory = sqlite3.Row
self._wal_active = apply_wal_with_fallback(conn, db_label="state.db") == "wal"
mode = apply_wal_with_fallback(conn, db_label="state.db")
# "wal" is also the *assumed* mode when the on-disk probe was blocked by a concurrent opener
# (#86515): the lock-free mode=ro read pool needs a confirmed WAL header, so confirm it here.
# Unknown -> reads queue on the writer lock (slow but correct) instead of racing SQLITE_BUSY
# on a file that may really be in rollback-journal mode.
self._wal_active = mode == "wal" and _on_disk_journal_mode(conn) == "wal"
apply_database_pragmas(conn, db_label="state.db")
conn.execute("PRAGMA foreign_keys=ON")
self._fts_cjk_loaded = load_fts5_cjk_extension(conn)
Expand Down Expand Up @@ -810,10 +817,20 @@ def _execute_write(
ioerr_begin_retried = False
while True:
self._raise_if_db_corrupt()
self._raise_if_db_replaced()
# NOTE: the replaced/generation live probe runs INSIDE the lock below,
# not here. close() mutates _conn and _db_sidecar_identity under that
# same lock, ending the WAL generation (SQLite unlinks the -wal/-shm
# sidecars on a clean close). A lock-free probe that races close() can
# observe the mid-teardown state — sidecars already unlinked while
# _db_sidecar_identity is not yet cleared — and misclassify this
# process's OWN clean close as an externally deleted generation,
# raising a sticky DeletedWalGenerationError that permanently refuses
# later writes (#105567). Inside the lock the probe only ever sees the
# stable post-close state (identity cleared → adopt / reopen path).
fn_started = False
try:
with self._lock:
self._raise_if_db_replaced()
if self._conn is None: # close() raced this writer
self._reopen_after_close_locked(context="write")
self._conn.execute("BEGIN IMMEDIATE")
Expand Down Expand Up @@ -914,13 +931,30 @@ def _do(conn):

def _read_one(self, sql: str, params: Any = ()) -> Optional[sqlite3.Row]:
"""``fetchone()`` of one read-only statement via ``_read_ctx``."""
with self._read_ctx() as conn:
return conn.execute(sql, params).fetchone()
return self._read_retrying_ioerr(lambda conn: conn.execute(sql, params).fetchone())

def _read_all(self, sql: str, params: Any = ()) -> List[sqlite3.Row]:
"""``fetchall()`` of one read-only statement via ``_read_ctx``."""
with self._read_ctx() as conn:
return conn.execute(sql, params).fetchall()
return self._read_retrying_ioerr(lambda conn: conn.execute(sql, params).fetchall())

def _read_retrying_ioerr(self, fn: Callable[[sqlite3.Connection], T]) -> T:
"""Run an idempotent SELECT through ``_read_ctx``, retrying a transient SQLITE_IOERR.

A warm ``mode=ro`` pooled reader can hit the same millisecond-wide WAL transition window as a
read-only OPEN (#100436) when its statement executes or steps: a sibling process's checkpoint /
WAL reset / frame flush surfaces ``disk I/O error`` because a read-only connection cannot rewrite
the -shm index (#100871, WSL2 ext4-on-vhdx, multi-process). The window closes on its own, so
the statement is replayed on the SAME connection within the read-only IOERR budget -- never
closed and reopened (close() cancels this process's POSIX locks for every sibling connection),
never quarantined (busy is not broken). A persistent IOERR exhausts the budget and propagates."""
for attempt in range(_READ_ONLY_IOERR_RETRY_ATTEMPTS + 1):
try:
with self._read_ctx() as conn:
return fn(conn)
except sqlite3.OperationalError as exc:
if attempt >= _READ_ONLY_IOERR_RETRY_ATTEMPTS or _DISK_IO_ERROR_MARKER not in str(exc).lower():
raise
time.sleep(_READ_ONLY_IOERR_RETRY_BACKOFF_S)

def _ensure_db_file_generation(self) -> None:
"""Mint a once-per-file generation stamp (state_meta + application_id). First opener wins (INSERT
Expand Down Expand Up @@ -998,8 +1032,10 @@ def _wal_generation_was_lost(self) -> bool:
if sys.platform.startswith("linux"):
watched = _watched_sqlite_sidecar_paths(self.db_path)
try:
for target in _proc_fd_targets(os.getpid()):
if " (deleted)" in target and _canonical_sqlite_path(target) in watched:
for target, fd_path in _proc_fd_targets(os.getpid()):
canonical = _canonical_sqlite_path(target)
if (" (deleted)" in target and canonical in watched
and _fd_is_truly_unlinked(fd_path, watched[canonical])):
return True
except OSError:
return False
Expand Down
46 changes: 38 additions & 8 deletions hermes_state_dbfile.py
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@
import threading
import time
from pathlib import Path
from typing import Any, Callable, Dict, List, Optional, Set, Tuple
from typing import Any, Callable, Dict, List, Optional, Tuple

from hermes_state_common import (
FTS_REBUILD_DEFERRAL_KEY, stat_db_file_identity as _stat_db_file_identity
Expand Down Expand Up @@ -136,13 +136,40 @@ def _canonical_sqlite_path(path: str) -> str:
return os.path.normcase(os.path.abspath(path.removesuffix(" (deleted)")))


def _watched_sqlite_sidecar_paths(db_path) -> Set[str]:
def _watched_sqlite_sidecar_paths(db_path) -> Dict[str, str]:
"""Map each sidecar's canonical (/proc-comparable) form to its literal, still-named path,
so a canonical match can be re-``stat``'d for identity rather than trusted as text."""
base = os.path.abspath(os.fspath(db_path))
return {_canonical_sqlite_path(base + "-wal"), _canonical_sqlite_path(base + "-shm")}
literal = (base + "-wal", base + "-shm")
return {_canonical_sqlite_path(path): path for path in literal}


def _fd_is_truly_unlinked(fd_path: str, watched_path: str) -> bool:
"""Confirm a `` (deleted)`` /proc fd target really names an orphaned generation, not the
CURRENT watched sidecar.

The suffix alone is not proof: on OpenZFS a live, still-linked file whose dentry was
unhashed is reported as deleted while it is still the very same inode the watched path
names. Conversely ``st_nlink == 0`` is not proof of the opposite — a stale generation can
keep a surviving hard link (a backup, an operator copy) after the watched path itself is
removed or replaced, leaving ``st_nlink >= 1`` on a truly orphaned inode. So compare
identity, not link count: only an exact ``(st_dev, st_ino)`` match between the fd and the
CURRENT watched path proves they are the same live file. A mismatch, or a watched path
that cannot be stat'd at all, means the fd holds a generation the watched path no longer
names — the guard keeps failing closed."""
try:
fd_stat = os.stat(fd_path)
except OSError:
return True
try:
watched_stat = os.stat(watched_path)
except OSError:
return True
return (fd_stat.st_dev, fd_stat.st_ino) != (watched_stat.st_dev, watched_stat.st_ino)


def _iter_proc_fd_targets():
"""Yield ``(pid, readlink target)`` for every readable ``/proc/<pid>/fd`` entry."""
"""Yield ``(pid, readlink target, fd path)`` for every readable ``/proc/<pid>/fd`` entry."""
for pid_str in os.listdir("/proc"):
if not pid_str.isdigit():
continue
Expand All @@ -153,7 +180,8 @@ def _iter_proc_fd_targets():
continue # process gone or not ours
for fd in fds:
with contextlib.suppress(OSError):
yield int(pid_str), os.readlink(f"{fd_dir}/{fd}")
fd_path = f"{fd_dir}/{fd}"
yield int(pid_str), os.readlink(fd_path), fd_path


def iter_deleted_sqlite_sidecar_holders(db_path) -> List[Tuple[int, str]]:
Expand All @@ -166,8 +194,10 @@ def iter_deleted_sqlite_sidecar_holders(db_path) -> List[Tuple[int, str]]:
holders: List[Tuple[int, str]] = []
watched = _watched_sqlite_sidecar_paths(db_path)
try:
for pid, target in _iter_proc_fd_targets():
if " (deleted)" in target and _canonical_sqlite_path(target) in watched:
for pid, target, fd_path in _iter_proc_fd_targets():
canonical = _canonical_sqlite_path(target)
if (" (deleted)" in target and canonical in watched
and _fd_is_truly_unlinked(fd_path, watched[canonical])):
holders.append((pid, target))
except Exception as exc:
logger.debug("deleted-WAL holder scan failed for %s: %s", db_path, exc)
Expand Down Expand Up @@ -614,7 +644,7 @@ def count_db_holders(db_path: Path) -> Optional[int]:
if not sys.platform.startswith("linux"):
return None
target = os.path.realpath(str(db_path))
return len({pid for pid, link in _iter_proc_fd_targets() if link == target})
return len({pid for pid, link, _fd_path in _iter_proc_fd_targets() if link == target})
except Exception:
return None

Expand Down
14 changes: 10 additions & 4 deletions hermes_state_errors.py
Original file line number Diff line number Diff line change
Expand Up @@ -74,8 +74,8 @@ def is_disk_full_error(exc: BaseException | str | None) -> bool:

# Every classify_persistence_error bucket; consumers enumerate this tuple.
PERSISTENCE_ERROR_CAUSES = (
"locked", "compression", "compression_closed", "turn_lease", "corrupt", "replaced", "disk",
"unknown",
"locked", "compression", "compression_closed", "turn_lease", "corrupt", "replaced",
"deleted_wal", "disk", "unknown",
)


Expand Down Expand Up @@ -182,14 +182,19 @@ class StateDbCorruptError(sqlite3.DatabaseError):
(SessionTurnLeaseLostError, "turn_lease"),
(CompressionSessionClosedError, "compression_closed"),
(CompressionSessionBusyError, "compression"),
# The WAL-generation error subclasses StateDbReplacedError so existing write diversion keeps
# working; classify it first because its recovery artifact and operator action are different.
(DeletedWalGenerationError, "deleted_wal"),
(StateDbReplacedError, "replaced"),
(StateDbCorruptError, "corrupt"),
)
_PERSISTENCE_CAUSE_BY_PHRASE = (
(("turn lease",), "turn_lease"),
(("closed by compression",), "compression_closed"),
(("being compressed", "compression lease"), "compression"),
(("was replaced underneath", "deleted state.db-wal", "deleted state.db-shm"), "replaced"),
# RPC-wrapped errors lose their exception type; retain the same sidecar/main-file split.
(("deleted state.db-wal", "deleted state.db-shm"), "deleted_wal"),
(("was replaced underneath",), "replaced"),
(_DB_CORRUPTION_MARKERS, "corrupt"),
(("locked", "busy"), "locked"),
)
Expand All @@ -200,7 +205,8 @@ def classify_persistence_error(exc_or_str) -> str:
matches: "locked" = busy, retry; "disk" = full/read-only/permissions;
"compression" = a live lease refused the write; "compression_closed" = adopt
the rotated session id; "turn_lease" = fencing, not storage; "corrupt" =
file damage (repair path, not disk space); "replaced" = stop writing."""
file damage (repair path, not disk space); "replaced" = main-file replacement;
"deleted_wal" = a retired sidecar generation requiring capture inspection."""
if exc_or_str is None:
return "unknown"
# Lease refusals contain neither "locked" nor "busy": match by type first,
Expand Down
9 changes: 5 additions & 4 deletions hermes_state_readpool.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,13 +69,14 @@
_fd_usage_cache: "tuple[float, Optional[int]]" = (0.0, None)


def _proc_fd_targets(pid: int) -> Iterator[str]:
"""readlink() of every entry in /proc/<pid>/fd (unreadable links skipped).
Raises OSError when the fd directory itself cannot be listed."""
def _proc_fd_targets(pid: int) -> "Iterator[tuple[str, str]]":
"""Yield ``(readlink target, fd path)`` for every entry in /proc/<pid>/fd (unreadable
links skipped). Raises OSError when the fd directory itself cannot be listed."""
fd_dir = f"/proc/{pid}/fd"
for fd in os.listdir(fd_dir):
fd_path = f"{fd_dir}/{fd}"
try:
yield os.readlink(f"{fd_dir}/{fd}")
yield os.readlink(fd_path), fd_path
except OSError:
continue

Expand Down
27 changes: 26 additions & 1 deletion hermes_state_wal.py
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,10 @@
_wal_fallback_warned_paths: set[str] = set()
_wal_fallback_warned_lock = threading.Lock()

# Dedup for the both-pragmas-failed WARNING (WAL *and* DELETE rejected, e.g. APFS external SSDs under contention).
_wal_delete_fallback_failed_paths: set[str] = set()
_wal_delete_fallback_failed_lock = threading.Lock()

# Dedup for the probe-unknown WARNING (on-disk journal mode unreadable, nothing touched).
_wal_probe_unknown_paths: set[str] = set()
_wal_probe_unknown_lock = threading.Lock()
Expand Down Expand Up @@ -291,7 +295,20 @@ def _wal_activated() -> str:
if require_wal:
raise WalUnsupportedError(str(exc)) from exc
_log_wal_fallback_once(db_label, exc)
_set_journal_mode_no_wait(conn, "DELETE")
try:
_set_journal_mode_no_wait(conn, "DELETE")
except sqlite3.OperationalError as delete_exc:
# Filesystems that reject BOTH pragmas (APFS external SSDs under heavy contention raise
# "disk I/O error" on DELETE too). The connection keeps its current mode — typically the
# SQLite default DELETE — and stays usable for reads/writes, so propagating would crash
# every DB-init caller (SessionDB, kanban_db, ResponseStore) for no benefit. Log once
# per db_label and return the mode actually in effect, read back rather than guessed.
_log_once("wal_delete_fallback_failed", db_label, exc, delete_exc)
try:
read_back = _mode_from_row(conn.execute("PRAGMA journal_mode").fetchone())
except sqlite3.OperationalError:
read_back = ""
return read_back or "delete"
return "delete"


Expand Down Expand Up @@ -398,6 +415,14 @@ def _wal_reset_repair_hint() -> str:
"%s: WAL journal_mode unsupported on this filesystem (%s) — falling back to journal_mode=DELETE (slower "
"rollback-journal mode; reduces concurrency but works on NFS/SMB/FUSE/ZFS). See "
"https://www.sqlite.org/wal.html for details. This message fires once per process per database."),
"wal_delete_fallback_failed": (_wal_delete_fallback_failed_lock, "_wal_delete_fallback_failed_paths", logging.WARNING,
# Both pragmas rejected (observed on APFS external SSDs under heavy contention): the connection keeps
# whatever mode it has (typically the SQLite default DELETE) and stays usable for reads/writes, so this
# is WARNING, not ERROR — propagating would crash every DB-init caller. Deduped per db_label because
# kanban_db.connect() runs on every kanban operation.
"%s: both WAL and DELETE journal_mode failed (WAL: %s; DELETE: %s) — continuing with the connection's "
"current journal mode (reads/writes still work). Typically seen on APFS external SSDs under heavy "
"contention. This message fires once per process per database."),
"delete_overridden": (_delete_overridden_warned_lock, "_delete_overridden_warned_paths", logging.ERROR,
# Never-live-downgrade keeps WAL; without this the operator never learns their delete had no effect.
"%s: database.journal_mode=delete is configured but the on-disk database is already WAL; keeping WAL (a live "
Expand Down
Loading
Loading