From 8d5d5f9ca296560fbff60bd7f9155dd7a4b00e33 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 21:14:22 -0700 Subject: [PATCH] refactor(gateway): third pass on scale_to_zero/wake/status_phrases/runtime_footer/systemd_notify --- gateway/runtime_footer.py | 72 +++++++------------- gateway/scale_to_zero.py | 121 ++++++++++++---------------------- gateway/status_phrases.py | 86 ++++++++---------------- gateway/systemd_notify.py | 37 ++++------- gateway/wake.py | 135 ++++++++++++++------------------------ 5 files changed, 161 insertions(+), 290 deletions(-) diff --git a/gateway/runtime_footer.py b/gateway/runtime_footer.py index 69362ed0da..090bdf0cc9 100644 --- a/gateway/runtime_footer.py +++ b/gateway/runtime_footer.py @@ -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..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.

.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) diff --git a/gateway/scale_to_zero.py b/gateway/scale_to_zero.py index 8f45f75958..a668d262a9 100644 --- a/gateway/scale_to_zero.py +++ b/gateway/scale_to_zero.py @@ -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 diff --git a/gateway/status_phrases.py b/gateway/status_phrases.py index 9070298056..22a42710ea 100644 --- a/gateway/status_phrases.py +++ b/gateway/status_phrases.py @@ -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: +, 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..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..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"]) diff --git a/gateway/systemd_notify.py b/gateway/systemd_notify.py index 9cabb6d46b..da66753533 100644 --- a/gateway/systemd_notify.py +++ b/gateway/systemd_notify.py @@ -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 diff --git a/gateway/wake.py b/gateway/wake.py index feaafa9a39..aa4bc84b98 100644 --- a/gateway/wake.py +++ b/gateway/wake.py @@ -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