diff --git a/gateway/run.py b/gateway/run.py index ec290de27f..c5345e0163 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -7,8 +7,7 @@ Run via ``python -m gateway.run`` or ``python cli.py --gateway``.""" try: import hermes_bootstrap # noqa: F401 except ModuleNotFoundError: - # A partial ``hermes update`` can leave the bootstrap unregistered; only Windows UTF-8 stdio suffers. - pass + pass # a partial ``hermes update`` can leave the bootstrap unregistered; only Windows UTF-8 stdio suffers import asyncio import concurrent.futures @@ -42,8 +41,7 @@ from agent.turn_context import compression_made_progress from hermes_cli.config import _is_ssh_remote_tilde_cwd, cfg_get from hermes_cli.fallback_config import get_fallback_chain -# Per-session AIAgent cache bounds (agents are heavy); enforced by -# _enforce_agent_cache_cap()/_session_expiry_watcher(). +# Per-session AIAgent cache bounds (agents are heavy); see _enforce_agent_cache_cap/_session_expiry_watcher. _AGENT_CACHE_MAX_SIZE = 128 _AGENT_CACHE_IDLE_TTL_SECS = 3600.0 # evict agents idle for >1h _PLATFORM_CONNECT_TIMEOUT_SECS_DEFAULT = 30.0 @@ -53,11 +51,9 @@ _TELEGRAM_CONNECT_TIMEOUT_SECS_DEFAULT = 180.0 _TELEGRAM_INITIAL_CONNECT_TIMEOUT_SECS_DEFAULT = 45.0 _ADAPTER_DISCONNECT_TIMEOUT_SECS_DEFAULT = 5.0 # End reasons meaning the USER deliberately closed this thread. Shared by _classify_completion_target and -# _resolve_async_delegation_session so they never disagree (a "delivered" reason the resolver drops is -# acked and then silently lost). +# _resolve_async_delegation_session so they never disagree (else a "delivered" reason is acked, then lost). _USER_BOUNDARY_END_REASONS = ("session_reset", "user_exit", "session_switch", "new_session") -# Bounds one stall-notify send so a wedged transport can't block the watcher; on timeout the next tick -# retries. +# Bounds one stall-notify send so a wedged transport can't block the watcher; on timeout the next tick retries. _STALL_NOTIFY_SEND_TIMEOUT_SECONDS = 15.0 _GATEWAY_PROXY_SSE_BUFFER_MAX_CHARS = 16 * 1024 * 1024 _TELEGRAM_COMMAND_MENTION_RE = re.compile(r"(?)\s+\d[\d,]*\s+messages,\s+retrying" r"|compressed\s+~[\d,]+\s+(?:→|->)\s+~[\d,]+\s+tokens,\s+retrying" @@ -96,11 +91,9 @@ _TELEGRAM_NOISY_STATUS_RE = re.compile( re.IGNORECASE | re.DOTALL) _HYGIENE_COOLDOWN_LADDER_MULTIPLIERS = (1, 3, 9) -# Ceiling on an escalated cooldown (cf. _RECONNECT_BACKOFF_CAP): an operator-raised base × ladder reaches -# 9h, indistinguishable from "compaction silently switched off". +# Ceiling on an escalated cooldown (cf. _RECONNECT_BACKOFF_CAP): base × ladder can reach 9h ≈ "compaction off". _HYGIENE_COOLDOWN_MAX_SECONDS = 3600.0 -# Flat retry-after when hygiene compression is ABANDONED by turn-hold expiry (not a failure, so -# outside the streak ladder); keeps sustained traffic from spawn/hold/cancelling one every turn. +# Flat retry-after when hygiene is ABANDONED by turn-hold expiry (not a failure: outside the streak ladder). _HYGIENE_TURNHOLD_RETRY_SECONDS = 60.0 @@ -111,10 +104,9 @@ def _gateway_session_db_inner(gateway): def _hygiene_cooldown_for_failure(gateway, session_key: str, base_cooldown_seconds: float) -> float: - """Bump the hygiene failure streak and return the escalated cooldown. + """Bump the hygiene failure streak and return the escalated cooldown (x1/x3/x9 over base, clamped). - Ladder (x1, x3, x9) over the base, clamped to the max. Hygiene's per-run ``AIAgent`` is fresh - (in-memory streak always 0), so the streak lives in SQLite keyed by rotation-stable ``session_key``. + Hygiene's per-run ``AIAgent`` is fresh, so the streak lives in SQLite keyed by rotation-stable session_key. """ streak, state = 1, None try: @@ -143,8 +135,7 @@ def _hygiene_cooldown_for_failure(gateway, session_key: str, base_cooldown_secon def _reset_hygiene_failure_streak(gateway, session_key: str) -> None: """Clear the hygiene failure streak after a compression that reduced context. - Peeks, never get-or-creates: a no-op 0 write must not create a never-evicted ``_sessions`` row. - """ + Peeks, never get-or-creates: a no-op 0 write must not create a never-evicted ``_sessions`` row.""" try: state = gateway._peek_session_state(session_key) if state is not None: @@ -164,9 +155,8 @@ def hygiene_compaction_recovered( approx_tokens: int, new_tokens: int) -> bool: """True when a hygiene run actually recovered the session (extracted to be unit testable). - Requires: no abort; transcript actually rewritten (the no-op path reuses pre-compression counts); - and material shrink per :func:`compression_made_progress` (a bare ``<`` misses row-count wins and - counts 30-50% estimate noise as one).""" + Requires no abort, a real rewrite (the no-op path reuses pre-compression counts) and material shrink per + :func:`compression_made_progress` (a bare ``<`` misses row-count wins and counts estimate noise).""" if aborted or not (rotated or in_place): return False return compression_made_progress(msg_count, new_count, approx_tokens, new_tokens) @@ -207,43 +197,37 @@ async def run_codex_hygiene_compaction( approx_tokens: int, timeout_seconds: float, failure_cooldown_seconds: float = 300.0) -> str: """Session hygiene for ``codex_app_server`` sessions. - The real context is the server-side thread; the local transcript is a never-replayed mirror, so - rewriting it shrinks nothing and evicting the live agent starts the next turn on an EMPTY thread. - So: compact the LIVE agent via ``thread/compact/start``, keep it cached, never build a detached - compressor. ``native``/``off`` skip without falling back to the local compressor. - Returns ``compacted``, ``skipped:`` or ``failed:``.""" + The real context is the server-side thread; the local transcript is a never-replayed mirror, so rewriting + it shrinks nothing and evicting the live agent starts the next turn on an EMPTY thread. So: compact the LIVE + agent via ``thread/compact/start``, keep it cached, never build a detached compressor. ``native``/``off`` + skip without local fallback. Returns ``compacted``, ``skipped:`` or ``failed:``.""" mode = str(auto_mode or "native").lower() if mode not in {"native", "hermes", "off"}: mode = "native" if mode != "hermes": - # native = app-server compacts itself; off = operator disabled it. A local transcript - # fallback cannot shrink the thread in any mode, so both skip cleanly with no eviction. + # native = app-server compacts itself; off = operator disabled. Local fallback can't shrink the thread. return f"skipped:mode={mode}" agent = _cached_agent_for_hygiene(gateway, session_key) if agent is None or agent is _AGENT_PENDING_SENTINEL: - # No live agent → no live thread; the detached mirror-only rewrite would be the very no-op - # this function exists to remove. + # No live agent → no live thread; a detached mirror-only rewrite is the no-op this exists to remove. return "skipped:no-cached-agent" if getattr(agent, "_codex_session", None) is None: return "skipped:no-live-thread" compressor = getattr(agent, "context_compressor", None) count_before = getattr(compressor, "compression_count", 0) - # copy_context keeps the caller's multiplexed profile secret scope and HERMES_HOME override in the - # worker (the default executor does not propagate ContextVars on the runtimes Hermes ships). + # copy_context carries profile secret scope / HERMES_HOME override (executors don't propagate ContextVars). worker_future = asyncio.get_running_loop().run_in_executor( None, copy_context().run, lambda: agent._compress_context(history, "", approx_tokens=approx_tokens)) track_worker = getattr(gateway, "_track_deferred_agent_worker", None) if callable(track_worker): - # ``wait_for`` only cancels the asyncio wrapper; keep the still-running executor thread - # visible to gateway shutdown. + # ``wait_for`` only cancels the asyncio wrapper; keep the running executor thread visible to shutdown. track_worker(worker_future, agent) try: await asyncio.wait_for(asyncio.shield(worker_future), timeout=max(float(timeout_seconds), 1.0)) except asyncio.TimeoutError: - # The executor thread keeps running (compact_thread has its own RPC timeouts); brake - # retries so a wedged app-server does not re-trigger compaction on every message. + # Executor thread keeps running (own RPC timeouts); brake retries so a wedged app-server isn't re-hit. if failure_cooldown_seconds >= 0: _record_hygiene_cooldown( gateway, session_id, failure_cooldown_seconds, "codex app-server thread compaction timed out") @@ -259,8 +243,7 @@ async def run_codex_hygiene_compaction( count_after = getattr(compressor, "compression_count", 0) if count_after > count_before: - # Native boundary recorded: compacted server-side; mirror intentionally NOT rewritten, - # agent stays cached. + # Native boundary recorded: compacted server-side; mirror NOT rewritten, agent stays cached. _reset_hygiene_failure_streak(gateway, session_key) return "compacted" # No boundary: internal skip or compaction error; the codex route already persisted its own cooldown. @@ -271,17 +254,15 @@ def hygiene_wait_should_extend( ) -> bool: """Whether the hygiene host should keep waiting for a slow summary. - A cancelled commit fence cannot produce a commit: extending the wait only queues inbound - messages behind a doomed attempt, so stop extending immediately.""" + A cancelled commit fence cannot commit: extending only queues inbound messages behind a doomed attempt.""" return not fence_cancelled and idle < timeout and waited < ceiling def _record_hygiene_cooldown( gateway, session_id: str, cooldown_seconds: float, error: Optional[str] = None) -> None: - """Persist a session-hygiene compression-failure cooldown to the state DB. + """Persist a session-hygiene compression-failure cooldown to the state DB (survives restarts). - Shares the in-conversation path's column/recorder so it survives restarts. ``error`` must be - forwarded: the recorder writes ``compression_failure_error`` UNCONDITIONALLY (else NULL clobber). + ``error`` must be forwarded: the recorder writes compression_failure_error UNCONDITIONALLY (NULL clobber). """ recorder = getattr(_gateway_session_db_inner(gateway), "record_compression_failure_cooldown", None) if recorder is None: @@ -295,8 +276,7 @@ def _record_hygiene_cooldown( def _status_template_to_regex(template: str) -> str: """Compile a compression status template constant into a regex source. - Literal text is escaped verbatim so wording drift cannot diverge from the matcher; - each ``{field}`` becomes a numeric-ish pattern.""" + Literal text is escaped verbatim (wording drift can't diverge from the matcher); ``{field}`` -> numeric.""" parts = re.split(r"\{[^{}]*\}", template) return r"[\d,]+".join(re.escape(part) for part in parts) @@ -317,8 +297,7 @@ _COMPRESSION_PROGRESS_STATUS_RE = re.compile( def _gateway_compression_progress_notices_enabled() -> bool: """True when ``compression.progress_notices`` is on (default False: chat is silent by design). - Read live (mtime-cached) so a config edit applies at the next status; fail-closed on read error. - """ + Read live (mtime-cached) so a config edit applies at the next status; fail-closed on read error.""" try: config = _load_gateway_config() compression_cfg = config.get("compression") if isinstance(config, dict) else None @@ -378,8 +357,7 @@ _GATEWAY_SECRET_PATTERNS = ( def _ensure_windows_gateway_venv_imports() -> None: """Make detached Windows gateway runs see the Hermes venv packages. - Patched before MCP discovery so tool injection does not depend on launchers preserving PYTHONPATH. - """ + Patched before MCP discovery so tool injection does not depend on launchers preserving PYTHONPATH.""" if sys.platform != "win32": return @@ -408,8 +386,7 @@ def _ensure_windows_gateway_venv_imports() -> None: site_entry = str(site_packages) if project_entry not in sys.path: sys.path.insert(0, project_entry) - # addsitepackages() semantics matter here: pywin32, used by the MCP - # SDK on Windows, relies on .pth processing to expose pywintypes. + # addsitedir semantics matter: pywin32 (MCP SDK on Windows) needs .pth processing for pywintypes. site.addsitedir(site_entry) if site_entry in sys.path: sys.path.remove(site_entry) @@ -442,9 +419,8 @@ def _non_conversational_metadata( def _interim_metadata(metadata: Optional[Dict[str, Any]] = None) -> Dict[str, Any]: """Mark a mid-turn status/advisory send as NOT the turn-final. - Stream-is-the-message adapters seal the live stream with the first unmarked send to an armed - (chat, turn) key, so every mid-turn gateway send MUST carry this marker or it seals the user's - answer stream with status text. Gateway-internal; adapters strip it before the wire.""" + Stream-is-the-message adapters seal the live stream with the first unmarked send to an armed (chat, turn) + key, so every mid-turn send MUST carry this marker. Gateway-internal; adapters strip it before the wire.""" merged = dict(metadata or {}) merged["_interim_send"] = True return merged @@ -453,9 +429,8 @@ def _interim_metadata(metadata: Optional[Dict[str, Any]] = None) -> Dict[str, An def _seed_hygiene_system_prompt(agent: Any, session_row: Optional[Dict[str, Any]]) -> bool: """Keep gateway hygiene from rebuilding a live session's system prompt. - Hygiene lacks the live session's initialized prompt environment and compression may persist a - system prompt, so a rebuild would strip external provider blocks. Seed the persisted prompt (or an - empty cache entry); the real turn rebuilds with fully initialized providers.""" + Hygiene lacks the live prompt environment, so a rebuild (persisted by compression) would strip external + provider blocks. Seed the persisted prompt (or an empty cache entry); the real turn rebuilds properly.""" stored_prompt = "" if isinstance(session_row, dict): raw_prompt = session_row.get("system_prompt") @@ -475,8 +450,7 @@ _TRANSIENT_NETWORK_ERROR_CLASS_NAMES = frozenset({ def _is_transient_network_error(exc: BaseException) -> bool: """True for transient network errors safe to log + swallow (the next poll recovers; never crash). - Walks the cause chain so wrapped errors (PTB ``NetworkError`` over ``httpx.ConnectError``) match. - """ + Walks the cause chain so wrapped errors (PTB ``NetworkError`` over ``httpx.ConnectError``) match.""" seen: set[int] = set() cur: Optional[BaseException] = exc depth = 0 @@ -496,8 +470,7 @@ def _gateway_loop_exception_handler( loop: "asyncio.AbstractEventLoop", context: Dict[str, Any]) -> None: """Loop-level safety net for transient network errors (installed once by ``start_gateway``). - Logs WARNING with traceback; non-transient errors go to the default handler so real bugs surface. - """ + Logs WARNING with traceback; non-transient errors go to the default handler so real bugs surface.""" exc = context.get("exception") if exc is not None and _is_transient_network_error(exc): task = context.get("future") or context.get("task") @@ -517,17 +490,15 @@ def _gateway_loop_exception_handler( def _redact_gateway_user_facing_secrets(text: str) -> str: """Secret redaction before text can leave the gateway. - Delegates to the shared ``redact_sensitive_text`` (full credential set) with ``force=True`` so - it holds even when ``security.redact_secrets`` is off; ``_GATEWAY_SECRET_PATTERNS`` is a second - pass so redaction degrades gracefully if that import fails.""" + Shared ``redact_sensitive_text`` with ``force=True`` (holds even when ``security.redact_secrets`` is off); + ``_GATEWAY_SECRET_PATTERNS`` is a second pass so redaction degrades gracefully if that import fails.""" redacted = str(text or "") try: from agent.redact import redact_sensitive_text redacted = redact_sensitive_text(redacted, force=True) except Exception: - # Fail-soft: the local pattern pass below still runs rather than leaking raw text to chat. - pass + pass # fail-soft: the local pattern pass below still runs rather than leaking raw text to chat for pattern in _GATEWAY_SECRET_PATTERNS: redacted = pattern.sub(lambda m: (m.group(1) if m.lastindex else "") + "[REDACTED]", redacted) return redacted @@ -536,8 +507,7 @@ def _redact_gateway_user_facing_secrets(text: str) -> str: def _redact_approval_command(cmd: "str | None") -> str: """Redact credentials from a command before it goes into an approval prompt. - The prompt is built from the raw command, so a Tirith-flagged credential would otherwise echo - verbatim to chat; ``force=True`` holds even when ``security.redact_secrets`` is off.""" + Else a Tirith-flagged credential echoes verbatim to chat; ``force=True`` holds even with redaction off.""" from agent.redact import redact_sensitive_text return redact_sensitive_text(str(cmd or ""), force=True) @@ -548,9 +518,8 @@ def _format_exec_approval_fallback( allow_session: bool = True, smart_denied: bool = False) -> str: """Render the text fallback from approval capabilities, not platform names.""" cmd_preview = command[:200] + "..." if len(command) > 200 else command - heading = "⚠️ **Dangerous command requires approval:**" - if smart_denied: - heading = "⚠️ **Smart DENY — owner override for one operation:**" + heading = ("⚠️ **Smart DENY — owner override for one operation:**" if smart_denied + else "⚠️ **Dangerous command requires approval:**") choices = [f"Reply `{command_prefix}approve` to execute this one operation"] if not smart_denied and allow_session: @@ -598,8 +567,7 @@ _GATEWAY_PROVIDER_ERROR_SHAPE_RE = re.compile( def _looks_like_gateway_provider_error(text: str) -> bool: """True when text is a provider failure envelope, not normal content. - Text must be short (envelopes are 1-3 lines) AND start with the marker, so prose that merely - mentions a status code does not match.""" + Must be short (envelopes are 1-3 lines) AND start with the marker, so prose citing a status code misses.""" if not text: return False body = str(text).strip() @@ -614,8 +582,7 @@ def _sanitize_gateway_final_response(platform: Any, text: str) -> str: if not text or _gateway_surface_passes_raw_text(platform): return text - # Lone UTF-16 surrogates make Telegram/Signal ``.encode()`` raise before any send. Last line of - # defense for legacy/plugin paths; the raw-text surfaces above pass through (JSON escapes safely). + # Lone UTF-16 surrogates make Telegram/Signal ``.encode()`` raise; last defense for legacy/plugin paths. from agent.message_sanitization import _sanitize_surrogates text = _sanitize_surrogates(str(text)) @@ -633,8 +600,7 @@ def _sanitize_gateway_final_response(platform: Any, text: str) -> str: def _prepare_gateway_status_message(platform: Any, event_type: str, message: str) -> Optional[str]: """Filter/sanitize agent status callbacks before platform delivery. - Local/CLI keep the raw diagnostic stream; messaging surfaces drop transient aux/compression noise. - """ + Local/CLI keep the raw diagnostic stream; messaging surfaces drop transient aux/compression noise.""" text = str(message or "").strip() if not text: return None @@ -642,8 +608,7 @@ def _prepare_gateway_status_message(platform: Any, event_type: str, message: str return text text = _redact_gateway_user_facing_secrets(text) - # Opt-in `compression.progress_notices` lets ROUTINE progress through; membership comes from the - # template constants, so other noise (aux failures, retry chatter) stays suppressed even when open. + # Opt-in `compression.progress_notices` lets ROUTINE (template-derived) progress through; other noise stays. if _TELEGRAM_NOISY_STATUS_RE.search(text) and not ( _gateway_compression_progress_notices_enabled() and _COMPRESSION_PROGRESS_STATUS_RE.search(text) ): @@ -656,8 +621,7 @@ def _prepare_gateway_status_message(platform: Any, event_type: str, message: str def render_notice_line(notice) -> str: """Render an AgentNotice to a single plaintext line (messaging has no status bar: one-shot push). - The notice policy already bakes the level glyph into the text — prepending one would DOUBLE it. - Fail-soft: a malformed/empty notice degrades to "" rather than raising.""" + The level glyph is already baked into the text (prepending would DOUBLE it); malformed/empty -> "".""" return str(getattr(notice, "text", "") or "").strip() @@ -673,9 +637,8 @@ async def _send_or_update_status_coro(adapter, chat_id, status_key, content, met def _approval_send_outcome(future, timeout: float) -> str: """Classify an approval prompt send as ``sent`` / ``failed`` / ``ambiguous``. - ``ambiguous`` = scheduling future timed out but the card may have posted (late connector ack): - keep the registration alive, do NOT re-send or fall back. Only a DEFINITIVE failure (error - result / non-timeout exception / no future) re-asks; those log their detail here.""" + ``ambiguous`` = future timed out but the card may have posted: keep the registration, do NOT re-send. + Only a DEFINITIVE failure (error result / non-timeout exception / no future) re-asks; logged here.""" if future is None: logger.warning("Prompt send failed: no scheduling future (loop unavailable)") return "failed" @@ -693,12 +656,10 @@ def _approval_send_outcome(future, timeout: float) -> str: def _clarify_send_disposition(fut, *, session_key: str, clarify_mod) -> "str | None": - """Decide whether a clarify prompt send aborts the wait, per the boundary rule. + """Decide whether a clarify prompt send aborts the wait; returns the abort sentinel or ``None``. - As with exec-approval, the scheduling future can time out while the card HAS posted. Only a - DEFINITIVE failure tears down the registration; ``ambiguous`` stays armed and proceeds to the - bounded wait (its response timeout covers a lost card). Returns the abort sentinel or ``None``. - """ + Only a DEFINITIVE failure tears down the registration; ``ambiguous`` (card may have posted) stays armed + and proceeds to the bounded wait, whose response timeout covers a lost card.""" outcome = _approval_send_outcome(fut, timeout=15) if outcome == "failed": # Undeliverable: clear the registration and return the sentinel so the agent falls back, not hangs. @@ -720,7 +681,6 @@ def _clarify_send_then_wait(fut, *, clarify_id: str, session_key: str, clarify_m timeout = clarify_mod.get_clarify_timeout() response = clarify_mod.wait_for_response(clarify_id, timeout=float(timeout)) if response is None or response == "": - # Timeout or session-boundary cancellation return f"[user did not respond within {int(timeout / 60)}m]" return response @@ -730,9 +690,9 @@ def _resolve_progress_thread_id( ) -> Optional[str]: """Return thread/root ID that progress/status bubbles should target. - ``reply_in_thread=False`` (Slack) disables the synthetic-thread fallback: progress messages - must not create a thread the final flat reply would inherit. A source.thread_id equal to the - event's own message id is the adapter's synthetic session-keying thread — treat as no thread.""" + ``reply_in_thread=False`` (Slack): no synthetic-thread fallback, else the final flat reply inherits a thread. + A source.thread_id equal to the event's message id is the adapter's synthetic session key: no thread. + """ platform_key = str(getattr(platform, "value", platform) or "").lower() if not reply_in_thread: if source_thread_id and event_message_id and str(source_thread_id) == str(event_message_id): @@ -762,8 +722,8 @@ def _resolve_gateway_display_bool( platform: Any = None, require_platform_override_for: set[Any] | None = None) -> bool: """Resolve a boolean display setting with optional platform-only opt-in. - Scratch-text features are too noisy for threaded surfaces (Mattermost) under a global opt-in, - so they require an explicit display.platforms.. override.""" + Scratch-text is too noisy for threaded surfaces (Mattermost): they need an explicit per-platform override. + """ current_platform = _gateway_platform_value(platform or platform_key) platform_only = {_gateway_platform_value(c) for c in (require_platform_override_for or set())} if ( @@ -784,9 +744,7 @@ def _resolve_gateway_display_bool( def _telegramize_command_mentions(text: str, platform: Any) -> str: - """Rewrite slash-command mentions to Telegram-valid names (lowercase, digits, underscores only). - - Other platforms' renderings are left unchanged.""" + """Rewrite slash-command mentions to Telegram-valid names (lowercase/digits/underscore); no-op elsewhere.""" platform_value = getattr(platform, "value", platform) if platform_value != "telegram": return text @@ -800,26 +758,22 @@ def _telegramize_command_mentions(text: str, platform: Any) -> str: return _TELEGRAM_COMMAND_MENTION_RE.sub(_replace, text) -# Auto-continue interrupted turns only while fresh (last transcript row timestamp), else stale -# tool-tail/resume_pending markers revive an unrelated old task after a restart. 1h covers -# ``agent.gateway_timeout`` (30 min) plus slack; override: ``agent.gateway_auto_continue_freshness``. +# Auto-continue interrupted turns only while fresh, else stale tool-tail/resume_pending markers revive an old +# task after a restart. 1h covers agent.gateway_timeout (30 min) + slack; cfg agent.gateway_auto_continue_freshness. _AUTO_CONTINUE_FRESHNESS_SECS_DEFAULT = 60 * 60 -# How long ``_finish_startup_restore`` waits on boot auto-resume turns before releasing the inbound -# gate. Override: ``agent.gateway_startup_restore_drain_timeout``. +# Boot auto-resume drain before the inbound gate opens. Override: agent.gateway_startup_restore_drain_timeout. _STARTUP_RESTORE_DRAIN_TIMEOUT_SECS_DEFAULT = 30.0 -# Bound on the boot warm-up run BEFORE the gate opens (so the first turn gets no skeleton system prompt); -# keeps a wedged init from wedging the gateway. Override: ``agent.gateway_startup_warmup_timeout`` -# (non-positive disables). +# Bound on the boot warm-up BEFORE the gate opens (no skeleton system prompt on turn one); keeps a wedged init +# from wedging the gateway. Override: ``agent.gateway_startup_warmup_timeout`` (non-positive disables). _STARTUP_WARMUP_TIMEOUT_SECS_DEFAULT = 20.0 def _coerce_gateway_timestamp(value: Any) -> Optional[float]: """Best-effort conversion of stored gateway timestamps to epoch seconds. - Missing/unparseable -> None, so legacy transcripts keep auto-continuing instead of being dropped. - """ + Missing/unparseable -> None, so legacy transcripts keep auto-continuing instead of being dropped.""" if value is None: return None if isinstance(value, datetime): @@ -846,18 +800,17 @@ def _coerce_gateway_timestamp(value: Any) -> Optional[float]: def _auto_continue_freshness_window() -> float: - """Return the configured auto-continue freshness window in seconds (non-positive disables the gate). + """Auto-continue freshness window in seconds (non-positive disables the gate). - Thin wrapper over ``gateway.session`` kept so ``gateway.run`` imports/test patches keep working. - """ + Thin wrapper over ``gateway.session`` kept so ``gateway.run`` imports/test patches keep working.""" from gateway.session import auto_continue_freshness_window return auto_continue_freshness_window() def _startup_restore_drain_timeout_secs() -> float: - """Max seconds ``_finish_startup_restore`` waits on boot auto-resume turns before opening the - inbound gate (all inbound QUEUED until then); non-positive disables. Duplicate-agent safety does - NOT depend on it: ``_schedule_resume_pending_sessions`` claims ``_running_agents`` SYNCHRONOUSLY. + """Max seconds ``_finish_startup_restore`` holds the inbound gate for boot auto-resume; <=0 disables. + + Duplicate-agent safety does NOT depend on it: ``_schedule_resume_pending_sessions`` claims SYNCHRONOUSLY. """ return _float_env("HERMES_STARTUP_RESTORE_DRAIN_TIMEOUT", _STARTUP_RESTORE_DRAIN_TIMEOUT_SECS_DEFAULT) @@ -865,17 +818,15 @@ def _startup_restore_drain_timeout_secs() -> float: def _startup_warmup_timeout_secs() -> float: """Max seconds the boot warm-up (``_warm_turn_prerequisites``) may hold the inbound gate shut. - On timeout the gate opens anyway and the warm-up finishes in the background, so a wedged - import/probe cannot wedge the gateway. Non-positive disables it.""" + On timeout the gate opens and the warm-up finishes in the background. Non-positive disables it.""" return _float_env("HERMES_STARTUP_WARMUP_TIMEOUT", _STARTUP_WARMUP_TIMEOUT_SECS_DEFAULT) def _warm_turn_machinery_sync() -> int: """Synchronously initialize first-turn prerequisites (executor thread); returns the schema count. - Covers exactly the lazy init seen in skeleton turns: the ``run_agent`` import graph, - ``get_tool_definitions`` (materializes schemas, primes the ``check_fn`` TTL cache), context files. - """ + Covers the lazy init seen in skeleton turns: ``run_agent`` import graph, tool schemas (+ ``check_fn`` + TTL cache), context files.""" import run_agent # noqa: F401 # heavy import graph, cached in sys.modules import model_tools @@ -892,17 +843,14 @@ def _warm_turn_machinery_sync() -> int: def _as_thread_info(info: Any) -> Optional[Tuple[str, str]]: """*info* as a (thread_id, initial_name) pair, or None if it isn't one. - The pair crosses the relay connector boundary, so its shape is the connector's word, not ours. - """ + The pair crosses the relay connector boundary, so its shape is the connector's word, not ours.""" if isinstance(info, tuple) and len(info) == 2 and all(isinstance(x, str) for x in info): return cast(Tuple[str, str], info) return None def _float_env(name: str, default: float) -> float: - """Read an env var as float; unset/empty/malformed fall back to ``default``. - - A misconfigured env var (``HERMES_AGENT_TIMEOUT=abc``) must not crash the gateway or a turn.""" + """Read an env var as float; unset/empty/malformed fall back to ``default`` (never crash the gateway).""" raw = os.environ.get(name) if raw is None or raw == "": return float(default) @@ -923,9 +871,7 @@ def _stamp_hygiene_compression_provenance( def _is_fresh_gateway_interruption( value: Any, *, now: Optional[float] = None, window_secs: Optional[float] = None) -> bool: - """True when an interruption marker is fresh enough to auto-continue. - - Unknown timestamps count as fresh (legacy transcripts, in-memory test scaffolding).""" + """True when an interruption marker is fresh enough to auto-continue (unknown timestamps count as fresh).""" window = float(window_secs) if window_secs is not None else float(_AUTO_CONTINUE_FRESHNESS_SECS_DEFAULT) if window <= 0: return True @@ -938,11 +884,9 @@ def _is_fresh_gateway_interruption( def build_resume_recovery_note( reason: Optional[str], message: str = "", *, interactive: bool = True) -> str: - """Build the resume-pending recovery system note for an interrupted turn. + """Build the resume-pending recovery system note for an interrupted turn (empty ``message`` = auto-resume). - Empty ``message`` = startup auto-resume. Interactive platforms report the restore and ask what - next; on non-interactive ones (``interactive_resume = False``) nobody can answer: finish the work. - """ + Interactive platforms report the restore and ask what next; non-interactive ones finish the work.""" reason_phrase = ( "a gateway restart" if reason == "restart_timeout" else "a gateway shutdown" if reason == "shutdown_timeout" else "a gateway interruption") @@ -980,18 +924,14 @@ def _prepare_resume_pending_message( reason: Optional[str], message: Optional[str], *, interactive: bool = True) -> tuple[str, str]: """Return the recovery message and the user text to persist. - Empty original (synthesized auto-resume): persist the note — a "" user row trips the pre-call - sanitizer every call. Real user text: persist clean words so the transcript stays scaffold-free. - """ + Empty original: persist the note (a "" user row trips the pre-call sanitizer). Real text: persist clean.""" recovery_message = build_resume_recovery_note(reason, message or "", interactive=interactive) persist_message = message if isinstance(message, str) and message.strip() else recovery_message return recovery_message, persist_message -# Assistant fields that must survive replay for CLI parity (reasoning continuity, prefix-cache hits, -# provider echo). ``reasoning``/``reasoning_content``: unreconstructable thinking text (DeepSeek/Kimi/ -# Moonshot). ``reasoning_details``: opaque signatures. ``codex_*_items``: Codex blobs (resent or -# caching degrades). +# Assistant fields that must survive replay for CLI parity (reasoning continuity, prefix-cache hits, provider +# echo): unreconstructable thinking text (DeepSeek/Kimi), opaque signatures, Codex blobs (caching degrades). _ASSISTANT_REPLAY_FIELDS: tuple[str, ...] = ( "reasoning", "reasoning_content", "reasoning_details", "codex_reasoning_items", "codex_message_items", "finish_reason") @@ -1002,12 +942,10 @@ def _build_replay_entry( ) -> Dict[str, Any]: """Build a replay entry for a non-tool-calling message, preserving ``_ASSISTANT_REPLAY_FIELDS``. - ``preserve_timestamp``: only user messages need it (the stale-dangerous-confirmation stripper - reads it). Falsy fields are dropped EXCEPT ``reasoning_content``: DeepSeek/Kimi treat "" as a - sentinel; dropping it can 400.""" + ``preserve_timestamp``: only user rows need it (stale-dangerous-confirmation stripper). Falsy fields are + dropped EXCEPT ``reasoning_content``: DeepSeek/Kimi treat "" as a sentinel; dropping it can 400.""" entry: Dict[str, Any] = {"role": role, "content": content} - # api_content sidecar: forward the exact bytes previously sent so the request prefix stays - # byte-stable — ONLY if this pipeline did not rewrite the content (else we resend what was stripped). + # api_content sidecar keeps the request prefix byte-stable — ONLY if this pipeline did not rewrite content. _sidecar = msg.get("api_content") if ( role in ("user", "assistant") @@ -1020,17 +958,11 @@ def _build_replay_entry( if _rkey not in msg: continue _rval = msg.get(_rkey) - if _rkey == "reasoning_content": - # Preserve empty-string sentinel for thinking-mode replay. - if _rval is None: - continue - elif not _rval: + if (_rval is None) if _rkey == "reasoning_content" else (not _rval): continue entry[_rkey] = _rval - if preserve_timestamp: - ts = msg.get("timestamp") - if ts: - entry["timestamp"] = ts + if preserve_timestamp and msg.get("timestamp"): + entry["timestamp"] = msg["timestamp"] return entry @@ -1042,8 +974,7 @@ _CURRENT_ADDRESSED_MESSAGE_HEADER = "[Current addressed message - answer only th def _uses_telegram_observed_group_context(channel_prompt: Optional[str]) -> bool: """Return True for Telegram group turns that may include observed chatter. - Observe-unmentioned mode persists skipped group chatter for later @mentions; those rows must - not replay as ordinary user turns or a weak wake word makes old chatter look like pending work. + Observed rows must not replay as ordinary user turns, or a weak wake word makes old chatter look like work. """ return bool(channel_prompt and _TELEGRAM_OBSERVED_CONTEXT_PROMPT_MARKER in channel_prompt) @@ -1054,17 +985,13 @@ def _csv_or_list_to_set(raw: Any) -> set[str]: return set() if isinstance(raw, list): return {str(part).strip() for part in raw if str(part).strip()} - s = str(raw).strip() - if not s: - return set() - return {part.strip() for part in s.split(",") if part.strip()} + return {part.strip() for part in str(raw).split(",") if part.strip()} def _slack_ignored_channels_from_gateway_config(config: Any) -> set[str]: """Return Slack channels that the generic gateway must never dispatch. - Deliberately duplicates the adapter's first-line drop as a fail-safe: even if a code path or - test hook bypasses the adapter, ignored channels cannot reach auth, pairing or sessions.""" + Duplicates the adapter's drop as a fail-safe so bypasses can't reach auth, pairing or sessions.""" platform_cfg = getattr(config, "platforms", {}).get(Platform.SLACK) raw = None if platform_cfg is not None: @@ -1088,9 +1015,7 @@ def _is_slack_ignored_channel(config: Any, chat_id: Any) -> bool: def _message_timestamps_enabled(user_config: Optional[dict]) -> bool: - """True when gateway.message_timestamps.enabled is opted in. - - Default OFF: a timestamp prefix on every user message changes what the model sees.""" + """True when gateway.message_timestamps.enabled is opted in (default OFF: changes what the model sees).""" if not isinstance(user_config, dict): return False gw = user_config.get("gateway") @@ -1108,8 +1033,7 @@ def _build_gateway_agent_history( inject_timestamps: bool = False) -> tuple[List[Dict[str, Any]], Optional[str]]: """Convert stored gateway transcript rows into agent replay messages. - Keeping that context out of ``conversation_history`` stops consecutive-user repair merging it - with the live turn and hiding the current message behind ``history_offset`` on persistence.""" + Observed context stays out of ``conversation_history`` so consecutive-user repair can't merge it in.""" from hermes_time import get_timezone as _get_msg_tz from gateway.message_timestamps import render_user_content_with_timestamp as _render_msg_ts @@ -1136,8 +1060,7 @@ def _build_gateway_agent_history( clean_msg = {k: v for k, v in msg.items() if k not in {"timestamp", "observed"}} agent_history.append(clean_msg) elif content: - # Strip persisted auto-continue notes from user messages (interrupted turns): keep the - # user's real text but never replay the recovery instruction — it caused infinite loops. + # Strip persisted auto-continue notes: keep the real user text, never replay the recovery note. if role == "user": content = _strip_auto_continue_noise(content) if not content: @@ -1152,12 +1075,10 @@ def _build_gateway_agent_history( # Strip interrupted tool-call tails so the LLM doesn't re-execute tools killed mid-flight. agent_history = strip_interrupted_tool_tails(agent_history) - # Strip a dangling assistant(tool_calls) tail with no tool answers — the SIGKILL-mid-tool-call signature - # (e.g. the tool ran `docker restart`); else the model re-issues the unanswered call on resume, forever. + # Strip a dangling assistant(tool_calls) tail (SIGKILL-mid-tool-call); else the model re-issues it forever. agent_history = strip_dangling_tool_call_tail(agent_history) - # Strip expired dangerous-confirmation phrases (e.g. "confirm forced restart") from user text: - # replayed, an unrelated follow-up could read as a fresh confirmation and re-trigger the action. + # Strip expired dangerous-confirmation phrases; replayed, a follow-up could read as a fresh confirmation. agent_history = strip_stale_dangerous_confirmations(agent_history, now=time.time()) observed_context = "\n".join(observed_group_context).strip() or None @@ -1166,12 +1087,10 @@ def _build_gateway_agent_history( def _select_cached_agent_history( persisted_history: List[Dict[str, Any]], live_history: Any) -> List[Dict[str, Any]]: - """Prefer the cached live transcript only when it is longer AND has a real, non-ephemeral - unpersisted row; otherwise return ``persisted_history`` unchanged. + """Prefer the cached live transcript only when it is longer AND has a real, non-ephemeral unpersisted row. - Guards the FTS write-corruption case: silent write failures reload a stale ``conversation_history`` - while the cached ``AIAgent`` still holds unpersisted real rows (same-session amnesia). Length alone - is not enough: a longer all-durable list can be an expected replay-filtering delta.""" + Guards FTS write-corruption amnesia (stale reload while the cached agent holds unpersisted rows). Length + alone is not enough: a longer all-durable list can be an expected replay-filtering delta.""" if isinstance(live_history, list) and len(live_history) > len(persisted_history): from run_agent import _is_ephemeral_scaffolding @@ -1205,10 +1124,9 @@ def _wrap_current_message_with_observed_context(message: Any, observed_context: def _last_transcript_timestamp(history: Optional[List[Dict[str, Any]]]) -> Any: - """Return the ``timestamp`` of the last usable transcript row, if any. + """Return the ``timestamp`` of the last usable (non-metadata) transcript row, if any. - Skips metadata-only rows dropped before reaching the agent. ``None`` when no usable row has - a timestamp — callers treat that as "fresh" for backward compatibility.""" + ``None`` when the last usable row has no timestamp — callers treat that as "fresh" (legacy rows).""" if not history: return None for msg in reversed(history): @@ -1220,17 +1138,14 @@ def _last_transcript_timestamp(history: Optional[List[Dict[str, Any]]]) -> Any: ts = msg.get("timestamp") if ts is not None: return ts - # First non-meta row lacks a timestamp — legacy row; None lets the caller take the legacy-fresh path. return None return None -# Tool output may hold literal MEDIA: examples (docs, logs); only tools that intentionally create -# deliverable media are eligible for auto-append when the model omits them from the final reply. +# Tool output may hold literal MEDIA: examples (docs, logs); only deliberate media producers may auto-append. _AUTO_APPEND_MEDIA_TOOL_NAMES = {"text_to_speech", "text_to_speech_tool", "image_generate"} -# Replay-tail sanitization lives in agent/replay_cleanup.py so every resume surface (this messaging -# gateway AND the TUI/WebUI gateway) shares one implementation. +# Replay-tail sanitization lives in agent/replay_cleanup.py so every resume surface shares one implementation. from agent.replay_cleanup import ( # noqa: E402 strip_interrupted_tool_tails, strip_dangling_tool_call_tail, strip_stale_dangerous_confirmations) @@ -1240,18 +1155,13 @@ _AUTO_CONTINUE_FALLBACK_PREFIX = "[System note: A new message" def _is_auto_continue_noise(content: Any) -> bool: - """Return True if this user-message content is a gateway-injected - auto-continue note that should NOT be replayed as a real user turn.""" - if not isinstance(content, str): - return False - return content.startswith((_AUTO_CONTINUE_NOTE_PREFIX, _AUTO_CONTINUE_FALLBACK_PREFIX)) + """Return True if this user-message content is a gateway-injected auto-continue note (never replay it).""" + return isinstance(content, str) and content.startswith( + (_AUTO_CONTINUE_NOTE_PREFIX, _AUTO_CONTINUE_FALLBACK_PREFIX)) def _strip_auto_continue_noise(content: Any) -> Any: - """Strip one or more leading persisted auto-continue note prefixes from user text. - - A row may hold both the note and the user's real question; the trailing real text is preserved. - """ + """Strip leading persisted auto-continue notes from user text; the trailing real question is preserved.""" if not _is_auto_continue_noise(content): return content text = str(content) @@ -1262,13 +1172,11 @@ def _strip_auto_continue_noise(content: Any) -> Any: text = text[end + 1 :].lstrip() return text -# Tools whose deliverable is a JSON payload with a local-file path field rather than a literal -# ``MEDIA:`` tag (e.g. image_generate -> ``{"success": true, "image": "/abs/path.png"}``). +# Tools whose deliverable is a JSON payload with a local-file path field rather than a literal ``MEDIA:`` tag. _JSON_MEDIA_TOOL_PATH_FIELDS = ("host_image", "image", "agent_visible_image") -# Extension-anchored MEDIA: matcher for tool results. Mirrors the dispatch-site pattern so a bare -# ``MEDIA:`` token in prose (no deliverable extension) is never auto-appended. +# Extension-anchored MEDIA: matcher (mirrors the dispatch site); a bare ``MEDIA:`` in prose never auto-appends. _TOOL_MEDIA_RE = re.compile( r'MEDIA:((?:[A-Za-z]:[/\\]|/|~\/)\S+\.(?:png|jpe?g|gif|webp|' r'mp4|mov|avi|mkv|webm|ogg|opus|mp3|wav|m4a|' @@ -1277,8 +1185,7 @@ _TOOL_MEDIA_RE = re.compile( re.IGNORECASE) -# Shared with cron delivery and gateway background tasks — the repair must run on every surface -# that feeds a final response into media extraction; canonical names live in gateway.media_repair. +# Shared with cron delivery and gateway background tasks; canonical names live in gateway.media_repair. from gateway.media_repair import tool_name_by_call_id as _tool_name_by_call_id # noqa: E402 @@ -3206,8 +3113,7 @@ class GatewayRunner( GatewayAgentCacheMixin): """Main gateway controller: manages adapter lifecycles, routes messages to/from the agent.""" - # Class-level defaults so partial construction in tests doesn't - # blow up on attribute access. + # Class-level defaults so partial construction in tests doesn't blow up on attribute access. _busy_input_mode: str = "interrupt" _busy_text_mode: str = "interrupt" _restart_drain_timeout: float = DEFAULT_GATEWAY_RESTART_DRAIN_TIMEOUT @@ -3230,8 +3136,7 @@ class GatewayRunner( _startup_restore_in_progress: bool = False _startup_warmup_task: Optional[asyncio.Task] = None - # Legacy per-session dict adapters: all per-session state lives in ``self._sessions``; these expose - # the old dict attrs as LIVE MutableMapping views. New code: ``self._session_state(key)``. + # Legacy per-session dict attrs as LIVE views over ``self._sessions``; new code: _session_state(key) _running_agents = legacy_dict_property("_running_agents") _running_agents_ts = legacy_dict_property("_running_agents_ts") _active_session_leases = legacy_dict_property("_active_session_leases") @@ -3239,11 +3144,9 @@ class GatewayRunner( _turn_lease_tokens = legacy_lease_token_property() _session_run_generation = legacy_dict_property("_session_run_generation") _session_model_overrides = legacy_dict_property("_session_model_overrides") - _pending_one_turn_model_restores = legacy_dict_property( - "_pending_one_turn_model_restores") + _pending_one_turn_model_restores = legacy_dict_property("_pending_one_turn_model_restores") _session_reasoning_overrides = legacy_dict_property("_session_reasoning_overrides") - _session_service_tier_overrides = legacy_dict_property( - "_session_service_tier_overrides") + _session_service_tier_overrides = legacy_dict_property("_session_service_tier_overrides") _last_resolved_model = legacy_dict_property("_last_resolved_model") _queued_events = legacy_dict_property("_queued_events") _pending_turn_sidecar_notes = legacy_dict_property("_pending_turn_sidecar_notes") @@ -3275,9 +3178,7 @@ class GatewayRunner( def _peek_session_state(self, session_key: str) -> Optional["SessionState"]: """Return the SessionState for ``session_key`` without creating one.""" sessions = self.__dict__.get("_sessions") - if not sessions: - return None - return sessions.get(session_key) + return sessions.get(session_key) if sessions else None def _is_session_running(self, session_key: str) -> bool: """True when the session holds a running-turn slot (agent or sentinel).""" @@ -3285,12 +3186,9 @@ class GatewayRunner( return state is not None and state.turn.agent is not None def _running_agent_items(self) -> List[tuple]: - """(session_key, agent) pairs for sessions with a running turn (incl. pending sentinels), - matching the old ``_running_agents`` dict contents.""" - return [ - (key, state.turn.agent) - for key, state in self._sessions_map().items() - if state.turn.agent is not None] + """(session_key, agent) pairs for sessions with a running turn (incl. pending sentinels).""" + return [(key, state.turn.agent) for key, state in self._sessions_map().items() + if state.turn.agent is not None] # Loop-liveness / watchdog handles; class-level defaults so partially constructed test runners work. _loop_heartbeat_task: Optional["asyncio.Task"] = None _loop_floor_timer_handle: Optional[Any] = None @@ -3316,8 +3214,7 @@ class GatewayRunner( # Non-None means SessionDB init failed — the gateway broadcasts a one-time warning to the home # channel(s) after connecting so the user learns persistence is broken before /resume fails. self._session_db_init_error: Optional[str] = None - # Adapters for NON-default profiles, keyed by profile name then Platform; self.adapters stays - # the default profile's map so existing sites are untouched when multiplexing is off (empty). + # Non-default profiles' adapters by profile then Platform; self.adapters stays the default's map. self._profile_adapters: Dict[str, Dict[Platform, BasePlatformAdapter]] = {} self._warn_if_docker_media_delivery_is_risky() _gateway_runner_ref = _weakref.ref(self) @@ -3332,15 +3229,13 @@ class GatewayRunner( def _init_runtime_settings(self) -> None: """Load ephemeral per-call config (prefill, reasoning, busy modes, timeouts, routing).""" - # Ephemeral config (config.yaml / env): injected at API-call time only, never persisted. self._prefill_messages = self._load_prefill_messages() self._reasoning_config = self._load_reasoning_config() self._service_tier = self._load_service_tier() self._show_reasoning = self._load_show_reasoning() self._busy_input_mode = self._load_busy_input_mode() self._busy_text_mode = self._load_busy_text_mode() - # Secondary-profile busy modes, snapshotted at multiplex startup; busy-message handlers consult - # them by routed source without rereading config or mutating process-global environment. + # Secondary-profile busy modes snapshotted at multiplex startup; handlers never reread config. self._busy_input_modes_by_profile: Dict[str, str] = {} self._busy_text_modes_by_profile: Dict[str, str] = {} self._restart_drain_timeout = self._load_restart_drain_timeout() @@ -3355,8 +3250,7 @@ class GatewayRunner( # Reset guard: a background process older than session_reset.bg_process_max_age_hours (24h # default) is stale and no longer blocks idle/daily reset (NOT killed, only ignored). from tools.process_registry import process_registry - _bg_max_age_hours = getattr( - self.config.default_reset_policy, "bg_process_max_age_hours", 24) + _bg_max_age_hours = getattr(self.config.default_reset_policy, "bg_process_max_age_hours", 24) _bg_max_age_seconds = ( _bg_max_age_hours * 3600 if _bg_max_age_hours and _bg_max_age_hours > 0 else None) self.session_store = SessionStore( @@ -3430,9 +3324,8 @@ class GatewayRunner( def _init_runtime_caches(self) -> None: """Agent cache, profile identity, Teams runtime, failed-platform tracking, slash-confirm counter.""" - # AIAgent per session preserves prompt caching (fresh agent per message breaks the prefix cache, - # ~10x cost on Anthropic). Value: (AIAgent, config_signature). LRU cap in _enforce_agent_cache_cap, - # idle TTL from _session_expiry_watcher. + # AIAgent per session preserves prompt caching (fresh agent per message ~10x cost on Anthropic). + # Value: (AIAgent, config_signature); LRU cap in _enforce_agent_cache_cap, TTL in expiry watcher. self._agent_cache: "OrderedDict[str, tuple]" = OrderedDict() self._agent_cache_lock = threading.Lock() # Launch-time identity of the profile that owns ``self.adapters``; ``_authorization_adapter`` @@ -3453,9 +3346,8 @@ class GatewayRunner( def _init_startup_checks(self) -> None: """Ensure tirith is installed and warn when manual approvals have no automated assessor.""" def _ensure_tirith() -> None: - # Downloads if needed; non-fatal — fail-open at scan time if unavailable. from tools.tirith_security import ensure_installed - ensure_installed(log_failures=False) + ensure_installed(log_failures=False) # downloads if needed; fail-open at scan time _best_effort(_ensure_tirith) @@ -3488,7 +3380,6 @@ class GatewayRunner( self._session_db_handles: Dict[Path, Any] = {} self._session_db_handles_lock = threading.Lock() from gateway.session_db_recovery import RecoverableHandleCache - self._session_db_handle_cache = RecoverableHandleCache( handles=self._session_db_handles, lock=self._session_db_handles_lock) try: @@ -3496,8 +3387,7 @@ class GatewayRunner( except Exception as e: # WARNING (not DEBUG) so it lands in errors.log; else an NFS HERMES_HOME silently loses /resume etc. logger.warning("SQLite session store not available: %s", e) - # Surfaced on the home channel(s) once connected; otherwise persistence degrades silently. - self._session_db_init_error = str(e) + self._session_db_init_error = str(e) # surfaced on the home channel(s) once connected # Opportunistic state.db maintenance (prune + optional VACUUM), at most once per min_interval_hours. # A few blocking seconds per day is fine for a long-lived gateway; failures log, never raise. @@ -3566,15 +3456,11 @@ class GatewayRunner( self._scale_to_zero_no_suspend_logged: bool = False def _open_session_db_for_active_scope(self, raise_on_error: bool = False) -> Any: - """Return the AsyncSessionDB for the profile scope active on this task. - - Resolved per access (not in ``__init__``) because ``SessionDB()`` reads the context-local - HERMES_HOME, letting a multiplexed profile use its own store; one handle cached per path. - Construction failure enters bounded backoff; ``raise_on_error=True`` (priming) propagates it. - """ + """AsyncSessionDB for the active profile scope, resolved per access (not in ``__init__``) since + ``SessionDB()`` reads the context-local HERMES_HOME; one handle cached per path. Construction + failure enters bounded backoff; ``raise_on_error=True`` (priming) propagates it.""" from hermes_state import AsyncSessionDB, _default_db_path, get_shared_session_db from gateway.session_db_recovery import RecoverableHandleCache - path = Path(_default_db_path()) cache = getattr(self, "_session_db_handle_cache", None) if cache is None: @@ -3584,9 +3470,8 @@ class GatewayRunner( self._session_db_handle_cache = cache def _open(): - # Borrow the SessionStore's handle (same path) rather than opening a second one — otherwise two - # writer connections + two read pools per state.db. The store owns/sweeps it at shutdown; this - # cache holds only the async wrapper and cannot go stale (store drops handles only in close_all). + # Borrow the SessionStore's handle (same path) so state.db doesn't get two writers/pools. + # The store owns/sweeps it at shutdown; this cache holds only the async wrapper (close_all). store = getattr(self, "session_store", None) borrowed = getattr(store, "_db", None) if store is not None else None if borrowed is not None: @@ -3607,8 +3492,7 @@ class GatewayRunner( self._session_db_init_error = None logger.info("SQLite session store recovered") - return cache.get( - path, _open, raise_on_error=raise_on_error, on_recovered=_recovered) + return cache.get(path, _open, raise_on_error=raise_on_error, on_recovered=_recovered) @property def _session_db(self) -> Any: @@ -3626,8 +3510,7 @@ class GatewayRunner( """Close every per-profile AsyncSessionDB this runner opened. Drained under the lock, closed outside it; a pinned handle is the pinner's to close. Wrappers - BORROWED from ``session_store`` are not closed: the store's sweep (runs first) closes them. - """ + BORROWED from ``session_store`` are skipped: the store's sweep (runs first) closes them.""" def _close(db) -> None: if getattr(db, "__dict__", {}).get("_hermes_borrowed_handle"): return @@ -3670,7 +3553,6 @@ class GatewayRunner( gateway process, so model-emitted paths like `/output/report.txt` must be host-readable.""" if os.getenv("TERMINAL_ENV", "").strip().lower() != "docker": return - connected = self.config.get_connected_platforms() messaging_platforms = [p for p in connected if p not in {Platform.LOCAL, Platform.API_SERVER, Platform.WEBHOOK}] if not messaging_platforms: @@ -3686,19 +3568,10 @@ class GatewayRunner( except Exception: logger.debug("Could not parse TERMINAL_DOCKER_VOLUMES for gateway media warning", exc_info=True) - has_explicit_output_mount = False for spec in volumes: match = _DOCKER_VOLUME_SPEC_RE.match(spec) - if not match: - continue - container_path = match.group("container") - if container_path in _DOCKER_MEDIA_OUTPUT_CONTAINER_PATHS: - has_explicit_output_mount = True - break - - if has_explicit_output_mount: - return - + if match and match.group("container") in _DOCKER_MEDIA_OUTPUT_CONTAINER_PATHS: + return logger.warning( "Docker backend is enabled for the messaging gateway but no explicit host-visible " "output mount (for example '/home/user/.hermes/cache/documents:/output') is configured. " @@ -3741,24 +3614,17 @@ class GatewayRunner( # Telegram General topic in forum-enabled private chats: clients omit message_thread_id or send "1"; both = root. _TELEGRAM_GENERAL_TOPIC_IDS = frozenset({"", "1"}) - _TELEGRAM_LOBBY_REMINDER_COOLDOWN_S = 30.0 - def _normalize_source_for_session_key( - self, source: SessionSource) -> SessionSource: - """Apply Telegram DM topic recovery to a source for session-key purposes. - - ``_handle_message_with_agent`` rewrites ``thread_id`` before deriving the session key, so - handlers keying off the raw ``event.source`` would store overrides under a key the next turn - never reads. Always derive override storage keys from the result so storage and read match. - """ + def _normalize_source_for_session_key(self, source: SessionSource) -> SessionSource: + """Apply Telegram DM topic recovery to a source for session-key purposes. Always derive override + storage keys from the result: ``_handle_message_with_agent`` rewrites ``thread_id`` before + deriving the session key, so keys from the raw ``event.source`` are never read next turn.""" try: recovered = self._recover_telegram_topic_thread_id(source) except Exception: return source - if recovered is None: - return source - return dataclasses.replace(source, thread_id=recovered) + return source if recovered is None else dataclasses.replace(source, thread_id=recovered) def _resolve_session_key_or_none(self, source, session_key: Optional[str]) -> Optional[str]: """``session_key`` if given, else the key for ``source`` (None when it cannot be derived).""" @@ -3785,18 +3651,15 @@ class GatewayRunner( def _persist_active_agents(self) -> None: """Persist the live in-flight agent count to ``gateway_state.json`` at every turn boundary. - - Passes ONLY ``active_agents`` so the read-merge-write preserves lifecycle state (``gateway_state=None`` + Passes ONLY ``active_agents`` so the read-merge-write keeps lifecycle state (gateway_state=None would clobber it). Best-effort: a failed write must never disrupt a turn.""" _write_runtime_status_quiet(active_agents=self._active_work_count()) def _running_agent_ids(self) -> set: """``id()`` of every agent mid-turn — identity-keyed so the lookup is O(1) and independent of ``AIAgent.__eq__`` (MagicMock overrides it in tests).""" - return { - id(a) - for _, a in self._running_agent_items() - if a is not None and a is not _AGENT_PENDING_SENTINEL} + return {id(a) for _, a in self._running_agent_items() + if a is not None and a is not _AGENT_PENDING_SENTINEL} def _snapshot_running_agents(self) -> Dict[str, Any]: return {k: a for k, a in self._running_agent_items() if a is not _AGENT_PENDING_SENTINEL} @@ -3815,8 +3678,7 @@ class GatewayRunner( steered: bool redirected: bool - # Bound for off-loop agent-resource cleanup: _cleanup_agent_resources is synchronous and can block - # long (subprocess teardown, memory-provider IO); inline it wedges the loop, so it runs in a worker. + # Worker bound for _cleanup_agent_resources: sync, can block long (subprocess teardown, memory IO). _CLEANUP_TIMEOUT_S = 30.0 # Budget for one finalize_session() dispatch (plugin on_session_finalize hooks + Relay close): @@ -3831,14 +3693,11 @@ class GatewayRunner( _AUTO_RESUME_REASONS = frozenset({"restart_timeout", "shutdown_timeout", "restart_interrupted"}) _MAX_SUPERVISED_RESTARTS = 5 - # A task that ran at least this long before crashing is HEALTHY: an isolated crash, not a - # crash-loop; the consecutive-restart counter resets so a long-lived daemon isn't abandoned. + # Ran this long before crashing = HEALTHY (isolated crash, not a crash-loop); restart counter resets. _SUPERVISED_HEALTHY_SECS = 300 - # Slow respawn tier after the reconnect watcher exhausts its restart budget. Long on purpose: the - # watcher is crashing on contact, so "check back later" beats a tight loop. + # Slow respawn tier once the watcher's restart budget is spent; long on purpose (crashes on contact). _RECONNECT_WATCHER_SLOW_RETRY_SECS = 300 - # Slow-tier respawns to attempt while work is still queued. Bounded: if half an hour of - # five-minute retries cannot keep a watcher alive, the fault is not transient — fail loudly. + # Slow-tier respawns while work is queued; if 30 min of 5-min retries can't keep it up, fail loudly. _MAX_SLOW_WATCHER_RESPAWNS = 6 _TELEGRAM_CAPABILITY_HINT_COOLDOWN_S = 300.0 _APPROVAL_TIMEOUT_SECONDS = 300 # 5 minutes @@ -3904,9 +3763,7 @@ class GatewayRunner( return facade def _get_cached_session_source(self, session_key: str): - if not session_key: - return None - cached_sources = getattr(self, "_session_sources", None) + cached_sources = getattr(self, "_session_sources", None) if session_key else None if not cached_sources: return None source = cached_sources.get(session_key) @@ -3918,7 +3775,6 @@ class GatewayRunner( @dataclasses.dataclass class _HygieneSettings: """Resolved session-hygiene configuration for one inbound turn.""" - model: str threshold_pct: float compression_enabled: bool @@ -3937,7 +3793,6 @@ class GatewayRunner( class _HygieneAttempt: """One detached hygiene compression attempt. ``cleanup_deferred`` is shared mutable state: wait handlers set it on raise paths; the owning ``finally`` reads it to decide on cleanup now.""" - agent: Any meta: Any commit_fence: Any = None @@ -3954,9 +3809,8 @@ class GatewayRunner( getattr(source, "thread_id", None), chat_type=getattr(source, "chat_type", None), reply_to_message_id=reply_to_message_id or getattr(source, "message_id", None)) if getattr(source, "platform", None) == Platform.SLACK: - # Per-turn egress identity: Slack chat.startStream needs recipient_user_id/team_id, and the relay - # adapter's _with_scope fallback reads them from per-chat caches a CONCURRENT turn overwrites. - # Stamping authentic values here reduces the cache to a restart/synthetic-send fallback. + # Per-turn egress identity: Slack chat.startStream needs recipient_user_id/team_id; the relay + # adapter's _with_scope fallback reads per-chat caches a CONCURRENT turn overwrites. team_id = getattr(source, "scope_id", None) user_id = getattr(source, "user_id", None) if team_id or user_id: @@ -4034,8 +3888,7 @@ class GatewayRunner( from gateway.session_context import set_session_vars # Async-delivery capability tells async tools whether this channel can wake a later turn. Default # True keeps CLI/unknown paths working; stateless adapters (api_server) declare False. - _adapters = getattr(self, "adapters", None) or {} - _adapter = _adapters.get(context.source.platform) + _adapter = (getattr(self, "adapters", None) or {}).get(context.source.platform) _async_delivery = getattr(_adapter, "supports_async_delivery", True) return set_session_vars( platform=context.source.platform.value, @@ -4062,8 +3915,7 @@ class GatewayRunner( """Run blocking work in the thread pool while preserving session contextvars.""" loop = asyncio.get_running_loop() ctx = copy_context() - return await loop.run_in_executor( - self._get_executor(), ctx.run, func, *args) + return await loop.run_in_executor(self._get_executor(), ctx.run, func, *args) def _get_executor(self) -> concurrent.futures.ThreadPoolExecutor: """Return the gateway-owned executor for blocking agent work.""" @@ -4071,7 +3923,6 @@ class GatewayRunner( if lock is None: lock = threading.Lock() self._executor_lock = lock - with lock: if getattr(self, "_executor_closing", False): raise RuntimeError("Gateway is shutting down; executor unavailable") @@ -4084,23 +3935,18 @@ class GatewayRunner( def _shutdown_executor(self, drain_timeout: float = 0.0) -> int: """Stop the gateway-owned executor; returns the number of worker threads still running. - ``drain_timeout=0`` is fire-and-forget; shutdown passes a bounded budget so blocking DB work - cannot outlive ``SessionDB.close()``. ``cancel_futures`` only drops unstarted work and a cancelled - ``run_in_executor`` awaitable does not stop its thread, so running workers are joined explicitly. - """ + cannot outlive ``SessionDB.close()``. ``cancel_futures`` only drops unstarted work and cancelling + a ``run_in_executor`` awaitable does not stop its thread, so running workers are joined.""" lock = getattr(self, "_executor_lock", None) if lock is None: return 0 - with lock: self._executor_closing = True executor = getattr(self, "_executor", None) self._executor = None - if executor is None: return 0 - try: executor.shutdown(wait=False, cancel_futures=True) except TypeError: @@ -4144,13 +3990,11 @@ class GatewayRunner( @staticmethod def _init_cached_agent_for_turn(agent: Any, interrupt_depth: int) -> None: """Reset per-turn state on a cached agent before a new turn starts. - The activity ts/desc/provenance triple resets together and only at depth 0 — else a session idle 29 min trips the watchdog before the first call; interrupt-recursive turns keep it so stuck-turn idle time accumulates to the 30-min timeout.""" if interrupt_depth == 0: from agent.session_activity import ActivityProvenance - agent._last_activity_ts = time.time() agent._last_activity_desc = "starting new turn (cached)" agent._last_activity_provenance = ActivityProvenance.UNKNOWN @@ -4162,10 +4006,9 @@ class GatewayRunner( def _profile_name_for_source(self, source: SessionSource) -> Optional[str]: """Resolve the profile name for an inbound source via configured routes (most specific wins). - - ``None`` = default/active profile. Gated on ``multiplex_profiles``: the scoped run only activates - under multiplexing, else keys would be profile-namespaced while the agent ran in ``agent:main``. - """ + ``None`` = default/active profile. Gated on ``multiplex_profiles``, since the scoped run only + activates under multiplexing; otherwise keys would be profile-namespaced while the agent ran in + ``agent:main``.""" config = getattr(self, "config", None) if not getattr(config, "multiplex_profiles", False): return None @@ -4189,14 +4032,12 @@ class GatewayRunner( except Exception as exc: logger.warning( "Rejecting profile route %r because the served-profile set could not be resolved", - matched.name, - exc_info=True) + matched.name, exc_info=True) raise ProfileRouteRejected(matched.name) from exc if matched.profile not in served: logger.warning( "Rejecting profile route %r: target profile %r is not served", - matched.name, - matched.profile) + matched.name, matched.profile) raise ProfileRouteRejected(matched.name) return matched.profile logger.debug( @@ -4209,25 +4050,20 @@ class GatewayRunner( """Resolve which profile's HERMES_HOME serves this source: ``source.profile``, then ``_profile_name_for_source`` (sources bypassing ``build_source``), then the active profile.""" from gateway.profile_routing import ProfileRouteRejected - from hermes_cli.profiles import ( - get_active_profile_name, get_profile_dir, profile_exists) + from hermes_cli.profiles import get_active_profile_name, get_profile_dir, profile_exists from hermes_constants import get_hermes_home - explicit_profile = None # explicitly requested (source or routing) vs. default fallback try: name = (source.profile or "").strip() or self._profile_name_for_source(source) explicit_profile = name or None if not name: name = get_active_profile_name() or "default" - profile_dir = get_profile_dir(name) if explicit_profile and not profile_exists(name): logger.warning( "Profile %r does not exist for source %s/%s (guild_id=%s), " "falling back to global HERMES_HOME", - explicit_profile, - source.platform.value, - source.chat_id, + explicit_profile, source.platform.value, source.chat_id, getattr(source, "guild_id", None)) return get_hermes_home() return profile_dir @@ -4237,17 +4073,13 @@ class GatewayRunner( logger.warning( "Failed to resolve profile directory for source %s/%s (guild_id=%s), " "falling back to global HERMES_HOME: %s", - source.platform.value, - source.chat_id, - getattr(source, "guild_id", None), - explicit_profile or "(no profile)", - exc_info=True) + source.platform.value, source.chat_id, getattr(source, "guild_id", None), + explicit_profile or "(no profile)", exc_info=True) return get_hermes_home() @dataclasses.dataclass class _RunAgentDisplay: """Per-turn display / progress settings resolved by ``_run_agent_display_settings``.""" - user_config: Any = None platform_key: Any = None enabled_toolsets: Any = None @@ -4270,7 +4102,6 @@ class GatewayRunner( @dataclasses.dataclass class _RunAgentWorker: """Executor future + inactivity-watchdog handles for one ``_run_agent_inner`` turn.""" - executor_task: Any = None agent_timeout: Optional[float] = None agent_warning: Optional[float] = None @@ -4285,12 +4116,9 @@ class GatewayRunner( def _run_planned_stop_watcher( stop_event: threading.Event, runner, loop: asyncio.AbstractEventLoop, shutdown_handler, *, poll_interval: float = 0.5) -> None: - """Poll for the planned-stop marker and trigger graceful shutdown. - - On Windows ``add_signal_handler`` is unavailable, so ``hermes gateway stop`` would never drain; - this cheap watcher (runs everywhere) turns the marker into the same shutdown-handler call a SIGTERM - would. On POSIX the signal handler consumes the marker first; ``_running``/``_draining`` guard re-triggers. - """ + """Poll for the planned-stop marker and trigger graceful shutdown (Windows lacks + ``add_signal_handler``, so ``hermes gateway stop`` would never drain). Runs everywhere; on POSIX + the signal handler consumes the marker first and ``_running``/``_draining`` guard re-triggers.""" from gateway.status import ( _get_planned_stop_marker_path, planned_stop_marker_targets_self) marker_path = _get_planned_stop_marker_path() @@ -4300,9 +4128,8 @@ def _run_planned_stop_watcher( marker_path.exists() and not getattr(runner, "_draining", False) and getattr(runner, "_running", False)): - # A marker may target a PREVIOUS instance (different PID) that exited before stop() cleaned - # up; firing on it means an "UNKNOWN" exit and a watchdog crash-loop. The probe unlinks - # stale/malformed markers. + # A marker may target a PREVIOUS instance that exited before stop() cleaned up; + # firing on it means an "UNKNOWN" exit and a watchdog crash-loop; probe unlinks stale. if not planned_stop_marker_targets_self(): stop_event.wait(poll_interval) continue @@ -4427,11 +4254,9 @@ def _housekeeping_memory_trim() -> None: def _start_gateway_housekeeping(stop_event: threading.Event, adapters=None, loop=None, interval: int = 60, cron_provider=None): - """Background thread for gateway-only periodic chores (NOT cron). - - Separate from the cron trigger so chores run regardless of ``CronScheduler`` provider (an external - scale-to-zero provider has no 60s loop). Cadences are in ticks of ``interval``; inner gates own the - real cadence.""" + """Background thread for gateway-only periodic chores (NOT cron). Separate from the cron trigger + so chores run under any ``CronScheduler`` provider (external scale-to-zero has no 60s loop). + Cadences are ticks of ``interval``; inner gates own the real cadence.""" chores: list[tuple[int, str, Any]] = [ (5, "Channel directory refresh", lambda: adapters and _housekeeping_channel_directory(adapters, loop)), (60, "Media cache cleanup", _housekeeping_media_caches), @@ -4475,22 +4300,18 @@ def _stop_cron_provider(provider) -> None: logger.debug("Cron provider stop() error: %s", exc) -# Upper bound for cooperatively draining the cron ticker on shutdown: the cron thread blocks on -# ``future.result(timeout=60)`` (cron/scheduler.py::_deliver_result), so a delivery unblocks in ~60s. +# Cron thread blocks on future.result(timeout=60) (cron/scheduler.py::_deliver_result) + margin. _CRON_SHUTDOWN_DRAIN_TIMEOUT = 65.0 -# Upper bound for draining the housekeeping ticker on shutdown: the channel-directory refresh blocks -# on ``fut.result(timeout=30)``, so cover that 30s plus margin or an in-flight refresh is abandoned. +# Housekeeping's channel-directory refresh blocks on fut.result(timeout=30); cover that + margin. _HOUSEKEEPING_SHUTDOWN_DRAIN_TIMEOUT = 35.0 async def _await_thread_exit( thread: Optional[threading.Thread], timeout: float, poll: float = 0.1) -> bool: """Wait for a daemon thread to exit WITHOUT blocking the event loop; True if it exited in time. - - A synchronous ``join()`` freezes the loop — fatal for the cron ticker, whose in-flight delivery is - a coroutine scheduled onto *this* loop: it could never run, so the join timed out and the message dropped. - """ + A synchronous ``join()`` freezes the loop — fatal for the cron ticker, whose in-flight delivery is a + coroutine on *this* loop: it could never run, so the join timed out and the message dropped.""" if thread is None: return True deadline = asyncio.get_running_loop().time() + max(0.0, timeout) @@ -4501,15 +4322,12 @@ async def _await_thread_exit( async def _shutdown_mcp_servers_nonblocking(timeout: float = 5.0) -> bool: """Close MCP servers off-loop with a bounded wait; True when done within ``timeout``. - - ``shutdown_mcp_servers()`` can block ~15s; on the loop thread that lets short-grace supervisors (s6 - 3s) SIGKILL us before ``mark_exited()`` runs, so every later boot reports a phantom unclean death. + ``shutdown_mcp_servers()`` can block ~15s; on the loop thread short-grace supervisors (s6 3s) + SIGKILL us before ``mark_exited()`` runs, so every later boot reports a phantom unclean death. On timeout shutdown proceeds and the daemon thread is left to finish or die.""" - def _do() -> None: try: from tools.mcp_tool import shutdown_mcp_servers - shutdown_mcp_servers() except Exception: logger.debug("MCP shutdown raised", exc_info=True) @@ -4520,8 +4338,7 @@ async def _shutdown_mcp_servers_nonblocking(timeout: float = 5.0) -> bool: if not done: logger.warning( "MCP shutdown did not finish within %.1fs; continuing gateway " - "teardown (background thread will be reaped at process exit)", - timeout) + "teardown (background thread will be reaped at process exit)", timeout) return done @@ -4540,31 +4357,26 @@ def _shutdown_gateway_health_export(runner: Any) -> None: def _gateway_stderr_formatter() -> logging.Formatter: """Return the redacting formatter used by the gateway stderr stream.""" from agent.redact import RedactingFormatter - return RedactingFormatter("%(asctime)s %(levelname)s %(name)s: %(message)s") def _replace_target_belongs_to_other_profile(existing_pid: int) -> bool: """Return True when ``--replace`` must refuse to signal ``existing_pid``. - A poisoned/stale PID record can point at another profile's LIVE gateway (cross-profile SIGTERM - restart loop). Ownership is decided by the persisted identity record ALONE, only while bound to the - live target by exact PID + start-time. Live argv can never PROVE ownership (no HERMES_HOME); it is - only a consistency check. Missing, legacy, conflicting or unprovable identity → refuse (fail closed). - """ + restart loop). Ownership is decided by the persisted identity record ALONE, bound to the live target + by exact PID + start-time; live argv can never PROVE ownership (no HERMES_HOME), it is only a + consistency check. Missing, legacy, conflicting or unprovable identity → refuse (fail closed).""" try: from gateway.status import ( _get_pid_path, _get_process_hermes_home, _get_process_start_time, _pid_from_record, _read_pid_record, _record_looks_like_gateway, _read_process_cmdline, _same_hermes_home) - our_home = _get_process_hermes_home() def refuse(msg: str, *args, level=logging.WARNING) -> bool: logger.log(level, "Refusing --replace: " + msg, *args) return True - # Authorize from the persisted identity record — bound claim: the record must describe THIS pid - # with THIS live start time, otherwise it is stale/poisoned and proves nothing. + # Bound claim: the record must name THIS pid with THIS live start time, else it proves nothing. record = _read_pid_record(_get_pid_path()) if not isinstance(record, dict) or not _record_looks_like_gateway(record): return refuse("no valid gateway pid record to prove ownership of PID %s.", existing_pid) @@ -4585,8 +4397,7 @@ def _replace_target_belongs_to_other_profile(existing_pid: int) -> bool: return refuse("pid record belongs to a different HERMES_HOME (%s, ours %s). Remove the stale PID " "record or stop the owning profile explicitly.", recorded_home, our_home, level=logging.ERROR) - # Argv consistency check (never authority): an explicit profile flag / HERMES_HOME= contradicting - # our home refuses even if the record agreed; bare argv adds nothing; on probe failure the record decides. + # Argv never proves ownership; an explicit contradicting --profile / HERMES_HOME= still refuses. live_cmdline = _best_effort(lambda: _read_process_cmdline(existing_pid)) if live_cmdline and _looks_like_profile_conflict_from_cmdline(live_cmdline, our_home): return refuse("target PID %s command line explicitly advertises a different profile than " @@ -4594,20 +4405,16 @@ def _replace_target_belongs_to_other_profile(existing_pid: int) -> bool: return False except Exception: # Destructive action + unknown ownership => fail closed. - logger.warning( - "cross-profile --replace ownership probe failed for PID %s; refusing to signal", - existing_pid, - exc_info=True) + logger.warning("cross-profile --replace ownership probe failed for PID %s; refusing to signal", + existing_pid, exc_info=True) return True def _looks_like_profile_conflict_from_cmdline(command: str, our_home) -> bool: """Token-exact contradiction check between a target argv and our home (authority is the pid record). - - Substring matching is not identity: ``--profile timothy`` must NOT read as profile ``tim``. - Returns False whenever the argv does not clearly contradict our home.""" + Substring matching is not identity: ``--profile timothy`` must NOT read as profile ``tim``. Returns + False whenever the argv does not clearly contradict our home.""" from gateway.status import _profile_name_for_home - profile_name = _profile_name_for_home(our_home) try: tokens = shlex.split(command) @@ -4674,7 +4481,6 @@ async def _wait_for_pid_exit(pid: int, attempts: int, delay: float) -> bool: async def _start_gateway_replace_existing_instance(existing_pid: int, replace: bool) -> bool: """Handle a live gateway PID under this HERMES_HOME: replace it (``--replace``) or refuse. - Returns False when startup must abort (refused, permission denied, target still alive).""" from gateway.status import get_process_start_time, remove_pid_file, terminate_pid if not replace: @@ -4690,29 +4496,24 @@ async def _start_gateway_replace_existing_instance(existing_pid: int, replace: b f" Or use 'hermes gateway run --replace' to auto-replace.\n") return False - # Never signal a live process we cannot prove belongs to this HERMES_HOME: a poisoned PID record - # steering --replace at another profile's gateway is exactly the restart-loop shape to avoid. + # Never signal a process not provably ours (a poisoned PID record → cross-profile restart loop). if _replace_target_belongs_to_other_profile(existing_pid): from gateway.status import _get_process_hermes_home - logger.error( "Refusing --replace: PID %d cannot be proven to belong " "to this profile's gateway (HERMES_HOME %s). Remove the " "stale PID record or stop the owning profile explicitly.", - existing_pid, - _get_process_hermes_home()) + existing_pid, _get_process_hermes_home()) return False existing_start_time = get_process_start_time(existing_pid) logger.info("Replacing existing gateway instance (PID %d) with --replace.", existing_pid) - # Takeover marker: the target's shutdown handler recognises its SIGTERM as a planned takeover and - # exits 0 (exit 1 would trigger systemd's Restart=on-failure and a flap loop against us). + # Takeover marker: target exits 0 on our SIGTERM (exit 1 → systemd Restart=on-failure flap loop). try: from gateway.status import write_takeover_marker write_takeover_marker(existing_pid) except Exception as e: logger.debug("Could not write takeover marker: %s", e) - # Snapshot children BEFORE signalling: once it exits, orphans are reparented and invisible to a parent - # walk; surviving adapter subprocesses hold scoped token locks and block us (Windows already tree-kills). + # Snapshot children BEFORE signalling: reparented orphans are invisible yet hold scoped token locks. try: from gateway.status import _snapshot_gateway_children _old_gateway_children = _snapshot_gateway_children(existing_pid) @@ -4736,13 +4537,11 @@ async def _start_gateway_replace_existing_instance(existing_pid: int, replace: b old_gateway_exited = True except (PermissionError, OSError): pass - # Confirm SIGKILL actually took (uninterruptible sleep, zombie) before clearing PID file / scoped - # locks, or two live gateways fight over the same token. + # Confirm SIGKILL took (D-state/zombie) before clearing PID/locks, or two gateways share a token. if not old_gateway_exited and not await _wait_for_pid_exit(existing_pid, 20, 0.25): logger.error( "Old gateway (PID %d) still appears alive after SIGKILL; " - "aborting replacement to avoid a duplicate gateway.", - existing_pid) + "aborting replacement to avoid a duplicate gateway.", existing_pid) _clear_takeover_marker_quiet() return False # Reap orphaned children (POSIX; mirrors Windows taskkill /T) so they stop holding scoped token locks. @@ -4776,8 +4575,7 @@ def _start_gateway_configure_logging(verbosity: Optional[int]) -> None: _best_effort(_sync_skills) - # Centralized logging — agent.log (INFO+), errors.log (WARNING+), gateway.log (INFO+, gateway - # records only). Idempotent, so repeated calls from AIAgent.__init__ don't duplicate. + # Centralized logging (agent.log INFO+, errors.log WARNING+, gateway.log gateway-only); idempotent. from hermes_logging import setup_logging, _safe_stderr setup_logging(hermes_home=_hermes_home, mode="gateway") @@ -4808,20 +4606,17 @@ def _start_gateway_configure_logging(verbosity: Optional[int]) -> None: def _start_gateway_make_shutdown_signal_handler(runner, _signal_initiated_shutdown: list): """Build the SIGINT/SIGTERM handler; ``_signal_initiated_shutdown[0]`` records an unplanned signal.""" def shutdown_signal_handler(received_signal=None): - # Planned --replace takeover (sibling wrote a marker naming this PID before SIGTERM): exit 0 so - # systemd's Restart=on-failure doesn't revive us to flap-fight the replacer. + # Planned --replace takeover (sibling marked this PID): exit 0 so systemd won't revive us. def _takeover() -> bool: from gateway.status import consume_takeover_marker_for_self return consume_takeover_marker_for_self() - # Planned stop: service managers and `hermes gateway stop` also send SIGTERM, indistinguishable - # from an external kill unless the CLI marks it first. SIGINT is an interactive Ctrl+C stop. + # Planned stop: CLI marks first, else its SIGTERM looks like an external kill. SIGINT = Ctrl+C. def _planned_stop() -> bool: from gateway.status import consume_planned_stop_marker_for_self return consume_planned_stop_marker_for_self() - # Fast (<10ms) snapshot of who's asking us to shut down — runs synchronously inside the asyncio - # signal handler: stdlib + /proc only, no subprocesses (a sync `ps aux` here once blocked ~3s). + # Fast (<10ms) sync snapshot: stdlib + /proc, no subprocesses (`ps aux` here once blocked ~3s). def _snapshot(): from gateway.shutdown_forensics import snapshot_shutdown_context return snapshot_shutdown_context(received_signal) @@ -4869,20 +4664,17 @@ def _start_gateway_claim_pid_file() -> bool: remove_pid_file, write_pid_file) _current_pid = get_running_pid() if _current_pid is not None and _current_pid != os.getpid(): - logger.error( - "Another gateway instance (PID %d) started during our startup. " - "Exiting to avoid double-running.", _current_pid) + logger.error("Another gateway instance (PID %d) started during our startup. " + "Exiting to avoid double-running.", _current_pid) return False if not acquire_gateway_runtime_lock(): - logger.error( - "Gateway runtime lock is already held by another instance. Exiting.") + logger.error("Gateway runtime lock is already held by another instance. Exiting.") return False try: write_pid_file() except FileExistsError: release_gateway_runtime_lock() - logger.error( - "PID file race lost to another gateway instance. Exiting.") + logger.error("PID file race lost to another gateway instance. Exiting.") return False atexit.register(remove_pid_file) atexit.register(release_gateway_runtime_lock) @@ -4895,16 +4687,13 @@ async def _start_gateway_start_control_socket(runner): _control_server = None try: from gateway.control_socket import GatewayControlServer - - # pause-for-update: the updater asks us to drain and exit cleanly (releasing venv file handles) - # instead of being tree-killed — same path as SIGUSR1/service restarts. The handler runs on the - # socket's executor thread, so the request is marshalled onto the loop; the ACK returns the budget. + # pause-for-update: the updater asks us to drain + exit (freeing venv handles) vs. a tree-kill + # (same path as SIGUSR1). Handler runs on the socket executor thread, so marshal onto the loop. _main_loop = asyncio.get_running_loop() def _pause_for_update_handler() -> dict: try: from hermes_cli.gateway import _get_restart_drain_timeout - _drain = float(_get_restart_drain_timeout()) except Exception: _drain = 30.0 @@ -4913,8 +4702,7 @@ async def _start_gateway_start_control_socket(runner): def _request() -> None: try: - accepted_box.append( - runner.request_restart(detached=False, via_service=True)) + accepted_box.append(runner.request_restart(detached=False, via_service=True)) finally: _done.set() @@ -4922,10 +4710,8 @@ async def _start_gateway_start_control_socket(runner): _done.wait(timeout=5.0) accepted = bool(accepted_box and accepted_box[0]) return { - "pausing": accepted, - "already_stopping": not accepted, - "pid": os.getpid(), - "drain_timeout": _drain} + "pausing": accepted, "already_stopping": not accepted, + "pid": os.getpid(), "drain_timeout": _drain} _control_server = GatewayControlServer( verb_handlers={"pause-for-update": _pause_for_update_handler}) @@ -4940,9 +4726,8 @@ async def _start_gateway_start_control_socket(runner): def _start_gateway_start_cron_and_housekeeping(runner): - """Start the cron scheduler thread + gateway housekeeping thread. - - Returns ``(cron_stop, cron_provider, cron_thread, housekeeping_thread)``.""" + """Start the cron scheduler thread + gateway housekeeping thread; returns + ``(cron_stop, cron_provider, cron_thread, housekeeping_thread)``.""" # The event loop is passed so cron delivery can use live adapters (E2EE support). from cron.scheduler_provider import ( InProcessCronScheduler, resolve_cron_scheduler, scheduler_for_profile_mode) @@ -4952,8 +4737,7 @@ def _start_gateway_start_cron_and_housekeeping(runner): resolve_cron_scheduler(), multiplex_profiles=multiplex_cron) cron_start_kwargs: Dict[str, Any] = {"adapters": runner.adapters, "loop": asyncio.get_running_loop()} - # Multiplex: tell the built-in ticker which profile homes to tick, else secondary profiles' cron jobs - # show as "scheduled" but never execute because no ticker owns that store. + # Multiplex: tell the ticker which profile homes to tick, else secondary profiles' jobs never run. if isinstance(cron_provider, InProcessCronScheduler) and multiplex_cron: try: profile_homes = _multiplex_profile_homes(runner.config) @@ -4965,15 +4749,12 @@ def _start_gateway_start_cron_and_housekeeping(runner): # cron through the default bot (even before that profile's adapter connects). cron_start_kwargs["default_profile"] = "default" logger.info( - "Cron scheduler will tick %d profile(s) under multiplex: %s", - len(profile_homes), + "Cron scheduler will tick %d profile(s) under multiplex: %s", len(profile_homes), [p[0] if isinstance(p, tuple) else p for p in profile_homes]) except Exception as exc: - logger.warning( - "Could not resolve profile homes for multiplex cron: %s", exc) + logger.warning("Could not resolve profile homes for multiplex cron: %s", exc) - # External cron providers own their remote scheduling contract; only the in-process ticker polls - # local due jobs, so only it receives the local external-drain dispatch gate. + # Only the in-process ticker polls local due jobs, so only it gets the external-drain dispatch gate. if isinstance(cron_provider, InProcessCronScheduler): cron_start_kwargs["can_dispatch"] = lambda: not ( runner._draining or runner._external_drain_active) @@ -5002,14 +4783,10 @@ def _start_gateway_start_cron_and_housekeeping(runner): # Gateway-only housekeeping runs independently of the cron provider; shares cron_stop for shutdown. housekeeping_thread = threading.Thread( - target=_start_gateway_housekeeping, - args=(cron_stop,), - kwargs={ - "adapters": runner.adapters, - "loop": asyncio.get_running_loop(), - "cron_provider": cron_provider}, - daemon=True, - name="gateway-housekeeping") + target=_start_gateway_housekeeping, args=(cron_stop,), + kwargs={"adapters": runner.adapters, "loop": asyncio.get_running_loop(), + "cron_provider": cron_provider}, + daemon=True, name="gateway-housekeeping") housekeeping_thread.start() return cron_stop, cron_provider, cron_thread, housekeeping_thread @@ -5045,16 +4822,13 @@ async def _start_gateway_shutdown_tail( if _exit_with_failure_verdict(runner): return False - # Cooperative stop, never join()ed: an in-flight cron delivery is a coroutine on THIS loop while the - # ticker blocks on future.result(); a sync join would starve it and drop the message. + # Never join(): an in-flight cron delivery is a coroutine on THIS loop; a sync join would drop it. cron_stop.set() _stop_cron_provider(cron_provider) if not await _await_thread_exit(cron_thread, timeout=_CRON_SHUTDOWN_DRAIN_TIMEOUT): - logger.warning( - "Cron ticker did not exit within %.0fs of shutdown — an in-flight " - "delivery may have been dropped.", _CRON_SHUTDOWN_DRAIN_TIMEOUT) - await _await_thread_exit( - housekeeping_thread, timeout=_HOUSEKEEPING_SHUTDOWN_DRAIN_TIMEOUT) + logger.warning("Cron ticker did not exit within %.0fs of shutdown — an in-flight " + "delivery may have been dropped.", _CRON_SHUTDOWN_DRAIN_TIMEOUT) + await _await_thread_exit(housekeeping_thread, timeout=_HOUSEKEEPING_SHUTDOWN_DRAIN_TIMEOUT) # Stop the planned-stop watcher (daemon=True so this is belt-and-suspenders). _planned_stop_watcher_stop.set() @@ -5066,19 +4840,16 @@ async def _start_gateway_shutdown_tail( if runner.exit_code is not None: raise SystemExit(runner.exit_code) - # An unexpected SIGTERM that wasn't a planned restart exits non-zero so systemd's Restart=on-failure - # revives the process; `hermes gateway stop` / Ctrl+C are planned stops and must not trigger revival. + # Unplanned SIGTERM exits non-zero so systemd Restart=on-failure revives us; planned stops must not. if _signal_initiated_shutdown[0] and not runner._restart_requested: - logger.info( - "Exiting with code 1 (signal-initiated shutdown without restart " - "request) so systemd Restart=on-failure can revive the gateway.") + logger.info("Exiting with code 1 (signal-initiated shutdown without restart " + "request) so systemd Restart=on-failure can revive the gateway.") return False # → sys.exit(1) in the caller # Older restart paths may reach here without ``runner.exit_code``; keep the non-zero fallback. if runner._restart_via_service: - logger.info( - "Exiting with code 75 (service-restart requested) so the service " - "manager relaunches the gateway.") + logger.info("Exiting with code 75 (service-restart requested) so the service " + "manager relaunches the gateway.") raise SystemExit(75) return True @@ -5087,26 +4858,21 @@ async def _start_gateway_shutdown_tail( async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = False, verbosity: Optional[int] = 0) -> bool: """Start the gateway and run until interrupted; False if it failed to start (non-zero exit so systemd can auto-restart). ``replace`` kills any existing instance first (avoids restart-loop deadlocks).""" - # Enable interactive exec approval on messaging platforms. Set here (not at module import) so - # incidental imports of gateway.run from CLI/tool code don't poison HERMES_EXEC_ASK. + # Set here (not at import) so incidental gateway.run imports from CLI code don't poison it. os.environ["HERMES_EXEC_ASK"] = "1" from hermes_cli.resource_limits import apply_nofile_soft_limit - apply_nofile_soft_limit() - # Snapshot the checkout revision while sys.modules still matches disk, so a later `git pull` under - # this long-lived process is detected (risky work refused) instead of crashing on stale modules. + # Snapshot the revision while sys.modules matches disk so a later `git pull` is detected safely. from gateway.code_skew import record_boot_fingerprint record_boot_fingerprint() - # Duplicate-instance guard: no two gateways under one HERMES_HOME. The PID file is scoped to - # HERMES_HOME, so multi-profile setups (distinct HERMES_HOME each) run concurrently untripped. + # Duplicate-instance guard scoped to HERMES_HOME; distinct-home multi-profile setups coexist. from gateway.status import get_running_pid existing_pid = get_running_pid() - if ( - existing_pid is not None and existing_pid != os.getpid() - and not await _start_gateway_replace_existing_instance(existing_pid, replace)): + if (existing_pid is not None and existing_pid != os.getpid() + and not await _start_gateway_replace_existing_instance(existing_pid, replace)): return False _start_gateway_configure_logging(verbosity) @@ -5119,8 +4885,7 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = # it to cold adapter connects and clears it before the background reconnect watcher starts. runner._platform_lock_takeover_on_start = bool(replace) - # Track whether an unexpected signal initiated shutdown: an unexpected SIGTERM exits non-zero so - # service managers revive us; planned stop paths write a marker first so they exit cleanly. + # Unexpected signals exit non-zero so service managers revive us; planned stops write a marker first. _signal_initiated_shutdown = [False] shutdown_signal_handler = _start_gateway_make_shutdown_signal_handler( @@ -5131,8 +4896,7 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = loop = asyncio.get_running_loop() - # Swallows transient network errors from background tasks (an unhandled telegram/httpx error in any - # awaited coroutine would kill the gateway). Deliberately narrow — everything else hits the default. + # Swallow transient network errors from background tasks; one unhandled httpx error would kill us. loop.set_exception_handler(_gateway_loop_exception_handler) if threading.current_thread() is threading.main_thread(): @@ -5146,9 +4910,8 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = else: logger.info("Skipping signal handlers (not running in main thread).") - # Windows has no add_signal_handler, so `hermes gateway stop`'s SIGTERM would never drain; this thread - # polls the planned-stop marker written BEFORE the kill and drives the same shutdown path. Runs - # everywhere (cheap) so environments masking SIGTERM still drain cleanly. + # Windows has no add_signal_handler, so `hermes gateway stop`'s SIGTERM would never drain; poll the + # planned-stop marker (written BEFORE the kill) instead. Runs everywhere so masked-SIGTERM drains. _planned_stop_watcher_stop = threading.Event() _planned_stop_watcher_thread = threading.Thread( target=_run_planned_stop_watcher, @@ -5156,13 +4919,11 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = name="planned-stop-watcher") _planned_stop_watcher_thread.start() - # Claim the PID file BEFORE any adapters: two concurrent `run --replace` both pass the termination - # wait, but only the O_EXCL winner ever opens Telegram polling, Discord sockets, etc. + # PID file BEFORE adapters: of two concurrent `run --replace`, only the O_EXCL winner opens sockets. if not _start_gateway_claim_pid_file(): return False - # Control socket right after the PID-file claim (winning that race makes us authoritative). Non-fatal: - # a bind failure leaves consumers on the process-scan/state-file layer. + # Right after the PID claim (which makes us authoritative); non-fatal — consumers fall back to scan. _control_server = await _start_gateway_start_control_socket(runner) def _lifecycle_record_startup() -> None: @@ -5179,8 +4940,7 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = _best_effort(_start_keepalive, "Nous auth keepalive did not start: %s") _ensure_windows_gateway_venv_imports() - # MCP discovery in an executor: discover_mcp_tools() blocks up to 120s, which on the loop thread would - # freeze platform heartbeats (Discord shard, Telegram polling). + # discover_mcp_tools() blocks up to 120s; on the loop thread it would freeze platform heartbeats. try: await _discover_gateway_mcp_tools(runner.config) except Exception as e: @@ -5206,8 +4966,7 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = _shutdown_gateway_health_export(runner) if runner.exit_reason: logger.error("Gateway exiting cleanly: %s", runner.exit_reason) - # An explicit exit code (e.g. GATEWAY_FATAL_CONFIG_EXIT_CODE) must propagate so the s6 finish - # script can translate it (78 → 125) and stop the restart loop instead of exiting 0. + # Explicit exit codes (GATEWAY_FATAL_CONFIG_EXIT_CODE) must propagate so s6 finish maps 78 → 125. if runner.exit_code is not None: raise SystemExit(runner.exit_code) return True @@ -5228,8 +4987,7 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = cron_stop, cron_provider, cron_thread, housekeeping_thread = ( _start_gateway_start_cron_and_housekeeping(runner)) - # READY is emitted only after adapters, cron and housekeeping reach their running boundary; - # missing config/systemd runtime state leaves the watchdog disabled without changing behavior. + # READY only once adapters, cron and housekeeping run; missing systemd state just disables watchdog. runner._start_systemd_watchdog() await runner.wait_for_shutdown() @@ -5283,50 +5041,38 @@ def main(): _best_effort(_step) import argparse - parser = argparse.ArgumentParser(description="Hermes Gateway - Multi-platform messaging") parser.add_argument("--config", "-c", help="Path to gateway config file") parser.add_argument("--verbose", "-v", action="store_true", help="Verbose output") - args = parser.parse_args() config = None if args.config: import yaml with open(args.config, encoding="utf-8") as f: - data = yaml.safe_load(f) or {} - config = GatewayConfig.from_dict(data) + config = GatewayConfig.from_dict(yaml.safe_load(f) or {}) - # start_gateway() finishes graceful teardown before returning OR raising SystemExit; force-exit after so - # a wedged non-daemon worker can't block Py_FinalizeEx's thread join. SystemExit is caught so EVERY - # exit path hits the os._exit backstop. + # start_gateway() completes teardown before returning/raising SystemExit; force-exit after so a + # wedged non-daemon worker can't block Py_FinalizeEx's join. SystemExit caught so EVERY path exits. try: success = asyncio.run(start_gateway(config)) exit_code = 0 if success else 1 except SystemExit as e: # e.code may be None (→ 0), an int, or a str (→ 1, like CPython). - if e.code is None: - exit_code = 0 - elif isinstance(e.code, int): - exit_code = e.code - else: - exit_code = 1 + exit_code = 0 if e.code is None else e.code if isinstance(e.code, int) else 1 _exit_after_graceful_shutdown(exit_code) def _exit_after_graceful_shutdown(exit_code: int) -> None: """Flush stdio, release the PID file + runtime lock, then hard-exit. - ``os._exit`` (not ``sys.exit``): SystemExit runs ``Py_FinalizeEx``, which joins every non-daemon thread — exactly the hang a wedged worker causes. It bypasses ``atexit``, so PID/lock release and the - bounded log drain (file handlers sit behind a ``QueueListener`` thread) are done here explicitly. - """ + bounded log drain (file handlers sit behind a ``QueueListener`` thread) are done here explicitly.""" for stream in (sys.stdout, sys.stderr): with suppress(Exception): stream.flush() def _release_locks() -> None: - # BEFORE the log drain: the drain is bounded but could take its full timeout on a wedged disk, - # and these locks must never be stranded. Idempotent (early SystemExit paths never ran _stop_impl). + # BEFORE the log drain (bounded, but could take its full timeout on a wedged disk); idempotent. from gateway.status import remove_pid_file, release_gateway_runtime_lock remove_pid_file() release_gateway_runtime_lock()