diff --git a/gateway/run.py b/gateway/run.py index b235253e3f..96d799750f 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -13481,6 +13481,17 @@ 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 + + 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 @@ -33421,7 +33432,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: @@ -33435,6 +33446,17 @@ 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. + try: + from gateway.startup_watchdog import arm_startup_watchdog + arm_startup_watchdog() + except Exception: + pass + # Force UTF-8 stdio on Windows — gateway logs and startup banner would # otherwise UnicodeEncodeError on cp1252 consoles. No-op on POSIX. try: diff --git a/gateway/startup_watchdog.py b/gateway/startup_watchdog.py new file mode 100644 index 0000000000..eed5fa04eb --- /dev/null +++ b/gateway/startup_watchdog.py @@ -0,0 +1,274 @@ +"""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 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. +""" + +from __future__ import annotations + +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() diff --git a/hermes_cli/gateway.py b/hermes_cli/gateway.py index 351d03b86e..ec0a66e8ed 100644 --- a/hermes_cli/gateway.py +++ b/hermes_cli/gateway.py @@ -6441,6 +6441,19 @@ 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 + # process-conflict guards: a --replace loser exiting above must not have + # armed a watchdog first. + try: + from gateway.startup_watchdog import arm_startup_watchdog + arm_startup_watchdog() + except Exception: + pass + # Detached Windows gateway runs must ignore console-control broadcasts # from sibling CLI processes, but foreground `hermes gateway run` still # needs to obey the banner's "Press Ctrl+C to stop" contract. diff --git a/tests/gateway/test_startup_watchdog.py b/tests/gateway/test_startup_watchdog.py new file mode 100644 index 0000000000..192c37061f --- /dev/null +++ b/tests/gateway/test_startup_watchdog.py @@ -0,0 +1,206 @@ +"""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. +""" + +from __future__ import annotations + +import json +import threading +import time +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 ( + StartupWatchdogHandle, + arm_startup_watchdog, + disarm_startup_watchdog, + get_startup_watchdog_dump_path, + resolve_startup_watchdog_timeout, + startup_watchdog_disabled, +) + + +@pytest.fixture(autouse=True) +def _isolate(tmp_path, monkeypatch): + """Every test gets a fresh singleton and its own HERMES_HOME.""" + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + monkeypatch.delenv(sw.ENV_STARTUP_WATCHDOG, raising=False) + monkeypatch.delenv(sw.ENV_STARTUP_WATCHDOG_TIMEOUT_S, raising=False) + sw._reset_for_tests() + yield + sw._reset_for_tests() + + +class _ExitCapture: + """Replaces StartupWatchdogHandle._exit so _fire() cannot kill pytest.""" + + def __init__(self): + self.codes: list[int] = [] + self.fired = threading.Event() + + def __call__(self, code: int) -> None: + self.codes.append(code) + self.fired.set() + + +@pytest.fixture +def exit_capture(monkeypatch): + capture = _ExitCapture() + monkeypatch.setattr(StartupWatchdogHandle, "_exit", staticmethod(capture)) + return capture + + +class TestConfigResolution: + def test_default_timeout(self): + assert resolve_startup_watchdog_timeout() == sw.DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S + + def test_env_override(self, monkeypatch): + monkeypatch.setenv(sw.ENV_STARTUP_WATCHDOG_TIMEOUT_S, "120") + assert resolve_startup_watchdog_timeout() == 120.0 + + def test_env_override_clamped_to_floor(self, monkeypatch): + monkeypatch.setenv(sw.ENV_STARTUP_WATCHDOG_TIMEOUT_S, "5") + assert resolve_startup_watchdog_timeout() == sw._MIN_TIMEOUT_S + + def test_garbage_env_falls_back_to_default(self, monkeypatch): + monkeypatch.setenv(sw.ENV_STARTUP_WATCHDOG_TIMEOUT_S, "soon") + assert resolve_startup_watchdog_timeout() == sw.DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S + + def test_nonpositive_env_falls_back_to_default(self, monkeypatch): + monkeypatch.setenv(sw.ENV_STARTUP_WATCHDOG_TIMEOUT_S, "-1") + assert resolve_startup_watchdog_timeout() == sw.DEFAULT_STARTUP_WATCHDOG_TIMEOUT_S + + @pytest.mark.parametrize("raw", ["0", "false", "no", "off", "FALSE", "Off"]) + def test_disabled_values(self, monkeypatch, raw): + monkeypatch.setenv(sw.ENV_STARTUP_WATCHDOG, raw) + assert startup_watchdog_disabled() is True + + @pytest.mark.parametrize("raw", ["", "1", "true", "yes"]) + def test_enabled_values(self, monkeypatch, raw): + monkeypatch.setenv(sw.ENV_STARTUP_WATCHDOG, raw) + assert startup_watchdog_disabled() is False + + +class TestArmDisarm: + def test_arm_returns_live_handle(self): + handle = arm_startup_watchdog(timeout_s=60) + assert handle is not None + assert handle.is_alive() + assert not handle.disarmed + disarm_startup_watchdog() + handle.join(timeout=2) + assert not handle.is_alive() + + def test_arm_is_idempotent(self): + first = arm_startup_watchdog(timeout_s=60) + second = arm_startup_watchdog(timeout_s=60) + assert first is second + disarm_startup_watchdog() + + def test_disarm_prevents_fire(self, exit_capture): + handle = arm_startup_watchdog(timeout_s=0.2) + assert handle is not None + disarm_startup_watchdog() + handle.join(timeout=2) + assert not exit_capture.fired.is_set() + assert exit_capture.codes == [] + + def test_disarm_without_arm_is_safe(self): + disarm_startup_watchdog() # must not raise + + def test_disarm_is_idempotent(self): + arm_startup_watchdog(timeout_s=60) + disarm_startup_watchdog() + disarm_startup_watchdog() # must not raise + + def test_disabled_via_env(self, monkeypatch): + monkeypatch.setenv(sw.ENV_STARTUP_WATCHDOG, "0") + assert arm_startup_watchdog(timeout_s=60) is None + + def test_rearm_after_disarm_starts_fresh_thread(self): + first = arm_startup_watchdog(timeout_s=60) + disarm_startup_watchdog() + first.join(timeout=2) + second = arm_startup_watchdog(timeout_s=60) + assert second is not None + assert second is not first + assert second.is_alive() + disarm_startup_watchdog() + + +class TestFire: + def test_fires_after_deadline_with_restart_code(self, exit_capture, tmp_path): + 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] + + def test_fire_writes_dump_record(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]) + assert record["tag"] == "startup_watchdog.fired" + assert record["exit_code"] == GATEWAY_SERVICE_RESTART_EXIT_CODE + assert record["timeout_s"] == pytest.approx(0.1) + + def test_fire_marks_lifecycle_exit(self, exit_capture, tmp_path, monkeypatch): + marked = {} + + def _fake_mark_exited(code, reason=None): + marked["code"] = code + marked["reason"] = reason + + import gateway.lifecycle_ledger as ledger + + 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. + assert marked == { + "code": GATEWAY_SERVICE_RESTART_EXIT_CODE, + "reason": "startup_liveness_watchdog", + } + + def test_custom_exit_code(self, exit_capture): + arm_startup_watchdog(timeout_s=0.1, exit_code=42) + assert exit_capture.fired.wait(timeout=5) + 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) == ( + tmp_path / "logs" / "gateway-startup-watchdog.log" + ) + + 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") + ) + sw._write_dump_record({"tag": "x"})