mirror of
https://github.com/NousResearch/hermes-agent.git
synced 2026-07-21 16:18:55 +00:00
fix(gateway): launchd ThrottleInterval + portable respawn-storm circuit breaker
Re-applied against v0.18.2. Still unconditional KeepAlive <true/> with no ThrottleInterval, and upstream's new restart_loop_guard only skips session auto-resume (never throttles the boot), so the launchd respawn storm is unguarded. - launchd plist: add ThrottleInterval=30 + ExitTimeOut=25. - status: record_start_and_check_storm — portable breaker (atomic gateway-starts.log, backs off when too many starts land in a window); run_gateway sleeps it before asyncio.run. Env HERMES_GATEWAY_MAX_STARTS (<=0 disables) / _START_WINDOW_S. Separate file from upstream's restart_loop.json — no collision. - tests: breaker (threshold/prune/atomic) + plist throttle keys. Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
71f4de3cd8
commit
d920954747
3 changed files with 185 additions and 2 deletions
|
|
@ -13,6 +13,7 @@ concurrently under distinct configurations).
|
|||
|
||||
import hashlib
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import shlex
|
||||
import signal
|
||||
|
|
@ -23,7 +24,7 @@ import time
|
|||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from hermes_constants import get_hermes_home, _get_platform_default_hermes_home
|
||||
from typing import Any, Optional
|
||||
from typing import Any, NamedTuple, Optional
|
||||
from utils import atomic_json_write
|
||||
|
||||
if sys.platform == "win32":
|
||||
|
|
@ -46,6 +47,84 @@ _GATEWAY_RUNNING_PID_CACHE_TTL_SECONDS = 1.0
|
|||
_gateway_running_pid_cache_lock = threading.Lock()
|
||||
_gateway_running_pid_cache: dict[tuple[str, bool, bool], tuple[float, tuple[Any, ...], Optional[int]]] = {}
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class StormInfo(NamedTuple):
|
||||
"""Result of a respawn-storm check: how many starts, over what window, and
|
||||
the backoff the caller should sleep to break the storm."""
|
||||
|
||||
count: int
|
||||
window_s: float
|
||||
backoff_s: float
|
||||
|
||||
|
||||
def _get_starts_log_path() -> Path:
|
||||
"""Path to the append-only gateway-start ledger used by the respawn-storm
|
||||
breaker. Distinct from ``restart_loop.json`` (the auto-resume guard) — no
|
||||
collision."""
|
||||
return get_hermes_home() / "gateway-starts.log"
|
||||
|
||||
|
||||
def record_start_and_check_storm(
|
||||
max_starts: int = 5, window_s: float = 120.0, *, backoff_cap_s: float = 300.0
|
||||
) -> Optional[StormInfo]:
|
||||
"""Record this gateway start and report whether a respawn storm is underway.
|
||||
|
||||
Appends the current UTC timestamp to the starts-log, prunes entries older
|
||||
than ``window_s``, and ring-buffers the file so it can't grow unbounded.
|
||||
Returns a :class:`StormInfo` when more than ``max_starts`` starts landed in
|
||||
the window (with an exponential backoff capped at ``backoff_cap_s``), else
|
||||
``None``.
|
||||
|
||||
Best-effort: any bookkeeping failure is logged and swallowed so a broken
|
||||
ledger can never crash gateway startup.
|
||||
"""
|
||||
try:
|
||||
path = _get_starts_log_path()
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
now = datetime.now(timezone.utc).timestamp()
|
||||
|
||||
existing: list[float] = []
|
||||
if path.exists():
|
||||
for line in path.read_text(encoding="utf-8").splitlines():
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
existing.append(float(line))
|
||||
except ValueError:
|
||||
continue
|
||||
|
||||
existing.append(now)
|
||||
|
||||
# Keep only starts within the sliding window for the storm decision.
|
||||
recent = [ts for ts in existing if now - ts <= window_s]
|
||||
|
||||
# Ring-buffer what we persist so the file stays bounded even if the
|
||||
# window is wide or starts are frequent.
|
||||
keep = max(max_starts * 4, 40)
|
||||
to_write = existing[-keep:]
|
||||
|
||||
tmp = path.with_suffix(".tmp")
|
||||
tmp.write_text(
|
||||
"\n".join(repr(ts) for ts in to_write) + "\n", encoding="utf-8"
|
||||
)
|
||||
os.replace(tmp, path)
|
||||
|
||||
if len(recent) > max_starts:
|
||||
backoff = min(
|
||||
backoff_cap_s, 5.0 * (2 ** min(len(recent) - max_starts, 6))
|
||||
)
|
||||
return StormInfo(count=len(recent), window_s=window_s, backoff_s=backoff)
|
||||
return None
|
||||
except Exception as _e:
|
||||
logger.debug(
|
||||
"respawn-storm breaker bookkeeping failed (non-fatal): %s", _e
|
||||
)
|
||||
return None
|
||||
|
||||
|
||||
def _get_process_hermes_home() -> Path:
|
||||
"""Return the process-level HERMES_HOME, skipping context-local overrides.
|
||||
|
|
|
|||
|
|
@ -3982,7 +3982,17 @@ def generate_launchd_plist() -> str:
|
|||
|
||||
<key>KeepAlive</key>
|
||||
<true/>
|
||||
|
||||
|
||||
<!-- ThrottleInterval raises launchd's default 10s minimum respawn interval
|
||||
to 30s so a crash-looping gateway can't hammer launchd into a rapid
|
||||
respawn storm; ExitTimeOut gives the gateway 25s of graceful-drain
|
||||
headroom before launchd escalates from SIGTERM to SIGKILL on stop. -->
|
||||
<key>ThrottleInterval</key>
|
||||
<integer>30</integer>
|
||||
|
||||
<key>ExitTimeOut</key>
|
||||
<integer>25</integer>
|
||||
|
||||
<key>StandardOutPath</key>
|
||||
<string>{log_dir}/gateway.log</string>
|
||||
|
||||
|
|
@ -4856,6 +4866,39 @@ def run_gateway(verbose: int = 0, quiet: bool = False, replace: bool = False, fo
|
|||
|
||||
_atexit.register(_atexit_hook)
|
||||
|
||||
# Portable, app-level respawn-storm circuit breaker. launchd/systemd have
|
||||
# their own throttles, but this backstop works on every platform (and covers
|
||||
# supervisors that lack a respawn floor). Set HERMES_GATEWAY_MAX_STARTS<=0 to
|
||||
# disable. Kept best-effort: a bookkeeping failure must never block startup.
|
||||
try:
|
||||
import time as _time
|
||||
|
||||
from gateway.status import record_start_and_check_storm
|
||||
|
||||
try:
|
||||
_max_starts = int(os.getenv("HERMES_GATEWAY_MAX_STARTS", "5"))
|
||||
except ValueError:
|
||||
_max_starts = 5
|
||||
try:
|
||||
_win = float(os.getenv("HERMES_GATEWAY_START_WINDOW_S", "120"))
|
||||
except ValueError:
|
||||
_win = 120.0
|
||||
_storm = (
|
||||
record_start_and_check_storm(max_starts=_max_starts, window_s=_win)
|
||||
if _max_starts > 0
|
||||
else None
|
||||
)
|
||||
if _storm is not None:
|
||||
logger.warning(
|
||||
"Gateway (re)started %d times in %.0fs — backing off %.0fs to break a respawn storm.",
|
||||
_storm.count,
|
||||
_storm.window_s,
|
||||
_storm.backoff_s,
|
||||
)
|
||||
_time.sleep(_storm.backoff_s)
|
||||
except Exception as _be:
|
||||
logger.debug("respawn-storm breaker check failed (non-fatal): %s", _be)
|
||||
|
||||
success = False
|
||||
try:
|
||||
success = asyncio.run(start_gateway(replace=replace, verbosity=verbosity))
|
||||
|
|
|
|||
|
|
@ -3,6 +3,7 @@
|
|||
import json
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
from pathlib import Path
|
||||
from types import SimpleNamespace
|
||||
|
||||
|
|
@ -1660,3 +1661,63 @@ class TestGatewayBusyDerivation:
|
|||
assert status.derive_gateway_drainable(
|
||||
gateway_running=True, gateway_state=state
|
||||
) is False, state
|
||||
|
||||
|
||||
class TestRespawnStormBreaker:
|
||||
def test_no_storm_under_threshold(self, tmp_path, monkeypatch):
|
||||
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
|
||||
for _ in range(5):
|
||||
result = status.record_start_and_check_storm(
|
||||
max_starts=5, window_s=120.0
|
||||
)
|
||||
assert result is None
|
||||
|
||||
def test_storm_detected_over_threshold(self, tmp_path, monkeypatch):
|
||||
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
|
||||
results = [
|
||||
status.record_start_and_check_storm(max_starts=5, window_s=120.0)
|
||||
for _ in range(7)
|
||||
]
|
||||
last = results[-1]
|
||||
assert last is not None
|
||||
assert last.count >= 6
|
||||
assert last.backoff_s > 0
|
||||
|
||||
def test_old_starts_pruned_outside_window(self, tmp_path, monkeypatch):
|
||||
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
|
||||
log_path = status._get_starts_log_path()
|
||||
old = time.time() - 10000
|
||||
log_path.write_text(
|
||||
"\n".join(repr(old) for _ in range(10)) + "\n", encoding="utf-8"
|
||||
)
|
||||
# All seeded entries are far outside the window, so a single new start
|
||||
# cannot exceed the threshold.
|
||||
result = status.record_start_and_check_storm(max_starts=5, window_s=120.0)
|
||||
assert result is None
|
||||
|
||||
def test_starts_log_written_atomically(self, tmp_path, monkeypatch):
|
||||
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
|
||||
for _ in range(6):
|
||||
status.record_start_and_check_storm(max_starts=5, window_s=120.0)
|
||||
|
||||
# No leftover temp file from the atomic write.
|
||||
assert not list(tmp_path.glob("*.tmp"))
|
||||
|
||||
log_path = status._get_starts_log_path()
|
||||
assert log_path.exists()
|
||||
for line in log_path.read_text(encoding="utf-8").splitlines():
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
float(line) # every persisted line must parse as a float
|
||||
|
||||
|
||||
class TestLaunchdPlistRespawnGovernance:
|
||||
def test_plist_has_throttle_interval(self, tmp_path, monkeypatch):
|
||||
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
|
||||
from hermes_cli.gateway import generate_launchd_plist
|
||||
|
||||
plist = generate_launchd_plist()
|
||||
assert "<key>ThrottleInterval</key>" in plist
|
||||
assert "<key>ExitTimeOut</key>" in plist
|
||||
assert "<key>KeepAlive</key>" in plist
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue