refactor(gateway): tier tables for memory/disk pressure, shared ISO parser for drain markers, hand-compacted docstrings

This commit is contained in:
Teknium
2026-09-02 18:54:58 -07:00
parent 5fa7066054
commit 3238af61c1
15 changed files with 256 additions and 431 deletions

View File

@@ -1,12 +1,11 @@
"""Memory-pressure bounds for the gateway's per-session AIAgent cache.
Each cached ``AIAgent`` pins ``_session_messages`` (the full live transcript,
tens of MB on a tool-heavy session). The cache's LRU cap counts entries, not
bytes, and the idle TTL defers eviction for busy sessions, so neither sees
actual memory use. This module supplies that signal: own anonymous RSS against
a budget derived from the cgroup limit; ``GatewayRunner`` sheds LRU transcripts
via soft eviction (rebuilt from the persisted session on the next turn).
Everything here is pure or read-only. Config lives under ``agent.agent_cache``.
Each cached ``AIAgent`` pins its full live transcript (tens of MB on a tool-heavy
session); the LRU cap counts entries, not bytes, and the idle TTL defers eviction
for busy sessions, so neither sees actual memory use. This module supplies that
signal — own anonymous RSS against a budget derived from the cgroup limit — and
``GatewayRunner`` sheds LRU transcripts via soft eviction (rebuilt from the
persisted session next turn). Pure/read-only; config under ``agent.agent_cache``.
"""
from __future__ import annotations
@@ -33,12 +32,8 @@ _OFF_WORDS = frozenset({"", "off", "none", "false", "disabled"})
@dataclass(frozen=True)
class AgentCacheBounds:
"""Operator-facing bounds for the per-session agent cache.
``max_size`` / ``idle_ttl_secs`` are ``None`` when unset so ``gateway/run.py``
keeps its module defaults; ``memory_high_mb`` is ``None`` when pressure
eviction is off.
"""
"""Operator-facing bounds. ``max_size``/``idle_ttl_secs`` are ``None`` when unset
so ``gateway/run.py`` keeps its defaults; ``memory_high_mb`` ``None`` = pressure eviction off."""
max_size: Optional[int] = None
idle_ttl_secs: Optional[float] = None
@@ -65,11 +60,10 @@ def _positive(value: Any, cast: Callable[[Any], Any] = int) -> Any:
def _cgroup_limit_bytes() -> Optional[int]:
"""Memory limit this process runs under, if cgroup-capped.
Prefers cgroup v2 ``memory.high`` (the throttling point) over ``memory.max``,
then cgroup v1. Checks the process's *own* cgroup first (where a systemd
unit's ``MemoryHigh=``/``MemoryMax=`` lands — the root files read ``max``
there), then the root for container-style limits. ``max`` and the v1
near-2^63 sentinel mean unlimited.
Prefers v2 ``memory.high`` (the throttling point) over ``memory.max``, then v1.
Own cgroup first (where a systemd unit's ``MemoryHigh=``/``MemoryMax=`` lands —
root reads ``max`` there), then root for container-style limits. ``max`` and
the v1 near-2^63 sentinel mean unlimited.
"""
if sys.platform != "linux":
return None
@@ -106,11 +100,8 @@ def _total_memory_bytes() -> Optional[int]:
def resolve_memory_high_mb(setting: Any) -> Optional[int]:
"""Resolve ``memory_high_mb`` into an absolute MB budget.
``"auto"`` derives it from the cgroup limit (or total RAM when uncapped);
a positive number is literal; anything falsy/off disables the pass.
"""
"""Absolute MB budget: ``"auto"`` derives from the cgroup limit (or total RAM when
uncapped); a positive number is literal; anything falsy/off disables the pass."""
if isinstance(setting, str):
normalized = setting.strip().lower()
if normalized != "auto":
@@ -129,11 +120,8 @@ def resolve_memory_high_mb(setting: Any) -> Optional[int]:
def resolve_agent_cache_bounds(config: Any) -> AgentCacheBounds:
"""Read ``agent.agent_cache`` out of the *raw* config mapping.
The gateway's loader does not deep-merge ``DEFAULT_CONFIG``, so an absent
key stays absent and callers can tell "operator chose 128" from "unset".
"""
"""Read ``agent.agent_cache`` from the *raw* config: the gateway loader does not
deep-merge ``DEFAULT_CONFIG``, so callers can tell "operator chose 128" from "unset"."""
section = (config.get("agent") or {}).get("agent_cache") if isinstance(config, dict) else None
if not isinstance(section, dict):
section = {}
@@ -154,12 +142,8 @@ def resolve_agent_cache_bounds(config: Any) -> AgentCacheBounds:
def read_anon_rss_mb() -> Optional[int]:
"""Process anonymous resident memory in MB, or None.
Anonymous pages are where cached transcripts live; file-backed pages are
noise. ``collect_memory_snapshot`` reads ``/proc/self/status`` without a
dependency; psutil covers other platforms (total RSS only).
"""
"""Anonymous RSS in MB (where cached transcripts live; file-backed pages are noise),
or None. ``/proc/self/status`` first; psutil covers other platforms (total RSS only)."""
try:
from hermes_cli.mem_trim import collect_memory_snapshot
@@ -179,13 +163,12 @@ def read_anon_rss_mb() -> Optional[int]:
def transcript_persistence_caught_up(agent: Any) -> bool:
"""True when the agent's live transcript is fully on disk.
"""True when the live transcript is fully on disk.
Soft eviction drops ``_session_messages`` and rebuilds from the persisted
session, so it is only safe once ``_last_flushed_db_idx`` (advanced to
``len(messages)`` only on a fully successful write) has caught up. Unknown
shapes are *not* caught up: a skipped eviction costs memory, a wrong one
costs the conversation.
Soft eviction rebuilds from the persisted session, so it is only safe once
``_last_flushed_db_idx`` (advanced only on a fully successful write) has caught
up. Unknown shapes are *not* caught up: a skipped eviction costs memory, a
wrong one costs the conversation.
"""
messages = getattr(agent, "_session_messages", None)
flushed = getattr(agent, "_last_flushed_db_idx", None)
@@ -193,19 +176,15 @@ def transcript_persistence_caught_up(agent: Any) -> bool:
def plan_pressure_evictions(
ordered_entries: Iterable[Tuple[str, Any]],
*,
is_evictable: Callable[[str, Any], bool],
max_evictions: int,
protect_recent: int = 0,
ordered_entries: Iterable[Tuple[str, Any]], *, is_evictable: Callable[[str, Any], bool],
max_evictions: int, protect_recent: int = 0,
) -> List[Tuple[str, Any]]:
"""Choose which cached sessions to shed, least-recently-used first.
``ordered_entries`` must be LRU→MRU (the cache OrderedDict is kept that way
by ``move_to_end`` on every hit). The batch is capped so one pass cannot
stall the gateway. ``protect_recent`` is clamped to half the cache: a few
huge transcripts can exhaust the budget alone, and a fixed guard would then
protect the whole cache with nothing left to shed.
``ordered_entries`` must be LRU→MRU (the cache OrderedDict ``move_to_end``s on
every hit). The batch is capped so one pass cannot stall the gateway.
``protect_recent`` is clamped to half the cache: a few huge transcripts can
exhaust the budget alone, and a fixed guard would leave nothing to shed.
"""
entries = list(ordered_entries)
if max_evictions <= 0 or not entries:

View File

@@ -23,43 +23,31 @@ logger = logging.getLogger(__name__)
# on a tiny volume is one download from write failures. Percent triggers are
# gated on absolute headroom also being low, and a hard absolute floor applies
# regardless of size (below it SQLite journaling / config writes are at risk).
_CRITICAL_FREE_MB = 256 # < 256 MB free: critical on any volume
_CRITICAL_PERCENT = 95.0 # >= 95% used AND < 1 GB free: critical
_CRITICAL_HEADROOM_MB = 1024
_ELEVATED_FREE_MB = 512 # < 512 MB free: elevated on any volume
_ELEVATED_PERCENT = 85.0 # >= 85% used AND < 4 GB free: elevated
_ELEVATED_HEADROOM_MB = 4096
# (level, free-MB floor, used-% trigger, headroom MB the % trigger is gated on); worst first.
_PRESSURE_TIERS = (
("critical", 256, 95.0, 1024), # < 256 MB free, or >= 95% used AND < 1 GB free
("elevated", 512, 85.0, 4096), # < 512 MB free, or >= 85% used AND < 4 GB free
)
_BYTES_PER_MB = 1024 * 1024
def classify_disk_pressure(free_mb: Any, total_mb: Any) -> str:
"""Map free/total MB to ``ok``/``elevated``/``critical``.
``unknown`` when the sample is missing or malformed — the caller must
not treat "we could not read it" as "disk is fine".
"""
"""``ok``/``elevated``/``critical`` from free/total MB; ``unknown`` when the sample
is missing/malformed — "could not read it" must never read as "fine"."""
free = _nonneg_int(free_mb)
total = _nonneg_int(total_mb)
if free is None or not total:
return "unknown"
used_percent = (1 - free / total) * 100.0
for level, free_floor, percent_floor, headroom in (
("critical", _CRITICAL_FREE_MB, _CRITICAL_PERCENT, _CRITICAL_HEADROOM_MB),
("elevated", _ELEVATED_FREE_MB, _ELEVATED_PERCENT, _ELEVATED_HEADROOM_MB),
):
for level, free_floor, percent_floor, headroom in _PRESSURE_TIERS:
if free < free_floor or (used_percent >= percent_floor and free < headroom):
return level
return "ok"
def collect_disk_status(home: Optional[Path] = None) -> Dict[str, Any]:
"""Build the ``disk`` block for ``/api/status``.
``home`` scopes the sample to a profile's HERMES_HOME (same contract as the
``memory`` block). Always returns a dict and never raises — an unreadable
or unmounted filesystem yields ``{"pressure": "unknown", ...}``.
"""
"""``disk`` block for ``/api/status`` (same ``home`` contract as ``memory``).
Never raises — an unreadable/unmounted filesystem yields ``pressure="unknown"``."""
status: Dict[str, Any] = {"pressure": "unknown", "total_mb": None, "free_mb": None, "used_percent": None}
try:
if home is None:

View File

@@ -1,19 +1,14 @@
"""External drain-control marker contract (dashboard → gateway).
There is no control channel into a running gateway, so begin/cancel-drain
writes (or removes) ``{HERMES_HOME}/.drain_request.json`` and a gateway watcher
reacts; this module owns the contract so writer and reader never disagree.
Presence of an ACTIVE marker means "external drain" (``gateway_state ->
"draining"``); absence or a stale marker means "not draining".
Staleness — two independent, individually-lenient signals (either suffices):
epoch mismatch (HERMES_HOME is a durable volume on Hermes Cloud, so a marker
survives the machine restart a drain-gated action ends in and would park the
fresh gateway in ``draining`` forever) and expiry (a same-epoch orphan is
ignored past :data:`DRAIN_REQUEST_MAX_AGE_SECONDS`; re-writing refreshes it).
Reading never raises: a malformed file reads as ``{}``, still drain-active
(fail-safe toward quiescing). Staleness rejects only on a *definite* verdict —
no epoch/timestamp, or no ``/proc``, degrades to presence-only, never fail-closed.
No control channel exists into a running gateway, so begin/cancel-drain writes
(or removes) ``{HERMES_HOME}/.drain_request.json`` and a gateway watcher reacts;
an ACTIVE marker means ``gateway_state -> "draining"``. Two lenient staleness
signals (either suffices): epoch mismatch (HERMES_HOME is a durable volume on
Hermes Cloud, so a marker survives the restart a drain-gated action ends in and
would park the fresh gateway in ``draining`` forever) and expiry (same-epoch
orphan past :data:`DRAIN_REQUEST_MAX_AGE_SECONDS`; re-writing refreshes it).
Reading never raises: a malformed file reads as ``{}`` — still drain-active
(fail-safe toward quiescing). Staleness rejects only on a *definite* verdict.
"""
from __future__ import annotations
@@ -25,6 +20,7 @@ from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Optional
from gateway.memory_status import _parse_iso
from hermes_constants import get_hermes_home
from utils import atomic_json_write
@@ -43,11 +39,10 @@ _expiry_logged_for: Optional[str] = None
def current_instantiation_epoch() -> str:
"""Identity of THIS container / VM instantiation ("<boot_id>:<pid1_start>").
Stable for the life of PID 1 (an s6 respawn of just the gateway, or a host
``hermes gateway restart``, keeps honouring an in-flight drain) but changes
whenever the machine is recreated: boot_id on a VM reboot, PID 1's start
time on a plain ``docker restart``. ``""`` when neither source is readable
(non-Linux, no ``/proc``), which disables the epoch check — never fail-closed.
Stable for the life of PID 1 (a gateway-only respawn keeps honouring an
in-flight drain) but changes when the machine is recreated: boot_id on a VM
reboot, PID 1's start time on ``docker restart``. ``""`` when neither is
readable (non-Linux, no ``/proc``) disables the epoch check — never fail-closed.
"""
boot_id = pid1_start = ""
with contextlib.suppress(OSError):
@@ -67,21 +62,17 @@ def drain_request_path(home: Optional[Path] = None) -> Path:
def write_drain_request(
*, principal: str = "drain-control", suppress_notification: bool = False, home: Optional[Path] = None
) -> dict[str, Any]:
"""Write the begin-drain marker atomically. Returns the payload written.
"""Write the begin-drain marker atomically; returns the payload.
Idempotent: re-writing refreshes ``requested_at`` (keep-alive past the
max-age). ``suppress_notification`` asks the shutdown ending this drain to
skip ONLY the home-channel "gateway shutting down" broadcast (the per-session
interrupt ping is never suppressed); which drains are quiet is the caller's
policy. Stamped with :func:`current_instantiation_epoch` so a copy surviving
a machine restart on the durable volume is recognised as stale.
Re-writing refreshes ``requested_at`` (keep-alive past the max-age).
``suppress_notification`` skips ONLY the home-channel "gateway shutting down"
broadcast (the per-session interrupt ping is never suppressed); which drains
are quiet is the caller's policy. Stamped with the instantiation epoch so a
copy surviving a machine restart on the durable volume reads as stale.
"""
payload = {
"action": "drain",
"requested_at": datetime.now(timezone.utc).isoformat(),
"principal": principal,
"epoch": current_instantiation_epoch(),
"suppress_notification": bool(suppress_notification),
"action": "drain", "requested_at": datetime.now(timezone.utc).isoformat(), "principal": principal,
"epoch": current_instantiation_epoch(), "suppress_notification": bool(suppress_notification),
}
atomic_json_write(drain_request_path(home), payload)
return payload
@@ -104,29 +95,21 @@ def _marker_is_expired(body: dict[str, Any]) -> bool:
"""True iff ``requested_at`` parses AND is older than the max-age.
Missing/unparseable and future-dated (clock skew) timestamps are honoured.
Logged once per marker, not per poll — the warning is the operator's
breadcrumb for a writer that leaked a marker.
Logged once per marker, not per poll — the operator's breadcrumb for a leak.
"""
global _expiry_logged_for
raw = body.get("requested_at")
if not isinstance(raw, str) or not raw:
requested_at = _parse_iso(raw)
if requested_at is None:
return False
try:
requested_at = datetime.fromisoformat(raw)
except ValueError:
return False
if requested_at.tzinfo is None:
requested_at = requested_at.replace(tzinfo=timezone.utc)
age = (datetime.now(timezone.utc) - requested_at).total_seconds()
if age <= DRAIN_REQUEST_MAX_AGE_SECONDS:
return False
if _expiry_logged_for != raw:
_expiry_logged_for = raw
_log.warning(
"drain-control: ignoring expired drain marker (requested_at=%s, "
"age=%.0fs > max %.0fs, principal=%s) — the drain that wrote it "
"was never cancelled; treating as stale so the gateway keeps "
"accepting turns.",
"drain-control: ignoring expired drain marker (requested_at=%s, age=%.0fs > max %.0fs, principal=%s) "
"— the drain that wrote it was never cancelled; treating as stale so the gateway keeps accepting turns.",
raw, age, DRAIN_REQUEST_MAX_AGE_SECONDS, body.get("principal"),
)
return True
@@ -149,11 +132,10 @@ def drain_requested(*, home: Optional[Path] = None) -> bool:
def drain_notification_suppressed(*, home: Optional[Path] = None) -> bool:
"""True iff an ACTIVE drain marker explicitly asks to suppress the shutdown broadcast.
"""True iff an ACTIVE marker asks to suppress the shutdown broadcast.
Same activeness rule as :func:`drain_requested`, so an orphaned marker can
never silence a fresh gateway's broadcast. A legacy marker without the
field or a contentless ``{}`` reads as False (fail toward the louder behaviour).
Same activeness rule as :func:`drain_requested`, so an orphan can never silence
a fresh gateway; a marker without the field reads False (fail toward louder).
"""
body = _active_drain_body(home)
return bool(body and body.get("suppress_notification"))

View File

@@ -1,14 +1,12 @@
"""Event hook system: fires handlers at gateway lifecycle points.
Hooks live in ~/.hermes/hooks/<name>/ with HOOK.yaml (name, description, events)
and handler.py (``def handle(event_type, context)``, sync or async). Handler
errors are logged and never block the pipeline. Events: gateway:startup,
session:start/end/reset, agent:start, agent:step (each tool-loop turn),
agent:end, command:* (wildcard). ``agent:start``/``agent:end`` context: platform,
user_id, chat_id, thread_id (forum-topic/thread root as str, "" outside a thread),
chat_type ("dm"|"group"|"forum"|""), session_id, message (500 chars); ``agent:end``
adds response (500 chars), model, provider. Telegram forum follow-ups should pass
``message_thread_id=int(thread_id)`` when ``chat_type == "forum"`` and thread_id set.
Hooks live in ~/.hermes/hooks/<name>/ with HOOK.yaml (name, description, events) and
handler.py (``def handle(event_type, context)``, sync or async); errors never block
the pipeline. Events: gateway:startup, session:start/end/reset, agent:start,
agent:step (each tool-loop turn), agent:end, command:* (wildcard). agent:* context:
platform, user_id, chat_id, thread_id ("" outside a thread), chat_type
("dm"|"group"|"forum"|""), session_id, message (500 chars); agent:end adds response,
model, provider. Forum follow-ups pass ``message_thread_id=int(thread_id)``.
"""
import asyncio

View File

@@ -23,10 +23,6 @@ from typing import Any, Dict, Optional
logger = logging.getLogger(__name__)
_LIFECYCLE_RELATIVE = ("state", "gateway.lifecycle.json")
_EXIT_DIAG_RELATIVE = ("logs", "gateway-exit-diag.log")
_STATE_DB_RELATIVE = ("state.db",)
def _process_hermes_home() -> Path:
"""HERMES_HOME for process-level identity files (ignore task overrides)."""
@@ -38,13 +34,13 @@ def _process_hermes_home() -> Path:
return get_hermes_home()
def _home_path(home: Optional[Path], relative: tuple) -> Path:
def _home_path(home: Optional[Path], *relative: str) -> Path:
return (_process_hermes_home() if home is None else home).joinpath(*relative)
def get_lifecycle_sentinel_path(home: Optional[Path] = None) -> Path:
"""Return ``<HERMES_HOME>/state/gateway.lifecycle.json``."""
return _home_path(home, _LIFECYCLE_RELATIVE)
return _home_path(home, "state", "gateway.lifecycle.json")
def _now_iso() -> str:
@@ -74,10 +70,8 @@ def sample_memory() -> Dict[str, Any]:
heartbeat so OOM crash cycles are classifiable from the volume alone.
"""
sample = _proc_fields("/proc/self/status", {"VmRSS": "rss_kib"})
mem = _proc_fields("/proc/meminfo", {
"MemTotal": "mem_total_kib", "MemAvailable": "mem_available_kib",
"SwapTotal": "SwapTotal", "SwapFree": "SwapFree",
})
mem = _proc_fields("/proc/meminfo", {"MemTotal": "mem_total_kib", "MemAvailable": "mem_available_kib",
"SwapTotal": "SwapTotal", "SwapFree": "SwapFree"})
swap_total, swap_free = mem.pop("SwapTotal", None), mem.pop("SwapFree", None)
sample.update(mem)
if swap_total is not None and swap_free is not None:
@@ -106,7 +100,7 @@ def _write_sentinel(payload: Dict[str, Any], home: Optional[Path]) -> None:
def _append_exit_diag(record: Dict[str, Any], home: Optional[Path]) -> None:
"""Append a JSON line to gateway-exit-diag.log (same format as the CLI's ``_exit_diag``)."""
path = _home_path(home, _EXIT_DIAG_RELATIVE)
path = _home_path(home, "logs", "gateway-exit-diag.log")
try:
path.parent.mkdir(parents=True, exist_ok=True)
with path.open("a", encoding="utf-8") as fh:
@@ -139,6 +133,18 @@ def _pid_alive_with_start_time(pid: Any, start_time: Any) -> bool:
return True
def _suspected_oom(mem: Dict[str, Any]) -> bool:
"""Heuristic only (classification stays with the reader); thresholds are
memory_status' "critical" tier so a live warning and a post-mortem verdict agree."""
from gateway.memory_status import _CRITICAL_AVAILABLE_FRACTION, _CRITICAL_AVAILABLE_KIB
total, avail = mem.get("mem_total_kib"), mem.get("mem_available_kib")
return isinstance(avail, int) and (
avail < _CRITICAL_AVAILABLE_KIB
or (isinstance(total, int) and total > 0 and avail / total < _CRITICAL_AVAILABLE_FRACTION)
)
def detect_unclean_exit(home: Optional[Path] = None) -> Optional[Dict[str, Any]]:
"""Evidence dict when the previous life died uncleanly, else ``None``. Read-only."""
sentinel = _read_json(get_lifecycle_sentinel_path(home))
@@ -146,10 +152,8 @@ def detect_unclean_exit(home: Optional[Path] = None) -> Optional[Dict[str, Any]]
return None
if _pid_alive_with_start_time(sentinel.get("pid"), sentinel.get("start_time")):
return None # live owner — planned takeover in flight, not a death
evidence: Dict[str, Any] = {
"prior_pid": sentinel.get("pid"),
"prior_started_at": sentinel.get("started_at"),
"prior_pid": sentinel.get("pid"), "prior_started_at": sentinel.get("started_at"),
"prior_start_time": sentinel.get("start_time"),
}
# Enrich with the last heartbeat: last proven liveness and memory at that moment.
@@ -164,30 +168,20 @@ def detect_unclean_exit(home: Optional[Path] = None) -> Optional[Dict[str, Any]]
mem = hb.get("mem")
if isinstance(mem, dict):
evidence["last_heartbeat_mem"] = mem
# OOM suspicion is only a hint (classification stays with the reader);
# thresholds are memory_status' "critical" tier so a live warning and a
# post-mortem verdict can never disagree.
from gateway.memory_status import _CRITICAL_AVAILABLE_FRACTION, _CRITICAL_AVAILABLE_KIB
total, avail = mem.get("mem_total_kib"), mem.get("mem_available_kib")
if isinstance(avail, int) and (
avail < _CRITICAL_AVAILABLE_KIB
or (isinstance(total, int) and total > 0 and avail / total < _CRITICAL_AVAILABLE_FRACTION)
):
if _suspected_oom(mem):
evidence["suspected_oom"] = True
return evidence
def check_state_db_integrity(home: Optional[Path] = None) -> str:
"""Return ``"ok"``, ``"absent"``, or the first ``quick_check`` complaint.
"""``"ok"``, ``"absent"``, or the first ``quick_check`` complaint. Never raises.
Called only after an unclean death — a SIGKILL mid-WAL-checkpoint can leave
half-written b-tree pages. ``quick_check(1)`` stops at the first problem
(~2s on a healthy 500MB store): cheap once per unclean boot, too costly every
boot. Opened normally, not read-only: a WAL store needs its -shm sidecar for
a read-only open, and the PRAGMA writes nothing. Never raises.
Only after an unclean death — SIGKILL mid-WAL-checkpoint can leave half-written
b-tree pages. ``quick_check(1)`` stops at the first problem (~2s on a healthy
500MB store): cheap once per unclean boot, too costly every boot. Opened
normally: a WAL store needs its -shm sidecar for read-only, and the PRAGMA writes nothing.
"""
path = _home_path(home, _STATE_DB_RELATIVE)
path = _home_path(home, "state.db")
if not path.exists():
return "absent"
try:
@@ -198,6 +192,25 @@ def check_state_db_integrity(home: Optional[Path] = None) -> str:
return "check-failed: no result" if not row or row[0] is None else str(row[0])
def _report_unclean_exit(evidence: Dict[str, Any], home: Optional[Path]) -> None:
"""Integrity-check the store, persist the exit-diag record, log at WARNING."""
# The death may have torn the store; this is the only moment we know to look.
verdict = evidence["state_db_integrity"] = check_state_db_integrity(home=home)
if verdict not in ("ok", "absent"):
logger.error(
"state.db FAILED integrity check after an unclean gateway exit: %s — sessions may read as "
"missing until it is repaired. Run `hermes doctor`.",
verdict,
)
_append_exit_diag({"ts": _now_iso(), "tag": "gateway.previous_unclean_exit", "pid": os.getpid(), **evidence}, home)
logger.warning(
"Previous gateway life (pid=%s, started_at=%s) exited UNCLEANLY (no exit path ran — SIGKILL / OOM / "
"VM death). last_heartbeat_at=%s last_mem=%s suspected_oom=%s",
evidence.get("prior_pid"), evidence.get("prior_started_at"), evidence.get("last_heartbeat_at"),
evidence.get("last_heartbeat_mem"), evidence.get("suspected_oom", False),
)
def record_startup(home: Optional[Path] = None) -> Optional[Dict[str, Any]]:
"""Boot entry point: report any unclean previous exit, then claim the sentinel.
@@ -208,29 +221,9 @@ def record_startup(home: Optional[Path] = None) -> Optional[Dict[str, Any]]:
try:
evidence = detect_unclean_exit(home)
if evidence is not None:
# The death may have torn the store; this is the only moment we know to look.
verdict = check_state_db_integrity(home=home)
evidence["state_db_integrity"] = verdict
if verdict not in ("ok", "absent"):
logger.error(
"state.db FAILED integrity check after an unclean gateway "
"exit: %s — sessions may read as missing until it is "
"repaired. Run `hermes doctor`.",
verdict,
)
_append_exit_diag(
{"ts": _now_iso(), "tag": "gateway.previous_unclean_exit", "pid": os.getpid(), **evidence}, home
)
logger.warning(
"Previous gateway life (pid=%s, started_at=%s) exited UNCLEANLY "
"(no exit path ran — SIGKILL / OOM / VM death). "
"last_heartbeat_at=%s last_mem=%s suspected_oom=%s",
evidence.get("prior_pid"), evidence.get("prior_started_at"), evidence.get("last_heartbeat_at"),
evidence.get("last_heartbeat_mem"), evidence.get("suspected_oom", False),
)
_report_unclean_exit(evidence, home)
except Exception:
logger.debug("Unclean-exit detection failed", exc_info=True)
try:
claim: Dict[str, Any] = {"phase": "running", "pid": os.getpid(), "start_time": time.time(), "started_at": _now_iso()}
# Carry the verdict on the PREVIOUS life on the new sentinel: it is the only
@@ -249,19 +242,18 @@ def record_startup(home: Optional[Path] = None) -> Optional[Dict[str, Any]]:
def mark_exited(exit_code: Optional[int] = None, reason: str = "graceful_shutdown", home: Optional[Path] = None) -> None:
"""Mark the current life as cleanly exited. Idempotent, never raises.
Only rewrites a sentinel provably owned by this process: during a
``--replace`` takeover the replacement claims it before the old process
finishes teardown, and the old life must not clobber the new ``running``
phase. ``pid=None`` / malformed sentinels have unknown ownership → left alone.
Only rewrites a sentinel provably owned by this process: during ``--replace``
the replacement claims it before the old process finishes teardown and must not
be clobbered. ``pid=None`` / malformed sentinels have unknown ownership → left alone.
"""
try:
sentinel = _read_json(get_lifecycle_sentinel_path(home))
if sentinel is not None and sentinel.get("pid") != os.getpid():
return
_write_sentinel({
"phase": "exited", "pid": os.getpid(), "exit_code": exit_code,
"exit_reason": reason, "exited_at": _now_iso(),
}, home)
_write_sentinel(
{"phase": "exited", "pid": os.getpid(), "exit_code": exit_code, "exit_reason": reason, "exited_at": _now_iso()},
home,
)
except Exception:
logger.debug("Failed to mark lifecycle sentinel exited", exc_info=True)

View File

@@ -1,18 +1,12 @@
"""Shared config→env bridge for media-delivery policy.
``validate_media_delivery_path`` (gateway/platforms/base.py) reads its policy
from environment variables:
- ``HERMES_MEDIA_DELIVERY_STRICT`` <- gateway.strict
- ``HERMES_MEDIA_ALLOW_DIRS`` <- gateway.media_delivery_allow_dirs
- ``HERMES_MEDIA_TRUST_RECENT_FILES`` <- gateway.trust_recent_files
Every delivery entrypoint (gateway startup, ``hermes cron run``, ``hermes send``)
calls :func:`apply_media_policy_env` before filtering media paths, so standalone
paths filter under the same policy as the gateway instead of silently dropping
attachments in strict/allowlisted deployments. An explicitly-set environment
variable WINS over config.yaml, so a shell override (and gateway startup's own
earlier run) survives.
``validate_media_delivery_path`` reads ``HERMES_MEDIA_DELIVERY_STRICT`` (gateway.strict),
``HERMES_MEDIA_ALLOW_DIRS`` (gateway.media_delivery_allow_dirs) and
``HERMES_MEDIA_TRUST_RECENT_FILES`` (gateway.trust_recent_files). Every delivery
entrypoint (gateway startup, ``hermes cron run``, ``hermes send``) calls
:func:`apply_media_policy_env` first so standalone paths filter under the gateway's
policy instead of silently dropping attachments in strict/allowlisted deployments.
An explicitly-set env var WINS over config.yaml, so shell overrides survive.
"""
from __future__ import annotations

View File

@@ -1,13 +1,11 @@
"""Repair model-mangled ``computer_use`` screenshot paths in final responses.
``computer_use`` persists a screenshot into the image cache and tells the model
its absolute path. Some models rewrite a Windows path into a POSIX-looking one
(``C:\\Users\\Alice\\...`` -> ``/Users/Alice/...``) inside an explicit ``MEDIA:``
Some models rewrite the Windows path ``computer_use`` reported into a POSIX-looking
one (``C:\\Users\\Alice\\...`` -> ``/Users/Alice/...``) inside an explicit ``MEDIA:``
directive, so delivery-path validation rejects it and drops the attachment.
Deliberately narrow: only rewrites paths inside a response that *already*
carries a ``MEDIA:`` directive, and only when its ``computer_use_<uuid>``
basename exactly matches a canonical path returned by ``computer_use`` this
turn. Never auto-attaches; normal media path validation still runs afterwards.
Deliberately narrow: only rewrites paths in a response that *already* carries a
``MEDIA:`` directive whose ``computer_use_<uuid>`` basename exactly matches a
canonical path returned this turn. Never auto-attaches; validation still runs.
"""
from __future__ import annotations
@@ -41,8 +39,7 @@ def tool_name_by_call_id(messages: List[Dict[str, Any]]) -> Dict[str, str]:
continue
for call in msg.get("tool_calls") or []:
call_id = call.get("id") or call.get("call_id")
fn = call.get("function") or {}
name = str(fn.get("name") or call.get("name") or "")
name = str((call.get("function") or {}).get("name") or call.get("name") or "")
if call_id and name:
mapping[str(call_id)] = name
return mapping
@@ -56,12 +53,9 @@ def _computer_use_capture_basename(path: Any) -> str:
def _iter_computer_use_capture_paths(content: Any) -> Iterator[str]:
"""Yield persisted screenshot paths from computer_use result content.
The tool can return JSON, a multimodal content list, or a text fallback; the
latter two keep the canonical path in the human-readable summary even though
the multimodal envelope's ``meta`` is not stored in the tool message.
"""
"""Yield persisted screenshot paths from computer_use result content (JSON, a
multimodal list, or text; the latter two keep the canonical path in the summary
line since the envelope's ``meta`` is not stored in the tool message)."""
if isinstance(content, str):
stripped = content.strip()
if stripped.startswith(("{", "[")):
@@ -114,13 +108,9 @@ def _current_turn_messages(messages: List[Dict[str, Any]], history_offset: int)
def repair_explicit_computer_use_media_paths(
response: str, messages: List[Dict[str, Any]], history_offset: int = 0
) -> str:
"""Recover model-mangled paths for explicitly requested screenshots.
Repairs only an already-explicit ``MEDIA:`` directive whose generated
basename case-insensitively matches a canonical screenshot path from this
turn. Fail-open: the repair is cosmetic, so any unexpected error returns
the response unchanged rather than aborting delivery.
"""
"""Recover model-mangled paths in explicit ``MEDIA:`` directives whose basename
matches (case-insensitively) a canonical screenshot path from this turn.
Fail-open: the repair is cosmetic, so any error returns the response unchanged."""
try:
return _repair_explicit_computer_use_media_paths_inner(response, messages, history_offset)
except Exception:

View File

@@ -50,11 +50,8 @@ def log_memory_usage(prefix: str = "") -> None:
"""Log ``[MEMORY] [<prefix> ]rss=... gc=... threads=... uptime=...``; safe from any thread."""
rss = _get_rss_mb()
logger.info(
"[MEMORY] %srss=%s gc=%s threads=%d uptime=%ds",
f"{prefix} " if prefix else "",
"unavailable" if rss is None else f"{rss}MB",
gc.get_count(), # (gen0, gen1, gen2)
threading.active_count(),
"[MEMORY] %srss=%s gc=%s threads=%d uptime=%ds", f"{prefix} " if prefix else "",
"unavailable" if rss is None else f"{rss}MB", gc.get_count(), threading.active_count(),
int(time.monotonic() - _start_time) if _start_time else 0,
)
@@ -81,16 +78,15 @@ def start_memory_monitoring(interval_seconds: float = 300.0) -> bool:
return False
if _get_rss_mb() is None:
logger.warning(
"[MEMORY] Memory monitoring unavailable: neither resource.getrusage "
"nor psutil could read process RSS — skipping periodic logging.",
"[MEMORY] Memory monitoring unavailable: neither resource.getrusage nor psutil could read process RSS "
"— skipping periodic logging.",
)
return False
_start_time = time.monotonic()
_stop_event = threading.Event()
log_memory_usage(prefix="baseline")
_monitor_thread = threading.Thread(
target=_monitor_loop, args=(_stop_event, float(interval_seconds)),
name="gateway-memory-monitor", daemon=True,
target=_monitor_loop, args=(_stop_event, float(interval_seconds)), name="gateway-memory-monitor", daemon=True
)
_monitor_thread.start()
logger.info("[MEMORY] Periodic memory monitoring started (interval: %ds)", int(interval_seconds))

View File

@@ -1,10 +1,9 @@
"""Memory status rollup for ``/api/status``.
Read side for signals the gateway already persists: the 30s
``state/gateway.heartbeat`` (RSS + MemAvailable/MemTotal + swap) and the
lifecycle sentinel's ``suspected_oom`` flag. Two small file reads, no IPC.
``/api/status`` is unauthenticated, so the block carries only coarse numbers
(MB), enums and booleans. Best-effort: a missing/corrupt file degrades to
Read side for signals the gateway already persists: the 30s ``state/gateway.heartbeat``
(RSS + MemAvailable/MemTotal + swap) and the lifecycle sentinel's ``suspected_oom``
flag — two small file reads, no IPC. ``/api/status`` is unauthenticated, so only
coarse numbers (MB), enums and booleans. A missing/corrupt file degrades to
``pressure="unknown"`` rather than raising into the status endpoint.
"""
@@ -24,6 +23,10 @@ _CRITICAL_AVAILABLE_KIB = 64 * 1024 # < 64 MiB available
_CRITICAL_AVAILABLE_FRACTION = 0.05 # < 5% of MemTotal
_ELEVATED_AVAILABLE_KIB = 128 * 1024 # < 128 MiB available
_ELEVATED_AVAILABLE_FRACTION = 0.15 # < 15% of MemTotal
_PRESSURE_TIERS = ( # order-sensitive: worst first
("critical", _CRITICAL_AVAILABLE_KIB, _CRITICAL_AVAILABLE_FRACTION),
("elevated", _ELEVATED_AVAILABLE_KIB, _ELEVATED_AVAILABLE_FRACTION),
)
# Writer cadence is 30s; 150s tolerates a briefly stalled loop without letting
# a long-dead gateway's last sample pose as current.
@@ -51,20 +54,14 @@ def _parse_iso(value: Any) -> Optional[datetime]:
def classify_pressure(available_kib: Any, total_kib: Any) -> str:
"""Map a MemAvailable/MemTotal pair to ``ok``/``elevated``/``critical``.
``unknown`` when the sample is missing or malformed — the caller must
not treat "we could not read it" as "memory is fine".
"""
"""``ok``/``elevated``/``critical`` from MemAvailable/MemTotal; ``unknown`` when the
sample is missing/malformed — "could not read it" must never read as "fine"."""
available = _nonneg_int(available_kib)
if available is None:
return "unknown"
total = _nonneg_int(total_kib)
fraction = available / total if total else None
for level, kib_floor, frac_floor in (
("critical", _CRITICAL_AVAILABLE_KIB, _CRITICAL_AVAILABLE_FRACTION),
("elevated", _ELEVATED_AVAILABLE_KIB, _ELEVATED_AVAILABLE_FRACTION),
):
for level, kib_floor, frac_floor in _PRESSURE_TIERS:
if available < kib_floor or (fraction is not None and fraction < frac_floor):
return level
return "ok"
@@ -86,13 +83,9 @@ def collect_memory_status(
*,
now: Optional[datetime] = None,
) -> Dict[str, Any]:
"""Build the ``memory`` block for ``/api/status``.
``home`` scopes the read to a profile's HERMES_HOME (``None`` = active
profile); ``now`` is injectable for tests. Always returns a dict and never
raises — a down gateway or corrupt files yield ``{"pressure": "unknown", ...}``
plus whatever fields could be recovered.
"""
"""``memory`` block for ``/api/status``; ``home`` scopes to a profile (``None`` =
active), ``now`` is injectable. Never raises — a down gateway or corrupt files
yield ``pressure="unknown"`` plus whatever fields could be recovered."""
moment = now or datetime.now(timezone.utc)
status: Dict[str, Any] = {
"pressure": "unknown", "gateway_rss_mb": None, "system_total_mb": None, "system_available_mb": None,

View File

@@ -56,11 +56,8 @@ def _parse_timestamp_match(match: re.Match, tz=None) -> Optional[float]:
def coerce_message_timestamp(ts_value: Any, tz=None) -> Optional[float]:
"""Coerce a timestamp-like value to Unix epoch seconds.
Accepts epoch numbers, datetime objects, ISO strings, and the gateway's
bracketed human-readable format. Returns ``None`` when uninterpretable.
"""
"""Epoch seconds from a number, datetime, ISO string, or the gateway's bracketed
format; ``None`` when uninterpretable."""
if ts_value is None:
return None
if isinstance(ts_value, (int, float)):
@@ -96,12 +93,9 @@ def format_message_timestamp(ts_value: Any, tz=None) -> str:
def strip_leading_message_timestamps(content: str, tz=None) -> Tuple[str, Optional[float]]:
"""Strip one or more leading gateway timestamp prefixes from ``content``.
Returns ``(clean_content, embedded_epoch)``. With multiple prefixes the one
closest to the message text wins, preserving the original platform-send
time for legacy contaminated rows like ``[processing time] [platform time] [sender] message``.
"""
"""Strip leading gateway timestamp prefixes → ``(clean_content, embedded_epoch)``.
With several prefixes the one closest to the text wins, preserving the platform-send
time of legacy rows like ``[processing time] [platform time] [sender] message``."""
if not isinstance(content, str) or not content:
return content, None
text = content
@@ -115,11 +109,8 @@ def strip_leading_message_timestamps(content: str, tz=None) -> Tuple[str, Option
def render_user_content_with_timestamp(content: str, ts_value: Any = None, tz=None) -> str:
"""Render a user message for LLM context with exactly one timestamp prefix.
An existing leading prefix is stripped and its parsed time wins over
``ts_value``. If no timestamp is available the cleaned content is returned.
"""
"""Render a user message for LLM context with exactly one timestamp prefix; an
existing prefix is stripped and its parsed time wins over ``ts_value``."""
clean_content, embedded_epoch = strip_leading_message_timestamps(content, tz=tz)
effective_ts = embedded_epoch if embedded_epoch is not None else ts_value
prefix = format_message_timestamp(effective_ts, tz=tz)

View File

@@ -23,24 +23,19 @@ def _origin_user_id(entry: dict) -> str:
def mirror_to_session(
platform: str, chat_id: str, message_text: str, source_label: str = "cli",
thread_id: Optional[str] = None, user_id: Optional[str] = None,
role: str = "assistant", session_id: Optional[str] = None,
platform: str, chat_id: str, message_text: str, source_label: str = "cli", thread_id: Optional[str] = None,
user_id: Optional[str] = None, role: str = "assistant", session_id: Optional[str] = None,
) -> bool:
"""Append a delivery-mirror message to the target session's SQLite transcript.
``session_id``: pass it when the caller already holds the exact session
(e.g. the cron in_channel seed that just created the row) to skip the
origin scan, which refuses to guess on a populated chat (flat session + N
thread sessions sharing one chat_id) and would silently drop the mirror.
``role`` defaults to ``"assistant"`` (the agent's own outgoing reply). Text
that is NOT the agent speaking (e.g. a cron brief) must pass ``role="user"``:
``mirror``/``mirror_source`` metadata is dropped at the SQLite boundary, so an
assistant-role mirror replays as a real assistant turn and produces
assistant→assistant pairs that break strict-alternation providers; a
user-role mirror collapses safely via the consecutive-user merge.
Pass ``session_id`` when the caller already holds the exact session (e.g. the
cron in_channel seed) to skip the origin scan, which refuses to guess on a
populated chat (flat + N thread sessions per chat_id) and would drop the mirror.
Text that is NOT the agent speaking (e.g. a cron brief) must pass
``role="user"``: ``mirror`` metadata is dropped at the SQLite boundary, so an
assistant-role mirror replays as a real turn and yields assistant→assistant
pairs that break strict-alternation providers, while a user-role mirror
collapses safely via the consecutive-user merge.
Returns True if mirrored, False if no matching session or error. Never raises.
"""
try:
@@ -48,42 +43,29 @@ def mirror_to_session(
session_id = _find_session_id(platform, str(chat_id), thread_id=thread_id, user_id=user_id)
if not session_id:
logger.warning(
"Mirror: no session found for %s:%s thread=%s user=%s "
"(explicit_id=none, origin-scan bailed)",
"Mirror: no session found for %s:%s thread=%s user=%s (explicit_id=none, origin-scan bailed)",
platform, chat_id, thread_id, user_id,
)
return False
_append_to_sqlite(session_id, {
"role": role,
"content": message_text,
"timestamp": datetime.now().isoformat(),
"mirror": True,
"mirror_source": source_label,
"role": role, "content": message_text, "timestamp": datetime.now().isoformat(),
"mirror": True, "mirror_source": source_label,
})
logger.debug("Mirror: wrote to session %s (from %s)", session_id, source_label)
return True
except Exception as e:
# WARNING, not debug: a silent mirror drop is the cron continuation-amnesia bug.
logger.warning(
"Mirror failed for %s:%s thread=%s user=%s session=%s: %s",
platform, chat_id, thread_id, user_id, session_id, e,
)
logger.warning("Mirror failed for %s:%s thread=%s user=%s session=%s: %s", platform, chat_id, thread_id, user_id, session_id, e)
return False
def _find_session_id(
platform: str, chat_id: str, thread_id: Optional[str] = None, user_id: Optional[str] = None,
) -> Optional[str]:
"""Find the active session_id for a platform + chat_id pair.
def _find_session_id(platform: str, chat_id: str, thread_id: Optional[str] = None, user_id: Optional[str] = None) -> Optional[str]:
"""Active session_id for a platform + chat_id pair.
state.db gateway session rows are primary; sessions.json is the fallback
for pre-migration databases. DM session keys don't embed the chat_id
(e.g. "agent:main:telegram:dm"), so matching is on the persisted origin.
With *user_id*, exact sender matches win. If several same-chat candidates
exist and none matches the user, return None rather than guess and
contaminate another participant's session.
state.db is primary; sessions.json is the pre-migration fallback. DM keys
don't embed the chat_id ("agent:main:telegram:dm"), so match on the persisted
origin. With *user_id*, exact sender matches win; several same-chat candidates
with no user match → None rather than contaminate another participant's session.
"""
try:
from hermes_state import get_shared_session_db, release_or_close
@@ -105,21 +87,16 @@ def _find_session_id(
except Exception:
return None
platform_lower = platform.lower()
candidates = []
for _key, entry in data.items():
# Keys starting with "_" (e.g. the gateway's "_README") are metadata sentinels.
if str(_key).startswith("_") or not isinstance(entry, dict):
continue
def _matches(entry: dict) -> bool:
origin = entry.get("origin") or {}
if (
(origin.get("platform") or entry.get("platform", "")).lower() != platform_lower
or str(origin.get("chat_id", "")) != str(chat_id)
or (thread_id is not None and str(origin.get("thread_id") or "") != str(thread_id))
):
continue
candidates.append(entry)
return (
(origin.get("platform") or entry.get("platform", "")).lower() == platform.lower()
and str(origin.get("chat_id", "")) == str(chat_id)
and (thread_id is None or str(origin.get("thread_id") or "") == str(thread_id))
)
# Keys starting with "_" (e.g. the gateway's "_README") are metadata sentinels.
candidates = [e for k, e in data.items() if not str(k).startswith("_") and isinstance(e, dict) and _matches(e)]
if not candidates:
return None
if user_id:

View File

@@ -51,9 +51,7 @@ def _probe_config(home: Path) -> dict[str, Any]:
raw = yaml.safe_load(path.read_text(encoding="utf-8"))
except Exception as exc:
return _check("degraded", f"invalid config ({type(exc).__name__})")
if raw is not None and not isinstance(raw, dict):
return _check("degraded", "top level is not a mapping")
return _check("ok")
return _check("ok") if raw is None or isinstance(raw, dict) else _check("degraded", "top level is not a mapping")
def _probe_disk(home: Path) -> dict[str, Any]:
@@ -81,27 +79,19 @@ def _probe_gateway(runtime_status: dict[str, Any]) -> dict[str, Any]:
def _probe_session_store(runtime_status: dict[str, Any], state_db_probe: dict[str, Any]) -> dict[str, Any]:
"""Report the running gateway cache state, not an independent reopen."""
runtime_store = runtime_status.get("session_store")
if isinstance(runtime_store, dict):
state = str(runtime_store.get("status") or "unknown")
if state in {"ok", "unavailable", "retrying"}:
return _check(state)
state = str(runtime_store.get("status") or "unknown") if isinstance(runtime_store, dict) else ""
if state in {"ok", "unavailable", "retrying"}:
return _check(state)
# Older gateways publish no cache state: fall back to the state_db probe.
return _check("ok" if state_db_probe.get("status") == "ok" else "unavailable")
def collect_runtime_readiness(
*,
configured_model: str,
runtime_status: dict[str, Any] | None,
active_api_runs: int = 0,
process_completion_queue_depth: int = 0,
active_delegations: int = 0,
*, configured_model: str, runtime_status: dict[str, Any] | None, active_api_runs: int = 0,
process_completion_queue_depth: int = 0, active_delegations: int = 0,
) -> dict[str, Any]:
"""Return bounded readiness diagnostics without mutating runtime state.
Even on the authenticated endpoint, probes expose status and counts only:
never config values, credentials, paths, queue payloads, or exception messages.
"""
"""Bounded readiness diagnostics, no runtime mutation. Even authenticated, probes
expose status and counts only: never config values, credentials, paths, payloads."""
home = get_hermes_home()
runtime = runtime_status if isinstance(runtime_status, dict) else {}
state_db_probe = _probe_state_db(home)
@@ -113,8 +103,7 @@ def collect_runtime_readiness(
"disk": _probe_disk(home),
"gateway": _probe_gateway(runtime),
"background_queues": _check(
"ok",
active_api_runs=max(0, int(active_api_runs)),
"ok", active_api_runs=max(0, int(active_api_runs)),
process_completions=max(0, int(process_completion_queue_depth)),
active_delegations=max(0, int(active_delegations)),
),

View File

@@ -12,12 +12,7 @@ from typing import Any
# Exact whole-response markers meaning "the agent intentionally chose not to
# reply". Keep small and explicit; arbitrary empty output remains an
# error/empty-response path, not silence.
LIVE_GATEWAY_SILENT_MARKERS = frozenset({
"[SILENT]",
"SILENT",
"NO_REPLY",
"NO REPLY",
})
LIVE_GATEWAY_SILENT_MARKERS = frozenset({"[SILENT]", "SILENT", "NO_REPLY", "NO REPLY"})
# Longer than any marker could plausibly be, even with stray punctuation.
_MARKER_LENGTH_CAP = 64
@@ -68,12 +63,10 @@ def is_intentional_silence_response(response: Any) -> bool:
def is_autonomous_silence_response(response: Any) -> bool:
"""Loose silence matcher for autonomous lanes (cron, webhook).
Autonomous lanes ask for ``[SILENT]`` when a tick produced nothing worth
attention, and models reliably bracket the marker with a short note. Unlike
:func:`is_intentional_silence_response` (interactive rule: EXACTLY a marker),
this suppresses when a marker is the whole response, sits on its own first or
last line, or the bracketed sentinel opens the response (``[SILENT] No
changes detected``). A token buried mid-sentence is still delivered.
Models reliably bracket ``[SILENT]`` with a short note, so unlike the
interactive EXACT rule this also suppresses when a marker sits on its own
first/last line or the bracketed sentinel opens the response (``[SILENT] No
changes detected``). A token buried mid-sentence is still delivered.
Shares :data:`LIVE_GATEWAY_SILENT_MARKERS` so the two sets cannot drift.
"""
if not isinstance(response, str):
@@ -104,13 +97,10 @@ def is_intentional_silence_agent_result(agent_result: dict | None, response: Any
def is_partial_silence_marker(text: Any) -> bool:
"""True while streamed ``text`` could still resolve to a silence marker.
The streaming path must decide, before the whole response is known, whether
to show its buffer. A buffer whose canonical form is a non-empty *prefix* of
a marker (``"NO"`` on the way to ``"NO_REPLY"``, or an exact marker not yet
terminated by stream-end) is held back so a raw marker is never shown and
then retracted. Anything that has diverged from every marker, or exceeds the
marker cap, returns False so normal streaming resumes. Shares the marker set
and canonicalization with :func:`is_intentional_silence_response`.
A buffer whose canonical form is a non-empty *prefix* of a marker (``"NO"`` on
the way to ``"NO_REPLY"``, or an exact marker not yet terminated by stream-end)
is held back so a raw marker is never shown and then retracted. Divergence
from every marker, or exceeding the cap, resumes normal streaming.
"""
return any(
c and any(marker.startswith(c) for marker in LIVE_GATEWAY_SILENT_MARKERS)

View File

@@ -6,11 +6,10 @@ from collections.abc import Mapping
from hermes_cli.config import DEFAULT_CONFIG
# EX_TEMPFAIL (sysexits.h): ask the service manager to restart the gateway
# after a graceful drain/reload path completes.
# EX_TEMPFAIL (sysexits.h): ask the service manager to restart after a graceful drain/reload.
GATEWAY_SERVICE_RESTART_EXIT_CODE = 75
# EX_CONFIG (sysexits.h): fatal configuration error (token collision, no
# platforms); the s6 finish script maps it to exit 125 so the supervisor stops restarting.
# EX_CONFIG (sysexits.h): fatal configuration error (token collision, no platforms);
# the s6 finish script maps it to exit 125 so the supervisor stops restarting.
GATEWAY_FATAL_CONFIG_EXIT_CODE = 78
# Set by ``hermes gateway run --external-supervisor``. Unlike systemd's INVOCATION_ID
@@ -22,9 +21,8 @@ DEFAULT_GATEWAY_RESTART_DRAIN_TIMEOUT = float(DEFAULT_CONFIG["agent"]["restart_d
DEFAULT_GATEWAY_SIGNAL_INTERRUPT_GRACE_TIMEOUT = float(DEFAULT_CONFIG["gateway"]["signal_interrupt_grace_timeout"])
DEFAULT_GATEWAY_POST_INTERRUPT_GRACE_TIMEOUT = 5.0
# In-band restart (``/restart``, SIGUSR1, self-restart) waits for active turns to
# finish *before* ``stop()`` begins. Distinct from ``restart_drain_timeout``, the
# force-interrupt budget once ``stop()`` runs (must stay short under systemd TimeoutStopSec).
# In-band restart waits for active turns to finish *before* ``stop()`` begins; distinct from
# ``restart_drain_timeout``, the force-interrupt budget once ``stop()`` runs (short under TimeoutStopSec).
DEFAULT_GATEWAY_RESTART_AFTER_TURN_TIMEOUT = float(DEFAULT_CONFIG["agent"]["restart_after_turn_timeout"])
# Cron-only floor under the ``stop()`` drain. ``restart_drain_timeout`` defaults to 0
@@ -123,11 +121,7 @@ def parse_signal_interrupt_grace_timeout(raw: object) -> float:
def resolve_cron_drain_budget(
drain_timeout: float,
cron_drain_timeout: float,
*,
watchdog_delay: float,
elapsed: float = 0.0,
drain_timeout: float, cron_drain_timeout: float, *, watchdog_delay: float, elapsed: float = 0.0,
cleanup_reserve_s: float = CRON_DRAIN_CLEANUP_RESERVE_S,
) -> float:
"""Seconds the shutdown drain may spend waiting on in-flight cron work.
@@ -147,11 +141,8 @@ def resolve_cron_drain_budget(
def resolve_systemd_timeout_stop_sec(
drain_timeout: float,
cron_drain_timeout: float = DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT,
*,
cleanup_reserve_s: float = CRON_DRAIN_CLEANUP_RESERVE_S,
headroom_s: float = SYSTEMD_STOP_HEADROOM_S,
drain_timeout: float, cron_drain_timeout: float = DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT, *,
cleanup_reserve_s: float = CRON_DRAIN_CLEANUP_RESERVE_S, headroom_s: float = SYSTEMD_STOP_HEADROOM_S,
floor_s: float = SYSTEMD_TIMEOUT_STOP_SEC_FLOOR,
) -> int:
"""Seconds systemd ``TimeoutStopSec`` must cover the full stop budget.

View File

@@ -1,15 +1,12 @@
"""Auto-resume restart-loop breaker (defense-3).
Defenses 1 and 2 (the ``_HERMES_GATEWAY`` guard on ``hermes gateway stop|restart``
+ ``terminal_tool``, and the cron-creation lifecycle filter) stop the agent
scheduling its own restart but not every SIGTERM source (``launchctl kickstart``,
a bad external monitor, any repeated crash): the supervisor respawns, the gateway
auto-resumes the restart-interrupted session, whose next turn re-runs the
offending logic. Each such boot is persisted to ``<HERMES_HOME>/gateway/
restart_loop.json``; boots CHAIN while consecutive gaps stay within
``max_gap_seconds`` (a ~150s watchdog-kill cycle trips like a ~10s respawn loop).
Tripped → caller SKIPS auto-resume; inbound messages still served. Any I/O
failure fails OPEN — a broken breaker must never wedge a healthy gateway.
Defenses 1-2 (``_HERMES_GATEWAY`` guard on ``hermes gateway stop|restart`` /
``terminal_tool``, cron lifecycle filter) stop the agent scheduling its own restart
but not every SIGTERM source: the supervisor respawns, the gateway auto-resumes the
restart-interrupted session, whose next turn re-runs the offending logic. Boots are
persisted to ``<HERMES_HOME>/gateway/restart_loop.json`` and CHAIN while gaps stay
within ``max_gap_seconds`` (a ~150s watchdog-kill cycle trips like a ~10s loop).
Tripped → caller SKIPS auto-resume. Any I/O failure fails OPEN, never wedging.
"""
from __future__ import annotations
@@ -28,15 +25,12 @@ logger = logging.getLogger("gateway.run")
# within a few cycles.
DEFAULT_MAX_RESTARTS = 3
DEFAULT_WINDOW_SECONDS = 60
# Longest gap between consecutive restart-interrupted boots that still counts
# them as the SAME loop. A fixed-window prune only sees cycles faster than the
# window (a slower loop drops its own history every boot and never trips);
# chaining on the inter-boot gap is period-agnostic, and real quiet resets it.
# Longest gap between consecutive restart-interrupted boots still counted as the
# SAME loop. A fixed-window prune only sees cycles faster than the window (a slower
# loop drops its history every boot and never trips); chaining on the inter-boot
# gap is period-agnostic, and real quiet resets it.
DEFAULT_MAX_GAP_SECONDS = 300
# Cap the persisted chain; only the newest ``max_restarts`` entries can change
# a verdict, the rest are forensics.
# Only the newest ``max_restarts`` entries can change a verdict; the rest are forensics.
_MAX_STORED_BOOTS = 50
@@ -66,12 +60,10 @@ def _chain_gap(window_seconds: int, max_gap_seconds: int) -> float:
def _chain_ending_at(boots: List[float], ts: float, gap: float) -> List[float]:
"""Unbroken chain of boots leading up to ``ts`` (oldest first).
Walks backwards while each successive gap stays within ``gap``; the first
wider gap ends the chain (older boots belong to a resolved episode).
Nothing recent enough -> empty list: how a healthy gateway forgets a loop.
"""
"""Unbroken chain of boots leading up to ``ts`` (oldest first): walks backwards
while each gap stays within ``gap``; the first wider gap ends the chain (older
boots are a resolved episode). Empty when nothing is recent — how a healthy
gateway forgets a loop."""
chain: List[float] = []
prev = ts
for t in sorted(boots, reverse=True):
@@ -89,15 +81,11 @@ def _chain_ending_at(boots: List[float], ts: float, gap: float) -> List[float]:
def record_restart_interrupted_boot(
window_seconds: int = DEFAULT_WINDOW_SECONDS,
*,
now: Optional[float] = None,
window_seconds: int = DEFAULT_WINDOW_SECONDS, *, now: Optional[float] = None,
max_gap_seconds: int = DEFAULT_MAX_GAP_SECONDS,
) -> List[float]:
"""Record a restart-interrupted boot; return the pruned chain + now (most recent last).
Best-effort — a persistence failure returns the in-memory list without raising.
"""
"""Record a restart-interrupted boot; return the pruned chain + now (most recent
last). A persistence failure returns the in-memory list without raising."""
ts = time.time() if now is None else now
boots = _chain_ending_at(_load_boots(), ts, _chain_gap(window_seconds, max_gap_seconds))
boots.append(ts)
@@ -112,32 +100,19 @@ def clear() -> None:
def check_and_record(
max_restarts: int = DEFAULT_MAX_RESTARTS,
window_seconds: int = DEFAULT_WINDOW_SECONDS,
*,
now: Optional[float] = None,
max_gap_seconds: int = DEFAULT_MAX_GAP_SECONDS,
max_restarts: int = DEFAULT_MAX_RESTARTS, window_seconds: int = DEFAULT_WINDOW_SECONDS, *,
now: Optional[float] = None, max_gap_seconds: int = DEFAULT_MAX_GAP_SECONDS,
) -> bool:
"""Record this boot and return True when auto-resume should be SKIPPED.
The single entry point the gateway calls: appends the current boot, then
checks whether the updated chain has reached ``max_restarts``.
"""
boots = record_restart_interrupted_boot(
window_seconds, now=now, max_gap_seconds=max_gap_seconds
)
"""Gateway entry point: record this boot; True when the chain reached
``max_restarts`` and auto-resume should be SKIPPED."""
boots = record_restart_interrupted_boot(window_seconds, now=now, max_gap_seconds=max_gap_seconds)
tripped = max_restarts > 0 and len(boots) >= max_restarts
if tripped:
logger.warning(
"Restart-loop breaker TRIPPED: %d chained restart-interrupted "
"gateway boots (no gap wider than %ds; threshold %d). Skipping "
"auto-resume to break a suspected SIGTERM-respawn loop (#30719, "
"#81642). Restart-interrupted sessions stay resume-pending and "
"will continue on the next real user message. If this is a false "
"positive, delete %s.",
len(boots),
int(_chain_gap(window_seconds, max_gap_seconds)),
max_restarts,
_state_path(),
"Restart-loop breaker TRIPPED: %d chained restart-interrupted gateway boots (no gap wider than %ds; "
"threshold %d). Skipping auto-resume to break a suspected SIGTERM-respawn loop (#30719, #81642). "
"Restart-interrupted sessions stay resume-pending and will continue on the next real user message. "
"If this is a false positive, delete %s.",
len(boots), int(_chain_gap(window_seconds, max_gap_seconds)), max_restarts, _state_path(),
)
return tripped