diff --git a/cli.py b/cli.py index 0b96e16252..d0ff845579 100644 --- a/cli.py +++ b/cli.py @@ -21704,6 +21704,13 @@ def main( # Handle gateway mode (messaging + cron) if gateway: import asyncio + # Startup-liveness watchdog (OOF-298): this legacy entry point must + # be covered too — arm before importing the gateway graph. + try: + from hermes_startup_watchdog import arm_startup_watchdog + arm_startup_watchdog() + except Exception: + pass from gateway.run import start_gateway print("Starting Hermes Gateway (messaging platforms)...") asyncio.run(start_gateway()) diff --git a/gateway/run.py b/gateway/run.py index 96d799750f..177592bbc3 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -13481,17 +13481,20 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew self._gateway_loop = None if self._gateway_loop is not None: self._start_loop_liveness_guards(self._gateway_loop) - # The event loop is confirmed live: the startup-liveness watchdog's - # job is done and the loop-liveness watchdog (armed just above) - # takes over from here (OOF-298). Disarm even when the loop guards - # are config-disabled — the startup watchdog only covers the - # pre-loop window, never adapter connects or steady-state. - try: - from gateway.startup_watchdog import disarm_startup_watchdog + # The event loop is confirmed live: the startup-liveness + # watchdog's job is done and the loop-liveness watchdog (armed + # just above) takes over from here (OOF-298). Disarm even when + # the loop guards are config-disabled — the startup watchdog + # only covers the pre-loop window, never adapter connects or + # steady-state. Deliberately inside the loop-confirmed branch: + # if the loop somehow isn't live, startup has NOT reached the + # milestone and the watchdog must stay armed. + try: + from gateway.startup_watchdog import disarm_startup_watchdog - disarm_startup_watchdog() - except Exception: - logger.debug("Startup watchdog disarm failed", exc_info=True) + disarm_startup_watchdog() + except Exception: + logger.debug("Startup watchdog disarm failed", exc_info=True) logger.info("Session storage: %s", self.config.sessions_dir) # Sanity-check that systemd's TimeoutStopSec covers our drain @@ -33432,7 +33435,7 @@ def main(): os.environ.setdefault("AI_AGENT", "hermes-agent") os.environ.setdefault("HERMES_AGENT", "true") -# Positive process identity: ledger registration + Windows job-object + # Positive process identity: ledger registration + Windows job-object # self-attach, so update-time reapers can identify this gateway (and its # child tree dies with it on Windows). Best-effort — never blocks startup. try: @@ -33446,11 +33449,13 @@ def main(): except Exception: pass - # Startup-liveness watchdog (OOF-298): armed before ANY other startup - # work — config load, imports, DB opens — so a deadlock anywhere in the - # pre-event-loop window still gets the process respawned by the service - # supervisor instead of wedging as a live-PID zombie. Disarmed by - # GatewayRunner once the event loop is confirmed live. + # Startup-liveness watchdog (OOF-298): armed before config load, DB + # opens, and the rest of pre-loop startup so a deadlock in that window + # still gets the process respawned by the service supervisor instead of + # wedging as a live-PID zombie. (Import-time coverage for the standard + # ``hermes gateway run`` path is provided even earlier, by the argv + # fast-path in hermes_cli.main.) Disarmed by GatewayRunner once the + # event loop is confirmed live. try: from gateway.startup_watchdog import arm_startup_watchdog arm_startup_watchdog() diff --git a/gateway/startup_watchdog.py b/gateway/startup_watchdog.py index eed5fa04eb..3b114c7b27 100644 --- a/gateway/startup_watchdog.py +++ b/gateway/startup_watchdog.py @@ -1,274 +1,40 @@ -"""Startup-liveness watchdog — respawn a gateway that wedges before its loop runs (OOF-298). +"""Compatibility shim — the real implementation is ``hermes_startup_watchdog``. -The existing liveness backstops all assume startup succeeded: +The startup-liveness watchdog (OOF-298) must be armable *before* the +``gateway`` package is imported: ``gateway/__init__`` eagerly pulls in the +config/session/delivery graph, and an import-time deadlock is squarely inside +the watchdog's coverage mandate. The implementation therefore lives at the +repository top level as a stdlib-only module. -* 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. - -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 the simplest thing that works: 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. - -Deadline rationale: a healthy startup reaches the disarm point in seconds. -The slowest legitimate pre-loop work is MCP tool discovery (bounded 120s -internal wait), so the 300s default leaves comfortable headroom. Platform -adapter connects — which can genuinely take minutes (WhatsApp pairing, npm -cold installs) — happen *after* the disarm point and are never covered by -this watchdog. - -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. +This shim keeps the intuitive ``gateway.startup_watchdog`` import path +working for code that runs after the package is loaded (the disarm site in +``gateway.run``, tests, operators poking at a REPL). """ -from __future__ import annotations +from hermes_startup_watchdog import ( # noqa: F401 + DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S, + ENV_STARTUP_WATCHDOG, + ENV_STARTUP_WATCHDOG_TIMEOUT_S, + SERVICE_RESTART_EXIT_CODE, + StartupWatchdogHandle, + arm_startup_watchdog, + disarm_startup_watchdog, + get_startup_watchdog_dump_path, + kick_startup_watchdog, + resolve_startup_watchdog_timeout, + startup_watchdog_disabled, +) -import faulthandler -import json -import logging -import os -import threading -import time -from datetime import datetime, timezone -from pathlib import Path -from typing import Any, Dict, Optional - -from gateway.restart import GATEWAY_SERVICE_RESTART_EXIT_CODE - -logger = logging.getLogger(__name__) - -DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S = 300.0 -_MIN_TIMEOUT_S = 30.0 - -ENV_STARTUP_WATCHDOG = "HERMES_STARTUP_WATCHDOG" -ENV_STARTUP_WATCHDOG_TIMEOUT_S = "HERMES_STARTUP_WATCHDOG_TIMEOUT_S" - -_DUMP_RELATIVE = ("logs", "gateway-startup-watchdog.log") - -_FALSEY = frozenset({"0", "false", "no", "off"}) - -# Module-level singleton: the arm sites (gateway.run.main and the -# hermes_cli.gateway CLI wrapper) and the disarm site -# (GatewayRunner._start_loop_liveness_guards) have no shared object to hand a -# handle through, 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 (ignore task overrides).""" - val = os.environ.get("HERMES_HOME", "").strip() - if val: - return Path(val) - from hermes_constants import get_hermes_home - - return get_hermes_home() - - -def get_startup_watchdog_dump_path(home: Optional[Path] = None) -> Path: - """Return ``/logs/gateway-startup-watchdog.log``.""" - base = home if home is not None else _process_hermes_home() - return base.joinpath(*_DUMP_RELATIVE) - - -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 - - -def resolve_startup_watchdog_timeout() -> float: - """Deadline in seconds; env override, floor-clamped, default on garbage.""" - raw = os.environ.get(ENV_STARTUP_WATCHDOG_TIMEOUT_S, "").strip() - if not raw: - return DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S - try: - value = float(raw) - except ValueError: - logger.warning( - "Ignoring non-numeric %s=%r; using default %.0fs", - ENV_STARTUP_WATCHDOG_TIMEOUT_S, - raw, - DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S, - ) - return DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S - if value <= 0: - return DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S - 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.""" - 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") - except Exception: - logger.debug("Failed to write startup watchdog dump record", exc_info=True) - - -class StartupWatchdogHandle: - """Disarm/inspect handle for the armed startup watchdog thread.""" - - def __init__(self, timeout_s: float, exit_code: int): - self.timeout_s = timeout_s - self.exit_code = exit_code - self.armed_at = time.monotonic() - self._disarmed = threading.Event() - self._thread: Optional[threading.Thread] = None - - def disarm(self) -> None: - """Startup reached a live event loop — stand down. Idempotent.""" - self._disarmed.set() - - @property - def disarmed(self) -> bool: - return self._disarmed.is_set() - - def is_alive(self) -> bool: - return self._thread is not None and self._thread.is_alive() - - def join(self, timeout: Optional[float] = None) -> None: - if self._thread is not None: - self._thread.join(timeout=timeout) - - # ── internals ──────────────────────────────────────────────────────── - - def _fire(self) -> None: - elapsed = time.monotonic() - self.armed_at - try: - logger.critical( - "Gateway startup did not reach a live event loop within %.0fs " - "(elapsed %.0fs); dumping all thread stacks and exiting with " - "code %d so the service supervisor can restart it (OOF-298).", - self.timeout_s, - elapsed, - self.exit_code, - ) - except Exception: - pass - _write_dump_record( - { - "ts": datetime.now(timezone.utc).isoformat(), - "tag": "startup_watchdog.fired", - "pid": os.getpid(), - "timeout_s": self.timeout_s, - "elapsed_s": round(elapsed, 3), - "exit_code": self.exit_code, - } - ) - try: - faulthandler.dump_traceback(all_threads=True) - except Exception: - logger.debug("Startup watchdog faulthandler dump failed", exc_info=True) - # Record the exit in the lifecycle sentinel so the next boot reports - # "startup watchdog hard-exit" instead of misclassifying this as an - # unclean SIGKILL/OOM death (NS-608). - try: - from gateway.lifecycle_ledger import mark_exited - - mark_exited(self.exit_code, reason="startup_liveness_watchdog") - except Exception: - pass - self._exit(self.exit_code) - - @staticmethod - def _exit(code: int) -> None: - """Seam for tests; production is a bare ``os._exit``.""" - os._exit(code) - - def _run(self) -> None: - if self._disarmed.wait(timeout=self.timeout_s): - return - if self._disarmed.is_set(): - return - self._fire() - - def _start(self) -> bool: - thread = threading.Thread( - target=self._run, - daemon=True, - name="gateway-startup-watchdog", - ) - try: - thread.start() - except Exception: - logger.debug("Failed to start gateway startup watchdog", exc_info=True) - return False - self._thread = thread - return True - - -def arm_startup_watchdog( - timeout_s: Optional[float] = None, - *, - exit_code: int = GATEWAY_SERVICE_RESTART_EXIT_CODE, -) -> Optional[StartupWatchdogHandle]: - """Arm the process-wide startup watchdog. Idempotent; never raises. - - Returns the (possibly pre-existing) handle, or ``None`` when disabled via - ``HERMES_STARTUP_WATCHDOG=0`` or when the thread could not be started. - """ - global _handle - try: - if startup_watchdog_disabled(): - return None - with _handle_lock: - if _handle is not None and _handle.is_alive(): - return _handle - resolved = ( - float(timeout_s) - if timeout_s is not None and float(timeout_s) > 0 - else resolve_startup_watchdog_timeout() - ) - handle = StartupWatchdogHandle(resolved, exit_code) - if not handle._start(): - return None - _handle = handle - return handle - except Exception: - logger.debug("Failed to arm gateway startup watchdog", exc_info=True) - return None - - -def disarm_startup_watchdog() -> None: - """Disarm the process-wide startup watchdog, if armed. Never raises.""" - global _handle - try: - with _handle_lock: - handle = _handle - _handle = None - if handle is not None: - handle.disarm() - except Exception: - logger.debug("Failed to disarm gateway startup watchdog", exc_info=True) - - -def _reset_for_tests() -> None: - """Drop the module singleton (test isolation only).""" - global _handle - with _handle_lock: - handle = _handle - _handle = None - if handle is not None: - handle.disarm() +__all__ = [ + "DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S", + "ENV_STARTUP_WATCHDOG", + "ENV_STARTUP_WATCHDOG_TIMEOUT_S", + "SERVICE_RESTART_EXIT_CODE", + "StartupWatchdogHandle", + "arm_startup_watchdog", + "disarm_startup_watchdog", + "get_startup_watchdog_dump_path", + "kick_startup_watchdog", + "resolve_startup_watchdog_timeout", + "startup_watchdog_disabled", +] diff --git a/hermes_cli/gateway.py b/hermes_cli/gateway.py index ec0a66e8ed..bf111ac3d7 100644 --- a/hermes_cli/gateway.py +++ b/hermes_cli/gateway.py @@ -6441,15 +6441,15 @@ def run_gateway(verbose: int = 0, quiet: bool = False, replace: bool = False, fo _guard_existing_gateway_process_conflict(replace=replace) sys.path.insert(0, str(PROJECT_ROOT)) - # Startup-liveness watchdog (OOF-298): armed before config load, imports, - # and DB opens so a deadlock anywhere in the pre-event-loop window still - # gets the process respawned by the service supervisor instead of wedging - # as a live-PID zombie that s6/systemd will never restart. Disarmed by - # GatewayRunner once the event loop is confirmed live. Placed after the + # Startup-liveness watchdog (OOF-298), idempotent backstop: normal + # ``hermes gateway run`` invocations already armed in hermes_cli.main's + # argv fast-path (before the heavy import graph), but programmatic + # callers can enter run_gateway() directly. Placed after the # process-conflict guards: a --replace loser exiting above must not have - # armed a watchdog first. + # armed a watchdog first. Disarmed by GatewayRunner once the event loop + # is confirmed live. try: - from gateway.startup_watchdog import arm_startup_watchdog + from hermes_startup_watchdog import arm_startup_watchdog arm_startup_watchdog() except Exception: pass @@ -6636,6 +6636,15 @@ def run_gateway(verbose: int = 0, quiet: bool = False, replace: bool = False, fo _storm.window_s, _storm.backoff_s, ) + # The backoff sleep is intentional idle time — tell the startup + # watchdog (OOF-298) so it isn't mistaken for a parked deadlock + # and hard-exited mid-backoff (which would defeat the breaker). + try: + from gateway.startup_watchdog import kick_startup_watchdog + + kick_startup_watchdog(extra_s=_storm.backoff_s) + except Exception: + pass _time.sleep(_storm.backoff_s) except Exception as _be: logger.debug("respawn-storm breaker check failed (non-fatal): %s", _be) diff --git a/hermes_cli/main.py b/hermes_cli/main.py index bab6cf6f7e..d1f084d5b7 100644 --- a/hermes_cli/main.py +++ b/hermes_cli/main.py @@ -102,6 +102,24 @@ try: except Exception: pass +# Startup-liveness watchdog (OOF-298): for gateway runs, arm BEFORE the heavy +# module-level import graph below — an import-time deadlock (native-extension +# init, contended import lock) is exactly the "wedged before the event loop, +# no logs, live PID" class this watchdog exists for. ``hermes_startup_watchdog`` +# is stdlib-only, so importing it here cannot itself wedge on application +# code. argv sniffing is deliberately crude: over-arming is harmless (any +# non-gateway command either exits well inside the deadline or... should be +# covered anyway if it wedges), while under-arming recreates OOF-298. +# GatewayRunner disarms once the event loop is confirmed live. +if "gateway" in sys.argv[1:] and "run" in sys.argv[1:]: + try: + from hermes_startup_watchdog import arm_startup_watchdog as _arm_sw + + _arm_sw() + del _arm_sw + except Exception: + pass + def _exit_after_oneshot(rc: object) -> None: """Exit one-shot mode without letting late native finalizers change rc. diff --git a/hermes_startup_watchdog.py b/hermes_startup_watchdog.py new file mode 100644 index 0000000000..1ce307e9a6 --- /dev/null +++ b/hermes_startup_watchdog.py @@ -0,0 +1,471 @@ +"""Startup-liveness watchdog — respawn a gateway that wedges before its loop runs (OOF-298). + +The existing liveness backstops all assume startup succeeded: + +* 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. + +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. Before firing, the watchdog checks +whether the process consumed meaningful CPU time during the expired window +(``time.process_time()`` is process-wide). Long-but-legitimate synchronous +startup work — most importantly large ``state.db`` schema migrations, which +run inside ``SessionDB.__init__`` well before the loop starts and can +genuinely exceed any fixed deadline on multi-GB databases — burns CPU +continuously, so the deadline is extended (with a log line each time) for as +long as progress continues. The OOF-298 deadlock class parks every thread in +futex waits and accrues ~zero CPU, so it still fires on schedule. Known +limitation, documented deliberately: a *spinning* (busy-wait) startup +deadlock reads as CPU progress and will not fire; the observed incident class +is parked-thread deadlocks, which this catches. + +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. +""" + +from __future__ import annotations + +import faulthandler +import json +import logging +import os +import sys +import threading +import time +from datetime import datetime, timezone +from pathlib import Path +from typing import Any, Dict, Optional + +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. +SERVICE_RESTART_EXIT_CODE = 75 + +ENV_STARTUP_WATCHDOG = "HERMES_STARTUP_WATCHDOG" +ENV_STARTUP_WATCHDOG_TIMEOUT_S = "HERMES_STARTUP_WATCHDOG_TIMEOUT_S" + +_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. +_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 an extension. A parked +# futex deadlock accrues microseconds; a schema migration accrues orders of +# magnitude more than this per window even on slow disks. +_CPU_PROGRESS_MIN_S = 1.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). +_LEDGER_JOIN_TIMEOUT_S = 5.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. +_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. +_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. + """ + val = os.environ.get("HERMES_HOME", "").strip() + if val: + return Path(val) + if sys.platform == "win32": + local_appdata = os.environ.get("LOCALAPPDATA", "").strip() + base = Path(local_appdata) if local_appdata else Path.home() / "AppData" / "Local" + return base / "hermes" + return Path.home() / ".hermes" + + +def get_startup_watchdog_dump_path(home: Optional[Path] = None) -> Path: + """Return ``/logs/gateway-startup-watchdog.log``.""" + base = home if home is not None else _process_hermes_home() + return base.joinpath(*_DUMP_RELATIVE) + + +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 + + +def resolve_startup_watchdog_timeout() -> float: + """Deadline in seconds; env override, floor-clamped, default on garbage.""" + raw = os.environ.get(ENV_STARTUP_WATCHDOG_TIMEOUT_S, "").strip() + if not raw: + return DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S + try: + value = float(raw) + except ValueError: + logger.warning( + "Ignoring non-numeric %s=%r; using default %.0fs", + ENV_STARTUP_WATCHDOG_TIMEOUT_S, + raw, + DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S, + ) + return DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S + if value <= 0: + return DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S + 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.""" + 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") + except Exception: + logger.debug("Failed to write startup watchdog dump record", exc_info=True) + + +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. + """ + try: + from gateway.lifecycle_ledger import mark_exited + + mark_exited(exit_code, reason="startup_liveness_watchdog") + except Exception: + pass + + +class StartupWatchdogHandle: + """Disarm/inspect handle for the armed startup watchdog thread.""" + + def __init__(self, timeout_s: float, exit_code: int): + self.timeout_s = timeout_s + self.exit_code = exit_code + self.armed_at = time.monotonic() + self._state = _ARMED + self._state_lock = threading.Lock() + self._deadline = self.armed_at + timeout_s + self._disarmed_event = threading.Event() + self._thread: Optional[threading.Thread] = None + self._extensions = 0 + + 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). + """ + with self._state_lock: + if self._state == _ARMED: + self._state = _DISARMED + self._disarmed_event.set() + + 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. + """ + try: + extra = max(0.0, float(extra_s)) + except (TypeError, ValueError): + extra = 0.0 + with self._state_lock: + self._deadline = time.monotonic() + self.timeout_s + extra + + @property + def disarmed(self) -> bool: + return self._state == _DISARMED + + def is_alive(self) -> bool: + return self._thread is not None and self._thread.is_alive() + + def join(self, timeout: Optional[float] = None) -> None: + if self._thread is not None: + self._thread.join(timeout=timeout) + + # ── internals ──────────────────────────────────────────────────────── + + @staticmethod + def _process_cpu_seconds() -> Optional[float]: + """Process-wide CPU time (user+system, all threads); None on failure.""" + try: + return time.process_time() + except Exception: + return None + + def _fire(self) -> None: + elapsed = time.monotonic() - self.armed_at + try: + logger.critical( + "Gateway startup did not reach a live event loop within %.0fs " + "(elapsed %.0fs, %d extension(s)) and shows no CPU progress; " + "dumping all thread stacks and exiting with code %d so the " + "service supervisor can restart it (OOF-298).", + self.timeout_s, + elapsed, + self._extensions, + self.exit_code, + ) + except Exception: + pass + _write_dump_record( + { + "ts": datetime.now(timezone.utc).isoformat(), + "tag": "startup_watchdog.fired", + "pid": os.getpid(), + "timeout_s": self.timeout_s, + "elapsed_s": round(elapsed, 3), + "extensions": self._extensions, + "exit_code": self.exit_code, + } + ) + try: + 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). + try: + ledger_thread = threading.Thread( + target=_mark_lifecycle_exit, + args=(self.exit_code,), + daemon=True, + name="gateway-startup-watchdog-ledger", + ) + ledger_thread.start() + ledger_thread.join(timeout=_LEDGER_JOIN_TIMEOUT_S) + except Exception: + pass + self._exit(self.exit_code) + + @staticmethod + def _exit(code: int) -> None: + """Seam for tests; production is a bare ``os._exit``.""" + os._exit(code) + + def _run(self) -> None: + last_cpu = self._process_cpu_seconds() + while True: + with self._state_lock: + if self._state != _ARMED: + return + deadline = self._deadline + remaining = deadline - time.monotonic() + if remaining > 0: + if self._disarmed_event.wait(timeout=min(remaining, _POLL_SLICE_S)): + return + continue + # Deadline expired. A slow-but-alive startup (large state.db + # schema migration inside SessionDB.__init__) burns CPU + # continuously; a parked futex deadlock accrues ~none. Extend + # for the former, fire only for the latter. + cpu = self._process_cpu_seconds() + if ( + cpu is not None + and last_cpu is not None + and (cpu - last_cpu) >= _CPU_PROGRESS_MIN_S + ): + 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 + try: + logger.warning( + "Gateway startup exceeded %.0fs but is consuming CPU " + "(%.1fs this window) — likely a long schema migration; " + "extending the startup watchdog deadline (extension #%d).", + self.timeout_s, + window_delta, + self._extensions, + ) + except Exception: + pass + continue + # No progress: claim the fire transition atomically so a disarm + # racing this exact moment can still win if it gets there first. + with self._state_lock: + if self._state != _ARMED: + return + self._state = _FIRING + self._fire() + return + + def _start(self) -> bool: + thread = threading.Thread( + target=self._run, + daemon=True, + name="gateway-startup-watchdog", + ) + try: + thread.start() + except Exception: + logger.debug("Failed to start gateway startup watchdog", exc_info=True) + return False + self._thread = thread + return True + + +def arm_startup_watchdog( + timeout_s: Optional[float] = None, + *, + exit_code: int = SERVICE_RESTART_EXIT_CODE, +) -> Optional[StartupWatchdogHandle]: + """Arm the process-wide startup watchdog. Idempotent; never raises. + + Returns the (possibly pre-existing) handle, or ``None`` when disabled via + ``HERMES_STARTUP_WATCHDOG=0`` or when the thread could not be started. + """ + global _handle + try: + if startup_watchdog_disabled(): + return None + with _handle_lock: + if _handle is not None and _handle.is_alive(): + return _handle + resolved = ( + float(timeout_s) + if timeout_s is not None and float(timeout_s) > 0 + else resolve_startup_watchdog_timeout() + ) + handle = StartupWatchdogHandle(resolved, exit_code) + if not handle._start(): + return None + _handle = handle + return handle + except Exception: + logger.debug("Failed to arm gateway startup watchdog", exc_info=True) + return None + + +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. + """ + global _handle + try: + with _handle_lock: + handle = _handle + _handle = None + if handle is not None: + handle.disarm() + except Exception: + 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. + """ + try: + with _handle_lock: + handle = _handle + if handle is not None: + handle.kick(extra_s) + except Exception: + logger.debug("Failed to kick gateway startup watchdog", exc_info=True) + + +def _reset_for_tests() -> None: + """Drop the module singleton (test isolation only).""" + global _handle + with _handle_lock: + handle = _handle + _handle = None + if handle is not None: + handle.disarm() diff --git a/pyproject.toml b/pyproject.toml index d0b1dba14d..f455c40ed5 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -576,6 +576,7 @@ py-modules = [ "hermes_state_portability", "hermes_state_schema", "hermes_state_search", + "hermes_startup_watchdog", "hermes_time", "hermes_logging", "utils", diff --git a/scripts/hermes-gateway b/scripts/hermes-gateway index b0d45810e3..2d39cfbe9d 100755 --- a/scripts/hermes-gateway +++ b/scripts/hermes-gateway @@ -294,6 +294,13 @@ def is_windows() -> bool: def run_gateway(): """Run the gateway in foreground.""" + # Startup-liveness watchdog (OOF-298): arm before importing the gateway + # graph so an import-time or pre-loop deadlock still gets respawned. + try: + from hermes_startup_watchdog import arm_startup_watchdog + arm_startup_watchdog() + except Exception: + pass from gateway.run import start_gateway print("Starting Hermes Gateway...") print("Press Ctrl+C to stop.") diff --git a/tests/gateway/test_startup_watchdog.py b/tests/gateway/test_startup_watchdog.py index 192c37061f..d4fb6671c4 100644 --- a/tests/gateway/test_startup_watchdog.py +++ b/tests/gateway/test_startup_watchdog.py @@ -1,10 +1,14 @@ """Startup-liveness watchdog tests (OOF-298). -The watchdog covers the pre-event-loop window: armed at process entry, -disarmed once the gateway's asyncio loop is confirmed live. If neither -happens within the deadline it must dump diagnostics, record a lifecycle -exit, and hard-exit with the service-restart code so the supervisor -respawns the process instead of babysitting a live-PID zombie. +The watchdog covers the pre-event-loop window: armed at process entry +(before the gateway package imports — the implementation is the stdlib-only +top-level module ``hermes_startup_watchdog``; ``gateway.startup_watchdog`` +is a re-export shim), disarmed once the gateway's asyncio loop is confirmed +live. If neither happens within the deadline — and the process shows no CPU +progress, so slow-but-alive schema migrations are exempt — it must dump +diagnostics, record a lifecycle exit, and hard-exit with the service-restart +code so the supervisor respawns the process instead of babysitting a +live-PID zombie. """ from __future__ import annotations @@ -16,13 +20,14 @@ from pathlib import Path import pytest -import gateway.startup_watchdog as sw -from gateway.restart import GATEWAY_SERVICE_RESTART_EXIT_CODE -from gateway.startup_watchdog import ( +import hermes_startup_watchdog as sw +from hermes_startup_watchdog import ( + SERVICE_RESTART_EXIT_CODE, StartupWatchdogHandle, arm_startup_watchdog, disarm_startup_watchdog, get_startup_watchdog_dump_path, + kick_startup_watchdog, resolve_startup_watchdog_timeout, startup_watchdog_disabled, ) @@ -39,6 +44,18 @@ def _isolate(tmp_path, monkeypatch): sw._reset_for_tests() +@pytest.fixture(autouse=True) +def _no_cpu_progress(monkeypatch): + """Freeze process CPU time so fire tests never see 'progress'. + + Individual tests that exercise the CPU-progress extension override this + with their own sequence. + """ + monkeypatch.setattr( + StartupWatchdogHandle, "_process_cpu_seconds", staticmethod(lambda: 0.0) + ) + + class _ExitCapture: """Replaces StartupWatchdogHandle._exit so _fire() cannot kill pytest.""" @@ -58,6 +75,56 @@ def exit_capture(monkeypatch): return capture +class TestContracts: + def test_restart_code_parity_with_gateway_restart(self): + """The stdlib-only module duplicates the exit-code constant; keep it + in lockstep with the canonical gateway.restart definition.""" + from gateway.restart import GATEWAY_SERVICE_RESTART_EXIT_CODE + + assert SERVICE_RESTART_EXIT_CODE == GATEWAY_SERVICE_RESTART_EXIT_CODE + + def test_gateway_shim_reexports_same_objects(self): + import gateway.startup_watchdog as shim + + assert shim.arm_startup_watchdog is arm_startup_watchdog + assert shim.disarm_startup_watchdog is disarm_startup_watchdog + assert shim.kick_startup_watchdog is kick_startup_watchdog + + def test_implementation_module_is_stdlib_only(self): + """Import-lightness is a correctness property (arm-before-imports, + no import-lock dependence at fire time): the implementation module + must not import the gateway/agent/hermes_cli graphs at module level.""" + import ast + import inspect + + source = inspect.getsource(sw) + tree = ast.parse(source) + forbidden_roots = { + "gateway", + "agent", + "hermes_cli", + "hermes_state", + "hermes_constants", + "tools", + "plugins", + } + offenders = [] + for node in ast.walk(tree): + # Only module-level and unconditional imports matter; function- + # bodied imports (the ledger helper) are deliberate and guarded. + if isinstance(node, ast.Import): + names = [alias.name for alias in node.names] + elif isinstance(node, ast.ImportFrom): + names = [node.module or ""] + else: + continue + for name in names: + root = name.split(".")[0] + if root in forbidden_roots and node.col_offset == 0: + offenders.append(name) + assert offenders == [] + + class TestConfigResolution: def test_default_timeout(self): assert resolve_startup_watchdog_timeout() == sw.DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S @@ -135,29 +202,123 @@ class TestArmDisarm: assert second.is_alive() disarm_startup_watchdog() + def test_disarm_after_deadline_expiry_wins_over_fire(self, exit_capture, monkeypatch): + """The P2 race from review: deadline expires, but disarm lands before + the fire transition claims the state — the disarm must win. We force + the interleaving by blocking the watchdog thread inside the CPU probe + (which runs after deadline expiry, before the fire transition). The + probe is also called once at thread start for the baseline, so only + the second call blocks.""" + in_probe = threading.Event() + release_probe = threading.Event() + calls = {"n": 0} + + def _blocking_probe(): + calls["n"] += 1 + if calls["n"] >= 2: + in_probe.set() + release_probe.wait(timeout=10) + return 0.0 + + monkeypatch.setattr( + StartupWatchdogHandle, + "_process_cpu_seconds", + staticmethod(_blocking_probe), + ) + handle = arm_startup_watchdog(timeout_s=0.1) + assert handle is not None + # Wait until the deadline has expired and the thread is inside the + # probe (post-expiry, pre-fire-transition). + assert in_probe.wait(timeout=5) + disarm_startup_watchdog() + release_probe.set() + handle.join(timeout=5) + assert not exit_capture.fired.is_set() + assert exit_capture.codes == [] + + +class TestKick: + def test_kick_extends_deadline(self, exit_capture): + handle = arm_startup_watchdog(timeout_s=0.3) + assert handle is not None + # Kick far enough out that the original 0.3s deadline can't fire + # while we watch. + kick_startup_watchdog(extra_s=60) + time.sleep(0.6) + assert not exit_capture.fired.is_set() + disarm_startup_watchdog() + + def test_kick_without_arm_is_safe(self): + kick_startup_watchdog(extra_s=30) # must not raise + + def test_kick_with_garbage_extra_is_safe(self): + arm_startup_watchdog(timeout_s=60) + kick_startup_watchdog(extra_s="nonsense") # type: ignore[arg-type] + disarm_startup_watchdog() + + +class TestCpuProgressExtension: + def test_cpu_progress_extends_instead_of_firing(self, exit_capture, monkeypatch): + """A long schema migration burns CPU: the watchdog must extend, not + fire (the P1 false-fire/restart-loop case from review).""" + # Each probe call reports +10s CPU — always 'progress'. + counter = {"cpu": 0.0} + + def _busy_probe(): + counter["cpu"] += 10.0 + return counter["cpu"] + + monkeypatch.setattr( + StartupWatchdogHandle, "_process_cpu_seconds", staticmethod(_busy_probe) + ) + handle = arm_startup_watchdog(timeout_s=0.1) + assert handle is not None + time.sleep(0.6) + assert not exit_capture.fired.is_set() + assert handle._extensions >= 1 + disarm_startup_watchdog() + + def test_no_cpu_progress_fires(self, exit_capture): + # autouse fixture pins CPU time at 0.0 — no progress. + arm_startup_watchdog(timeout_s=0.1) + assert exit_capture.fired.wait(timeout=5) + assert exit_capture.codes == [SERVICE_RESTART_EXIT_CODE] + + def test_probe_failure_fails_toward_firing(self, exit_capture, monkeypatch): + """If CPU time can't be read the watchdog must still fire on a real + deadlock rather than extending forever.""" + monkeypatch.setattr( + StartupWatchdogHandle, "_process_cpu_seconds", staticmethod(lambda: None) + ) + arm_startup_watchdog(timeout_s=0.1) + assert exit_capture.fired.wait(timeout=5) + class TestFire: - def test_fires_after_deadline_with_restart_code(self, exit_capture, tmp_path): + def test_fires_after_deadline_with_restart_code(self, exit_capture): handle = arm_startup_watchdog(timeout_s=0.1) assert handle is not None assert exit_capture.fired.wait(timeout=5) - assert exit_capture.codes == [GATEWAY_SERVICE_RESTART_EXIT_CODE] + assert exit_capture.codes == [SERVICE_RESTART_EXIT_CODE] - def test_fire_writes_dump_record(self, exit_capture, tmp_path): + def test_fire_writes_dump_record_and_stacks(self, exit_capture, tmp_path): arm_startup_watchdog(timeout_s=0.1) assert exit_capture.fired.wait(timeout=5) dump_path = get_startup_watchdog_dump_path(tmp_path) - # The record write happens before _exit; poll briefly for the file. deadline = time.monotonic() + 2 while not dump_path.exists() and time.monotonic() < deadline: time.sleep(0.02) assert dump_path.exists() - record = json.loads(dump_path.read_text(encoding="utf-8").splitlines()[0]) + content = dump_path.read_text(encoding="utf-8") + record = json.loads(content.splitlines()[0]) assert record["tag"] == "startup_watchdog.fired" - assert record["exit_code"] == GATEWAY_SERVICE_RESTART_EXIT_CODE + assert record["exit_code"] == SERVICE_RESTART_EXIT_CODE assert record["timeout_s"] == pytest.approx(0.1) + # File-based faulthandler dump follows the JSON record (stderr may be + # absent on detached runs). + assert "Thread" in content or "Current thread" in content - def test_fire_marks_lifecycle_exit(self, exit_capture, tmp_path, monkeypatch): + def test_fire_marks_lifecycle_exit(self, exit_capture, monkeypatch): marked = {} def _fake_mark_exited(code, reason=None): @@ -169,10 +330,10 @@ class TestFire: monkeypatch.setattr(ledger, "mark_exited", _fake_mark_exited) arm_startup_watchdog(timeout_s=0.1) assert exit_capture.fired.wait(timeout=5) - # mark_exited runs just before _exit on the same thread; once fired - # is set the _exit stub has returned, so mark_exited already ran. + # The ledger write runs on a helper thread joined (with timeout) + # before _exit; once fired is set the join already happened. assert marked == { - "code": GATEWAY_SERVICE_RESTART_EXIT_CODE, + "code": SERVICE_RESTART_EXIT_CODE, "reason": "startup_liveness_watchdog", } @@ -182,16 +343,6 @@ class TestFire: assert exit_capture.codes == [42] -class TestFireTimeoutClamp: - def test_explicit_timeout_below_floor_still_used_directly(self, exit_capture): - """arm_startup_watchdog(timeout_s=...) is a trusted caller/test seam — - it bypasses the env floor clamp so tests stay fast. Only env-provided - values are clamped (they come from operators).""" - handle = arm_startup_watchdog(timeout_s=0.1) - assert handle is not None - assert handle.timeout_s == pytest.approx(0.1) - - class TestDumpPath: def test_dump_path_under_home(self, tmp_path): assert get_startup_watchdog_dump_path(tmp_path) == ( @@ -199,7 +350,6 @@ class TestDumpPath: ) def test_dump_write_failure_is_swallowed(self, monkeypatch): - # Point the dump at an unwritable location; must not raise. monkeypatch.setattr( sw, "get_startup_watchdog_dump_path", lambda home=None: Path("/dev/null/nope") )