refactor(gateway): third pass on scale_to_zero/wake/status_phrases/runtime_footer/systemd_notify

This commit is contained in:
Teknium
2026-09-02 21:14:22 -07:00
parent 4dc4a0e6d9
commit 8d5d5f9ca2
5 changed files with 161 additions and 290 deletions

View File

@@ -1,22 +1,11 @@
"""Gateway runtime-metadata footer (model · context % · cwd), off by default
to keep replies minimal.
Config (``display.runtime_footer``, per-platform override under
``display.platforms.<platform>.runtime_footer``; toggled by ``/footer on|off``)::
display:
runtime_footer:
enabled: true
fields: [model, context_pct, cwd] # order shown; drop any to hide
Fields: ``model`` (vendor prefix dropped), ``context_pct`` (last-call
occupancy), ``latency`` (turn wall-clock, opt-in — NOT in the default set so an
unset ``fields`` renders exactly as before), ``cwd`` (home-relative).
``gateway/run.py`` appends the footer to the final response only (never to
tool-progress or streaming partials); when streaming already delivered the text,
it goes out as a trailing message via ``send_trailing_footer()``.
"""
"""Gateway runtime-metadata footer (model · context % · cwd), off by default to keep replies
minimal. Config: ``display.runtime_footer: {enabled: bool, fields: [model, context_pct, cwd]}``
(order shown; drop any to hide), per-platform override ``display.platforms.<p>.runtime_footer``,
toggled by ``/footer on|off``. Fields: ``model`` (vendor prefix dropped), ``context_pct`` (last-call
occupancy), ``latency`` (turn wall-clock, opt-in — NOT in the default set so an unset ``fields``
renders exactly as before), ``cwd`` (home-relative). ``gateway/run.py`` appends the footer to the
final response only (never to tool-progress or streaming partials); when streaming already
delivered the text, it goes out as a trailing message via ``send_trailing_footer()``."""
from __future__ import annotations
@@ -82,16 +71,12 @@ def _format_latency(seconds: float) -> str:
return f"{m}m{sec:02d}s"
def format_runtime_footer(
*, model: Optional[str], context_tokens: int, context_length: Optional[int],
cwd: Optional[str] = None, turn_seconds: Optional[float] = None,
fields: Iterable[str] = _DEFAULT_FIELDS,
) -> str:
"""Render the footer line, or "" if no fields have data.
Fields whose data is missing (and unknown field names) are skipped silently —
a partial footer beats ``?%`` or empty slots.
"""
def format_runtime_footer(*, model: Optional[str], context_tokens: int,
context_length: Optional[int], cwd: Optional[str] = None,
turn_seconds: Optional[float] = None,
fields: Iterable[str] = _DEFAULT_FIELDS) -> str:
"""Render the footer line, or "" if no fields have data. Fields whose data is missing (and
unknown field names) are skipped silently — a partial footer beats ``?%`` or empty slots."""
def context_pct() -> str:
if context_length and context_length > 0 and context_tokens >= 0:
return f"{max(0, min(100, round((context_tokens / context_length) * 100)))}%"
@@ -104,26 +89,19 @@ def format_runtime_footer(
"latency": lambda: _format_latency(turn_seconds) if turn_seconds is not None and turn_seconds >= 0 else "",
"cwd": lambda: _home_relative_cwd(cwd or _env_cwd()),
}
parts = [value for field in fields if (render := renderers.get(field)) and (value := render())]
return _SEP.join(parts)
return _SEP.join(v for field in fields if (render := renderers.get(field)) and (v := render()))
def build_footer_line(
*, user_config: dict[str, Any] | None, platform_key: str | None, model: Optional[str],
context_tokens: int, context_length: Optional[int], cwd: Optional[str] = None,
turn_seconds: Optional[float] = None,
) -> str:
"""Entry point for gateway/run.py: footer text, or "" when disabled / no data.
Callers append it to the final response themselves, preserving a single
blank line of separation.
``turn_seconds`` is the caller-measured (``time.monotonic()``) run duration;
``None`` skips the ``latency`` field.
"""
def build_footer_line(*, user_config: dict[str, Any] | None, platform_key: str | None,
model: Optional[str], context_tokens: int, context_length: Optional[int],
cwd: Optional[str] = None, turn_seconds: Optional[float] = None) -> str:
"""Entry point for gateway/run.py: footer text, or "" when disabled / no data. Callers append it
to the final response themselves, preserving a single blank line of separation.
``turn_seconds`` is the caller-measured (``time.monotonic()``) run duration; ``None`` skips the
``latency`` field."""
cfg = resolve_footer_config(user_config, platform_key)
if not cfg.get("enabled"):
return ""
return format_runtime_footer(
model=model, context_tokens=context_tokens, context_length=context_length,
cwd=cwd, turn_seconds=turn_seconds, fields=cfg.get("fields") or _DEFAULT_FIELDS,
)
return format_runtime_footer(model=model, context_tokens=context_tokens,
context_length=context_length, cwd=cwd, turn_seconds=turn_seconds,
fields=cfg.get("fields") or _DEFAULT_FIELDS)

View File

@@ -1,14 +1,12 @@
"""Scale-to-zero idle detection + dormant-quiesce for the gateway.
Owns the *decision* to go idle, drives the relay transport's ``go_dormant()``,
then SUSPENDS the machine via the local Fly Machines API socket; wake stays
platform-side (autostart-on-wakeUrl). The gateway self-suspends because Fly
Proxy only sees INBOUND proxied connections — it would suspend mid-turn or
before ``go_dormant()`` flipped the relay destination (buffered-event black hole).
Enable is gated SOLELY by the NAS "Labs" toggle env stamp (not config); the idle
timeout IS config.yaml. Quiesce uses ``go_dormant()`` (never disconnect/drain)
and ``mark_resume_pending`` is NOT called: suspend preserves RAM.
"""
Owns the *decision* to go idle, drives the relay transport's ``go_dormant()``, then SUSPENDS the
machine via the local Fly Machines API socket; wake stays platform-side (autostart on wakeUrl).
Self-suspend because Fly Proxy only sees INBOUND proxied connections: it would suspend mid-turn or
before ``go_dormant()`` flipped the relay destination (buffered-event black hole). Enable is gated
SOLELY by the NAS "Labs" toggle env stamp (not config); the idle timeout IS config.yaml. Quiesce
uses ``go_dormant()`` (never disconnect/drain); ``mark_resume_pending`` is NOT called: suspend
preserves RAM."""
from __future__ import annotations
@@ -27,8 +25,7 @@ FLY_APP_NAME_ENV = "FLY_APP_NAME" # Fly-injected identity; both needed for self
FLY_MACHINE_ID_ENV = "FLY_MACHINE_ID"
# Local flaps (Fly Machines API) socket; POST .../suspend freezes THIS machine.
FLY_API_SOCKET = "/.fly/api"
# Short is safe: real work always blocks the suspend and resume is sub-second; longer just bills
# idle RAM.
# Short is safe: real work always blocks the suspend, resume is sub-second; longer bills idle RAM.
DEFAULT_IDLE_TIMEOUT_MINUTES = 2
_TRUTHY = {"1", "true", "yes", "on"}
# Dashboard-client liveness marker, touched by the (separate) dashboard process on every /api/ws
@@ -47,26 +44,21 @@ def scale_to_zero_enabled(environ: Optional[dict] = None) -> bool:
return _env_str(environ, SCALE_TO_ZERO_ENV).lower() in _TRUTHY
def parse_idle_timeout_seconds(
cfg_value: Any, default_minutes: int = DEFAULT_IDLE_TIMEOUT_MINUTES
) -> float:
"""Coerce ``scale_to_zero.idle_timeout_minutes`` to seconds.
Non-numeric / non-positive degrades to the default (never <= 0: instant dormancy).
"""
def parse_idle_timeout_seconds(cfg_value: Any,
default_minutes: int = DEFAULT_IDLE_TIMEOUT_MINUTES) -> float:
"""Coerce ``scale_to_zero.idle_timeout_minutes`` to seconds. Non-numeric / non-positive
degrades to the default (never <= 0: instant dormancy)."""
try:
minutes = float(cfg_value)
except (TypeError, ValueError):
minutes = float(default_minutes)
minutes = 0.0
return (float(default_minutes) if minutes <= 0 else minutes) * 60.0
def messaging_is_relay_only_or_absent(platforms: Iterable[Any]) -> bool:
"""True iff the only connected platform is RELAY, or there is none.
A directly-connected platform holds a live socket and cannot scale to zero.
Compared by ``.value``/name so this module stays enum-import-free.
"""
"""True iff the only connected platform is RELAY, or there is none. A directly-connected
platform holds a live socket and cannot scale to zero. Compared by ``.value``/name so this
module stays enum-import-free."""
names = {str(getattr(p, "value", p)).strip().lower() for p in platforms}
names.discard("relay")
return not names
@@ -79,21 +71,14 @@ def should_arm(*, enabled: bool, relay_only_or_absent: bool, wake_url: Optional[
return bool(enabled) and bool(relay_only_or_absent) and bool(wake_url)
def is_idle(
*, active_work_count: int, seconds_since_last_inbound: float, idle_timeout_seconds: float,
has_live_background_work: bool,
) -> bool:
def is_idle(*, active_work_count: int, seconds_since_last_inbound: float,
idle_timeout_seconds: float, has_live_background_work: bool) -> bool:
"""Pure idle predicate: no active work, no inbound within the window, no live background work.
``active_work_count`` is the BROAD aggregate (agent turns + cron + API runs) — passing only
``len(_running_agents)`` reopens the mid-cron-job suspend hole. Callers that cannot read a
work source must fail AWAKE (pass a positive sentinel), never fail to 0.
"""
return (
active_work_count <= 0
and not has_live_background_work
and seconds_since_last_inbound >= idle_timeout_seconds
)
work source must fail AWAKE (pass a positive sentinel), never fail to 0."""
return (active_work_count <= 0 and not has_live_background_work
and seconds_since_last_inbound >= idle_timeout_seconds)
def dashboard_client_heartbeat_path(hermes_home: Optional[os.PathLike | str] = None):
@@ -109,8 +94,7 @@ def touch_dashboard_client_heartbeat(path: Optional[os.PathLike | str] = None) -
try:
p = dashboard_client_heartbeat_path() if path is None else path
os.makedirs(os.path.dirname(p), exist_ok=True)
with open(p, "a", encoding="utf-8"):
pass
open(p, "a", encoding="utf-8").close()
os.utime(p, None)
return True
except Exception: # noqa: BLE001 - liveness garnish must never break the WS
@@ -118,15 +102,12 @@ def touch_dashboard_client_heartbeat(path: Optional[os.PathLike | str] = None) -
return False
def dashboard_client_last_seen(
path: Optional[os.PathLike | str] = None, *, now: Optional[float] = None
) -> Optional[float]:
"""Epoch seconds a dashboard client last sent a WS frame, or None if never.
Missing marker -> None (steady state when nobody has the dashboard open — NOT fail-awake,
or no instance would ever sleep). Unreadable marker -> ``now`` (fail-awake, as in ``is_idle``).
Clamped to now: an NTP step-back can leave the mtime in the future.
"""
def dashboard_client_last_seen(path: Optional[os.PathLike | str] = None, *,
now: Optional[float] = None) -> Optional[float]:
"""Epoch seconds a dashboard client last sent a WS frame, or None if never. Missing marker ->
None (steady state when nobody has the dashboard open — NOT fail-awake, or no instance would
ever sleep). Unreadable marker -> ``now`` (fail-awake, as in ``is_idle``). Clamped to now: an
NTP step-back can leave the mtime in the future."""
current = time.time() if now is None else now
p = dashboard_client_heartbeat_path() if path is None else path
try:
@@ -139,46 +120,31 @@ def dashboard_client_last_seen(
def self_suspend_available(environ: Optional[dict] = None) -> bool:
"""True iff Fly machine identity is present AND the local Machines API socket exists.
Off-Fly the watcher skips the quiesce: the platform owns the freeze.
"""
return bool(
_env_str(environ, FLY_APP_NAME_ENV)
and _env_str(environ, FLY_MACHINE_ID_ENV)
and os.path.exists(FLY_API_SOCKET)
)
Off-Fly the watcher skips the quiesce: the platform owns the freeze."""
return bool(_env_str(environ, FLY_APP_NAME_ENV) and _env_str(environ, FLY_MACHINE_ID_ENV)
and os.path.exists(FLY_API_SOCKET))
def suspend_self(
environ: Optional[dict] = None, *, socket_path: str = FLY_API_SOCKET, timeout: float = 10.0
) -> bool:
def suspend_self(environ: Optional[dict] = None, *, socket_path: str = FLY_API_SOCKET,
timeout: float = 10.0) -> bool:
"""POST /v1/apps/{app}/machines/{id}/suspend on the local flaps socket (the socket is the
credential).
Returns True when flaps accepted (2xx); the kernel then freezes this process shortly after, so
treat as fire-and-forget. Never raises: a failed suspend leaves the machine running (fail-
awake). stdlib-only on purpose — a plain unix-socket HTTP/1.1 request, no async plumbing to
freeze mid-await.
"""
app = _env_str(environ, FLY_APP_NAME_ENV)
machine_id = _env_str(environ, FLY_MACHINE_ID_ENV)
credential). Returns True when flaps accepted (2xx); the kernel then freezes this process
shortly after, so treat as fire-and-forget. Never raises: a failed suspend leaves the machine
running (fail-awake). stdlib-only on purpose — a plain unix-socket HTTP/1.1 request, no async
plumbing to freeze mid-await."""
app, machine_id = _env_str(environ, FLY_APP_NAME_ENV), _env_str(environ, FLY_MACHINE_ID_ENV)
if not app or not machine_id:
logger.warning("scale-to-zero: suspend_self called without Fly machine identity")
return False
request = (
f"POST /v1/apps/{app}/machines/{machine_id}/suspend HTTP/1.1\r\n"
"Host: flaps\r\nContent-Length: 0\r\nConnection: close\r\n\r\n"
)
request = (f"POST /v1/apps/{app}/machines/{machine_id}/suspend HTTP/1.1\r\n"
"Host: flaps\r\nContent-Length: 0\r\nConnection: close\r\n\r\n")
try:
with socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) as sock:
sock.settimeout(timeout)
sock.connect(socket_path)
sock.sendall(request.encode("ascii"))
response = b""
while len(response) < 65536:
chunk = sock.recv(4096)
if not chunk:
break
while len(response) < 65536 and (chunk := sock.recv(4096)):
response += chunk
except OSError as exc:
logger.warning("scale-to-zero: flaps suspend request failed: %s", exc)
@@ -190,7 +156,6 @@ def suspend_self(
logger.info("scale-to-zero: machine suspend accepted by flaps (%s)", status_line)
else:
body = response.split(b"\r\n\r\n", 1)[-1][:500].decode("utf-8", "replace")
logger.warning(
"scale-to-zero: flaps suspend rejected: %s %s", status_line, json.dumps(body)[:500]
)
logger.warning("scale-to-zero: flaps suspend rejected: %s %s", status_line,
json.dumps(body)[:500])
return ok

View File

@@ -1,21 +1,10 @@
"""Human-friendly generic gateway status phrases.
Turns Hermes' long-running gateway status surface into short chat-safe lines
without relaying raw model scratch text: only configured phrase strings are
used — tool args, commands, previews, and reasoning are never interpolated.
Built-in defaults live in ``gateway/assets/status_phrases.yaml``. Users can add
profile-relative catalogs under ``HERMES_HOME`` via the conventional paths
``status_phrases.yaml`` / ``status_phrases/*.yaml`` or via config::
display:
status_phrases:
path: status_phrases/whatsapp.yaml # relative to HERMES_HOME
mode: append # append (default) or replace
Absolute paths and ``..`` escapes are ignored on purpose so config stays
profile-portable and cannot read arbitrary files.
"""
"""Human-friendly generic gateway status phrases: short chat-safe lines for the long-running
status surface without relaying raw model scratch text — only configured phrase strings are used;
tool args, commands, previews, and reasoning are never interpolated. Built-in defaults live in
``gateway/assets/status_phrases.yaml``; users add profile-relative catalogs under ``HERMES_HOME``
via ``status_phrases.yaml`` / ``status_phrases/*.yaml`` or ``display.status_phrases: {path:
<HERMES_HOME-relative>, mode: append|replace}``. Absolute paths and ``..`` escapes are ignored on
purpose so config stays profile-portable and cannot read arbitrary files."""
from __future__ import annotations
@@ -28,8 +17,8 @@ import yaml
from hermes_constants import get_hermes_home
# Hermes UI surfaces, not app/vendor buckets. Long-running-only: regular
# tool/thinking/interim chatter is deliberately not rewritten (too noisy in chat).
# Hermes UI surfaces, not app/vendor buckets. Long-running-only: regular tool/thinking/interim
# chatter is deliberately not rewritten (too noisy in chat).
_STATUS_SURFACES = ("status", "generic")
_MAX_CUSTOM_PHRASES_PER_SURFACE = 80
_MAX_PHRASE_CHARS = 160
@@ -44,19 +33,16 @@ _FALLBACK_PHRASES: dict[str, list[str]] = {
def _clean_phrase_list(value: Any) -> list[str]:
if not isinstance(value, list):
return []
cleaned: list[str] = []
for item in value[:_MAX_CUSTOM_PHRASES_PER_SURFACE]:
for item in value[:_MAX_CUSTOM_PHRASES_PER_SURFACE] if isinstance(value, list) else ():
phrase = str(item or "").strip()
if phrase and len(phrase) <= _MAX_PHRASE_CHARS and phrase not in cleaned:
cleaned.append(phrase)
return cleaned
def _merge_phrase_mapping(
catalog: dict[str, list[str]], section: Mapping[str, Any], *, inherited_mode: str | None = None
) -> None:
def _merge_phrase_mapping(catalog: dict[str, list[str]], section: Mapping[str, Any], *,
inherited_mode: str | None = None) -> None:
replace = str(section.get("mode") or inherited_mode or "append").strip().lower() == "replace"
phrase_map = section.get("phrases") if isinstance(section.get("phrases"), Mapping) else section
for surface in _STATUS_SURFACES:
@@ -77,10 +63,8 @@ def _merge_phrase_file(catalog: dict[str, list[str]], path: Path, *, inherited_m
def _iter_phrase_files(base_dir: Path, raw_path: Any) -> list[Path]:
"""YAML files under ``base_dir/raw_path``; [] for absolute / ``..`` / escaping paths."""
raw = str(raw_path or "").strip()
if not raw:
return []
candidate = Path(raw).expanduser()
if candidate.is_absolute() or ".." in candidate.parts:
if not raw or candidate.is_absolute() or ".." in candidate.parts:
return []
base = base_dir.resolve()
path = (base / candidate).resolve()
@@ -95,9 +79,8 @@ def _iter_phrase_files(base_dir: Path, raw_path: Any) -> list[Path]:
return []
def _merge_phrase_paths(
catalog: dict[str, list[str]], paths: Any, *, base_dir: Path, inherited_mode: str | None = None
) -> None:
def _merge_phrase_paths(catalog: dict[str, list[str]], paths: Any, *, base_dir: Path,
inherited_mode: str | None = None) -> None:
if paths is None:
return
for raw_path in paths if isinstance(paths, list) else [paths]:
@@ -125,24 +108,19 @@ def _merge_phrase_config(catalog: dict[str, list[str]], section: Any, *, base_di
_merge_phrase_mapping(catalog, section)
def resolve_status_phrase_catalog(
user_config: Mapping[str, Any] | None, platform_key: str | None = None
) -> dict[str, list[str]]:
"""Resolve built-in + user-configured generic status phrases.
Resolution order mirrors gateway display settings: built-ins, conventional
profile-relative user files, global ``display.status_phrases`` (or legacy
alias ``generic_status_phrases``), then
``display.platforms.<platform>.status_phrases``.
"""
def resolve_status_phrase_catalog(user_config: Mapping[str, Any] | None,
platform_key: str | None = None) -> dict[str, list[str]]:
"""Resolve built-in + user-configured generic status phrases. Order mirrors gateway display
settings: built-ins, conventional profile-relative user files, global
``display.status_phrases`` (or legacy alias ``generic_status_phrases``), then
``display.platforms.<platform>.status_phrases``."""
catalog = _copy_catalog(_DEFAULT_PHRASES)
hermes_home = get_hermes_home()
_merge_phrase_paths(catalog, list(_CONVENTIONAL_RELATIVE_PATHS), base_dir=hermes_home)
display = (user_config or {}).get("display") if isinstance(user_config, Mapping) else None
if not isinstance(display, Mapping):
return catalog
sections = [display]
platforms = display.get("platforms")
sections, platforms = [display], display.get("platforms")
if platform_key and isinstance(platforms, Mapping) and isinstance(platforms.get(platform_key), Mapping):
sections.append(platforms[platform_key])
for section in sections:
@@ -151,25 +129,19 @@ def resolve_status_phrase_catalog(
return catalog
def classify_status_context(
kind: str, *, tool_name: str | None = None, preview: str | None = None, args: Any = None,
) -> str:
def classify_status_context(kind: str, *, tool_name: str | None = None, preview: str | None = None,
args: Any = None) -> str:
"""Classify an internal gateway event into a Hermes UI-surface bucket."""
if str(kind or "").strip().lower() in {"heartbeat", "waiting", "long_running", "status"}:
return "status"
return "generic"
def choose_status_phrase(
kind: str, *, tool_name: str | None = None, preview: str | None = None, args: Any = None,
recent: MutableSequence[str] | None = None, rng: Any = None,
catalog: Mapping[str, list[str]] | None = None,
) -> str:
"""Pick a short generic status phrase, avoiding recent repeats.
``preview`` and ``args`` are accepted for callback compatibility, but their
raw contents are never embedded in the returned phrase.
"""
def choose_status_phrase(kind: str, *, tool_name: str | None = None, preview: str | None = None,
args: Any = None, recent: MutableSequence[str] | None = None,
rng: Any = None, catalog: Mapping[str, list[str]] | None = None) -> str:
"""Pick a short generic status phrase, avoiding recent repeats. ``preview`` and ``args`` are
accepted for callback compatibility, but their raw contents are never embedded in the result."""
phrase_catalog = catalog or _DEFAULT_PHRASES
category = classify_status_context(kind, tool_name=tool_name, preview=preview, args=args)
candidates = list(phrase_catalog.get(category) or phrase_catalog.get("generic") or _DEFAULT_PHRASES["generic"])

View File

@@ -10,34 +10,27 @@ import socket
def _notify_socket() -> str:
address = os.environ.get("NOTIFY_SOCKET", "").strip()
return address if address and hasattr(socket, "AF_UNIX") else ""
return os.environ.get("NOTIFY_SOCKET", "").strip() if hasattr(socket, "AF_UNIX") else ""
def notify(message: str) -> bool:
"""Send one nonblocking sd_notify datagram when systemd configured it.
Failures are non-fatal: a missing socket or an older platform must never prevent the gateway
from starting.
"""
"""Send an sd_notify datagram if systemd configured it; failures never block gateway startup."""
if not (address := _notify_socket()) or not isinstance(message, str) or not message:
return False
try:
payload = message.encode("utf-8")
with socket.socket(socket.AF_UNIX, socket.SOCK_DGRAM) as sender:
sender.setblocking(False) # a full receiver buffer must not stall the event loop
# systemd's ``@abstract`` notation -> Python's leading-NUL address form.
# systemd's ``@abstract`` notation -> Python's leading-NUL address form
sender.connect("\0" + address[1:] if address.startswith("@") else address)
sender.send(payload)
sender.send(message.encode("utf-8"))
return True
except (OSError, UnicodeError, ValueError):
return False
def watchdog_interval_seconds() -> float | None:
raw = os.environ.get("WATCHDOG_USEC", "").strip() if _notify_socket() else ""
try:
interval = float(raw) / 1_000_000.0
interval = float(os.environ.get("WATCHDOG_USEC", "") if _notify_socket() else "") / 1e6
except (TypeError, ValueError):
return None
return interval if math.isfinite(interval) and interval > 0 else None
@@ -46,9 +39,7 @@ def watchdog_interval_seconds() -> float | None:
class SystemdWatchdog:
"""Feed systemd while the asyncio event loop continues to make progress."""
def __init__(
self, *, config_enabled: bool = True, lag_tolerance_seconds: float | None = None
) -> None:
def __init__(self, *, config_enabled: bool = True, lag_tolerance_seconds: float | None = None):
self._config_enabled = bool(config_enabled)
self.interval_seconds = watchdog_interval_seconds()
self._lag_tolerance_seconds = lag_tolerance_seconds
@@ -61,11 +52,10 @@ class SystemdWatchdog:
def _lag_tolerance(self) -> float:
default = max(0.1, (self.interval_seconds or 0.0) * 0.25)
try:
with contextlib.suppress(TypeError, ValueError):
value = float(self._lag_tolerance_seconds)
except (TypeError, ValueError):
return default
return max(0.0, value) if math.isfinite(value) else default
return max(0.0, value) if math.isfinite(value) else default
return default
def start(self) -> bool:
if not self.enabled:
@@ -111,17 +101,16 @@ class SystemdWatchdog:
now = loop.time()
if not self.record_tick(scheduled_at=scheduled_at, now=now):
return
scheduled_at = (
scheduled_at + cadence if scheduled_at + cadence >= now else now + cadence
)
scheduled_at += cadence
if scheduled_at < now:
scheduled_at = now + cadence
async def stop(self) -> None:
"""Stop feeding systemd and emit ``STOPPING=1`` at most once."""
self._stopping = True
task = self._task
if task is not None and task is not asyncio.current_task():
if not task.done():
task.cancel()
task.cancel() # no-op on a finished task
with contextlib.suppress(asyncio.CancelledError, Exception):
await task
self._task = None

View File

@@ -1,17 +1,13 @@
"""Wake an existing agent session from a background completion event.
Delivery is selected by the adapter's ``supports_async_delivery`` flag. Push-capable
adapters get a synthetic ``MessageEvent(internal=True)`` via ``handle_message``.
Stateless adapters (API server) would run that under a ``build_session_key()`` key
that never matches the raw ``X-Hermes-Session-Id`` real turns use (invisible parallel
session), so we self-POST ``/v1/chat/completions`` with the raw session id header and
resume the REAL session. Async-delegation completions are the exception there: the
CLIENT owns the next turn, so a completion is never self-POSTed as a new ``role=user``
prompt (it could cross a pending human-confirmation gate); ``persist_delegation_delivery``
writes a durable DELIVERY row (``display_kind="async_delegation_complete"``, the shape the
TUI/desktop pollers read) instead. Failures RAISE (after bounded retries on transient
errors) so callers can rewind cursors / retry instead of silently losing the event.
"""
"""Wake an existing agent session from a background completion event. Push-capable adapters
(``supports_async_delivery``) get a synthetic ``MessageEvent(internal=True)`` via handle_message;
stateless adapters (API server) would run that under a ``build_session_key()`` key that never
matches the raw ``X-Hermes-Session-Id`` real turns use (invisible parallel session), so we self-POST
``/v1/chat/completions`` with the raw id header to resume the REAL session. Exception:
async-delegation completions: the CLIENT owns the next turn, so they are never self-POSTed as a
new ``role=user`` prompt (could cross a pending human-confirmation gate); instead
``persist_delegation_delivery`` writes a durable DELIVERY row (``display_kind=
"async_delegation_complete"``, read by TUI/desktop pollers). Failures RAISE (after bounded retries
on transient errors) so callers can rewind cursors / retry instead of silently losing the event."""
from __future__ import annotations
@@ -21,34 +17,29 @@ from typing import Any, Optional
logger = logging.getLogger(__name__)
# A wake self-post runs the whole agent turn synchronously (stream=false);
# generous ceiling so long tool-using turns aren't killed mid-flight.
# A wake self-post runs the whole agent turn synchronously (stream=false); generous ceiling so long
# tool-using turns aren't killed mid-flight.
WAKE_TURN_TIMEOUT_SECONDS = 600.0
# Backoff between retries on transient failures. The API server has no
# per-session lock (concurrent turns are last-writer-wins) but DOES enforce a
# global max_concurrent_runs cap via HTTP 429, which is worth waiting out.
# Backoff between retries on transient failures. The API server has no per-session lock (concurrent
# turns are last-writer-wins) but DOES enforce a global max_concurrent_runs cap via HTTP 429, which
# is worth waiting out.
_RETRY_DELAYS_SECONDS = (2.0, 5.0, 10.0)
def adapter_supports_push(adapter: Any) -> bool:
"""Whether this adapter can push a message to the user after a turn ends.
Reads ``supports_async_delivery`` off the adapter class rather than the
request-scoped contextvar — background watchers run outside any bound
session context. Adapters that don't declare the flag are push-capable.
"""
"""Whether this adapter can push a message to the user after a turn ends. Reads
``supports_async_delivery`` off the adapter class rather than the request-scoped contextvar
(background watchers run outside any bound session context). Adapters that don't declare
the flag are push-capable."""
return bool(getattr(adapter, "supports_async_delivery", True))
async def deliver_wake(adapter: Any, *, text: str, session_id: str = "", source: Any = None) -> None:
"""Deliver a wake turn to the session behind ``adapter``.
``session_id`` is the RAW session id (``X-Hermes-Session-Id`` / state.db
key) — required for non-push adapters. ``source`` is the ``SessionSource``
for the synthetic event — required for push-capable adapters. Raises on
failure so the caller can rewind/retry.
"""
"""Deliver a wake turn to the session behind ``adapter``. ``session_id`` is the RAW session id
(``X-Hermes-Session-Id`` / state.db key) — required for non-push adapters. ``source`` is the
``SessionSource`` for the synthetic event — required for push-capable adapters. Raises on
failure so the caller can rewind/retry."""
if adapter_supports_push(adapter):
if source is None:
raise ValueError("deliver_wake: push-capable adapter requires a SessionSource")
@@ -57,30 +48,23 @@ async def deliver_wake(adapter: Any, *, text: str, session_id: str = "", source:
await adapter.handle_message(synth_event)
return
if not session_id:
raise ValueError(
"deliver_wake: non-push adapter (supports_async_delivery=False) "
"requires the raw session id to self-post the wake turn"
)
raise ValueError("deliver_wake: non-push adapter (supports_async_delivery=False) "
"requires the raw session id to self-post the wake turn")
await _self_post_chat_completion(adapter, text=text, session_id=session_id)
def _delegation_display_metadata(evt: dict) -> dict:
"""Display-only metadata for a persisted delegation delivery row.
Mirrors ``tui_gateway.server._async_delegation_display_metadata`` (same
``display_kind`` consumer contract) without importing the TUI stack.
"""
"""Display-only metadata for a persisted delegation delivery row. Mirrors
``tui_gateway.server._async_delegation_display_metadata`` (same ``display_kind`` consumer
contract) without importing the TUI stack."""
raw_results = evt.get("results")
results = [r for r in raw_results if isinstance(r, dict)] if isinstance(raw_results, list) else []
task_count = len(results) or 1
completed_count = sum(1 for r in results if r.get("status") in {"completed", "success"})
failed_count = sum(1 for r in results if r.get("status") in {"failed", "error"})
metadata = {
"delegation_id": str(evt.get("delegation_id") or ""),
"task_count": task_count,
"completed_count": completed_count or task_count - failed_count,
"failed_count": failed_count,
}
metadata = {"delegation_id": str(evt.get("delegation_id") or ""), "task_count": task_count,
"completed_count": completed_count or task_count - failed_count,
"failed_count": failed_count}
duration = evt.get("total_duration_seconds") or evt.get("duration_seconds")
if isinstance(duration, (int, float)):
metadata["duration_seconds"] = duration
@@ -88,23 +72,17 @@ def _delegation_display_metadata(evt: dict) -> dict:
async def persist_delegation_delivery(adapter: Any, *, text: str, session_id: str, evt: Optional[dict] = None) -> None:
"""Persist an async-delegation completion as a durable DELIVERY row
(see module docstring) WITHOUT running any agent turn.
Raises on failure so the caller can release the durable claim and retry.
"""
"""Persist an async-delegation completion as a durable DELIVERY row (see module docstring)
WITHOUT running any agent turn. Raises on failure so the caller can release the durable claim
and retry."""
if not session_id:
raise ValueError(
"persist_delegation_delivery: raw session id required to persist "
"the completion on the api_server session transcript"
)
raise ValueError("persist_delegation_delivery: raw session id required to persist "
"the completion on the api_server session transcript")
ensure = getattr(adapter, "_ensure_session_db", None)
db: Any = await asyncio.to_thread(ensure) if callable(ensure) else None
if db is None:
raise RuntimeError(
"persist_delegation_delivery: api_server SessionDB unavailable — "
f"cannot persist completion for session {session_id}"
)
raise RuntimeError("persist_delegation_delivery: api_server SessionDB unavailable — "
f"cannot persist completion for session {session_id}")
await asyncio.to_thread(
db.append_message, session_id, "user", content=text,
display_kind="async_delegation_complete", display_metadata=_delegation_display_metadata(evt or {}),
@@ -115,12 +93,10 @@ async def persist_delegation_delivery(adapter: Any, *, text: str, session_id: st
async def _self_post_chat_completion(adapter: Any, *, text: str, session_id: str) -> None:
"""POST the wake text to the in-pod API server as a normal session turn.
Uses the adapter's own bind host/port/key. Session continuation via
``X-Hermes-Session-Id`` is 403-gated on ``API_SERVER_KEY``, so a missing
key is a hard error rather than a wake in a fresh session nobody watches.
"""
"""POST the wake text to the in-pod API server as a normal session turn, using the adapter's
own bind host/port/key. Session continuation via ``X-Hermes-Session-Id`` is 403-gated on
``API_SERVER_KEY``, so a missing key is a hard error rather than a wake in a fresh session
nobody watches."""
import aiohttp
host = str(getattr(adapter, "_host", "") or "127.0.0.1")
if host in ("0.0.0.0", "::", "*"):
@@ -128,20 +104,15 @@ async def _self_post_chat_completion(adapter: Any, *, text: str, session_id: str
port = int(getattr(adapter, "_port", 0) or 8642)
api_key = str(getattr(adapter, "_api_key", "") or "")
if not api_key:
raise RuntimeError(
"wake self-post requires API_SERVER_KEY: session continuation via "
"X-Hermes-Session-Id is rejected (403) on an unauthenticated API "
"server, so the wake cannot reach the target session"
)
raise RuntimeError("wake self-post requires API_SERVER_KEY: session continuation via "
"X-Hermes-Session-Id is rejected (403) on an unauthenticated API "
"server, so the wake cannot reach the target session")
if ":" in host and not host.startswith("["):
host = f"[{host}]" # bare IPv6 literal
url = f"http://{host}:{port}/v1/chat/completions"
headers = {"Authorization": f"Bearer {api_key}", "X-Hermes-Session-Id": session_id}
payload = {
"model": str(getattr(adapter, "_model_name", "") or "hermes-agent"),
"messages": [{"role": "user", "content": text}],
"stream": False,
}
payload = {"model": str(getattr(adapter, "_model_name", "") or "hermes-agent"),
"messages": [{"role": "user", "content": text}], "stream": False}
last_err: Optional[BaseException] = None
attempts = 1 + len(_RETRY_DELAYS_SECONDS)
for attempt in range(attempts):
@@ -151,15 +122,13 @@ async def _self_post_chat_completion(adapter: Any, *, text: str, session_id: str
timeout = aiohttp.ClientTimeout(total=WAKE_TURN_TIMEOUT_SECONDS)
async with aiohttp.ClientSession(timeout=timeout) as http:
async with http.post(url, json=payload, headers=headers) as resp:
if resp.status == 429:
# Global concurrency cap — transient; back off and retry.
if resp.status == 429: # concurrency cap — transient; back off and retry
last_err = RuntimeError(
f"wake self-post got HTTP 429 (concurrency cap) for session {session_id}"
)
logger.warning("%s; attempt %d/%d", last_err, attempt + 1, attempts)
continue
if resp.status >= 400:
# Non-transient (auth/validation) — fail immediately.
if resp.status >= 400: # non-transient (auth/validation): fail immediately
body = (await resp.text())[:300]
raise RuntimeError(
f"wake self-post failed for session {session_id}: HTTP {resp.status}: {body}"
@@ -169,10 +138,8 @@ async def _self_post_chat_completion(adapter: Any, *, text: str, session_id: str
return
except (aiohttp.ClientError, asyncio.TimeoutError, OSError) as exc:
last_err = exc
logger.warning(
"wake self-post transient failure for session %s (attempt %d/%d): %s",
session_id, attempt + 1, attempts, exc,
)
logger.warning("wake self-post transient failure for session %s (attempt %d/%d): %s",
session_id, attempt + 1, attempts, exc)
raise RuntimeError(
f"wake self-post gave up for session {session_id} after {attempts} attempts: {last_err}"
) from last_err