fix(gateway): address startup-watchdog review findings (OOF-298, PR #89750)
Independent review of the initial startup-liveness watchdog surfaced two
P1s and three P2s. All are addressed here.
P1 — legitimate slow startups (large state.db schema migrations inside
SessionDB.__init__, which run synchronously before the loop starts) could
exceed the fixed 300s deadline and restart-loop. The watchdog now checks
process CPU time (time.process_time(), process-wide) when the deadline
expires: continuous CPU consumption means a live migration, so the deadline
is extended (with a warning log per extension). The OOF-298 deadlock class
parks every thread in futex waits and accrues ~zero CPU, so it still fires
on schedule. Documented limitation: a spinning busy-wait deadlock reads as
progress and won't fire — the observed incident class is parked threads.
P1 — import-time deadlocks were outside coverage. The implementation moved
to a stdlib-only top-level module (hermes_startup_watchdog), and
hermes_cli/main.py arms it via an argv fast-path ("gateway" + "run" in
argv) BEFORE the heavy module-level import graph. gateway/startup_watchdog
remains as a re-export shim so the intuitive import path keeps working for
the disarm site, tests, and REPL use. Import-lightness is a correctness
property, tested via AST inspection: at fire time the wedged main thread
may hold the import lock, so the fire path performs no imports on its own
thread — the lifecycle-ledger write runs on a bounded-join helper thread
and os._exit happens regardless.
P2 — disarm/fire race: the handle now has an explicit state machine
(armed → disarmed | firing) guarded by a lock; whichever transition takes
the lock first wins, so a disarm landing after deadline expiry but before
the fire transition is honored. Regression test forces the exact
interleaving by blocking inside the CPU probe.
P2 — uncovered entry points: cli.py --gateway and scripts/hermes-gateway
run_gateway() now arm the watchdog before importing the gateway graph.
hermes_cli/gateway.py run_gateway() keeps an idempotent backstop arm for
programmatic callers.
P2 — respawn-storm backoff interaction: the storm breaker's intentional
backoff sleep (up to minutes, ~zero CPU — indistinguishable from a parked
deadlock) now calls kick_startup_watchdog(extra_s=backoff) so the deadline
is pushed past the sleep instead of firing mid-backoff.
Also: the faulthandler stack dump is now additionally written to
logs/gateway-startup-watchdog.log (stderr may be absent on detached/
windowless runs); the disarm site in gateway/run.py moved inside the
loop-confirmed branch (if the loop is NOT live, the milestone was not
reached and the watchdog must stay armed); hermes_startup_watchdog added
to pyproject py-modules so sealed venvs ship it; SERVICE_RESTART_EXIT_CODE
is duplicated in the stdlib-only module with a parity test against
gateway.restart.
Tests: 38 in tests/gateway/test_startup_watchdog.py (contracts incl.
stdlib-only AST check and shim re-export identity, config resolution,
arm/disarm/kick, CPU-progress extension vs no-progress fire, probe-failure
fails toward firing, disarm-vs-fire race, dump record + file stacks,
lifecycle ledger, custom exit code).
This commit is contained in:
@@ -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())
|
||||
|
||||
+21
-16
@@ -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()
|
||||
|
||||
+35
-269
@@ -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 ``<HERMES_HOME>/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",
|
||||
]
|
||||
|
||||
+16
-7
@@ -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)
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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 ``<HERMES_HOME>/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()
|
||||
@@ -576,6 +576,7 @@ py-modules = [
|
||||
"hermes_state_portability",
|
||||
"hermes_state_schema",
|
||||
"hermes_state_search",
|
||||
"hermes_startup_watchdog",
|
||||
"hermes_time",
|
||||
"hermes_logging",
|
||||
"utils",
|
||||
|
||||
@@ -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.")
|
||||
|
||||
@@ -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")
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user