refactor(gateway): unify /proc + coercer helpers, merge timestamp regexes, collapse defensive layers in lifecycle/monitoring/media modules

This commit is contained in:
Teknium
2026-09-02 18:23:35 -07:00
parent 113f04616b
commit 98dccb4ca8
18 changed files with 447 additions and 682 deletions
+37 -72
View File
@@ -1,15 +1,12 @@
"""Memory-pressure bounds for the gateway's per-session AIAgent cache.
Each cached ``AIAgent`` pins ``_session_messages`` (the full live transcript,
tool outputs included — 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:
the process's own anonymous RSS against a budget derived from the cgroup limit.
``GatewayRunner._sweep_agent_cache_under_pressure`` uses it to shed LRU
transcripts via the soft-eviction path (rebuilt from the persisted session on
the next turn).
Everything here is pure or read-only. Config lives under ``agent.agent_cache``.
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``.
"""
from __future__ import annotations
@@ -33,6 +30,7 @@ _DEFAULT_MAX_EVICTIONS_PER_PASS = 16
_DEFAULT_PROTECT_RECENT = 8
_BYTES_PER_MB = 1024 * 1024
_OFF_WORDS = frozenset({"", "off", "none", "false", "disabled"})
@dataclass(frozen=True)
@@ -51,7 +49,11 @@ class AgentCacheBounds:
protect_recent: int = _DEFAULT_PROTECT_RECENT
def _positive(value: Any, cast: Callable[[Any], Any]) -> Any:
def _is_int(value: Any) -> bool:
return isinstance(value, int) and not isinstance(value, bool)
def _positive(value: Any, cast: Callable[[Any], Any] = int) -> Any:
"""``cast(value)`` if it is a positive number (bools rejected), else None."""
if isinstance(value, bool) or value is None:
return None
@@ -62,21 +64,13 @@ def _positive(value: Any, cast: Callable[[Any], Any]) -> Any:
return parsed if parsed > 0 else None
def _positive_int(value: Any) -> Optional[int]:
return _positive(value, int)
def _positive_float(value: Any) -> Optional[float]:
return _positive(value, float)
def _cgroup_limit_bytes() -> Optional[int]:
"""Return the memory limit this process runs under, if cgroup-capped.
"""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
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
there), then the root for container-style limits. ``max`` and the v1
near-2^63 sentinel mean unlimited.
"""
if sys.platform != "linux":
@@ -97,18 +91,11 @@ def _cgroup_limit_bytes() -> Optional[int]:
]
for candidate in candidates:
try:
raw = Path(candidate).read_text(encoding="utf-8").strip()
except OSError:
limit = int(Path(candidate).read_text(encoding="utf-8").strip())
except (OSError, ValueError): # unreadable, empty, or "max"
continue
if not raw or raw == "max":
continue
try:
limit = int(raw)
except ValueError:
continue
if limit <= 0 or limit >= (1 << 62):
continue
return limit
if 0 < limit < (1 << 62):
return limit
return None
@@ -134,16 +121,12 @@ def resolve_memory_high_mb(setting: Any) -> Optional[int]:
if isinstance(setting, str):
normalized = setting.strip().lower()
if normalized != "auto":
return (
None
if normalized in ("", "off", "none", "false", "disabled")
else _positive_int(normalized)
)
return None if normalized in _OFF_WORDS else _positive(normalized)
elif isinstance(setting, bool):
if not setting:
return None
else:
return _positive_int(setting)
return _positive(setting)
limit = _cgroup_limit_bytes() or _total_memory_bytes()
if not limit:
@@ -158,46 +141,32 @@ def resolve_agent_cache_bounds(config: Any) -> AgentCacheBounds:
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".
"""
section: Any = None
if isinstance(config, dict):
agent_cfg = config.get("agent")
if isinstance(agent_cfg, dict):
section = agent_cfg.get("agent_cache")
section = (config.get("agent") or {}).get("agent_cache") if isinstance(config, dict) else None
if not isinstance(section, dict):
section = {}
max_evictions = _positive_int(section.get("max_evictions_per_pass"))
protect_recent = section.get("protect_recent")
protect_parsed = _positive_int(protect_recent)
# 0 means "shed anything" — distinct from unset. The bool guard keeps
protect_parsed = _positive(protect_recent)
# 0 means "shed anything" — distinct from unset. The bool guard keeps
# `protect_recent: false` (False == 0) on the default instead of silently
# disabling MRU protection.
if (
protect_parsed is None
and isinstance(protect_recent, int)
and not isinstance(protect_recent, bool)
and protect_recent == 0
):
if protect_parsed is None and _is_int(protect_recent) and protect_recent == 0:
protect_parsed = 0
return AgentCacheBounds(
max_size=_positive_int(section.get("max_size")),
idle_ttl_secs=_positive_float(section.get("idle_ttl_secs")),
max_size=_positive(section.get("max_size")),
idle_ttl_secs=_positive(section.get("idle_ttl_secs"), float),
memory_high_mb=resolve_memory_high_mb(section.get("memory_high_mb", "auto")),
max_evictions_per_pass=(
max_evictions if max_evictions is not None else _DEFAULT_MAX_EVICTIONS_PER_PASS
),
protect_recent=(
protect_parsed if protect_parsed is not None else _DEFAULT_PROTECT_RECENT
),
max_evictions_per_pass=_positive(section.get("max_evictions_per_pass")) or _DEFAULT_MAX_EVICTIONS_PER_PASS,
protect_recent=_DEFAULT_PROTECT_RECENT if protect_parsed is None else protect_parsed,
)
def read_anon_rss_mb() -> Optional[int]:
"""Return the process's anonymous resident memory in MB, or None.
"""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
noise. ``collect_memory_snapshot`` reads ``/proc/self/status`` without a
dependency; psutil covers other platforms (total RSS only).
"""
try:
@@ -224,17 +193,13 @@ def transcript_persistence_caught_up(agent: Any) -> bool:
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)`` by ``AIAgent._flush_messages_to_session_db`` 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.
``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.
"""
messages = getattr(agent, "_session_messages", None)
if not isinstance(messages, list):
return False
flushed = getattr(agent, "_last_flushed_db_idx", None)
if not isinstance(flushed, int) or isinstance(flushed, bool):
return False
return flushed >= len(messages)
return isinstance(messages, list) and _is_int(flushed) and flushed >= len(messages)
def plan_pressure_evictions(
@@ -247,8 +212,8 @@ def plan_pressure_evictions(
"""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
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.
"""
+3 -5
View File
@@ -3,11 +3,9 @@
Runs as ``ExecStopPost=`` after the gateway's main process has exited: the
safety net for long-lived helpers the gateway doesn't track (``adb``, platform
bridges) that would otherwise be orphaned in the cgroup and block
``Restart=always``.
Per-PID SIGKILLs over ``cgroup.procs`` are used deliberately instead of writing
``1`` to ``cgroup.kill``: the kernel has returned ``EINVAL`` on the cgroup-wide
kill while per-PID signal delivery still works.
``Restart=always``. Per-PID SIGKILLs over ``cgroup.procs`` are used instead of
writing ``1`` to ``cgroup.kill``: the kernel has returned ``EINVAL`` on the
cgroup-wide kill while per-PID signal delivery still works.
"""
from __future__ import annotations
+5 -10
View File
@@ -1,15 +1,11 @@
"""Detect when the gateway is running stale code after a hot ``git pull``.
The gateway's ``sys.modules`` is frozen at boot. If the checkout is updated
underneath it (manual ``git pull``, or the window before ``hermes update``'s
graceful restart), a first-time lazy import can resolve a freshly-pulled module
against a stale cached dependency -> ImportError
(``tests/test_stale_utils_module_import.py``). We snapshot the revision at
underneath it, a first-time lazy import can resolve a freshly-pulled module
against a stale cached dependency -> ImportError. We snapshot the revision at
startup so risky callers (e.g. ``/model`` switching) can refuse with a clear
"restart the gateway" message instead.
If the revision can't be read (non-git install, IO error) the boot snapshot
stays ``None`` and detection no-ops — never a false positive.
"restart the gateway" message. If the revision can't be read (non-git install,
IO error) the boot snapshot stays ``None`` and detection no-ops — never a false positive.
"""
from __future__ import annotations
@@ -47,8 +43,7 @@ def _short(fingerprint: str) -> str:
def detect_code_skew() -> tuple[str, str] | None:
"""Return ``(boot_rev, disk_rev)`` short labels if the checkout drifted
since boot, else ``None``."""
"""``(boot_rev, disk_rev)`` short labels if the checkout drifted since boot, else ``None``."""
if _boot_fingerprint is None:
return None
current = _fingerprint()
+3 -6
View File
@@ -22,12 +22,9 @@ def resolve_placeholder_terminal_cwd(
) -> str | None:
"""Return the ``TERMINAL_CWD`` value to set, or ``None`` to leave it unset.
Cases:
- **local** + placeholder → ``MESSAGING_CWD`` or ``home_fallback``
- **docker** + placeholder + mount on + host ``MESSAGING_CWD`` → host path
(for ``terminal_tool`` ``/workspace`` mapping)
- **docker** + placeholder + mount off → ``None`` (sandbox default)
- other non-local backends + placeholder → ``None``
local + placeholder → ``MESSAGING_CWD`` or ``home_fallback``; docker +
placeholder + mount on + host ``MESSAGING_CWD`` → that host path (for the
``/workspace`` mapping); any other non-local backend → ``None`` (sandbox default).
"""
if configured_cwd and configured_cwd not in CWD_PLACEHOLDERS:
return configured_cwd
+19 -37
View File
@@ -1,20 +1,10 @@
"""Disk-usage rollup for ``/api/status``.
Companion to :mod:`gateway.memory_status`: a hosted agent can fill its data
volume (SQLite writes failing, sessions and config saves lost) while its
dashboard and the NAS agent card look healthy. ``gateway/readiness.py``
already probes disk, but readiness is a component verdict, not user-facing
telemetry — this module produces the public block the dashboard SPA and the
NAS availability sweep consume.
Disk is sampled live via :func:`shutil.disk_usage` (one ``statvfs`` call), so
there is no staleness dimension and no ``sampled_at``.
Public-safety: ``/api/status`` is unauthenticated (``PUBLIC_API_PATHS``); this
block carries only coarse numbers (MB, one-decimal percent) and an enum — the same
disclosure class as the ``memory`` block.
Best-effort and read-only: an unreadable filesystem degrades to
volume (SQLite writes failing, config saves lost) while its dashboard looks
healthy. Sampled live via one ``statvfs`` call, so there is no ``sampled_at``.
``/api/status`` is unauthenticated: only coarse numbers (MB, one-decimal
percent) and an enum. Best-effort: an unreadable filesystem degrades to
``pressure="unknown"`` rather than raising into the status endpoint.
"""
@@ -25,10 +15,12 @@ import shutil
from pathlib import Path
from typing import Any, Dict, Optional
from gateway.memory_status import _nonneg_int
logger = logging.getLogger(__name__)
# Percent alone misleads both ways: 90% used on 100 GB leaves 10 GB, while 50%
# on a tiny volume is one download from write failures. So percent triggers are
# 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
@@ -41,21 +33,15 @@ _ELEVATED_HEADROOM_MB = 4096
_BYTES_PER_MB = 1024 * 1024
def _coerce_mb(value: Any) -> Optional[int]:
if isinstance(value, bool) or not isinstance(value, int) or value < 0:
return None
return value
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".
"""
free = _coerce_mb(free_mb)
total = _coerce_mb(total_mb)
if free is None or total is None or total <= 0:
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 (
@@ -71,16 +57,10 @@ 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; on hosted images every profile shares one volume).
Always returns a dict and never raises — an unreadable/unmounted filesystem
yields ``{"pressure": "unknown", ...}``.
``memory`` block). Always returns a dict and never raises — an unreadable
or unmounted filesystem yields ``{"pressure": "unknown", ...}``.
"""
status: Dict[str, Any] = {
"pressure": "unknown",
"total_mb": None,
"free_mb": None,
"used_percent": None,
}
status: Dict[str, Any] = {"pressure": "unknown", "total_mb": None, "free_mb": None, "used_percent": None}
try:
if home is None:
from hermes_constants import get_hermes_home
@@ -93,8 +73,10 @@ def collect_disk_status(home: Optional[Path] = None) -> Dict[str, Any]:
return status
total_mb = usage.total // _BYTES_PER_MB
free_mb = usage.free // _BYTES_PER_MB
status["total_mb"] = total_mb
status["free_mb"] = free_mb
status["used_percent"] = round((usage.used / usage.total) * 100, 1)
status["pressure"] = classify_disk_pressure(free_mb, total_mb)
status.update(
total_mb=total_mb,
free_mb=free_mb,
used_percent=round((usage.used / usage.total) * 100, 1),
pressure=classify_disk_pressure(free_mb, total_mb),
)
return status
+41 -72
View File
@@ -1,38 +1,27 @@
"""External drain-control marker contract (dashboard → gateway).
There is no control channel into a running gateway, so the dashboard's
begin/cancel-drain endpoint talks to it the same way ``.restart_notify.json``
does: it writes (or removes) a marker file and a gateway watcher reacts. This
module owns the contract so writer and reader can never disagree.
Contract (presence-based):
* begin-drain → write ``{HERMES_HOME}/.drain_request.json`` with
``{"action": "drain", "requested_at": <iso>, "principal": <str>,
"epoch": <instantiation-epoch>, "suppress_notification": <bool>}``.
* cancel-drain → remove the marker.
* The watcher treats presence of an ACTIVE marker as "external drain": flip
``gateway_state -> "draining"`` and stop accepting new turns. Absence or a
stale marker means "not draining" (revert to ``running``).
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 that the drain-gated action (update,
migrate, env edit) ends in; honouring it would park the fresh gateway in
``draining`` forever. A marker is stamped with this instantiation's epoch
and ignored only on a *definite* mismatch.
* expiry: a same-epoch orphan (action completed without a restart, writer
never cancelled) is ignored once ``requested_at`` is older than
:data:`DRAIN_REQUEST_MAX_AGE_SECONDS`. Re-calling :func:`write_drain_request`
refreshes the timestamp — the sanctioned keep-alive for a long drain.
marker survives the machine restart that a drain-gated action ends in;
honouring it would park the fresh gateway in ``draining`` forever.
* expiry: a same-epoch orphan (writer never cancelled) is ignored once
``requested_at`` is older than :data:`DRAIN_REQUEST_MAX_AGE_SECONDS`;
re-calling :func:`write_drain_request` refreshes it (long-drain keep-alive).
Reading never raises: a malformed/half-written file reads as "present but
contentless" (``{}``), which still counts as drain-active — fail-safe toward
quiescing. Both staleness checks only reject on a *definite* verdict; a marker
with no epoch/timestamp, or a host without ``/proc``, degrades to the
presence-only behaviour and never fails closed.
Reading never raises: a malformed/half-written file reads as ``{}``, which still
counts as drain-active (fail-safe toward quiescing). Both staleness checks
reject only on a *definite* verdict — a marker with no epoch/timestamp, or a
host without ``/proc``, degrades to presence-only and never fails closed.
"""
from __future__ import annotations
import contextlib
import functools
import json
import logging
@@ -42,7 +31,6 @@ from typing import Any, Optional
from hermes_constants import get_hermes_home
from utils import atomic_json_write
import contextlib
_log = logging.getLogger(__name__)
@@ -52,9 +40,8 @@ _DRAIN_REQUEST_FILENAME = ".drain_request.json"
# leaked marker can cause. Long drains refresh the marker instead of raising this.
DRAIN_REQUEST_MAX_AGE_SECONDS = 3600.0
# Dedup for the expired-marker warning: the watcher re-reads every second and an
# expired orphan sits on disk until removed. Keyed by ``requested_at`` so a
# keep-alive re-write that later expires again logs again.
# Dedup for the expired-marker warning (the watcher re-reads every second).
# Keyed by ``requested_at`` so a keep-alive re-write that later expires logs again.
_expiry_logged_for: Optional[str] = None
@@ -62,27 +49,22 @@ _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 (so an s6 respawn of just the gateway, or a
host ``hermes gateway restart`` under systemd/launchd, keeps honouring an
in-flight drain) but changes whenever the machine is recreated: boot_id
changes on a VM/microVM reboot, PID 1's start time on a plain
``docker restart`` (same host kernel, brand-new ``/init``).
Returns ``""`` when neither source is readable (non-Linux, no ``/proc``),
which disables the epoch check downstream — never fail-closed.
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.
"""
boot_id = ""
with contextlib.suppress(OSError):
boot_id = Path("/proc/sys/kernel/random/boot_id").read_text(encoding="utf-8").strip()
pid1_start = ""
try:
with contextlib.suppress(OSError, IndexError):
# "<pid> (<comm>) <state> ...": comm may contain spaces/parens, so split
# on the LAST ')'. starttime is field 22 (1-indexed) = tail index 19.
stat = Path("/proc/1/stat").read_text(encoding="utf-8")
pid1_start = stat.rsplit(")", 1)[1].split()[19]
except (OSError, IndexError):
pass
if not boot_id and not pid1_start:
return ""
@@ -91,8 +73,7 @@ def current_instantiation_epoch() -> str:
def drain_request_path(home: Optional[Path] = None) -> Path:
"""Absolute path to the drain-request marker, respecting HERMES_HOME."""
base = home if home is not None else get_hermes_home()
return Path(base) / _DRAIN_REQUEST_FILENAME
return Path(home if home is not None else get_hermes_home()) / _DRAIN_REQUEST_FILENAME
def write_drain_request(
@@ -104,13 +85,11 @@ def write_drain_request(
"""Write the begin-drain marker atomically. Returns the payload written.
Idempotent: re-writing refreshes ``requested_at`` (keep-alive past the
max-age). ``suppress_notification`` asks the shutdown that ends this drain
to skip ONLY the home-channel "gateway shutting down" broadcast (the
per-active-session interrupt ping is never suppressed); the policy of
which drains are quiet lives entirely in the caller. Defaults False so
legacy/operator drains behave exactly as before. The marker is stamped
with :func:`current_instantiation_epoch` so a copy surviving a machine
restart on the durable volume can be recognised as stale.
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.
"""
payload = {
"action": "drain",
@@ -136,20 +115,12 @@ def clear_drain_request(*, home: Optional[Path] = None) -> bool:
return False
def _marker_epoch_is_stale(body: dict[str, Any]) -> bool:
"""True iff ``body``'s epoch is a *definite* mismatch (both epochs known and differ)."""
current = current_instantiation_epoch()
marker_epoch = body.get("epoch")
return bool(current and marker_epoch and marker_epoch != current)
def _marker_is_expired(body: dict[str, Any]) -> bool:
"""True iff ``requested_at`` parses AND is older than the max-age.
Missing/unparseable timestamps and future-dated ones (clock skew) are
honoured. The expiry is logged once per marker, not per poll: this path
only fires when a writer leaked a marker, and the warning is the
operator's breadcrumb.
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.
"""
global _expiry_logged_for
raw = body.get("requested_at")
@@ -179,15 +150,14 @@ def _marker_is_expired(body: dict[str, Any]) -> bool:
return True
def _marker_is_stale(body: dict[str, Any]) -> bool:
"""True iff the marker is definitely from a drain that is already over."""
return _marker_epoch_is_stale(body) or _marker_is_expired(body)
def _active_drain_body(home: Optional[Path]) -> Optional[dict[str, Any]]:
"""Marker body if present AND not stale, else None."""
"""Marker body if present AND not stale (definite epoch mismatch or expired), else None."""
body = read_drain_request(home=home)
if body is None or _marker_is_stale(body):
if body is None:
return None
current = current_instantiation_epoch()
marker_epoch = body.get("epoch")
if (current and marker_epoch and marker_epoch != current) or _marker_is_expired(body):
return None
return body
@@ -200,10 +170,9 @@ 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.
Uses the same activeness rule as :func:`drain_requested`, so an orphaned
marker can never silence a fresh gateway's legitimate broadcast. A legacy
marker without the field or a contentless ``{}`` body reads as False
(fail toward the louder behaviour).
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).
"""
body = _active_drain_body(home)
return bool(body and body.get("suppress_notification"))
+58 -53
View File
@@ -25,6 +25,7 @@ Handlers posting a follow-up into the same Telegram forum-topic should pass
import asyncio
import importlib.util
import sys
from pathlib import Path
from typing import Any, Callable, Dict, List, Optional
import yaml
@@ -35,6 +36,49 @@ from hermes_cli.config import get_hermes_home
HOOKS_DIR = get_hermes_home() / "hooks"
def _load_hook_dir(hook_dir: Path) -> Optional[tuple]:
"""``(name, events, handle_fn, description)`` for a valid hook dir, else None (reason printed)."""
manifest_path = hook_dir / "HOOK.yaml"
handler_path = hook_dir / "handler.py"
if not manifest_path.exists() or not handler_path.exists():
return None
manifest = yaml.safe_load(manifest_path.read_text(encoding="utf-8"))
if not manifest or not isinstance(manifest, dict):
print(f"[hooks] Skipping {hook_dir.name}: invalid HOOK.yaml", flush=True)
return None
hook_name = manifest.get("name", hook_dir.name)
events = manifest.get("events", [])
if not events:
print(f"[hooks] Skipping {hook_name}: no events declared", flush=True)
return None
# Register in sys.modules BEFORE exec_module so Pydantic/dataclass forward
# references (from ``from __future__ import annotations``) resolve; otherwise
# a handler declaring a BaseModel fails at first dispatch with
# "TypeAdapter ... is not fully defined".
module_name = f"hermes_hook_{hook_name}"
spec = importlib.util.spec_from_file_location(module_name, handler_path)
if spec is None or spec.loader is None:
print(f"[hooks] Skipping {hook_name}: could not load handler.py", flush=True)
return None
module = importlib.util.module_from_spec(spec)
sys.modules[module_name] = module
try:
spec.loader.exec_module(module)
except Exception:
sys.modules.pop(module_name, None)
raise
handle_fn = getattr(module, "handle", None)
if handle_fn is None:
print(f"[hooks] Skipping {hook_name}: no 'handle' function found", flush=True)
return None
return hook_name, events, handle_fn, manifest.get("description", "")
class HookRegistry:
"""Discovers, loads, and fires event hooks."""
@@ -60,62 +104,23 @@ class HookRegistry:
for hook_dir in sorted(HOOKS_DIR.iterdir()):
if not hook_dir.is_dir():
continue
manifest_path = hook_dir / "HOOK.yaml"
handler_path = hook_dir / "handler.py"
if not manifest_path.exists() or not handler_path.exists():
continue
try:
manifest = yaml.safe_load(manifest_path.read_text(encoding="utf-8"))
if not manifest or not isinstance(manifest, dict):
print(f"[hooks] Skipping {hook_dir.name}: invalid HOOK.yaml", flush=True)
continue
hook_name = manifest.get("name", hook_dir.name)
events = manifest.get("events", [])
if not events:
print(f"[hooks] Skipping {hook_name}: no events declared", flush=True)
continue
# Register in sys.modules BEFORE exec_module so Pydantic/dataclass
# forward references (from ``from __future__ import annotations``)
# resolve; otherwise a handler declaring a BaseModel fails at first
# dispatch with "TypeAdapter ... is not fully defined".
module_name = f"hermes_hook_{hook_name}"
spec = importlib.util.spec_from_file_location(module_name, handler_path)
if spec is None or spec.loader is None:
print(f"[hooks] Skipping {hook_name}: could not load handler.py", flush=True)
continue
module = importlib.util.module_from_spec(spec)
sys.modules[module_name] = module
try:
spec.loader.exec_module(module)
except Exception:
sys.modules.pop(module_name, None)
raise
handle_fn = getattr(module, "handle", None)
if handle_fn is None:
print(f"[hooks] Skipping {hook_name}: no 'handle' function found", flush=True)
continue
for event in events:
self._handlers.setdefault(event, []).append(handle_fn)
self._loaded_hooks.append({
"name": hook_name,
"description": manifest.get("description", ""),
"events": events,
"path": str(hook_dir),
})
print(f"[hooks] Loaded hook '{hook_name}' for events: {events}", flush=True)
loaded = _load_hook_dir(hook_dir)
except Exception as e:
print(f"[hooks] Error loading hook {hook_dir.name}: {e}", flush=True)
continue
if loaded is None:
continue
hook_name, events, handle_fn, description = loaded
for event in events:
self._handlers.setdefault(event, []).append(handle_fn)
self._loaded_hooks.append({
"name": hook_name,
"description": description,
"events": events,
"path": str(hook_dir),
})
print(f"[hooks] Loaded hook '{hook_name}' for events: {events}", flush=True)
def _resolve_handlers(self, event_type: str) -> List[Callable]:
"""Exact-match handlers first, then ``<base>:*`` wildcard handlers.
+83 -138
View File
@@ -1,24 +1,12 @@
"""Gateway lifecycle ledger — durable termination-reason evidence.
Graceful shutdowns already leave forensics (:mod:`gateway.shutdown_forensics`,
``gateway-exit-diag.log``); an **unclean death** (SIGKILL, kernel OOM kill, VM
death) runs no handler, so the next boot would otherwise have no record of it.
A tiny state machine persisted to ``<HERMES_HOME>/state/gateway.lifecycle.json``
closes that gap:
* :func:`record_startup` reads the previous life's sentinel. ``phase ==
"running"`` means no exit path ran → unclean death. The finding (with the
last heartbeat's memory sample) is appended to ``gateway-exit-diag.log`` as a
``gateway.previous_unclean_exit`` record and logged at WARNING, then the
sentinel is rewritten as ``phase=running`` for the new life.
* :func:`mark_exited` rewrites the sentinel as ``phase=exited`` on every clean
exit path (``_exit_after_graceful_shutdown`` and the watchdog ``os._exit`` sites).
* :func:`sample_memory` is the cheap /proc snapshot the 30s loop heartbeat
embeds, so OOM crash cycles are classifiable from the volume alone.
Everything here is best-effort: a forensics failure must never affect the
gateway lifecycle it is observing.
Graceful shutdowns leave forensics (``gateway-exit-diag.log``); an unclean
death (SIGKILL, kernel OOM, VM death) runs no handler. A sentinel at
``<HERMES_HOME>/state/gateway.lifecycle.json`` closes the gap:
:func:`record_startup` finds ``phase == "running"`` from the previous life →
unclean death, appended to the exit-diag log as ``gateway.previous_unclean_exit``
and logged at WARNING; :func:`mark_exited` rewrites ``phase=exited`` on every
clean exit path. Best-effort: forensics must never affect the lifecycle.
"""
from __future__ import annotations
@@ -28,6 +16,7 @@ import logging
import os
import sqlite3
import time
from contextlib import closing
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Dict, Optional
@@ -38,11 +27,6 @@ _LIFECYCLE_RELATIVE = ("state", "gateway.lifecycle.json")
_EXIT_DIAG_RELATIVE = ("logs", "gateway-exit-diag.log")
_STATE_DB_RELATIVE = ("state.db",)
# Heuristic OOM-suspicion thresholds on the last heartbeat's memory sample.
# Deliberately conservative: only a hint; classification stays with the reader.
_LOW_MEM_AVAILABLE_KIB = 64 * 1024 # < 64 MiB available
_LOW_MEM_AVAILABLE_FRACTION = 0.05 # < 5% of MemTotal available
def _process_hermes_home() -> Path:
"""HERMES_HOME for process-level identity files (ignore task overrides)."""
@@ -55,8 +39,7 @@ def _process_hermes_home() -> Path:
def _home_path(home: Optional[Path], relative: tuple) -> Path:
base = home if home is not None else _process_hermes_home()
return base.joinpath(*relative)
return (home if home is not None else _process_hermes_home()).joinpath(*relative)
def get_lifecycle_sentinel_path(home: Optional[Path] = None) -> Path:
@@ -64,38 +47,41 @@ def get_lifecycle_sentinel_path(home: Optional[Path] = None) -> Path:
return _home_path(home, _LIFECYCLE_RELATIVE)
def sample_memory() -> Dict[str, Any]:
"""Cheap memory snapshot: own RSS + system availability + swap.
def _now_iso() -> str:
return datetime.now(timezone.utc).isoformat()
Pure ``/proc`` reads, Linux-only (returns ``{}`` elsewhere), never
raises. Values in KiB to match the kernel's units.
"""
sample: Dict[str, Any] = {}
def _proc_fields(path: str, wanted: Dict[str, str]) -> Dict[str, int]:
"""``{dst: int}`` for each ``src: dst`` key found in a ``Key: value`` /proc file."""
found: Dict[str, int] = {}
try:
with open("/proc/self/status", encoding="utf-8") as fh:
for line in fh:
if line.startswith("VmRSS:"):
sample["rss_kib"] = int(line.split()[1])
break
except (OSError, ValueError, IndexError):
pass
try:
meminfo: Dict[str, int] = {}
wanted = {"MemTotal", "MemAvailable", "SwapTotal", "SwapFree"}
with open("/proc/meminfo", encoding="utf-8") as fh:
with open(path, encoding="utf-8") as fh:
for line in fh:
key = line.split(":", 1)[0]
if key in wanted:
meminfo[key] = int(line.split()[1])
if len(meminfo) == len(wanted):
found[wanted[key]] = int(line.split()[1])
if len(found) == len(wanted):
break
for src, dst in (("MemTotal", "mem_total_kib"), ("MemAvailable", "mem_available_kib")):
if src in meminfo:
sample[dst] = meminfo[src]
if "SwapTotal" in meminfo and "SwapFree" in meminfo:
sample["swap_used_kib"] = meminfo["SwapTotal"] - meminfo["SwapFree"]
except (OSError, ValueError, IndexError):
pass
return {}
return found
def sample_memory() -> Dict[str, Any]:
"""Cheap /proc snapshot (KiB, kernel units): own RSS + MemTotal/MemAvailable + swap used.
Linux-only (``{}`` elsewhere), never raises; embedded in the 30s loop
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",
})
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:
sample["swap_used_kib"] = swap_total - swap_free
return sample
@@ -119,8 +105,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`` records so existing tooling greps both)."""
"""Append a JSON line to gateway-exit-diag.log (same format as the CLI's ``_exit_diag``)."""
path = _home_path(home, _EXIT_DIAG_RELATIVE)
try:
path.parent.mkdir(parents=True, exist_ok=True)
@@ -133,9 +118,8 @@ def _append_exit_diag(record: Dict[str, Any], home: Optional[Path]) -> None:
def _pid_alive_with_start_time(pid: Any, start_time: Any) -> bool:
"""True when ``pid`` is a live process matching ``start_time`` (±2s).
Guards the takeover race: during ``--replace`` the old gateway can still
be mid-teardown when the new one boots — a live matching owner is a
planned handover, not an unclean death.
Guards the ``--replace`` takeover race: a live matching owner mid-teardown
is a planned handover, not an unclean death.
"""
try:
pid_int = int(pid)
@@ -144,9 +128,8 @@ def _pid_alive_with_start_time(pid: Any, start_time: Any) -> bool:
if pid_int <= 0:
return False
try:
# NOT os.kill(pid, 0): on Windows that sends CTRL_C_EVENT to the
# target's console group. _pid_exists is the canonical no-kill probe.
from gateway.status import _pid_exists
# NOT os.kill(pid, 0): on Windows that sends CTRL_C_EVENT to the target's console group.
from gateway.status import _pid_exists, get_process_start_time
if not _pid_exists(pid_int):
return False
@@ -155,8 +138,6 @@ def _pid_alive_with_start_time(pid: Any, start_time: Any) -> bool:
if start_time is None:
return True # alive; can't disambiguate PID reuse — err on "alive"
try:
from gateway.status import get_process_start_time
actual = get_process_start_time(pid_int)
return actual is None or abs(float(actual) - float(start_time)) <= 2.0
except Exception:
@@ -164,9 +145,7 @@ def _pid_alive_with_start_time(pid: Any, start_time: Any) -> bool:
def detect_unclean_exit(home: Optional[Path] = None) -> Optional[Dict[str, Any]]:
"""Inspect the previous life's sentinel; return an evidence dict when it
died uncleanly, else ``None``. Read-only — does not rewrite the sentinel.
"""
"""Evidence dict when the previous life died uncleanly, else ``None``. Read-only."""
sentinel = _read_json(get_lifecycle_sentinel_path(home))
if not sentinel or sentinel.get("phase") != "running":
return None
@@ -178,7 +157,6 @@ def detect_unclean_exit(home: Optional[Path] = None) -> Optional[Dict[str, Any]]
"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.
try:
from gateway.shutdown_watchdog import get_loop_heartbeat_path
@@ -191,15 +169,15 @@ 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
total = mem.get("mem_total_kib")
avail = mem.get("mem_available_kib")
# 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 < _LOW_MEM_AVAILABLE_KIB
or (
isinstance(total, int)
and total > 0
and avail / total < _LOW_MEM_AVAILABLE_FRACTION
)
avail < _CRITICAL_AVAILABLE_KIB
or (isinstance(total, int) and total > 0 and avail / total < _CRITICAL_AVAILABLE_FRACTION)
):
evidence["suspected_oom"] = True
return evidence
@@ -209,22 +187,18 @@ def check_state_db_integrity(home: Optional[Path] = None) -> str:
"""Return ``"ok"``, ``"absent"``, or the first ``quick_check`` complaint.
Called only after an unclean death — a SIGKILL mid-WAL-checkpoint can leave
half-written b-tree pages (see ``_enforce_macos_synchronous_full`` in
:mod:`hermes_state`). ``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 itself writes nothing. Never raises.
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.
"""
path = _home_path(home, _STATE_DB_RELATIVE)
if not path.exists():
return "absent"
try:
conn = sqlite3.connect(str(path))
try:
with closing(sqlite3.connect(str(path))) as conn:
row = conn.execute("PRAGMA quick_check(1)").fetchone()
finally:
conn.close()
except Exception as exc: # sqlite3.Error, OSError, anything
except Exception as exc:
return f"check-failed: {exc}"
if not row or row[0] is None:
return "check-failed: no result"
@@ -232,12 +206,10 @@ def check_state_db_integrity(home: Optional[Path] = None) -> str:
def record_startup(home: Optional[Path] = None) -> Optional[Dict[str, Any]]:
"""Boot-time entry point: report any unclean previous exit, then claim
the sentinel for the current life.
"""Boot entry point: report any unclean previous exit, then claim the sentinel.
Returns the unclean-exit evidence dict (also persisted to
``gateway-exit-diag.log`` and logged at WARNING) or ``None``. Never
raises.
Returns the evidence dict (also persisted to ``gateway-exit-diag.log`` and
logged at WARNING) or ``None``. Never raises.
"""
evidence: Optional[Dict[str, Any]] = None
try:
@@ -253,36 +225,24 @@ def record_startup(home: Optional[Path] = None) -> Optional[Dict[str, Any]]:
"repaired. Run `hermes doctor`.",
verdict,
)
_append_exit_diag({
"ts": datetime.now(timezone.utc).isoformat(),
"tag": "gateway.previous_unclean_exit",
"pid": os.getpid(),
**evidence,
}, home)
_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),
evidence.get("prior_pid"), evidence.get("prior_started_at"), evidence.get("last_heartbeat_at"),
evidence.get("last_heartbeat_mem"), evidence.get("suspected_oom", False),
)
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": datetime.now(timezone.utc).isoformat(),
}
# Carry the verdict on the PREVIOUS life forward on the new sentinel: it
# is the only machine-readable copy (the exit-diag log is append-only
# prose) and /api/status reads it to report an OOM restart. Scoped to
# this life only — the next clean exit or boot rewrites the sentinel.
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
# machine-readable copy (/api/status reads it to report an OOM restart).
# Scoped to this life — the next clean exit or boot rewrites the sentinel.
if evidence is not None:
claim["prior_unclean_exit"] = True
if evidence.get("suspected_oom"):
@@ -293,49 +253,34 @@ def record_startup(home: Optional[Path] = None) -> Optional[Dict[str, Any]]:
return evidence
def mark_exited(
exit_code: Optional[int] = None,
reason: str = "graceful_shutdown",
home: Optional[Path] = None,
) -> None:
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 the sentinel when it is provably owned by this process:
during a ``--replace`` takeover the replacement claims the sentinel before
the old process finishes teardown, and the old life must not clobber the
new owner's ``running`` phase. A sentinel with ``pid=None`` (or malformed)
has unknown ownership and is likewise left alone.
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.
"""
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": datetime.now(timezone.utc).isoformat(),
"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)
def read_prior_exit_label(profile_home: Path) -> str:
"""Container-boot helper: ``clean`` / ``unclean`` / ``unknown`` summary of
how the profile's last gateway life ended. Read-only, exception-free;
used by ``hermes_cli.container_boot`` to annotate ``container-boot.log``.
"""``clean`` / ``unclean`` / ``unknown`` for how the profile's last gateway life ended.
Read-only, exception-free; annotates ``container-boot.log``. ``running`` is
unclean because the old PID namespace is gone at container boot.
"""
try:
sentinel = _read_json(get_lifecycle_sentinel_path(profile_home))
if not sentinel:
return "unknown"
phase = sentinel.get("phase")
if phase == "exited":
return "clean"
if phase == "running":
# Old PID namespace is gone at container boot — never exited cleanly.
return "unclean"
phase = (_read_json(get_lifecycle_sentinel_path(profile_home)) or {}).get("phase")
return {"exited": "clean", "running": "unclean"}.get(phase, "unknown")
except Exception:
pass
return "unknown"
return "unknown"
+6 -9
View File
@@ -7,15 +7,12 @@ from environment variables:
- ``HERMES_MEDIA_ALLOW_DIRS`` <- gateway.media_delivery_allow_dirs
- ``HERMES_MEDIA_TRUST_RECENT_FILES`` <- gateway.trust_recent_files
The translation used to run only in gateway startup, so standalone delivery
paths (``hermes cron run``, ``hermes send``, a standalone cron tick) filtered
MEDIA paths under a different policy and silently dropped attachments in
strict/allowlisted deployments (text is unaffected -- only media goes through
path validation). ``apply_media_policy_env()`` is the shared,
idempotent helper every delivery entrypoint calls before filtering media paths.
Precedence: an explicitly-set environment variable WINS over config.yaml, so a
shell-exported override (and gateway startup's own earlier run) survives.
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.
"""
from __future__ import annotations
+29 -48
View File
@@ -1,19 +1,13 @@
"""Repair model-mangled ``computer_use`` screenshot paths in final responses.
``computer_use`` persists a screenshot into the Hermes 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:`` directive, so delivery-path validation rejects it and drops the
attachment.
The repair is deliberately narrow: it only rewrites paths inside a response that
*already* carries an explicit ``MEDIA:`` directive, and only when the directive's
generated ``computer_use_<uuid>`` basename exactly matches a canonical screenshot
path returned by ``computer_use`` in the current turn. It never auto-attaches
captures; normal media path validation still runs afterwards.
Own module (like ``gateway/media_policy.py``) so the gateway turn path, gateway
background tasks and cron delivery share one implementation.
``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:``
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.
"""
from __future__ import annotations
@@ -29,15 +23,12 @@ logger = logging.getLogger(__name__)
# root, or UNC share. Shared string so the summary regex stays in sync.
_ABS_PATH_PREFIX_PATTERN = r"(?:[A-Za-z]:[/\\]|/|\\\\)"
_ABS_PATH_PREFIX_RE = re.compile(r"^" + _ABS_PATH_PREFIX_PATTERN)
_CAPTURE_BASENAME_PATTERN = r"computer_use_[0-9a-f]{32}\.(?:png|jpe?g)"
_COMPUTER_USE_CAPTURE_BASENAME_RE = re.compile(
r"^computer_use_[0-9a-f]{32}\.(?:png|jpe?g)$",
re.IGNORECASE,
)
_COMPUTER_USE_CAPTURE_BASENAME_RE = re.compile(r"^" + _CAPTURE_BASENAME_PATTERN + r"$", re.IGNORECASE)
_COMPUTER_USE_CAPTURE_SUMMARY_RE = re.compile(
r"\(shareable screenshot saved to "
r"(?P<path>" + _ABS_PATH_PREFIX_PATTERN + r"[^\r\n]*?"
r"computer_use_[0-9a-f]{32}\.(?:png|jpe?g))\)",
r"(?P<path>" + _ABS_PATH_PREFIX_PATTERN + r"[^\r\n]*?" + _CAPTURE_BASENAME_PATTERN + r")\)",
re.IGNORECASE,
)
@@ -86,29 +77,23 @@ def _iter_computer_use_capture_paths(content: Any) -> Iterator[str]:
return
for match in _COMPUTER_USE_CAPTURE_SUMMARY_RE.finditer(content):
yield match.group("path").strip()
return
if isinstance(content, list):
elif isinstance(content, list):
for part in content:
yield from _iter_computer_use_capture_paths(part)
return
if not isinstance(content, dict):
return
screenshot_path = content.get("screenshot_path")
if isinstance(screenshot_path, str):
yield screenshot_path
meta = content.get("meta")
if isinstance(meta, dict) and isinstance(meta.get("screenshot_path"), str):
yield meta["screenshot_path"]
# Producer shapes (tools/computer_use/tool.py::_capture_response):
# content/text = multimodal parts; text_summary/summary = the line carrying
# "(shareable screenshot saved to ...)".
for field in ("content", "text", "text_summary", "summary"):
nested = content.get(field)
if isinstance(nested, (str, dict, list)):
yield from _iter_computer_use_capture_paths(nested)
elif isinstance(content, dict):
screenshot_path = content.get("screenshot_path")
if isinstance(screenshot_path, str):
yield screenshot_path
meta = content.get("meta")
if isinstance(meta, dict) and isinstance(meta.get("screenshot_path"), str):
yield meta["screenshot_path"]
# Producer shapes (tools/computer_use/tool.py::_capture_response):
# content/text = multimodal parts; text_summary/summary = the line carrying
# "(shareable screenshot saved to ...)".
for field in ("content", "text", "text_summary", "summary"):
nested = content.get(field)
if isinstance(nested, (str, dict, list)):
yield from _iter_computer_use_capture_paths(nested)
def _current_turn_messages(messages: List[Dict[str, Any]], history_offset: int) -> List[Dict[str, Any]]:
@@ -135,13 +120,11 @@ def repair_explicit_computer_use_media_paths(
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
turn. Fail-open: the repair is cosmetic, so any unexpected error returns
the response unchanged rather than aborting delivery.
"""
try:
return _repair_explicit_computer_use_media_paths_inner(
response, messages, history_offset
)
return _repair_explicit_computer_use_media_paths_inner(response, messages, history_offset)
except Exception:
logger.debug("computer_use media path repair failed", exc_info=True)
return response
@@ -163,9 +146,7 @@ def _repair_explicit_computer_use_media_paths_inner(
if msg.get("role") not in {"tool", "function"}:
continue
call_id = str(msg.get("tool_call_id") or msg.get("call_id") or "")
tool_name = str(
msg.get("name") or msg.get("tool_name") or call_id_names.get(call_id) or ""
)
tool_name = str(msg.get("name") or msg.get("tool_name") or call_id_names.get(call_id) or "")
if tool_name != "computer_use":
continue
for path in _iter_computer_use_capture_paths(msg.get("content")):
+13 -23
View File
@@ -1,15 +1,15 @@
"""Periodic process memory usage logging for the gateway.
Ported from cline/cline#10343 (src/standalone/memory-monitor.ts). Emits one
grep-friendly ``[MEMORY] ...`` line every N seconds (default 300) from a daemon
thread so slow leaks in the long-lived gateway show up as an RSS time series in
``agent.log`` / ``gateway.log``. A baseline snapshot is logged on start and a
final one on stop. Uses stdlib ``resource`` first, ``psutil`` as fallback
(Windows); if neither works the monitor warns once and stays disabled.
Emits one grep-friendly ``[MEMORY] ...`` line every N seconds (default 300)
from a daemon thread so slow leaks show up as an RSS time series in the logs.
A baseline snapshot is logged on start and a final one on stop. Uses stdlib
``resource`` first, ``psutil`` as fallback (Windows); if neither works the
monitor warns once and stays disabled.
"""
from __future__ import annotations
import contextlib
import gc
import logging
import os
@@ -17,7 +17,6 @@ import sys
import threading
import time
from typing import Optional
import contextlib
logger = logging.getLogger(__name__)
@@ -49,25 +48,16 @@ def _get_rss_mb() -> Optional[int]:
def log_memory_usage(prefix: str = "") -> None:
"""Log current memory usage as ``[MEMORY] [<prefix> ]rss=... gc=... threads=... uptime=...``.
Safe to call on-demand from any thread at lifecycle moments.
"""
"""Log ``[MEMORY] [<prefix> ]rss=... gc=... threads=... uptime=...``; safe from any thread."""
rss = _get_rss_mb()
uptime = int(time.monotonic() - _start_time) if _start_time else 0
try:
gc_counts = gc.get_count() # (gen0, gen1, gen2)
except Exception:
gc_counts = (0, 0, 0)
try:
thread_count = threading.active_count()
except Exception:
thread_count = 0
tag = f"{prefix} " if prefix else ""
rss_text = "unavailable" if rss is None else f"{rss}MB"
logger.info(
"[MEMORY] %srss=%s gc=%s threads=%d uptime=%ds", tag, rss_text, gc_counts, thread_count, uptime
"[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(),
uptime,
)
+21 -41
View File
@@ -1,17 +1,10 @@
"""Memory status rollup for ``/api/status``.
Read side for memory-pressure signals the gateway already persists but only
logs: the 30s ``state/gateway.heartbeat`` (RSS + MemAvailable/MemTotal + swap,
from :func:`gateway.shutdown_watchdog.write_loop_heartbeat`) and the lifecycle
sentinel's ``suspected_oom`` flag (:func:`gateway.lifecycle_ledger.record_startup`).
Distills them into a compact block the dashboard SPA and the NAS availability
sweep consume — no new sampling, no IPC, two small file reads.
Public-safety: ``/api/status`` is unauthenticated (``PUBLIC_API_PATHS``), so this
block carries only coarse numbers (MB granularity), enums, and booleans — the
same disclosure class as ``active_agents`` / ``nous_session_valid``.
Best-effort and read-only: 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 the block carries only coarse numbers
(MB), enums and booleans. Best-effort: a missing/corrupt file degrades to
``pressure="unknown"`` rather than raising into the status endpoint.
"""
@@ -24,21 +17,18 @@ from typing import Any, Dict, Optional
logger = logging.getLogger(__name__)
# Thresholds on system MemAvailable. ``critical`` mirrors the lifecycle ledger's
# OOM-suspicion heuristics (_LOW_MEM_AVAILABLE_KIB / _LOW_MEM_AVAILABLE_FRACTION):
# a level that would make a later unclean death "suspected OOM" should already
# warn while the process is alive.
# Thresholds on system MemAvailable. ``critical`` doubles as the lifecycle
# ledger's OOM-suspicion heuristic: a level that makes a later unclean death
# "suspected OOM" already warns while the process is alive.
_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
# Writer cadence is 30s (DEFAULT_HEARTBEAT_INTERVAL_S); 150s tolerates a briefly
# stalled loop without letting a long-dead gateway's last sample pose as current.
# Writer cadence is 30s; 150s tolerates a briefly stalled loop without letting
# a long-dead gateway's last sample pose as current.
_HEARTBEAT_FRESH_TTL_S = 150.0
_KIB_PER_MB = 1024
def _nonneg_int(value: Any) -> Optional[int]:
"""Return *value* if it is a non-negative int (bools rejected), else None."""
@@ -49,7 +39,7 @@ def _nonneg_int(value: Any) -> Optional[int]:
def _mb(kib: Any) -> Optional[int]:
kib = _nonneg_int(kib)
return None if kib is None else kib // _KIB_PER_MB
return None if kib is None else kib // 1024
def _parse_iso(value: Any) -> Optional[datetime]:
@@ -82,23 +72,15 @@ def classify_pressure(available_kib: Any, total_kib: Any) -> str:
return "ok"
def _read_heartbeat(home: Optional[Path]) -> Optional[Dict[str, Any]]:
try:
from gateway.lifecycle_ledger import _read_json
from gateway.shutdown_watchdog import get_loop_heartbeat_path
return _read_json(get_loop_heartbeat_path(home))
except Exception:
return None
def _read_sentinel(home: Optional[Path]) -> Optional[Dict[str, Any]]:
def _read_state_files(home: Optional[Path]) -> tuple:
"""``(heartbeat, sentinel)`` dicts, each ``None`` when unreadable."""
try:
from gateway.lifecycle_ledger import _read_json, get_lifecycle_sentinel_path
from gateway.shutdown_watchdog import get_loop_heartbeat_path
return _read_json(get_lifecycle_sentinel_path(home))
return _read_json(get_loop_heartbeat_path(home)), _read_json(get_lifecycle_sentinel_path(home))
except Exception:
return None
return None, None
def collect_memory_status(
@@ -110,8 +92,8 @@ def collect_memory_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/never-started gateway or corrupt files yield
``{"pressure": "unknown", ...}`` plus whatever fields could be recovered.
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] = {
@@ -123,13 +105,12 @@ def collect_memory_status(
"sampled_at": None,
"last_boot_unclean": False,
"last_boot_suspected_oom": False,
# Identity of the CURRENT gateway life (sentinel started_at) — changes on
# every boot, so the dashboard can key banner dismissal on it and
# acknowledging one OOM restart does not mute the NEXT one.
# Identity of the CURRENT life (sentinel started_at): the dashboard keys
# banner dismissal on it so acknowledging one OOM restart does not mute the NEXT.
"boot_id": None,
}
heartbeat = _read_heartbeat(home)
heartbeat, sentinel = _read_state_files(home)
if heartbeat:
sampled_at = _parse_iso(heartbeat.get("updated_at"))
mem = heartbeat.get("mem")
@@ -146,7 +127,6 @@ def collect_memory_status(
if 0 <= (moment - sampled_at).total_seconds() <= _HEARTBEAT_FRESH_TTL_S:
status["pressure"] = classify_pressure(mem.get("mem_available_kib"), mem.get("mem_total_kib"))
sentinel = _read_sentinel(home)
if sentinel:
status["last_boot_unclean"] = bool(sentinel.get("prior_unclean_exit"))
status["last_boot_suspected_oom"] = bool(sentinel.get("prior_suspected_oom"))
+24 -27
View File
@@ -12,17 +12,19 @@ from datetime import datetime
from typing import Any, Optional, Tuple
# Current gateway format: [Tue 2026-04-28 13:40:53 CEST]
_HUMAN_TIMESTAMP_RE = re.compile(
r"^\[(?P<dow>[A-Z][a-z]{2}) "
# Leading timestamp prefix, either the current human format
# ``[Tue 2026-04-28 13:40:53 CEST]`` or the older ISO one
# ``[2026-04-13T17:02:06+0200]`` / ``[...+02:00]`` (human tried first).
_TIMESTAMP_PREFIX_RE = re.compile(
r"^\[(?:"
r"(?P<dow>[A-Z][a-z]{2}) "
r"(?P<date>\d{4}-\d{2}-\d{2}) "
r"(?P<time>\d{2}:\d{2}:\d{2})"
r"(?: (?P<tz>[A-Za-z0-9_+\-/:]+))?\]\s*"
r"(?: (?P<tz>[A-Za-z0-9_+\-/:]+))?"
r"|(?P<iso>\d{4}-\d{2}-\d{2}T[^\]]+)"
r")\]\s*"
)
# Older gateway format: [2026-04-13T17:02:06+0200] or [+02:00]
_ISO_TIMESTAMP_RE = re.compile(r"^\[(?P<iso>\d{4}-\d{2}-\d{2}T[^\]]+)\]\s*")
def _localize(dt: datetime, tz) -> float:
"""Epoch for ``dt``; naive values take ``tz`` if given, else the local zone."""
@@ -44,7 +46,7 @@ def _parse_iso(text: str, tz=None) -> Optional[float]:
def _parse_timestamp_match(match: re.Match, tz=None) -> Optional[float]:
if match.groupdict().get("iso"):
if match.group("iso"):
return _parse_iso(match.group("iso"), tz)
try:
dt = datetime.strptime(f"{match.group('date')} {match.group('time')}", "%Y-%m-%d %H:%M:%S")
@@ -53,10 +55,6 @@ def _parse_timestamp_match(match: re.Match, tz=None) -> Optional[float]:
return _localize(dt, tz)
def _match_timestamp_prefix(text: str) -> Optional[re.Match]:
return _HUMAN_TIMESTAMP_RE.match(text) or _ISO_TIMESTAMP_RE.match(text)
def coerce_message_timestamp(ts_value: Any, tz=None) -> Optional[float]:
"""Coerce a timestamp-like value to Unix epoch seconds.
@@ -72,21 +70,20 @@ def coerce_message_timestamp(ts_value: Any, tz=None) -> Optional[float]:
return float(ts_value.timestamp())
except Exception:
return None
if isinstance(ts_value, str):
text = ts_value.strip()
if not text:
return None
match = _match_timestamp_prefix(text)
if match is not None:
parsed = _parse_timestamp_match(match, tz=tz)
if parsed is not None:
return parsed
try:
return float(text)
except (TypeError, ValueError):
pass
if not isinstance(ts_value, str):
return None
text = ts_value.strip()
if not text:
return None
match = _TIMESTAMP_PREFIX_RE.match(text)
if match is not None:
parsed = _parse_timestamp_match(match, tz=tz)
if parsed is not None:
return parsed
try:
return float(text)
except (TypeError, ValueError):
return _parse_iso(text, tz)
return None
def format_message_timestamp(ts_value: Any, tz=None) -> str:
@@ -109,7 +106,7 @@ def strip_leading_message_timestamps(content: str, tz=None) -> Tuple[str, Option
return content, None
text = content
embedded_epoch: Optional[float] = None
while (match := _match_timestamp_prefix(text)) is not None:
while (match := _TIMESTAMP_PREFIX_RE.match(text)) is not None:
parsed = _parse_timestamp_match(match, tz=tz)
if parsed is not None:
embedded_epoch = parsed
+17 -20
View File
@@ -33,17 +33,15 @@ def mirror_to_session(
``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. ``_find_session_id`` refuses to guess on a populated chat
(flat session + N thread sessions sharing one chat_id) and would silently
drop the mirror.
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"``, correct when the mirrored text is
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 ``repair_message_sequence``'s consecutive-user merge.
``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.
Returns True if mirrored, False if no matching session or error. Never raises.
"""
@@ -118,11 +116,11 @@ def _find_session_id(
if str(_key).startswith("_") or not isinstance(entry, dict):
continue
origin = entry.get("origin") or {}
if (origin.get("platform") or entry.get("platform", "")).lower() != platform_lower:
continue
if str(origin.get("chat_id", "")) != str(chat_id):
continue
if thread_id is not None and str(origin.get("thread_id") or "") != str(thread_id):
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)
@@ -144,13 +142,12 @@ def _find_session_id(
def _append_to_sqlite(session_id: str, message: dict) -> None:
"""Append a message to the SQLite session database."""
db = None
try:
from hermes_state import get_shared_session_db, release_or_close
db = get_shared_session_db()
db.append_message(session_id=session_id, role=message.get("role", "assistant"), content=message.get("content"))
try:
db.append_message(session_id=session_id, role=message.get("role", "assistant"), content=message.get("content"))
finally:
release_or_close(db)
except Exception as e:
logger.debug("Mirror SQLite write failed: %s", e)
finally:
if db is not None:
release_or_close(db)
+9 -10
View File
@@ -14,6 +14,7 @@ from hermes_constants import get_hermes_home
_DISK_DEGRADED_PERCENT = 90.0
_CONNECTED_STATES = {"connected", "running", "ok"}
def _check(status: str, detail: str | None = None, **extra: Any) -> dict[str, Any]:
@@ -48,34 +49,32 @@ def _probe_config(home: Path) -> dict[str, Any]:
return _check("ok", "using defaults")
try:
raw = yaml.safe_load(path.read_text(encoding="utf-8"))
if raw is not None and not isinstance(raw, dict):
return _check("degraded", "top level is not a mapping")
return _check("ok")
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")
def _probe_disk(home: Path) -> dict[str, Any]:
try:
usage = shutil.disk_usage(home)
used_pct = round((usage.used / usage.total) * 100, 1) if usage.total else 0.0
status = "degraded" if used_pct >= _DISK_DEGRADED_PERCENT else "ok"
return _check(status, used_percent=used_pct, free_bytes=usage.free)
except Exception as exc:
return _check("degraded", type(exc).__name__)
used_pct = round((usage.used / usage.total) * 100, 1) if usage.total else 0.0
status = "degraded" if used_pct >= _DISK_DEGRADED_PERCENT else "ok"
return _check(status, used_percent=used_pct, free_bytes=usage.free)
def _probe_gateway(runtime_status: dict[str, Any]) -> dict[str, Any]:
state = str(runtime_status.get("gateway_state") or "unknown")
platforms = runtime_status.get("platforms")
if not isinstance(platforms, dict):
platforms = {}
platforms = platforms if isinstance(platforms, dict) else {}
connected = sum(
1
for value in platforms.values()
if isinstance(value, dict)
and str(value.get("state") or value.get("status") or "").lower()
in {"connected", "running", "ok"}
and str(value.get("state") or value.get("status") or "").lower() in _CONNECTED_STATES
)
status = "ok" if state in {"running", "draining"} else "degraded"
return _check(status, state=state, connected_platforms=connected, platforms=len(platforms))
+13 -19
View File
@@ -42,20 +42,18 @@ def _strip_edge_silence_punctuation(text: str) -> str:
return text[start:end].strip()
def _canonical_silence_candidates(text: str) -> tuple[str, ...]:
exact = _canonical_silence_candidate(text)
stripped = _strip_edge_silence_punctuation(text.strip())
if stripped == text.strip():
return (exact,)
return (exact, _canonical_silence_candidate(stripped))
def _short_stripped(text: Any) -> str:
"""Stripped text if it is a non-empty string within the marker cap, else ''."""
def _canonical_silence_candidates(text: Any) -> tuple[str, ...]:
"""Canonical forms of a short marker-sized response; ``()`` when not a candidate at all."""
if not isinstance(text, str):
return ""
return ()
stripped = text.strip()
return stripped if 0 < len(stripped) <= _MARKER_LENGTH_CAP else ""
if not 0 < len(stripped) <= _MARKER_LENGTH_CAP:
return ()
exact = _canonical_silence_candidate(stripped)
depunctuated = _strip_edge_silence_punctuation(stripped)
if depunctuated == stripped:
return (exact,)
return (exact, _canonical_silence_candidate(depunctuated))
def is_intentional_silence_response(response: Any) -> bool:
@@ -64,10 +62,7 @@ def is_intentional_silence_response(response: Any) -> bool:
Prose that merely mentions ``NO_REPLY`` must be delivered normally. A blank
response is not silence either — that is the empty-response failure path.
"""
stripped = _short_stripped(response)
return bool(stripped) and any(
c in LIVE_GATEWAY_SILENT_MARKERS for c in _canonical_silence_candidates(stripped)
)
return any(c in LIVE_GATEWAY_SILENT_MARKERS for c in _canonical_silence_candidates(response))
def is_autonomous_silence_response(response: Any) -> bool:
@@ -117,8 +112,7 @@ def is_partial_silence_marker(text: Any) -> bool:
marker cap, returns False so normal streaming resumes. Shares the marker set
and canonicalization with :func:`is_intentional_silence_response`.
"""
stripped = _short_stripped(text)
return bool(stripped) and any(
return any(
c and any(marker.startswith(c) for marker in LIVE_GATEWAY_SILENT_MARKERS)
for c in _canonical_silence_candidates(stripped)
for c in _canonical_silence_candidates(text)
)
+44 -62
View File
@@ -15,28 +15,12 @@ GATEWAY_SERVICE_RESTART_EXIT_CODE = 75
# failure) so the supervisor stops restarting the gateway.
GATEWAY_FATAL_CONFIG_EXIT_CODE = 78
def is_global_startup_conflict(error_code: str | None) -> bool:
"""True when an adapter's fatal error is a single-writer ownership conflict.
Adapters emit ``{scope}_lock`` with ``retryable=True`` so a *mid-run*
reconnect can recover once the holder exits or a stale record is cleared.
At startup a live foreign
holder is a configuration conflict (two gateways cannot poll one token), so
the startup router must not treat it as a transient blip. Matches by error
CODE only (``lock_conflict`` / ``*_lock``), never by message text.
"""
code = (error_code or "").strip().lower()
return bool(code) and (code == "lock_conflict" or code.endswith("_lock"))
# Set by ``hermes gateway run --external-supervisor``. Unlike systemd's
# INVOCATION_ID and launchd's XPC_SERVICE_NAME, this survives wrappers that
# replace the child environment (e.g. ``sudo env -i``).
EXTERNAL_GATEWAY_SUPERVISOR_ENV = "HERMES_GATEWAY_EXTERNAL_SUPERVISOR"
DEFAULT_GATEWAY_RESTART_DRAIN_TIMEOUT = float(
DEFAULT_CONFIG["agent"]["restart_drain_timeout"]
)
DEFAULT_GATEWAY_RESTART_DRAIN_TIMEOUT = float(DEFAULT_CONFIG["agent"]["restart_drain_timeout"])
DEFAULT_GATEWAY_SIGNAL_INTERRUPT_GRACE_TIMEOUT = float(
DEFAULT_CONFIG["gateway"]["signal_interrupt_grace_timeout"]
)
@@ -53,12 +37,10 @@ DEFAULT_GATEWAY_RESTART_AFTER_TURN_TIMEOUT = float(
# Cron-only floor under the ``stop()`` drain. ``restart_drain_timeout`` defaults
# to 0 because interrupting a *chat* turn is cheap and recoverable (user is
# told, session pre-marked resume_pending). An interrupted *cron* run has
# neither property — nobody is waiting on it, it lands in jobs.json as a
# permanent failure and a recurring job just waits for its next schedule — so a
# zero-second drain silently destroys work.
DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT = float(
DEFAULT_CONFIG["agent"]["cron_drain_timeout"]
)
# neither property — it lands in jobs.json as a permanent failure and a
# recurring job just waits for its next schedule — so a zero-second drain
# silently destroys work.
DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT = float(DEFAULT_CONFIG["agent"]["cron_drain_timeout"])
# Seconds of the shutdown watchdog leash held back for post-drain work
# (interrupt agents, kill tool subprocesses, mark jobs interrupted, disconnect
@@ -72,10 +54,23 @@ CRON_DRAIN_CLEANUP_RESERVE_S = 10.0
SYSTEMD_STOP_HEADROOM_S = 30.0
SYSTEMD_TIMEOUT_STOP_SEC_FLOOR = 60.0
_TRUTHY = {"1", "true", "yes", "on"}
def is_gateway_supervisor_process(
environ: Mapping[str, str] | None = None,
) -> bool:
def is_global_startup_conflict(error_code: str | None) -> bool:
"""True when an adapter's fatal error is a single-writer ownership conflict.
Adapters emit ``{scope}_lock`` with ``retryable=True`` so a *mid-run*
reconnect can recover once the holder exits. At startup a live foreign
holder is a configuration conflict (two gateways cannot poll one token), so
the startup router must not treat it as a transient blip. Matches by error
CODE only (``lock_conflict`` / ``*_lock``), never by message text.
"""
code = (error_code or "").strip().lower()
return bool(code) and (code == "lock_conflict" or code.endswith("_lock"))
def is_gateway_supervisor_process(environ: Mapping[str, str] | None = None) -> bool:
"""Return whether this gateway process is owned by a supervisor."""
env = os.environ if environ is None else environ
if env.get("INVOCATION_ID") or env.get("HERMES_S6_SUPERVISED_CHILD"):
@@ -83,19 +78,14 @@ def is_gateway_supervisor_process(
xpc_service = env.get("XPC_SERVICE_NAME", "")
if xpc_service and xpc_service != "0":
return True
return str(env.get(EXTERNAL_GATEWAY_SUPERVISOR_ENV, "")).strip().lower() in {
"1",
"true",
"yes",
"on",
}
return str(env.get(EXTERNAL_GATEWAY_SUPERVISOR_ENV, "")).strip().lower() in _TRUTHY
def is_container_restart_context() -> bool:
"""Whether the gateway runs in a container (Docker/Podman): the detached
setsid restart path dies with the cgroup, so exit-75 service restart is the
only viable path. Separate function so tests can mock container detection
(a real ``/.dockerenv`` on CI otherwise flips the routing)."""
"""Whether the gateway runs in a container (Docker/Podman): the detached setsid
restart path dies with the cgroup, so exit-75 service restart is the only viable
path. Separate function so tests can mock container detection (a real
``/.dockerenv`` on CI otherwise flips the routing)."""
return os.path.exists("/.dockerenv") or os.path.exists("/run/.containerenv")
@@ -107,22 +97,24 @@ def _seconds(value: object, fallback: float = 0.0) -> float:
return fallback
def _parse_timeout_keeping_zero(raw: object, default: float) -> float:
"""Parse a timeout where ``0`` is a deliberate disable (must NOT fall
through to ``default``), unlike None / blank / non-numeric input."""
def _parse_timeout_keeping_zero(raw: object, default: float, *, finite: bool = False) -> float:
"""Parse a timeout where ``0`` is a deliberate disable (must NOT fall through
to ``default``), unlike None / blank / non-numeric (/ non-finite) input."""
if raw is None or (isinstance(raw, str) and not raw.strip()):
return default
try:
value = float(raw)
value = float(raw) # type: ignore[arg-type]
except (TypeError, ValueError):
return default
if finite and not math.isfinite(value):
return default
return max(0.0, value)
def parse_restart_drain_timeout(raw: object) -> float:
"""Parse a configured drain timeout, falling back to the shared default."""
try:
value = float(raw) if str(raw or "").strip() else DEFAULT_GATEWAY_RESTART_DRAIN_TIMEOUT
value = float(raw) if str(raw or "").strip() else DEFAULT_GATEWAY_RESTART_DRAIN_TIMEOUT # type: ignore[arg-type]
except (TypeError, ValueError):
return DEFAULT_GATEWAY_RESTART_DRAIN_TIMEOUT
return max(0.0, value)
@@ -138,6 +130,11 @@ def parse_cron_drain_timeout(raw: object) -> float:
return _parse_timeout_keeping_zero(raw, DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT)
def parse_signal_interrupt_grace_timeout(raw: object) -> float:
"""Parse the unexpected-signal post-interrupt grace timeout."""
return _parse_timeout_keeping_zero(raw, DEFAULT_GATEWAY_SIGNAL_INTERRUPT_GRACE_TIMEOUT, finite=True)
def resolve_cron_drain_budget(
drain_timeout: float,
cron_drain_timeout: float,
@@ -149,12 +146,11 @@ def resolve_cron_drain_budget(
"""Seconds the shutdown drain may spend waiting on in-flight cron work.
The configured floor is clamped to what this process can honour: the
shutdown watchdog hard-exits at ``watchdog_delay`` (and TimeoutStopSec is
sized from the same budget), so waiting past that leash minus
``cleanup_reserve_s`` swaps a cleanly-interrupted job for a SIGKILL that
leaves it wedged. Never returns less than ``drain_timeout``: the cron floor
only ever extends the wait, so an operator who deliberately configured a
long ``restart_drain_timeout`` keeps it.
shutdown watchdog hard-exits at ``watchdog_delay`` (TimeoutStopSec is sized
from the same budget), so waiting past that leash minus ``cleanup_reserve_s``
swaps a cleanly-interrupted job for a SIGKILL that leaves it wedged. Never
returns less than ``drain_timeout``: the cron floor only ever extends the
wait, so a deliberately long ``restart_drain_timeout`` is kept.
"""
drain = _seconds(drain_timeout)
floor = _seconds(cron_drain_timeout)
@@ -181,7 +177,7 @@ def resolve_systemd_timeout_stop_sec(
``restart_drain_timeout`` is only the chat-turn interrupt budget (default 0);
the stop path may first wait ``cron_drain_timeout`` + ``cleanup_reserve_s``
for cron work, so sizing from drain alone lets systemd SIGKILL an in-budget
drain. A zero cron timeout is an opt-out and does not extend the budget.
drain. A zero cron timeout is an opt-out and does not extend the budget.
Non-numeric inputs degrade to 0.
"""
drain = _seconds(drain_timeout)
@@ -204,17 +200,3 @@ def resolve_restart_exit_wait_budget(
expiry must cover both phases.
"""
return _seconds(drain_timeout) + _seconds(after_turn_timeout) + _seconds(headroom)
def parse_signal_interrupt_grace_timeout(raw: object) -> float:
"""Parse the unexpected-signal post-interrupt grace timeout."""
try:
if raw is None or (isinstance(raw, str) and not raw.strip()):
value = DEFAULT_GATEWAY_SIGNAL_INTERRUPT_GRACE_TIMEOUT
else:
value = float(raw)
except (TypeError, ValueError):
return DEFAULT_GATEWAY_SIGNAL_INTERRUPT_GRACE_TIMEOUT
if not math.isfinite(value):
return DEFAULT_GATEWAY_SIGNAL_INTERRUPT_GRACE_TIMEOUT
return max(0.0, value)
+22 -30
View File
@@ -1,35 +1,31 @@
"""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 via cron/CLI. They do NOT cover
every SIGTERM source (raw ``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.
stop|restart`` + ``terminal_tool``, and the cron-creation lifecycle filter) stop
the agent scheduling its own restart. They do NOT cover every SIGTERM source
(raw ``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.
This module is the last-resort circuit breaker. Each boot with
restart-interrupted sessions pending is timestamped and persisted (each boot
is a fresh process, so in-memory state is useless). Boots CHAIN while
consecutive gaps stay within ``max_gap_seconds``, so slow crash cycles (a
wedged loop killed by the liveness watchdog every ~150s) trip exactly like the
fast ~10s respawn loop. When tripped, the caller SKIPS auto-resume for that
boot — the gateway still serves real inbound messages, it just stops replaying
the session that keeps killing it.
State lives in ``<HERMES_HOME>/gateway/restart_loop.json`` (profile-scoped).
Best-effort: any read/write failure fails OPEN (no false trip) because a
broken breaker must never wedge a healthy gateway.
Last-resort circuit breaker: each boot with restart-interrupted sessions pending
is timestamped and persisted to ``<HERMES_HOME>/gateway/restart_loop.json``
(each boot is a fresh process). Boots CHAIN while consecutive gaps stay within
``max_gap_seconds``, so a slow crash cycle (liveness watchdog every ~150s) trips
exactly like a fast ~10s respawn loop. When tripped, the caller SKIPS
auto-resume for that boot — real inbound messages are still served.
Best-effort: any read/write failure fails OPEN (a broken breaker must never
wedge a healthy gateway).
"""
from __future__ import annotations
import contextlib
import json
import logging
import time
from typing import List, Optional
from hermes_constants import get_hermes_home
import contextlib
logger = logging.getLogger("gateway.run")
@@ -39,10 +35,9 @@ 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
# 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 makes the breaker period-agnostic, and a
# single boot followed by real quiet resets the chain.
# 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
@@ -63,16 +58,14 @@ def _load_boots() -> List[float]:
def _save_boots(boots: List[float]) -> None:
try:
with contextlib.suppress(OSError):
path = _state_path()
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(json.dumps({"boots": boots}), encoding="utf-8")
except OSError:
pass
def _chain_gap(window_seconds: int, max_gap_seconds: int) -> float:
"""Inter-boot gap that still links two boots. Floored by ``window_seconds`` so
"""Inter-boot gap that still links two boots. Floored by ``window_seconds`` so
widening the window never makes the breaker *less* sensitive."""
return float(max(1, window_seconds, max_gap_seconds))
@@ -80,10 +73,9 @@ 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 keeping boots while each successive gap stays within
``gap``; the first wider gap ends the chain (older boots belong to an
already-resolved episode). Nothing recent enough -> empty list, which is how
a healthy gateway forgets an old loop.
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.
"""
chain: List[float] = []
prev = ts
@@ -139,7 +131,7 @@ def check_and_record(
boots = record_restart_interrupted_boot(
window_seconds, now=now, max_gap_seconds=max_gap_seconds
)
tripped = len(boots) >= max_restarts if max_restarts > 0 else False
tripped = max_restarts > 0 and len(boots) >= max_restarts
if tripped:
logger.warning(
"Restart-loop breaker TRIPPED: %d chained restart-interrupted "