From 857b43ef9299a4f2555e32b7ddfea0182d15937f Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 18:14:53 -0700 Subject: [PATCH] refactor(hermes_bootstrap,hermes_startup_watchdog,hermes_logging): collapse duplicated dump/handle plumbing, compact docs --- hermes_bootstrap.py | 206 +++++----------------- hermes_logging.py | 271 +++++++++++----------------- hermes_startup_watchdog.py | 352 ++++++++++++++----------------------- 3 files changed, 272 insertions(+), 557 deletions(-) diff --git a/hermes_bootstrap.py b/hermes_bootstrap.py index c0622bb0d9..38913cf241 100644 --- a/hermes_bootstrap.py +++ b/hermes_bootstrap.py @@ -1,50 +1,12 @@ -"""Windows UTF-8 bootstrap for Hermes entry points. +"""Windows UTF-8 bootstrap for Hermes entry points (no-op on POSIX). -Python on Windows has two long-standing text-encoding footguns: - -1. ``sys.stdout`` / ``sys.stderr`` are bound to the console code page - (``cp1252`` on US-locale installs), so ``print("café")`` crashes with - ``UnicodeEncodeError: 'charmap' codec can't encode character``. - -2. Child processes spawned via ``subprocess`` don't know to use UTF-8 - unless ``PYTHONUTF8`` and/or ``PYTHONIOENCODING`` are set in their - environment — so any Python subprocess (the execute_code sandbox, - delegation children, linter subprocesses, etc.) inherits the same - cp1252 defaults and hits the same UnicodeEncodeError. - -This module fixes both on Windows *only* — POSIX is untouched. It -should be imported at the very top of every Hermes entry point -(``hermes``, ``hermes-agent``, ``hermes-acp``, ``python -m gateway.run``, -``batch_runner.py``, ``cron/scheduler.py``) before any other imports -that might do file I/O or print to stdout. - -What this module does on Windows: - - - Sets ``os.environ["PYTHONUTF8"] = "1"`` (PEP 540 UTF-8 mode) so - every child process we spawn uses UTF-8 for ``open()`` and stdio. - - Sets ``os.environ["PYTHONIOENCODING"] = "utf-8"`` for belt-and- - suspenders — some tools read this instead of / in addition to - ``PYTHONUTF8``. - - Reconfigures ``sys.stdout`` / ``sys.stderr`` to UTF-8 in the current - process, using the ``reconfigure()`` API (Python 3.7+). This fixes - ``print("café")`` in the parent without a re-exec. - -What this module does NOT do: - - - It does not re-exec Python with ``-X utf8``, so ``open()`` calls in - the *current* process still default to locale encoding. Those need - an explicit ``encoding="utf-8"`` at the call site (lint rule - ``PLW1514`` / ``PYI058``). Ruff is the right tool for that sweep. - -What this module does on POSIX: - - - Nothing. POSIX systems are already UTF-8 by default in 99% of cases, - and we don't want to touch ``LANG``/``LC_*`` behavior that users may - have configured intentionally. If someone hits a C/POSIX locale on - Linux, they can export ``PYTHONUTF8=1`` themselves — we won't override. - -Idempotent: safe to call multiple times. ``_bootstrap_once`` guards -against double-reconfigure. +Windows binds stdio to the console code page (cp1252), so ``print("café")`` raises +``UnicodeEncodeError``, and Python children inherit the same default unless +``PYTHONUTF8``/``PYTHONIOENCODING`` are set. Import this module first in every entry +point (``hermes``, ``hermes-agent``, ``hermes-acp``, ``gateway.run``, ``batch_runner``, +``cron/scheduler``). It does NOT re-exec with ``-X utf8``: ``open()`` in the current +process still needs an explicit ``encoding="utf-8"`` (ruff ``PLW1514``). POSIX is left +alone deliberately — users' ``LANG``/``LC_*`` choices are respected. """ from __future__ import annotations @@ -57,97 +19,45 @@ _bootstrap_applied = False def apply_windows_utf8_bootstrap() -> bool: - """Apply the Windows UTF-8 bootstrap if we're on Windows. - - Returns True if bootstrap was applied (i.e. we're on Windows and - haven't already done this), False otherwise. The return value is - advisory — callers normally don't need it, but tests may want to - assert the path was taken. - - Idempotent: subsequent calls after the first are a no-op. - """ + """Apply the Windows UTF-8 bootstrap once; True only when it was applied this call.""" global _bootstrap_applied - if not _IS_WINDOWS: - return False - if _bootstrap_applied: + if not _IS_WINDOWS or _bootstrap_applied: return False - # 1. Child processes inherit these and run in UTF-8 mode. - # We use setdefault() rather than overwriting so the user can - # explicitly opt out by setting PYTHONUTF8=0 in their environment - # (or PYTHONIOENCODING=something-else) if they really want to. + # setdefault() so a user can opt out with PYTHONUTF8=0 / PYTHONIOENCODING=... os.environ.setdefault("PYTHONUTF8", "1") os.environ.setdefault("PYTHONIOENCODING", "utf-8") - # 2. Reconfigure the current process's stdio to UTF-8. Needed - # because os.environ changes don't retroactively rebind sys.stdout - # — those were bound at interpreter startup based on the console - # code page. ``reconfigure`` is a TextIOWrapper method since 3.7. - # - # errors="replace" means that if we ever *read* something from - # stdin that isn't UTF-8 (unlikely but possible with piped input - # from legacy tools), we'll get U+FFFD replacement chars rather - # than a crash. Output is pure UTF-8. - for stream_name in ("stdout", "stderr"): - stream = getattr(sys, stream_name, None) - if stream is None: - continue - reconfigure = getattr(stream, "reconfigure", None) + # os.environ changes don't rebind streams bound at interpreter startup, so + # reconfigure them in-process. errors="replace" keeps a non-UTF-8 legacy + # pipe on stdin from crashing us (U+FFFD instead of an exception). + # Non-TextIOWrapper streams (BytesIO in tests, embedded hosts) have no + # reconfigure(): skip — the env-var fix for children is the bigger win. + for stream_name in ("stdout", "stderr", "stdin"): + reconfigure = getattr(getattr(sys, stream_name, None), "reconfigure", None) if reconfigure is None: - # Not a TextIOWrapper (could be redirected to a BytesIO in - # tests, or a non-standard stream in some embedded cases). - # Skip silently — the env-var fix is still in effect for - # child processes, which is the bigger win. continue try: reconfigure(encoding="utf-8", errors="replace") except (OSError, ValueError): - # Already closed, or someone replaced it with something - # non-reconfigurable. Non-fatal. - pass - - # stdin is reconfigured separately with errors="replace" too — input - # from a legacy pipe shouldn't crash the process. - stdin = getattr(sys, "stdin", None) - if stdin is not None: - reconfigure = getattr(stdin, "reconfigure", None) - if reconfigure is not None: - try: - reconfigure(encoding="utf-8", errors="replace") - except (OSError, ValueError): - pass + pass # closed, or replaced with something non-reconfigurable _bootstrap_applied = True return True def suppress_platform_ver_console() -> None: - """Stub ``platform._syscmd_ver`` on Windows — decode-crash + flash guard. + """Stub ``platform._syscmd_ver`` on Windows — decode-crash + console-flash guard. - CPython's ``platform.win32_ver()`` (reached via ``platform.uname()`` / - ``platform.platform()``, which the OpenAI SDK touches for its - platform headers) shells out ``cmd /c ver``. Two failure modes: - - - **Console flash**: the ``check_output(..., shell=True)`` call has no - ``CREATE_NO_WINDOW``, so a windowless parent (pythonw gateway, slash - workers, kanban workers) flashes a visible console per call. - - **UnicodeDecodeError on Python 3.11.0/3.11.1**: those micros lack - CPython's ``encoding="locale"`` fix (added 3.11.2), so under PEP 540 - UTF-8 mode (which we enable above) the ``ver`` output — OEM code page - bytes on localized Windows — is strict-utf-8 decoded and raises, - crashing ``platform.platform()`` in any process that inherits - ``PYTHONUTF8=1`` (issue #69413). - - Stubbing ``_syscmd_ver`` to return its inputs makes ``win32_ver()`` hit - its documented fallback and read the version from - ``sys.getwindowsversion()`` — same data, in-process, no subprocess. - Mirrors ``hermes_cli._subprocess_compat.suppress_platform_ver_console`` - (kept there for callers that don't import bootstrap); double - application is harmless. Lives here so EVERY entry point gets it — - ``tui_gateway/slash_worker.py``, ``tui_gateway/entry.py``, - ``run_agent.py``, ``batch_runner.py``, and ``cli.py`` import only - ``hermes_bootstrap``, never ``hermes_cli.main``. + ``platform.win32_ver()`` (reached via ``platform.platform()``, which the OpenAI SDK + calls) shells out ``cmd /c ver`` with ``shell=True`` and no ``CREATE_NO_WINDOW``: a + windowless parent (pythonw gateway, slash/kanban workers) flashes a console per call, + and Python 3.11.0/3.11.1 (no ``encoding="locale"`` fix) strict-utf-8-decodes the OEM + code page output under PEP 540 mode and raises (#69413). Returning the inputs makes + ``win32_ver()`` fall back to ``sys.getwindowsversion()`` — same data, no subprocess. + Mirrors ``hermes_cli._subprocess_compat.suppress_platform_ver_console`` for callers + that never import ``hermes_cli.main``; double application is harmless. """ if not _IS_WINDOWS: return @@ -161,34 +71,19 @@ def suppress_platform_ver_console() -> None: platform._syscmd_ver = _quiet_syscmd_ver except Exception: - # Hardening only — never let it break an entry point. - pass + pass # hardening only — never break an entry point def harden_import_path(src_root: str | None = None) -> None: """Stop a package in the current directory from shadowing Hermes modules. - Hermes ships top-level modules with common names (``utils``, ``proxy``, - ``ui``). Python always seeds ``sys.path`` with the current directory, so - launching an entry point from a project that has its own ``utils/`` package - makes ``from utils import ...`` resolve to the *user's* package and crash - with an ImportError before the gateway can even start. - - The current directory reaches ``sys.path`` two ways, and a complete guard - has to handle both: - - - As the empty string ``""`` (or ``"."``) that Python inserts at - ``sys.path[0]`` for ``-m`` / script launches. - - As its own *absolute* path, when a venv activation or a project that - adds itself to ``PYTHONPATH`` puts the directory there explicitly. - - We drop the relative forms outright, then force the real Hermes source root - to the front — relocating it ahead of any absolute cwd entry rather than - only inserting when absent, so an absolute cwd path can't keep winning. - - ``src_root`` defaults to the directory this module lives in, which is the - repository root for every shipped entry point, so the guard is - self-sufficient and does not depend on the spawner exporting an env var. + Hermes ships top-level modules with common names (``utils``, ``proxy``, ``ui``); a + project with its own ``utils/`` launched from its directory would win the import. + The cwd reaches ``sys.path`` as ``""``/``"."`` (script/``-m`` launches) AND as an + absolute path (venv activation, PYTHONPATH), so both are handled: relative forms are + dropped and the Hermes root is *relocated* to the front, not merely inserted when + absent. ``src_root`` defaults to this module's directory (the repo root for every + shipped entry point), so no spawner env var is required. """ root = src_root or os.environ.get("HERMES_PYTHON_SRC_ROOT") or os.path.dirname( os.path.abspath(__file__) @@ -202,18 +97,12 @@ def harden_import_path(src_root: str | None = None) -> None: def activate_durable_lazy_target() -> None: - """Put the durable lazy-install dir on ``sys.path`` if one is configured. + """Put the durable lazy-install dir (``HERMES_LAZY_INSTALL_TARGET``) on ``sys.path``. - On immutable Docker images the agent venv is sealed and lazy installs - are redirected to a writable dir on the data volume - (``HERMES_LAZY_INSTALL_TARGET``, e.g. ``/opt/data/lazy-packages``). - Packages installed there on a previous run must be importable on this - run, so we activate the dir here — at the very first import, before any - backend module imports its SDK. - - The activation appends to the END of ``sys.path`` so the core venv - always wins name collisions (see ``tools.lazy_deps`` for the full - security rationale). Never raises; a missing/empty target is a no-op. + Immutable Docker images seal the venv and redirect lazy installs to the data volume; + packages installed there on a previous run must be importable before any backend + imports its SDK. Appends to the END of ``sys.path`` so the core venv always wins name + collisions (see ``tools.lazy_deps``). Never raises; unset target is a no-op. """ if not os.environ.get("HERMES_LAZY_INSTALL_TARGET", "").strip(): return @@ -221,19 +110,10 @@ def activate_durable_lazy_target() -> None: from tools import lazy_deps lazy_deps.activate_durable_lazy_target() except Exception: - # Bootstrap must never crash an entry point. If activation fails the - # backend simply reports itself unavailable, exactly as before. - pass + pass # a failed activation just leaves the backend reporting itself unavailable -# Apply on import — entry points just need ``import hermes_bootstrap`` -# (or ``from hermes_bootstrap import apply_windows_utf8_bootstrap``) at -# the very top of their module, before importing anything else. The -# import side effect does the right thing. +# Apply on import — entry points only need ``import hermes_bootstrap`` first. apply_windows_utf8_bootstrap() suppress_platform_ver_console() - -# Activate the durable lazy-install target (immutable Docker images) so -# packages installed into the data volume on a previous run are importable -# this run, before any backend module imports its SDK. No-op when unset. activate_durable_lazy_target() diff --git a/hermes_logging.py b/hermes_logging.py index 2f2b91b65d..2dfa09ad07 100644 --- a/hermes_logging.py +++ b/hermes_logging.py @@ -1,12 +1,9 @@ """Centralized logging setup for Hermes Agent. -Log files produced: agent.log — INFO+, all agent/tool/session activity (the main log) errors.log — -WARNING+, errors and warnings only (quick triage) gateway.log — INFO+, gateway-only events (created -when mode="gateway") gui.log — INFO+, dashboard/websocket/TUI-gateway events (created when -mode="gui") - -All files use ``RotatingFileHandler`` with ``RedactingFormatter`` so secrets are never written to -disk. +Log files: agent.log (INFO+, everything), errors.log (WARNING+), gateway.log (INFO+, +gateway components; ``mode="gateway"``), gui.log (INFO+, dashboard/TUI-gateway; +``mode="gui"``). All are rotating files driven through one async queue and formatted +with ``RedactingFormatter`` so secrets never reach disk. """ import atexit @@ -21,28 +18,15 @@ from logging.handlers import QueueHandler, QueueListener from pathlib import Path from typing import Optional, Sequence -# On Windows, stdlib ``RotatingFileHandler`` calls ``os.rename()`` in -# ``doRollover()`` and fails with ``PermissionError [WinError 32]`` whenever -# another process holds an append-mode handle on ``agent.log`` — which is -# essentially always in Hermes (TUI, gateway, ``hy_memory`` server, MCP -# servers, and on-demand CLI commands all log from separate processes), -# pinning ``agent.log`` at the 5 MiB threshold and spamming stderr with -# a traceback on every emit. ``concurrent-log-handler`` wraps the rename in a -# cross-process file lock (via ``portalocker``: pywin32 on Windows) so only -# one process rotates at a time and the others wait their turn. -# -# This swap is Windows-ONLY and deliberately so: -# * The bug (WinError 32 on rename-while-open) is specific to Windows file -# locking semantics — POSIX renames an open file fine, so stdlib already -# works correctly on Linux/macOS. -# * On POSIX, managed-mode (NixOS) relies on the exact ``_open()`` / -# ``doRollover()`` lifecycle of stdlib ``RotatingFileHandler`` (the -# ``_ManagedRotatingFileHandler`` subclass chmods 0660 after each). CLH -# opens lazily and rotates differently, which breaks the group-writable -# guarantee and the eager file-creation those paths depend on. -# Aliasing keeps every existing ``RotatingFileHandler`` reference in this -# module (class declaration, ``isinstance`` checks, docstring) working -# unchanged. See #44873. +# Windows-ONLY swap (#44873): stdlib ``RotatingFileHandler.doRollover()`` calls +# ``os.rename()``, which fails with ``PermissionError [WinError 32]`` whenever +# another process holds an append handle on ``agent.log`` — essentially always +# in Hermes (TUI, gateway, hy_memory, MCP servers, CLI commands all log) — +# pinning the file at the size threshold and spamming stderr on every emit. +# ``concurrent-log-handler`` serializes rollover with a cross-process lock. +# POSIX keeps stdlib: renames of open files work, and managed mode (NixOS) +# relies on stdlib's exact ``_open()``/``doRollover()`` lifecycle for the +# 0660 chmod and eager file creation; CLH opens lazily and rotates differently. if sys.platform == "win32": from concurrent_log_handler import ( # noqa: E402 ConcurrentRotatingFileHandler as RotatingFileHandler, @@ -53,17 +37,13 @@ else: from hermes_constants import get_config_path, get_hermes_home, mkdir_under_hermes_home -# Sentinel to track whether setup_logging() has already run. The function -# is idempotent — calling it twice is safe but the second call is a no-op -# unless ``force=True``. +# setup_logging() is idempotent: a second call is a no-op unless ``force=True``. _logging_initialized = False -# Thread-local storage for per-conversation session context. +# Thread-local per-conversation session context. _session_context = threading.local() -# Default log format — includes timestamp, level, optional session tag, -# logger name, and message. The ``%(session_tag)s`` field is guaranteed to -# exist on every LogRecord via _install_session_record_factory() below. +# ``%(session_tag)s`` exists on every LogRecord via _install_session_record_factory(). _LOG_FORMAT = "%(asctime)s %(levelname)s%(session_tag)s %(name)s: %(message)s" _LOG_FORMAT_VERBOSE = "%(asctime)s - %(name)s - %(levelname)s%(session_tag)s - %(message)s" @@ -71,17 +51,16 @@ _LOG_FORMAT_VERBOSE = "%(asctime)s - %(name)s - %(levelname)s%(session_tag)s - % def _safe_stderr(): # type: ignore[return] """Return a stderr stream that tolerates Unicode on all platforms. - We wrap ``sys.stderr`` in a ``TextIOWrapper`` with ``errors='replace'`` so log lines are never - lost — un-encodable characters are replaced with ``?`` instead of crashing the process. + Wraps ``sys.stderr`` with ``errors='replace'`` so un-encodable characters become + ``?`` instead of crashing the process. """ stream = sys.stderr encoding = getattr(stream, "encoding", None) or "utf-8" - # Already UTF-8 or surrogate-aware — no wrapping needed. if encoding.lower().replace("-", "") in ("utf8", "utf8surrogateescape"): return stream try: wrapped = io.TextIOWrapper(stream.buffer, encoding="utf-8", errors="replace", line_buffering=True) - # Prevent the wrapper from closing the underlying buffer when it is garbage-collected. + # Prevent the wrapper from closing the underlying buffer when garbage-collected. wrapped.close = lambda: None # type: ignore[assignment] return wrapped except Exception: @@ -89,12 +68,11 @@ def _safe_stderr(): # type: ignore[return] def _is_windows_concurrent_log_lock_timeout(exc: BaseException | None) -> bool: - """Return True for concurrent-log-handler's Windows lock timeout. + """True for concurrent-log-handler's Windows lock timeout. - On Windows Desktop, slash-command workers and the gateway can all write to the same rotating log - files. ``concurrent-log-handler`` serializes rollover with a cross-process lock, but when - another process holds that lock too long it raises this RuntimeError. Logging failures should - not escape into Desktop chat output. + Slash-command workers and the gateway share rotating files on Windows Desktop; + when another process holds the rollover lock too long CLH raises this + RuntimeError, which must not escape into Desktop chat output. """ return ( sys.platform == "win32" @@ -117,8 +95,6 @@ def _quiet_noisy_loggers() -> None: logging.getLogger(name).setLevel(logging.WARNING) -# Public session context API - def set_session_context(session_id: str) -> None: """Set the session ID for the current thread.""" _session_context.session_id = session_id @@ -129,27 +105,24 @@ def clear_session_context() -> None: _session_context.session_id = None -# Record factory — injects session_tag into every LogRecord at creation - def _install_session_record_factory() -> None: """Replace the global LogRecord factory with one that adds ``session_tag``. - Unlike a handler/logger ``Filter``, the record factory runs for EVERY record in the process, - including propagated and third-party-handled ones, so ``%(session_tag)s`` is always available - and never KeyErrors. Idempotent: a marker attribute prevents double-wrapping on reload. + Unlike a Filter, the record factory runs for EVERY record in the process (propagated + and third-party-handled ones included), so ``%(session_tag)s`` never KeyErrors. + Idempotent via a marker attribute. """ current_factory = logging.getLogRecordFactory() if getattr(current_factory, "_hermes_session_injector", False): - return # already installed + return def _session_record_factory(*args, **kwargs): record = current_factory(*args, **kwargs) sid = getattr(_session_context, "session_id", None) record.session_tag = f" [{sid}]" if sid else "" # type: ignore[attr-defined] - # QueueListener formats records on its own thread, after the - # profile-scoped ContextVar has gone out of scope. Keep the resolved - # home on the record so a multiplex desktop ticker can route the log - # to the job owner's files (#97489). + # QueueListener formats on its own thread, after the profile-scoped + # ContextVar is gone; keep the resolved home on the record so a + # multiplex desktop ticker can route to the job owner's files (#97489). try: record.hermes_home = str(get_hermes_home().resolve()) # type: ignore[attr-defined] except Exception: @@ -160,13 +133,10 @@ def _install_session_record_factory() -> None: logging.setLogRecordFactory(_session_record_factory) -# Install immediately on import — session_tag is available on all records -# from this point forward, even before setup_logging() is called. +# Install on import so session_tag exists on all records even before setup_logging(). _install_session_record_factory() -# Filters - class _ComponentFilter(logging.Filter): """Only pass records whose logger name starts with one of *prefixes*.""" @@ -178,13 +148,11 @@ class _ComponentFilter(logging.Filter): return record.name.startswith(self._prefixes) -# Logger name prefixes that belong to each component. -# Used by _ComponentFilter and exposed for ``hermes logs --component``. +# Logger name prefixes per component; used by _ComponentFilter and ``hermes logs --component``. COMPONENT_PREFIXES = { - # ``plugins.platforms`` covers messaging-platform adapters that migrated - # out of ``gateway/platforms/`` into bundled plugins (#41112) — they are - # still gateway components and their logs belong in gateway.log / match - # ``hermes logs --component gateway``. + # ``plugins.platforms``: messaging adapters that migrated out of + # ``gateway/platforms/`` into bundled plugins (#41112) are still gateway + # components and belong in gateway.log. "gateway": ("gateway", "hermes_plugins", "plugins.platforms"), "agent": ("agent", "run_agent", "model_tools", "batch_runner"), "tools": ("tools",), @@ -194,8 +162,6 @@ COMPONENT_PREFIXES = { } -# Main setup - def setup_logging( *, hermes_home: Optional[Path] = None, @@ -205,18 +171,16 @@ def setup_logging( mode: Optional[str] = None, force: bool = False, ) -> Path: - """Configure the Hermes logging subsystem. + """Configure the Hermes logging subsystem; returns the ``logs/`` directory. - Safe to call multiple times; the second call is a no-op unless *force* is ``True``. Level and - rotation defaults come from config.yaml ``logging.*``. ``mode="gateway"`` adds ``gateway.log`` - (gateway components only) and ``mode="gui"`` adds ``gui.log`` (dashboard / TUI-gateway). - Returns the ``logs/`` directory. + Safe to call multiple times; the second call is a no-op unless *force*. Level and + rotation defaults come from config.yaml ``logging.*``. ``mode="gateway"`` adds + ``gateway.log`` and ``mode="gui"`` adds ``gui.log``. """ global _logging_initialized home = hermes_home or get_hermes_home() log_dir = mkdir_under_hermes_home(home / "logs") - # Read config defaults (best-effort — config may not be loaded yet). cfg_level, cfg_max_size, cfg_backup = _read_logging_config() level_name = (log_level or cfg_level or "INFO").upper() @@ -224,13 +188,12 @@ def setup_logging( max_bytes = (max_size_mb or cfg_max_size or 5) * 1024 * 1024 backups = backup_count or cfg_backup or 3 - # Lazy import to avoid circular dependency at module load time. - from agent.redact import RedactingFormatter + from agent.redact import RedactingFormatter # lazy: circular at module load root = logging.getLogger() - # (filename, level, max_bytes, backup_count, component) — component gates on ``mode`` and - # restricts the file to that component's logger prefixes. + # (filename, level, max_bytes, backup_count, component) — a component gates + # the file on ``mode`` and restricts it to that component's logger prefixes. handler_specs = ( ("agent.log", level, max_bytes, backups, None), ("errors.log", logging.WARNING, 2 * 1024 * 1024, 2, None), @@ -249,7 +212,7 @@ def setup_logging( if _logging_initialized and not force: return log_dir - # Ensure root logger level is low enough for the handlers to fire. + # Root level must be low enough for the handlers to fire. if root.level == logging.NOTSET or root.level > level: root.setLevel(level) @@ -265,7 +228,6 @@ def setup_verbose_logging() -> None: root = logging.getLogger() - # Avoid adding duplicate stream handlers. if any(getattr(h, "_hermes_verbose", False) for h in root.handlers): return @@ -275,7 +237,6 @@ def setup_verbose_logging() -> None: handler._hermes_verbose = True # type: ignore[attr-defined] root.addHandler(handler) - # Lower root logger level so DEBUG records reach all handlers. if root.level > logging.DEBUG: root.setLevel(logging.DEBUG) @@ -284,8 +245,6 @@ def setup_verbose_logging() -> None: logging.getLogger("rex-deploy").setLevel(logging.INFO) -# Internal helpers - def _quietly(fn) -> None: """Call *fn* (a ``close``/``stop`` bound method) swallowing errors — teardown must never raise.""" try: @@ -295,21 +254,19 @@ def _quietly(fn) -> None: class _ManagedRotatingFileHandler(RotatingFileHandler): - """RotatingFileHandler that ensures group-writable perms in managed mode + """RotatingFileHandler with managed-mode perms and external-rotation detection. - In managed mode (NixOS) the setgid stateDir needs group-readable files, but ``_open()`` and - ``doRollover()`` honor the umask (0644), so ``chmod 0660`` is applied after both. Also, a - rotating handler holds an fd: if the file is rotated externally (logrotate, ``mv``) writes - silently go to the old inode, so before each emit the path's inode is compared to the open - stream's and the file reopened on mismatch (the ``WatchedFileHandler`` pattern). + In managed mode (NixOS) the setgid stateDir needs group-readable files, but + ``_open()``/``doRollover()`` honor the umask (0644), so ``chmod 0660`` follows both. + A rotating handler also holds an fd: if the file is rotated externally (logrotate, + ``mv``) writes silently go to the old inode, so each emit compares the path's inode + to the open stream's and reopens on mismatch (the ``WatchedFileHandler`` pattern). """ def __init__(self, *args, **kwargs): from hermes_cli.config import is_managed self._managed = is_managed() super().__init__(*args, **kwargs) - # Snapshot the inode of the currently open stream so emit() can - # detect external rotation without an extra fstat per write. self._record_stream_stat() def _chmod_if_managed(self): @@ -320,7 +277,7 @@ class _ManagedRotatingFileHandler(RotatingFileHandler): pass def _record_stream_stat(self, st: Optional[os.stat_result] = None) -> None: - """Snapshot dev/ino of ``baseFilename`` so we can detect external rotation.""" + """Snapshot dev/ino of ``baseFilename`` so emit() can detect external rotation.""" try: st = st or os.stat(self.baseFilename) self._stat_dev, self._stat_ino = st.st_dev, st.st_ino @@ -328,10 +285,10 @@ class _ManagedRotatingFileHandler(RotatingFileHandler): self._stat_dev, self._stat_ino = None, None def _reopen_stream(self, stat_result=None) -> None: - """Close the current stream and open ``baseFilename`` afresh (best-effort). + """Close and reopen ``baseFilename`` (best-effort). - On failure the stream is left ``None`` so the next emit bails rather than writing to a - stale inode. + On failure the stream is left ``None`` so the next emit bails rather than + writing to a stale inode. """ if self.stream is not None: _quietly(self.stream.close) @@ -343,18 +300,15 @@ class _ManagedRotatingFileHandler(RotatingFileHandler): self._record_stream_stat(stat_result) def _reopen_if_externally_rotated(self) -> None: - """Reopen the stream when ``baseFilename`` no longer matches our fd. + """Reopen when ``baseFilename`` was renamed, unlinked, or replaced by another inode. - Triggered when ``baseFilename`` was renamed (logrotate), unlinked, or replaced by a - different inode. Silent + best-effort: any error falls back to the existing (possibly stale) + Silent + best-effort: any error falls back to the existing (possibly stale) stream so logging keeps working instead of dying on a stat failure. """ try: st = os.stat(self.baseFilename) except FileNotFoundError: - # File was rotated/unlinked underneath us: reopen so a fresh inode - # is created at the expected path. - self._reopen_stream() + self._reopen_stream() # rotated/unlinked underneath us: recreate at the path return except OSError: return # transient — try again on the next emit @@ -362,22 +316,20 @@ class _ManagedRotatingFileHandler(RotatingFileHandler): if self._stat_dev is None or self._stat_ino is None: self._record_stream_stat(st) elif (st.st_dev, st.st_ino) != (self._stat_dev, self._stat_ino): - # baseFilename now points at a DIFFERENT inode than the one we hold open. self._reopen_stream(st) def emit(self, record: logging.LogRecord) -> None: - # Cheap-ish stat-per-record check; the kernel caches inode metadata - # so the syscall is sub-microsecond on a hot file. + # The kernel caches inode metadata, so this stat is sub-microsecond on a hot file. if self.stream is not None or os.path.exists(self.baseFilename): self._reopen_if_externally_rotated() super().emit(record) def handleError(self, record: logging.LogRecord) -> None: - """Suppress the known Windows ``concurrent-log-handler`` lock timeout + """Suppress the known Windows ``concurrent-log-handler`` lock timeout. - CLH's ``emit()`` catches the ``"Cannot acquire lock after N attempts"`` RuntimeError and - routes it here, so this override is the single point to silence it before stdlib prints to - stderr (which the Desktop slash-worker captures and surfaces into chat output). + CLH's ``emit()`` routes that RuntimeError here, so this is the single point to + silence it before stdlib prints to stderr (which the Desktop slash-worker + captures into chat output). """ if not _is_windows_concurrent_log_lock_timeout(sys.exc_info()[1]): super().handleError(record) @@ -390,8 +342,8 @@ class _ManagedRotatingFileHandler(RotatingFileHandler): def doRollover(self): super().doRollover() self._chmod_if_managed() - # Our own rollover writes a new baseFilename; refresh the snapshot - # so the next emit doesn't mistake it for external rotation. + # Our own rollover writes a new baseFilename; refresh the snapshot so + # the next emit doesn't mistake it for external rotation. self._record_stream_stat() @@ -411,9 +363,8 @@ def _new_file_handler( class _ProfileRoutingFileHandler(logging.Handler): """Route queued records to the log file for their Hermes home. - The handler itself is used only behind the existing QueueListener, so its small routing lock - never blocks an agent or dashboard event loop. The underlying handlers retain the existing - rotation, redaction, and managed permission behavior. + Used only behind the QueueListener, so its small routing lock never blocks an agent + or dashboard event loop. Per-home handlers keep rotation, redaction and managed perms. """ def __init__(self, existing: RotatingFileHandler, profile_homes: Sequence[Path]) -> None: @@ -465,15 +416,10 @@ class _ProfileRoutingFileHandler(logging.Handler): super().close() -# Asynchronous file logging — keep the cross-process rotation lock off the loop -# -# The rotating file handlers serialize rollover with a cross-process lock (see -# the module header): when several Hermes processes log to the same file, an -# ``emit`` can block while another process holds that lock. When the emitting -# thread is an asyncio event loop, that block stalls the loop and drops -# WebSocket clients. To keep file I/O off the hot path, every file handler is -# driven by a single ``QueueListener`` on a dedicated thread; loggers only touch -# an in-memory queue (a non-blocking enqueue). +# Asynchronous file logging: an ``emit`` can block on the cross-process +# rotation lock (module header); on an asyncio thread that stalls the loop and +# drops WebSocket clients. Every file handler is therefore driven by a single +# QueueListener thread; loggers only do a non-blocking enqueue. _log_queue: "Optional[queue.SimpleQueue]" = None _queue_listener: Optional[QueueListener] = None @@ -481,20 +427,19 @@ _queued_file_handlers: list = [] _queue_atexit_registered = False # Guards every read-modify-write of the four globals above. setup_logging() # holds no lock and its _logging_initialized guard runs AFTER handler -# registration, so _register_queued_handler() can run concurrently with a -# flush/reset from another thread (gateway init racing a plugin/CLI path). -# Without this, two threads can interleave listener.stop()/reassign/start() -# and leave the queue with two live listeners or an orphaned worker thread. +# registration, so _register_queued_handler() can race a flush/reset from +# another thread (gateway init vs a plugin/CLI path); without this, two +# threads can interleave stop()/reassign/start() and leave two live listeners. _queue_state_lock = threading.Lock() class _NonFormattingQueueHandler(QueueHandler): """``QueueHandler`` for an in-process queue. - Stdlib ``prepare()`` formats and strips ``args``/``exc_info`` for pickling across processes; - our queue is in-process, so the target handlers get an unformatted record and apply their own - ``RedactingFormatter`` on the listener thread. A shallow copy is returned because the emitting - thread's synchronous handlers may mutate ``record.message`` while the listener reads it. + Stdlib ``prepare()`` formats and strips ``args``/``exc_info`` for cross-process + pickling; ours is in-process, so targets get the unformatted record and apply their + own ``RedactingFormatter`` on the listener thread. A shallow copy is returned because + the emitting thread's synchronous handlers may mutate ``record.message`` meanwhile. """ def prepare(self, record: logging.LogRecord) -> logging.LogRecord: @@ -502,10 +447,7 @@ class _NonFormattingQueueHandler(QueueHandler): def _stop_queue_listener() -> None: - """Flush and stop the background log listener (idempotent, thread-safe). - - This is the atexit hook, so it must acquire the state lock itself. - """ + """Flush and stop the background log listener (idempotent; atexit hook, so it takes the lock).""" global _queue_listener with _queue_state_lock: listener, _queue_listener = _queue_listener, None @@ -516,8 +458,8 @@ def _stop_queue_listener() -> None: def _start_queue_listener_locked() -> None: """(Re)build + start a listener over the current handler set (``_queue_state_lock`` held). - A running listener is stopped first; this only happens while handlers are being added - (queue empty), so ``stop()`` returns immediately. + A running listener is stopped first; this only happens while handlers are being + added (queue empty), so ``stop()`` returns immediately. """ global _queue_listener if _queue_listener is not None: @@ -527,9 +469,10 @@ def _start_queue_listener_locked() -> None: def _register_queued_handler(handler: logging.Handler) -> None: - """Route *handler* through the shared async queue instead of attaching it to *root* directly, so - emitting threads never block on file I/O or the cross-process rotation lock. The - ``QueueListener`` applies each handler's own level and filters on its worker thread. + """Route *handler* through the shared async queue instead of attaching it to root. + + Emitting threads never block on file I/O or the rotation lock; the ``QueueListener`` + applies each handler's own level and filters on its worker thread. """ global _log_queue, _queue_atexit_registered with _queue_state_lock: @@ -537,9 +480,7 @@ def _register_queued_handler(handler: logging.Handler) -> None: _log_queue = queue.SimpleQueue() qh = _NonFormattingQueueHandler(_log_queue) qh._hermes_queue = True # type: ignore[attr-defined] - # Always funnel through the root logger so records from any logger - # (production passes root here; callers may pass a child) reach the - # queue via propagation. + # Always on the root logger so records from any logger reach the queue. logging.getLogger().addHandler(qh) _queued_file_handlers.append(handler) _start_queue_listener_locked() @@ -553,12 +494,10 @@ def _register_queued_handler(handler: logging.Handler) -> None: def flush_log_queue() -> None: """Block until all queued records have been written, then resume. - Draining is done by stopping the listener (which processes every pending record before joining) - and restarting it. Used by tests that read a log file right after emitting to it. - - NOTE: ``stop()`` joins the worker thread, so this blocks until the queue is empty. Do NOT call - this on a hard-exit path where the listener may be wedged on the rotation lock — use - ``drain_log_queue()`` there instead, which bounds the wait. + Stops the listener (which processes every pending record before joining) and + restarts it. ``stop()`` joins the worker thread — do NOT call this on a hard-exit + path where the listener may be wedged on the rotation lock; use + ``drain_log_queue()`` there, which bounds the wait. """ with _queue_state_lock: listener = _queue_listener @@ -570,10 +509,8 @@ def flush_log_queue() -> None: def drain_log_queue(timeout: float = 1.0) -> None: """Best-effort, time-bounded drain for hard-exit paths (no restart). - Unlike ``flush_log_queue()``, this stops the listener WITHOUT restarting it (the process is - about to exit) and bounds the drain: if the listener's worker thread is wedged on the cross- - process rotation lock — the very failure this async-logging change exists to survive — an - unbounded ``stop()``/join would re-freeze the shutdown path. + If the listener's worker is wedged on the cross-process rotation lock — the very + failure async logging exists to survive — an unbounded join would re-freeze shutdown. """ listener = _queue_listener if listener is None: @@ -586,12 +523,10 @@ def drain_log_queue(timeout: float = 1.0) -> None: def enable_profile_log_routing(profile_homes: Sequence[str | Path]) -> bool: """Make the queued file logs follow a desktop profile context. - ``setup_logging`` normally binds handlers to one process home. The desktop dashboard is the - exception: its embedded cron ticker may run jobs for every profile. Replace the existing static - file handlers with profile routers after that profile list is known. - - Returns ``True`` when routing is enabled or was already enabled. A single-profile caller is left - untouched because its existing handlers are already correctly scoped. + ``setup_logging`` binds handlers to one process home; the desktop dashboard's + embedded cron ticker may run jobs for every profile, so its static file handlers + are replaced with profile routers once the profile list is known. Returns ``True`` + when routing is (or already was) enabled; a single-profile caller is left untouched. """ global _queue_listener homes: list[Path] = [] @@ -654,9 +589,7 @@ def _add_rotating_handler( formatter: logging.Formatter, log_filter: Optional[logging.Filter] = None, ) -> None: - """Register a queued ``RotatingFileHandler`` for *path*, skipping if one already exists for the - same resolved file path (idempotent). - """ + """Register a queued ``RotatingFileHandler`` for *path*; idempotent per resolved path.""" resolved = path.resolve() for existing in _queued_file_handlers: # Already attached directly, or already covered by the profile router. @@ -671,19 +604,16 @@ def _add_rotating_handler( ) if log_filter is not None: handler.addFilter(log_filter) - # Route through the async queue instead of ``logger.addHandler(handler)`` so - # the rotation-lock wait never runs on the caller's (often event-loop) thread. + # Queue, not ``addHandler``: the rotation-lock wait never runs on the caller's thread. _register_queued_handler(handler) def _read_logging_config(): """Best-effort read of ``logging.*`` from config.yaml.""" try: - # Prefer the shared (mtime, size)-keyed raw-config cache so this read - # reuses the parse hermes_cli.main's early bridge already did (one - # config.yaml parse per process instead of 3-4). Fall back to a - # direct parse when hermes_cli.config isn't importable (bare - # hermes_logging consumers). + # Prefer the shared (mtime, size)-keyed raw-config cache so this reuses + # hermes_cli.main's early parse (one config.yaml parse per process); + # fall back to a direct parse for bare hermes_logging consumers. try: from hermes_cli.config import read_raw_config as _rrc cfg = _rrc() or {} @@ -696,8 +626,7 @@ def _read_logging_config(): cfg = fast_safe_load(f) or {} if not cfg: return (None, None, None) - # Managed scope: an administrator can pin logging.* too. Overlay via - # the shared helper (fail-open) since this reads config.yaml directly. + # Managed scope: an administrator can pin logging.* too (fail-open overlay). try: from hermes_cli import managed_scope cfg = managed_scope.apply_managed_overlay(cfg) diff --git a/hermes_startup_watchdog.py b/hermes_startup_watchdog.py index 52a3bd0556..9985689909 100644 --- a/hermes_startup_watchdog.py +++ b/hermes_startup_watchdog.py @@ -1,82 +1,29 @@ """Startup-liveness watchdog — respawn a gateway that wedges before its loop runs (OOF-298). -The existing liveness backstops all assume startup succeeded: +Every other liveness backstop (loop-liveness watchdog, shutdown watchdog, heartbeat file) +assumes startup succeeded; none can fire if the process deadlocks before the event loop +is alive (OOF-298: ~30h with every thread in ``futex_wait_queue``, zero logs, s6 saw a +live PID). A daemon thread armed at process entry and disarmed once the loop is confirmed +live dumps all-thread stacks (``faulthandler``), records the exit in the lifecycle ledger +(NS-608) and ``os._exit``\\ s with the service-restart code so s6/systemd respawn. -* the loop-liveness watchdog (:mod:`gateway.shutdown_watchdog`) is armed by - ``GatewayRunner._start_loop_liveness_guards`` — *inside* the running event - loop's startup path; -* the shutdown watchdog is armed at ``stop()``; -* the loop heartbeat file is written by an asyncio task. +Slow-but-alive startups are not killed. Order of authority: (1) phase-owned progress +leases (:func:`report_startup_progress`) — authoritative, prove the *startup path itself* +is alive, work for I/O-bound phases (schema migration, corruption repair in +``SessionDB.__init__``) with ~zero CPU; (2) process-wide CPU progress, capped at +``_MAX_CPU_EXTENSIONS`` since an unrelated thread burning CPU must not hide a parked +startup thread forever. Known limitation: a *spinning* startup deadlock earns the capped +extensions before firing; the observed class is parked-thread deadlocks. Idle-by-design +waits call :func:`kick_startup_watchdog` (respawn-storm backoff, up to 300s); MCP +discovery's 120s wait sits inside the 300s default. -None of them can fire if the process deadlocks **before the event loop comes -alive**. That failure mode is real: OOF-298 documents a hosted gateway whose -process sat for ~30 hours with every thread parked in ``futex_wait_queue``, -zero log lines written, ``/health`` unreachable — while s6 saw a live PID and -therefore never respawned it, and a stale ``gateway_state.json`` from the -*previous* life told every status surface the gateway was "draining". - -This module closes that gap with a plain daemon OS thread armed at process -entry, disarmed the moment the event loop is confirmed live (the point where -the existing loop-liveness watchdog takes over). If startup neither reaches -that milestone nor exits within the deadline, the watchdog dumps all-thread -stacks via ``faulthandler``, records the exit in the lifecycle ledger -(NS-608) so the next boot classifies it correctly, and ``os._exit``\\ s with -the service-restart code so s6/systemd revive the process instead of -babysitting a zombie. - -Slow-but-alive startups are NOT killed. Two mechanisms, in order of -authority: - -1. **Phase-owned progress leases** (:func:`report_startup_progress`): a - startup phase that is about to do legitimately long synchronous work - (large ``state.db`` schema migrations, corruption repair/backup — both - run inside ``SessionDB.__init__`` well before the loop starts, and both - can be I/O-bound with near-zero CPU) declares a lease for its honest - worst case. The lease is the authoritative signal: it proves the - *startup path itself* is alive, not merely that the process is warm. -2. **CPU progress, as a bounded fallback only**: if the deadline expires - but the process consumed meaningful CPU during the window - (``time.process_time()`` is process-wide), the deadline is extended — - at most ``_MAX_CPU_EXTENSIONS`` times. Process-wide CPU proves activity, - not startup progress (an unrelated daemon thread burning CPU must not - hide a parked startup thread forever), hence the cap. Phases that hold - a current lease are never subject to the cap. - -The OOF-298 deadlock class parks every thread in futex waits, accrues ~zero -CPU, and owns no lease — it fires on schedule. Known limitation, documented -deliberately: a *spinning* (busy-wait) startup deadlock reads as CPU -progress and gets the capped extensions before firing; the observed -incident class is parked-thread deadlocks, which fire immediately. - -Waits that are idle-by-design get explicit handling instead: - -* the respawn-storm breaker's intentional backoff sleep (up to 300s) calls - :func:`kick_startup_watchdog` with the sleep budget before sleeping; -* MCP tool discovery's internal wait is bounded at 120s, comfortably inside - the 300s default deadline. - -IMPORT-LIGHTNESS IS A CORRECTNESS PROPERTY of this module, not a style -preference. It lives at the repository top level (not inside the ``gateway`` -package) and imports **only stdlib** because: - -1. ``gateway/__init__`` eagerly imports the config/session/delivery graph — - hundreds of modules, DB-adjacent code included. Arming must happen - *before* that graph is imported, or an import-time deadlock (a plausible - shape of "wedged before the loop, no logs") sits outside the watchdog's - coverage. -2. At fire time the main thread may be wedged **holding the import lock**; - any import attempted on the watchdog thread could then block forever. - The fire path therefore performs no imports at all on its own thread — - the lifecycle-ledger write (which does import) runs on a short-lived - helper thread joined with a timeout, and ``os._exit`` happens regardless. - -Config surface is deliberately env-only (``HERMES_STARTUP_WATCHDOG=0`` to -disable, ``HERMES_STARTUP_WATCHDOG_TIMEOUT_S`` to tune): the watchdog must be -armed before config.yaml is loaded — a wedge during config parsing is exactly -in scope — so it cannot depend on config for its own enablement. - -Everything here is best-effort: a watchdog failure must never affect the -startup it is observing. +IMPORT-LIGHTNESS IS A CORRECTNESS PROPERTY: top-level module, stdlib only. Arming must +precede importing ``gateway`` (hundreds of modules; an import-time deadlock is in scope), +and at fire time the wedged main thread may hold the import lock — so the fire path does +no imports on its own thread; the ledger write runs on a helper thread joined with a +timeout. Config is env-only (``HERMES_STARTUP_WATCHDOG=0``, +``HERMES_STARTUP_WATCHDOG_TIMEOUT_S``) because config.yaml parsing is itself in scope. +Everything is best-effort: a watchdog failure must never affect the startup it observes. """ from __future__ import annotations @@ -97,9 +44,8 @@ logger = logging.getLogger(__name__) DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S = 300.0 _MIN_TIMEOUT_S = 30.0 -# Mirrors gateway.restart.GATEWAY_SERVICE_RESTART_EXIT_CODE. Duplicated here -# (with a parity test in tests/gateway/test_startup_watchdog.py) because this -# module must not import the gateway package — see module docstring. +# Mirrors gateway.restart.GATEWAY_SERVICE_RESTART_EXIT_CODE (parity test in +# tests/gateway/test_startup_watchdog.py) — this module must not import gateway. SERVICE_RESTART_EXIT_CODE = 75 ENV_STARTUP_WATCHDOG = "HERMES_STARTUP_WATCHDOG" @@ -109,62 +55,50 @@ _DUMP_RELATIVE = ("logs", "gateway-startup-watchdog.log") _FALSEY = frozenset({"0", "false", "no", "off"}) -# The waiter re-reads its deadline at most this often, so kick_/deadline -# extensions take effect promptly without busy-waiting. +# The waiter re-reads its deadline at most this often so kicks/extensions +# take effect promptly without busy-waiting. _POLL_SLICE_S = 5.0 -# Minimum process CPU-time delta (seconds) within one expired deadline window -# for startup to count as "making progress" and earn a fallback extension. A -# parked futex deadlock accrues microseconds; a schema migration accrues -# orders of magnitude more than this per window even on slow disks. +# Minimum process CPU-time delta within one expired window to count as +# progress. A parked futex deadlock accrues microseconds; a schema migration +# accrues orders of magnitude more per window even on slow disks. _CPU_PROGRESS_MIN_S = 1.0 -# Hard cap on CPU-fallback extensions. CPU is process-wide evidence and can -# be produced by threads unrelated to startup, so it may only stretch the -# runway to (1 + cap) x timeout; anything longer must hold an explicit -# phase lease (report_startup_progress). 3 x 300s default = 20min total. +# Hard cap on CPU-fallback extensions: CPU is process-wide evidence, so it may +# only stretch the runway to (1 + cap) x timeout (3 x 300s = 20min); anything +# longer must hold an explicit phase lease. _MAX_CPU_EXTENSIONS = 3 -# Per-call clamp on progress leases (report_startup_progress). A phase that -# genuinely needs longer renews its lease — the renewal is itself the -# liveness evidence. 15 minutes covers the observed worst-case single -# migration step on multi-GB state.db files with generous margin. +# Per-call clamp on progress leases; a phase that needs longer renews (the +# renewal is the liveness evidence). 15min covers the observed worst-case +# single migration step on multi-GB state.db files with margin. _MAX_LEASE_S = 900.0 -# How long the fire path waits for the lifecycle-ledger helper thread before -# exiting anyway (the import lock may be held by the wedged main thread). +# Bounded wait for the lifecycle-ledger helper thread (import lock may be +# held by the wedged main thread). _LEDGER_JOIN_TIMEOUT_S = 5.0 -# Upper bound on the ENTIRE forensic fire path (logging, dump record, -# faulthandler, ledger). A sibling escort thread — which touches no logging, -# no filesystem, and no application locks — hard-exits the process if the -# forensics wedge (e.g. the wedged main thread holds the logging handler -# lock, or the disk is full/hung). Must exceed _LEDGER_JOIN_TIMEOUT_S. +# Upper bound on the ENTIRE forensic fire path. A sibling escort thread that +# touches no logging/filesystem/locks hard-exits if forensics wedge (e.g. the +# wedged main thread holds the logging handler lock, disk full/hung). Must +# exceed _LEDGER_JOIN_TIMEOUT_S. _FIRE_EXIT_BOUND_S = 10.0 -# Handle lifecycle states. Transitions are guarded by the handle's state -# lock so a disarm and a fire can never both "win" (P2 race, PR #89750 -# review): armed -> disarmed (startup reached a live loop) or -# armed -> firing (deadline expired with no CPU progress) — never both. +# Handle states. Transitions are guarded by the handle's state lock so a +# disarm and a fire can never both "win": armed -> disarmed or armed -> firing. _ARMED = "armed" _DISARMED = "disarmed" _FIRING = "firing" -# Module-level singleton: the arm sites (hermes_cli.main / hermes_cli.gateway -# / gateway.run.main / cli.py --gateway / scripts/hermes-gateway) and the -# disarm site (GatewayRunner, once the loop is live) have no shared object to -# hand a handle through, and only one gateway startup ever runs per process. +# Module singleton: the arm sites (hermes_cli.main / hermes_cli.gateway / +# gateway.run.main / cli.py --gateway) and the disarm site (GatewayRunner) +# share no object, and only one gateway startup ever runs per process. _handle_lock = threading.Lock() _handle: Optional["StartupWatchdogHandle"] = None def _process_hermes_home() -> Path: - """HERMES_HOME for process-level diagnostic files. - - Stdlib-only replica of ``hermes_constants``' platform default — this - module must not import application code (see module docstring). Hosted - images always set ``HERMES_HOME`` explicitly. - """ + """HERMES_HOME for diagnostic files — stdlib-only replica of the hermes_constants default.""" val = os.environ.get("HERMES_HOME", "").strip() if val: return Path(val) @@ -183,8 +117,7 @@ def get_startup_watchdog_dump_path(home: Optional[Path] = None) -> Path: def startup_watchdog_disabled() -> bool: """True when ``HERMES_STARTUP_WATCHDOG`` opts out explicitly.""" - raw = os.environ.get(ENV_STARTUP_WATCHDOG, "").strip().lower() - return raw in _FALSEY + return os.environ.get(ENV_STARTUP_WATCHDOG, "").strip().lower() in _FALSEY def resolve_startup_watchdog_timeout() -> float: @@ -207,23 +140,30 @@ def resolve_startup_watchdog_timeout() -> float: return max(value, _MIN_TIMEOUT_S) -def _write_dump_record(record: Dict[str, Any]) -> None: - """Append a one-line JSON metadata record beside the faulthandler dump.""" +def _append_dump(write, failure_msg: str) -> None: + """Open the dump file for append and hand it to *write*; failures only log at DEBUG.""" try: path = get_startup_watchdog_dump_path() path.parent.mkdir(parents=True, exist_ok=True) with open(path, "a", encoding="utf-8") as fh: - fh.write(json.dumps(record, default=str) + "\n") + write(fh) except Exception: - logger.debug("Failed to write startup watchdog dump record", exc_info=True) + logger.debug(failure_msg, exc_info=True) + + +def _write_dump_record(record: Dict[str, Any]) -> None: + """Append a one-line JSON metadata record beside the faulthandler dump.""" + _append_dump( + lambda fh: fh.write(json.dumps(record, default=str) + "\n"), + "Failed to write startup watchdog dump record", + ) def _mark_lifecycle_exit(exit_code: int) -> None: """Record the watchdog exit in the NS-608 lifecycle sentinel. - Runs on a dedicated helper thread (see ``_fire``): the ``import`` below - can block indefinitely on the interpreter import lock if the wedged main - thread holds it, and the fire path must reach ``os._exit`` regardless. + Runs on a helper thread (see ``_fire``): the import can block on the + interpreter import lock, and the fire path must reach ``os._exit`` regardless. """ try: from gateway.lifecycle_ledger import mark_exited @@ -246,22 +186,19 @@ class StartupWatchdogHandle: self._disarmed_event = threading.Event() self._thread: Optional[threading.Thread] = None self._extensions = 0 - # Phase-owned progress lease (see lease()). monotonic deadline the + # Phase-owned progress lease (see lease()): monotonic deadline the # current startup phase has claimed for legitimately long sync work. self._lease_until = 0.0 self._lease_phase: Optional[str] = None self._lease_count = 0 - # Set by _fire() once forensics complete; the exit escort thread - # uses it to stand down when the normal exit path won the race. + # Set by _fire() once forensics complete so the exit escort stands down. self._fire_done = threading.Event() def disarm(self) -> None: """Startup reached a live event loop — stand down. Idempotent. - Atomic with respect to firing: whichever of disarm/fire takes the - state lock first wins, so a disarm that lands before the fire - sequence begins is always honored (never lost to a deadline that - expired concurrently). + Atomic with respect to firing: whichever of disarm/fire takes the state + lock first wins, so a disarm landing before the fire sequence is never lost. """ with self._state_lock: if self._state == _ARMED: @@ -271,9 +208,8 @@ class StartupWatchdogHandle: def kick(self, extra_s: float = 0.0) -> None: """Push the deadline out to ``now + timeout + extra_s``. - For call sites that are about to block intentionally with ~zero CPU - activity (the respawn-storm breaker's backoff sleep), which would - otherwise be indistinguishable from a parked deadlock. + For call sites about to block intentionally with ~zero CPU (the + respawn-storm backoff sleep), otherwise indistinguishable from a parked deadlock. """ try: extra = max(0.0, float(extra_s)) @@ -283,18 +219,13 @@ class StartupWatchdogHandle: self._deadline = time.monotonic() + self.timeout_s + extra def lease(self, expected_s: float, phase: str = "") -> None: - """Claim a progress lease: this startup phase is alive and expects - up to ``expected_s`` more seconds of legitimate synchronous work. + """Claim a progress lease: this phase expects up to ``expected_s`` more seconds of work. - This is the authoritative "still making progress" signal — unlike - process-wide CPU time it is owned by the startup path itself, so it - works for I/O-bound phases (corruption repair, backups) that accrue - almost no CPU, and it cannot be counterfeited by unrelated threads. - - Leases are clamped to ``_MAX_LEASE_S`` per call so a single buggy - caller cannot silence the watchdog indefinitely; genuinely long - phases renew periodically (renewal proves continued liveness). - Never raises.""" + The authoritative "still making progress" signal — owned by the startup path, + so it works for I/O-bound phases and cannot be counterfeited by unrelated + threads. Clamped to ``_MAX_LEASE_S`` per call so one buggy caller cannot + silence the watchdog indefinitely; long phases renew. Never raises. + """ try: expected = float(expected_s) except (TypeError, ValueError): @@ -332,14 +263,11 @@ class StartupWatchdogHandle: def _fire(self) -> None: """Forensics, then exit — with the exit itself independently bounded. - Everything in here that produces forensics (logging, the JSON dump - record, faulthandler, the lifecycle ledger) can in principle block: - the wedged main thread may hold the logging handler lock, the disk - may be full or hung. None of that may stop the respawn. An escort - thread is started FIRST; it touches no logging, no filesystem and - no application locks — it sleeps, checks whether the normal exit - happened, and otherwise calls the exit seam itself. ``os._exit`` - is async-signal-safe and lock-free by design.""" + Everything producing forensics (logging, dump record, faulthandler, ledger) + can block: the wedged main thread may hold the logging handler lock, the disk + may be hung. None of that may stop the respawn, so the escort thread starts + FIRST. ``os._exit`` is async-signal-safe and lock-free by design. + """ try: escort = threading.Thread( target=self._exit_escort, @@ -381,22 +309,15 @@ class StartupWatchdogHandle: faulthandler.dump_traceback(all_threads=True) except Exception: logger.debug("Startup watchdog faulthandler dump failed", exc_info=True) - # Also dump stacks into the log file: on detached/windowless runs - # (pythonw, some service managers) stderr may be absent, and the - # whole point of firing is to leave forensics behind. - try: - path = get_startup_watchdog_dump_path() - path.parent.mkdir(parents=True, exist_ok=True) - with open(path, "a", encoding="utf-8") as fh: - faulthandler.dump_traceback(file=fh, all_threads=True) - except Exception: - logger.debug( - "Startup watchdog file-based faulthandler dump failed", exc_info=True - ) - # Lifecycle-ledger write on a helper thread: it imports application - # code, and the wedged main thread may hold the import lock. Bounded - # join, then exit regardless (NS-608 classification is best-effort; - # the respawn is not). + # Also dump into the log file: detached/windowless runs (pythonw, some + # service managers) may have no stderr, and forensics are the point. + _append_dump( + lambda fh: faulthandler.dump_traceback(file=fh, all_threads=True), + "Startup watchdog file-based faulthandler dump failed", + ) + # Ledger write on a helper thread (it imports application code; the + # wedged main thread may hold the import lock). Bounded join, then exit + # regardless — NS-608 classification is best-effort; the respawn is not. try: ledger_thread = threading.Thread( target=_mark_lifecycle_exit, @@ -414,9 +335,9 @@ class StartupWatchdogHandle: def _exit_escort(self) -> None: """Hard-exit if the forensic fire path wedges (bounded-exit seam). - Deliberately free of log handlers, filesystem access, module loads - and any lock shared with application code: its only dependencies - are a monotonic sleep, an Event check, and the exit seam.""" + Deliberately free of log handlers, filesystem access, module loads and any + lock shared with application code: only a sleep, an Event check and the exit seam. + """ self._sleep(_FIRE_EXIT_BOUND_S) if self._fire_done.is_set(): return @@ -432,6 +353,14 @@ class StartupWatchdogHandle: """Seam for tests; production is a bare ``os._exit``.""" os._exit(code) + def _extend_if_armed(self, deadline: float) -> bool: + """Set a new deadline under the state lock; False when no longer armed.""" + with self._state_lock: + if self._state != _ARMED: + return False + self._deadline = deadline + return True + def _run(self) -> None: last_cpu = self._process_cpu_seconds() while True: @@ -444,28 +373,16 @@ class StartupWatchdogHandle: if self._disarmed_event.wait(timeout=min(remaining, _POLL_SLICE_S)): return continue - # Deadline expired. Order of authority: - # - # 1. Phase lease (report_startup_progress): the startup path - # itself declared long legitimate work — honor it outright. - # Works for I/O-bound phases with ~zero CPU (corruption - # repair, backups) and cannot be faked by unrelated threads. - # 2. CPU progress, bounded: process-wide CPU proves the process - # is doing *something*, not that startup is progressing (an - # unrelated daemon thread could burn CPU while the startup - # thread sits parked forever). Extend at most - # _MAX_CPU_EXTENSIONS times, then fire regardless. + # Deadline expired. Order of authority: (1) a phase lease is honored + # outright; (2) CPU progress extends at most _MAX_CPU_EXTENSIONS + # times (process-wide CPU proves activity, not startup progress). now = time.monotonic() with self._state_lock: lease_until = self._lease_until lease_phase = self._lease_phase if lease_until > now: - with self._state_lock: - if self._state != _ARMED: - return - self._deadline = max( - lease_until, now + min(_POLL_SLICE_S, self.timeout_s) - ) + if not self._extend_if_armed(max(lease_until, now + min(_POLL_SLICE_S, self.timeout_s))): + return try: logger.warning( "Gateway startup exceeded %.0fs but phase %r holds a " @@ -476,7 +393,7 @@ class StartupWatchdogHandle: ) except Exception: pass - # Leased work may be I/O-bound; reset the CPU baseline so a + # Leased work may be I/O-bound; reset the CPU baseline so the # post-lease window is judged on its own activity. last_cpu = self._process_cpu_seconds() continue @@ -490,10 +407,8 @@ class StartupWatchdogHandle: window_delta = cpu - last_cpu last_cpu = cpu self._extensions += 1 - with self._state_lock: - if self._state != _ARMED: - return - self._deadline = time.monotonic() + self.timeout_s + if not self._extend_if_armed(time.monotonic() + self.timeout_s): + return try: logger.warning( "Gateway startup exceeded %.0fs but is consuming CPU " @@ -568,10 +483,9 @@ def arm_startup_watchdog( def disarm_startup_watchdog() -> None: """Disarm the process-wide startup watchdog, if armed. Never raises. - The handle's ``disarm()`` is called while still holding the singleton - lock — it is non-blocking, and holding the lock closes the window where - a concurrent re-arm could swap in a new handle that the disarm then - misses. + ``disarm()`` runs while still holding the singleton lock — it is non-blocking, + and this closes the window where a concurrent re-arm could swap in a new + handle that the disarm then misses. """ global _handle try: @@ -584,44 +498,36 @@ def disarm_startup_watchdog() -> None: logger.debug("Failed to disarm gateway startup watchdog", exc_info=True) -def kick_startup_watchdog(extra_s: float = 0.0) -> None: - """Extend the armed watchdog's deadline. No-op when not armed; never raises. - - Call before intentionally blocking with ~zero CPU activity (e.g. the - respawn-storm breaker's backoff sleep) so the idle wait is not mistaken - for a parked deadlock. - """ +def _with_armed_handle(method: str, failure_msg: str, *args) -> None: + """Call ``handle.(*args)`` on the armed handle; no-op when unarmed, never raises.""" try: with _handle_lock: handle = _handle if handle is not None: - handle.kick(extra_s) + getattr(handle, method)(*args) except Exception: - logger.debug("Failed to kick gateway startup watchdog", exc_info=True) + logger.debug(failure_msg, exc_info=True) + + +def kick_startup_watchdog(extra_s: float = 0.0) -> None: + """Extend the armed watchdog's deadline. No-op when not armed; never raises. + + Call before intentionally blocking with ~zero CPU activity (e.g. the + respawn-storm backoff sleep) so the idle wait is not mistaken for a parked deadlock. + """ + _with_armed_handle("kick", "Failed to kick gateway startup watchdog", extra_s) def report_startup_progress(expected_s: float, phase: str = "") -> None: """Declare a phase-owned progress lease on the armed startup watchdog. - Call from startup phases about to perform legitimately long synchronous - work — most importantly ``state.db`` schema migrations and corruption - repair/backup inside ``SessionDB.__init__`` — passing an honest worst - case for the work about to be done, and renew periodically for - multi-step phases. Unlike CPU-time inference, a lease is owned by the - startup path itself: it works for I/O-bound work that accrues ~zero CPU - and cannot be counterfeited by unrelated busy threads. - - Per-call lease duration is clamped to ``_MAX_LEASE_S``; renewals prove - continued liveness. No-op when the watchdog is not armed; never raises — - safe to call unconditionally from application code. + Call from startup phases about to do legitimately long synchronous work — most + importantly ``state.db`` schema migrations and corruption repair/backup inside + ``SessionDB.__init__`` — with an honest worst case, renewing for multi-step phases. + Per-call duration is clamped to ``_MAX_LEASE_S``. No-op when not armed; never + raises — safe to call unconditionally from application code. """ - try: - with _handle_lock: - handle = _handle - if handle is not None: - handle.lease(expected_s, phase) - except Exception: - logger.debug("Failed to report startup progress", exc_info=True) + _with_armed_handle("lease", "Failed to report startup progress", expected_s, phase) def _reset_for_tests() -> None: