refactor(gateway): fold docstring closers in run.py (text parity via cleandoc)
This commit is contained in:
+62
-124
@@ -1,8 +1,7 @@
|
||||
"""Gateway runner - entry point for messaging platform integrations.
|
||||
|
||||
Provides ``start_gateway()`` (start all configured adapters) and ``GatewayRunner`` (lifecycle).
|
||||
Run via ``python -m gateway.run`` or ``python cli.py --gateway``.
|
||||
"""
|
||||
Run via ``python -m gateway.run`` or ``python cli.py --gateway``."""
|
||||
|
||||
# hermes_bootstrap must be the very first import (UTF-8 stdio on Windows; no-op on POSIX).
|
||||
try:
|
||||
@@ -167,8 +166,7 @@ def hygiene_compaction_recovered(
|
||||
|
||||
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).
|
||||
"""
|
||||
counts 30-50% estimate noise as one)."""
|
||||
if aborted or not (rotated or in_place):
|
||||
return False
|
||||
return compression_made_progress(msg_count, new_count, approx_tokens, new_tokens)
|
||||
@@ -213,8 +211,7 @@ async def run_codex_hygiene_compaction(
|
||||
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:<reason>`` or ``failed:<reason>``.
|
||||
"""
|
||||
Returns ``compacted``, ``skipped:<reason>`` or ``failed:<reason>``."""
|
||||
mode = str(auto_mode or "native").lower()
|
||||
if mode not in {"native", "hermes", "off"}:
|
||||
mode = "native"
|
||||
@@ -275,8 +272,7 @@ def hygiene_wait_should_extend(
|
||||
"""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.
|
||||
"""
|
||||
messages behind a doomed attempt, so stop extending immediately."""
|
||||
return not fence_cancelled and idle < timeout and waited < ceiling
|
||||
|
||||
|
||||
@@ -300,8 +296,7 @@ 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.
|
||||
"""
|
||||
each ``{field}`` becomes a numeric-ish pattern."""
|
||||
parts = re.split(r"\{[^{}]*\}", template)
|
||||
return r"[\d,]+".join(re.escape(part) for part in parts)
|
||||
|
||||
@@ -449,8 +444,7 @@ def _interim_metadata(metadata: Optional[Dict[str, Any]] = None) -> Dict[str, An
|
||||
|
||||
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.
|
||||
"""
|
||||
answer stream with status text. Gateway-internal; adapters strip it before the wire."""
|
||||
merged = dict(metadata or {})
|
||||
merged["_interim_send"] = True
|
||||
return merged
|
||||
@@ -461,8 +455,7 @@ def _seed_hygiene_system_prompt(agent: Any, session_row: Optional[Dict[str, Any]
|
||||
|
||||
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.
|
||||
"""
|
||||
empty cache entry); the real turn rebuilds with fully initialized providers."""
|
||||
stored_prompt = ""
|
||||
if isinstance(session_row, dict):
|
||||
raw_prompt = session_row.get("system_prompt")
|
||||
@@ -526,8 +519,7 @@ def _redact_gateway_user_facing_secrets(text: str) -> str:
|
||||
|
||||
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.
|
||||
"""
|
||||
pass so redaction degrades gracefully if that import fails."""
|
||||
redacted = str(text or "")
|
||||
try:
|
||||
from agent.redact import redact_sensitive_text
|
||||
@@ -545,8 +537,7 @@ 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.
|
||||
"""
|
||||
verbatim to chat; ``force=True`` holds even when ``security.redact_secrets`` is off."""
|
||||
from agent.redact import redact_sensitive_text
|
||||
|
||||
return redact_sensitive_text(str(cmd or ""), force=True)
|
||||
@@ -608,8 +599,7 @@ 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.
|
||||
"""
|
||||
mentions a status code does not match."""
|
||||
if not text:
|
||||
return False
|
||||
body = str(text).strip()
|
||||
@@ -667,8 +657,7 @@ 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.
|
||||
"""
|
||||
Fail-soft: a malformed/empty notice degrades to "" rather than raising."""
|
||||
return str(getattr(notice, "text", "") or "").strip()
|
||||
|
||||
|
||||
@@ -686,8 +675,7 @@ def _approval_send_outcome(future, timeout: float) -> str:
|
||||
|
||||
``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.
|
||||
"""
|
||||
result / non-timeout exception / no future) re-asks; those log their detail here."""
|
||||
if future is None:
|
||||
logger.warning("Prompt send failed: no scheduling future (loop unavailable)")
|
||||
return "failed"
|
||||
@@ -744,8 +732,7 @@ def _resolve_progress_thread_id(
|
||||
|
||||
``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.
|
||||
"""
|
||||
event's own message id is the adapter's synthetic session-keying thread — treat as 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):
|
||||
@@ -776,8 +763,7 @@ def _resolve_gateway_display_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.<platform>.<setting> override.
|
||||
"""
|
||||
so they require an explicit display.platforms.<platform>.<setting> 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 (
|
||||
@@ -916,8 +902,7 @@ def _as_thread_info(info: Any) -> Optional[Tuple[str, str]]:
|
||||
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.
|
||||
"""
|
||||
A misconfigured env var (``HERMES_AGENT_TIMEOUT=abc``) must not crash the gateway or a turn."""
|
||||
raw = os.environ.get(name)
|
||||
if raw is None or raw == "":
|
||||
return float(default)
|
||||
@@ -940,8 +925,7 @@ 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).
|
||||
"""
|
||||
Unknown timestamps count as fresh (legacy transcripts, in-memory test scaffolding)."""
|
||||
window = float(window_secs) if window_secs is not None else float(_AUTO_CONTINUE_FRESHNESS_SECS_DEFAULT)
|
||||
if window <= 0:
|
||||
return True
|
||||
@@ -1020,8 +1004,7 @@ def _build_replay_entry(
|
||||
|
||||
``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.
|
||||
"""
|
||||
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).
|
||||
@@ -1081,8 +1064,7 @@ 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.
|
||||
"""
|
||||
test hook bypasses the adapter, ignored channels cannot reach auth, pairing or sessions."""
|
||||
platform_cfg = getattr(config, "platforms", {}).get(Platform.SLACK)
|
||||
raw = None
|
||||
if platform_cfg is not None:
|
||||
@@ -1108,8 +1090,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.
|
||||
"""
|
||||
Default OFF: a timestamp prefix on every user message changes what the model sees."""
|
||||
if not isinstance(user_config, dict):
|
||||
return False
|
||||
gw = user_config.get("gateway")
|
||||
@@ -1128,8 +1109,7 @@ def _build_gateway_agent_history(
|
||||
"""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.
|
||||
"""
|
||||
with the live turn and hiding the current message behind ``history_offset`` on persistence."""
|
||||
from hermes_time import get_timezone as _get_msg_tz
|
||||
from gateway.message_timestamps import render_user_content_with_timestamp as _render_msg_ts
|
||||
|
||||
@@ -1191,8 +1171,7 @@ def _select_cached_agent_history(
|
||||
|
||||
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.
|
||||
"""
|
||||
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
|
||||
|
||||
@@ -1229,8 +1208,7 @@ def _last_transcript_timestamp(history: Optional[List[Dict[str, Any]]]) -> Any:
|
||||
"""Return the ``timestamp`` of the last usable 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.
|
||||
"""
|
||||
a timestamp — callers treat that as "fresh" for backward compatibility."""
|
||||
if not history:
|
||||
return None
|
||||
for msg in reversed(history):
|
||||
@@ -1362,8 +1340,7 @@ def _collect_auto_append_media_tags(
|
||||
def _collect_history_media_paths(agent_history: List[Dict[str, Any]]) -> set:
|
||||
"""Collect every media path already delivered in prior assistant/tool output (dedup set).
|
||||
|
||||
Both the JSON-payload and assistant-message shapes must be covered or delivery repeats.
|
||||
"""
|
||||
Both the JSON-payload and assistant-message shapes must be covered or delivery repeats."""
|
||||
paths: set = set()
|
||||
tool_name_by_call_id = _tool_name_by_call_id(agent_history)
|
||||
|
||||
@@ -1407,8 +1384,7 @@ def _ensure_ssl_certs() -> None:
|
||||
"""Set SSL_CERT_FILE when the system hides CA certs from Python (NixOS etc.).
|
||||
|
||||
Must run BEFORE any HTTP library (discord, aiohttp) is imported. A set-but-missing path makes
|
||||
every later httpx/OpenAI client fail in ssl.load_verify_locations(), so treat it as unset.
|
||||
"""
|
||||
every later httpx/OpenAI client fail in ssl.load_verify_locations(), so treat it as unset."""
|
||||
configured_cert = os.environ.get("SSL_CERT_FILE")
|
||||
if configured_cert:
|
||||
if os.path.exists(configured_cert):
|
||||
@@ -1506,8 +1482,7 @@ def _reload_runtime_env_preserving_config_authority() -> None:
|
||||
Long-lived gateways reload ~/.hermes/.env per turn for rotated keys; config.yaml stays
|
||||
authoritative for budgets (else stale HERMES_MAX_ITERATIONS wins). Multiplex mode never reloads
|
||||
.env globally (secrets come from the per-turn ``set_secret_scope``; mutating ``os.environ`` would
|
||||
leak the default profile's keys to every profile) but still honors the max_turns bridge.
|
||||
"""
|
||||
leak the default profile's keys to every profile) but still honors the max_turns bridge."""
|
||||
from agent.secret_scope import is_multiplex_active
|
||||
if is_multiplex_active():
|
||||
_bridge_max_turns_from_config(_hermes_home)
|
||||
@@ -1535,8 +1510,7 @@ def _current_max_iterations() -> int:
|
||||
"""Return the current per-turn iteration budget after runtime env refresh.
|
||||
|
||||
Uses ``resolve_turn_limit`` so ``agent.max_turns: none``/``unlimited`` (bridged as a string
|
||||
into ``HERMES_MAX_ITERATIONS``) yields the unlimited sentinel instead of an ``int()`` crash.
|
||||
"""
|
||||
into ``HERMES_MAX_ITERATIONS``) yields the unlimited sentinel instead of an ``int()`` crash."""
|
||||
_reload_runtime_env_preserving_config_authority()
|
||||
from hermes_cli.config import resolve_turn_limit as _resolve_turn_limit
|
||||
return _resolve_turn_limit(os.getenv("HERMES_MAX_ITERATIONS"))
|
||||
@@ -1549,24 +1523,21 @@ class MultiplexConfigError(RuntimeError):
|
||||
"""A profile multiplexer config is invalid.
|
||||
|
||||
Distinct from a transient adapter-connect failure: the operator must fix config.yaml, so it
|
||||
propagates to the startup guard instead of being treated as retryable adapter noise.
|
||||
"""
|
||||
propagates to the startup guard instead of being treated as retryable adapter noise."""
|
||||
|
||||
|
||||
class SecondaryPortBindingConfigError(MultiplexConfigError):
|
||||
"""A secondary profile enabled a port-binding platform.
|
||||
|
||||
The default profile owns the single shared listener (serving every profile via /p/<profile>/),
|
||||
so this is always a misconfiguration and is skipped rather than taking down the multiplexer.
|
||||
"""
|
||||
so this is always a misconfiguration and is skipped rather than taking down the multiplexer."""
|
||||
|
||||
|
||||
class HygieneTurnHoldExceeded(Exception):
|
||||
"""Hygiene-compression turn-hold budget elapsed while the summary model was still streaming.
|
||||
|
||||
Availability boundary, not a failure: must NOT take the idle-timeout failure path
|
||||
(AGENT_COMPRESSION_TIMEOUT, "no output" message, failure cooldown ladder).
|
||||
"""
|
||||
(AGENT_COMPRESSION_TIMEOUT, "no output" message, failure cooldown ladder)."""
|
||||
|
||||
|
||||
def _multiplex_profile_homes(config: object) -> list[tuple[str, "Path"]]:
|
||||
@@ -1600,8 +1571,7 @@ def _handoff_watch_scopes(runner: object) -> list:
|
||||
``/handoff`` writes into the store of the profile the CLI ran under; an unscoped watcher polls
|
||||
only the ROOT store, so a secondary profile's handoff would never be seen (CLI times out).
|
||||
``(None, None)`` = root poll, always first; secondary profiles follow (default not repeated).
|
||||
A raising resolver degrades to the root poll rather than silently disabling the watcher.
|
||||
"""
|
||||
A raising resolver degrades to the root poll rather than silently disabling the watcher."""
|
||||
scopes: list = [(None, None)]
|
||||
try:
|
||||
config = getattr(runner, "config", None)
|
||||
@@ -1620,8 +1590,7 @@ async def _reclaim_stale(runner: object) -> None:
|
||||
|
||||
Runs once per store at watcher startup. ``running`` is only set for one in-process dispatch, so
|
||||
a leftover row belongs to a dead process and blocks ``request_handoff`` for that session forever.
|
||||
Defensive: a raising reclaim would abort startup.
|
||||
"""
|
||||
Defensive: a raising reclaim would abort startup."""
|
||||
session_db = getattr(runner, "_session_db", None)
|
||||
if session_db is None:
|
||||
return
|
||||
@@ -1675,8 +1644,7 @@ def _profile_runtime_scope(
|
||||
|
||||
``set_hermes_home_override`` is a contextvar (reaches the agent worker via ``copy_context()``);
|
||||
``set_secret_scope`` makes the profile ``.env`` the credential source without mutating
|
||||
``os.environ``, so subprocesses never inherit cross-profile secrets.
|
||||
"""
|
||||
``os.environ``, so subprocesses never inherit cross-profile secrets."""
|
||||
from hermes_constants import set_hermes_home_override, reset_hermes_home_override
|
||||
from agent.secret_scope import set_secret_scope, reset_secret_scope
|
||||
|
||||
@@ -1717,8 +1685,7 @@ def load_gateway_config_for_runner() -> "GatewayConfig":
|
||||
With multiplexing on, reload under the default profile's ``_profile_runtime_scope`` so platform
|
||||
tokens in that profile's ``.env`` resolve via the secret scope (as secondary profiles do);
|
||||
unscoped, ``_getenv`` falls to ``os.environ``, which often lacks a token living only under
|
||||
``profiles/<name>/.env``. Off -> identical to ``load_gateway_config()``.
|
||||
"""
|
||||
``profiles/<name>/.env``. Off -> identical to ``load_gateway_config()``."""
|
||||
cfg = load_gateway_config()
|
||||
if not getattr(cfg, "multiplex_profiles", False):
|
||||
return cfg
|
||||
@@ -1759,8 +1726,7 @@ def _platform_has_bot_credential(platform: "Platform", platform_config: "Platfor
|
||||
"""Return True when a token-authenticated platform has a usable bot credential.
|
||||
|
||||
Platforms that do not use ``PlatformConfig.token`` always return True so we
|
||||
never skip them here (Signal session paths, port-binding HTTP adapters, etc.).
|
||||
"""
|
||||
never skip them here (Signal session paths, port-binding HTTP adapters, etc.)."""
|
||||
from gateway.config import PLATFORM_TOKEN_ENV_NAMES, Platform
|
||||
|
||||
if platform not in PLATFORM_TOKEN_ENV_NAMES:
|
||||
@@ -1828,8 +1794,7 @@ def _bridge_max_turns_to_env(agent_cfg: Any) -> None:
|
||||
"""Bridge ``agent.max_turns`` preserving its raw spelling ("none", "unlimited", "120").
|
||||
|
||||
Python None (`null` / bare `key:`) clears a stale bridge instead: str(None) -> "None" would map
|
||||
to the unlimited sentinel rather than "absent = default".
|
||||
"""
|
||||
to the unlimited sentinel rather than "absent = default"."""
|
||||
if not isinstance(agent_cfg, dict) or "max_turns" not in agent_cfg:
|
||||
return
|
||||
raw = agent_cfg["max_turns"]
|
||||
@@ -1965,8 +1930,7 @@ def _load_bridge_config(config_path: Path) -> dict:
|
||||
"""Raw config read for the presence-sensitive env bridge, with the managed overlay applied.
|
||||
|
||||
Raw (not defaults-merged) so only keys the user wrote are bridged, else all of DEFAULT_CONFIG
|
||||
would be exported. Managed overlay applies BEFORE bridging so pinned values win in env too.
|
||||
"""
|
||||
would be exported. Managed overlay applies BEFORE bridging so pinned values win in env too."""
|
||||
from hermes_cli.config import _expand_env_vars, read_user_config_raw
|
||||
cfg = _expand_env_vars(read_user_config_raw(config_path))
|
||||
if not isinstance(cfg, dict):
|
||||
@@ -2172,8 +2136,7 @@ def _resolve_runtime_agent_kwargs() -> dict:
|
||||
"""Resolve provider credentials for gateway-created AIAgent instances.
|
||||
|
||||
``resolve_runtime_provider()`` falls through to env vars for legacy compatibility, but the
|
||||
gateway never consults env vars for behavioral config — config.yaml is authoritative.
|
||||
"""
|
||||
gateway never consults env vars for behavioral config — config.yaml is authoritative."""
|
||||
from hermes_cli.runtime_provider import (
|
||||
resolve_runtime_provider, format_runtime_provider_error, _get_model_config)
|
||||
from hermes_cli.auth import AuthError, is_rate_limited_auth_error
|
||||
@@ -2225,8 +2188,7 @@ def _runtime_agent_kwargs(runtime: dict) -> dict:
|
||||
"""AIAgent constructor kwargs shared by every runtime-provider resolution.
|
||||
|
||||
``request_overrides`` is passed through as resolved (custom_providers ``extra_body`` etc.) so
|
||||
the provider's configured request body reaches the per-turn route on the gateway path.
|
||||
"""
|
||||
the provider's configured request body reaches the per-turn route on the gateway path."""
|
||||
return {
|
||||
"api_key": runtime.get("api_key"),
|
||||
"base_url": runtime.get("base_url"),
|
||||
@@ -2253,8 +2215,7 @@ class _GatewayModelContext:
|
||||
def _resolve_gateway_model_context(model: Optional[str] = None) -> _GatewayModelContext:
|
||||
"""Resolve the configured gateway route and its effective context window.
|
||||
|
||||
Shared by status/session banners and slash commands. Call off the event loop (may block).
|
||||
"""
|
||||
Shared by status/session banners and slash commands. Call off the event loop (may block)."""
|
||||
from agent.model_metadata import DEFAULT_FALLBACK_CONTEXT, get_model_context_length
|
||||
|
||||
resolved_model = model or _resolve_gateway_model()
|
||||
@@ -2430,8 +2391,7 @@ def _event_media_is_video(event, index: int) -> bool:
|
||||
def _build_media_placeholder(event) -> str:
|
||||
"""Text placeholder for media-only events (later replaced by vision enrichment).
|
||||
|
||||
Queued media is dequeued via .text only, so a caption-less event would otherwise be lost.
|
||||
"""
|
||||
Queued media is dequeued via .text only, so a caption-less event would otherwise be lost."""
|
||||
parts = []
|
||||
media_urls = getattr(event, "media_urls", None) or []
|
||||
for i, url in enumerate(media_urls):
|
||||
@@ -2451,8 +2411,7 @@ def _build_document_context_note(
|
||||
"""Context note prepended to a user turn when they attach a document.
|
||||
|
||||
``content_inlined=False`` = cached without content, so tell the agent to read it. Binary docs
|
||||
must say *extract* the text; "ask the user" made it punt.
|
||||
"""
|
||||
must say *extract* the text; "ask the user" made it punt."""
|
||||
if mtype.startswith("text/") and content_inlined:
|
||||
return (
|
||||
f"[The user sent a text document: '{display_name}'. Its content has been included below. "
|
||||
@@ -2523,8 +2482,7 @@ async def _probe_audio_duration(path: str) -> Optional[str]:
|
||||
def _dequeue_pending_event(adapter, session_key: str) -> MessageEvent | None:
|
||||
"""Consume and return the full pending event for a session.
|
||||
|
||||
Media metadata is kept so follow-ups re-enter normal image/STT/document preprocessing.
|
||||
"""
|
||||
Media metadata is kept so follow-ups re-enter normal image/STT/document preprocessing."""
|
||||
return adapter.get_pending_message(session_key)
|
||||
|
||||
|
||||
@@ -2543,8 +2501,7 @@ def _reap_gateway_turn_processes(
|
||||
|
||||
``task_id`` is session-scoped, so a *replacement* turn can spawn its own process mid-reap;
|
||||
``is_still_current`` (closure over the captured run_generation) lets the caller bail instead of
|
||||
killing it — that turn snapshots its own baseline, so nothing stays unreaped.
|
||||
"""
|
||||
killing it — that turn snapshots its own baseline, so nothing stays unreaped."""
|
||||
if not task_id:
|
||||
# Blank task_id (sessionless callers) would match and kill every unrelated empty-task process.
|
||||
return 0
|
||||
@@ -2585,8 +2542,7 @@ def _dump_wedged_turn_stacks(task_id: str) -> None:
|
||||
"""Log the stack of every thread that looks like turn work, at reap time.
|
||||
|
||||
The hard interrupt frees the wedged worker before a profiler can attach, so dump BEFORE it.
|
||||
Best-effort, bounded (turn-machinery threads only, capped output), never raises.
|
||||
"""
|
||||
Best-effort, bounded (turn-machinery threads only, capped output), never raises."""
|
||||
try:
|
||||
frames = sys._current_frames()
|
||||
names = {t.ident: t.name for t in threading.enumerate()}
|
||||
@@ -2686,8 +2642,7 @@ def _strip_response_attachments_for_direct_send(response: str, adapter) -> str:
|
||||
|
||||
Replays only explicit ``MEDIA:`` attachments; bare paths/URLs stay visible (the post-stream
|
||||
uploader ignores them). No broad ``MEDIA:`` regex after ``extract_media()``: it deliberately
|
||||
preserves protected code spans and unvalidated tags.
|
||||
"""
|
||||
preserves protected code spans and unvalidated tags."""
|
||||
_, cleaned = adapter.extract_media(response)
|
||||
return cleaned.replace("[[audio_as_voice]]", "").replace("[[as_document]]", "").strip()
|
||||
|
||||
@@ -2795,8 +2750,7 @@ def _load_gateway_config(config_path: "Path | None" = None) -> dict:
|
||||
"""Load and parse a gateway config.yaml, returning {} on any error (fail-open).
|
||||
|
||||
Defaults to the active gateway home (``_hermes_home`` monkeypatches apply); multiplexed callers
|
||||
may pass a profile path.
|
||||
"""
|
||||
may pass a profile path."""
|
||||
if config_path is None:
|
||||
config_path = _gateway_config_home() / 'config.yaml'
|
||||
raw: dict = {}
|
||||
@@ -2841,8 +2795,7 @@ def _load_gateway_config(config_path: "Path | None" = None) -> dict:
|
||||
def _checkpoint_agent_kwargs(config: dict | None) -> dict:
|
||||
"""Translate gateway checkpoint config into ``AIAgent`` constructor args.
|
||||
|
||||
Gateway bypasses ``load_config()``, so defaults are here; legacy ``checkpoints: true`` works.
|
||||
"""
|
||||
Gateway bypasses ``load_config()``, so defaults are here; legacy ``checkpoints: true`` works."""
|
||||
cp_cfg = config.get("checkpoints", {}) if isinstance(config, dict) else {}
|
||||
if isinstance(cp_cfg, bool):
|
||||
cp_cfg = {"enabled": cp_cfg}
|
||||
@@ -2862,8 +2815,7 @@ def _load_gateway_runtime_config() -> dict:
|
||||
"""Load gateway config for runtime reads, expanding supported ``${VAR}`` refs.
|
||||
|
||||
Built on ``_load_gateway_config()``. Expansion failures are deliberately NOT swallowed —
|
||||
returning the unexpanded dict would mask the very bug this helper fixes.
|
||||
"""
|
||||
returning the unexpanded dict would mask the very bug this helper fixes."""
|
||||
cfg = _load_gateway_config()
|
||||
if not isinstance(cfg, dict) or not cfg:
|
||||
return {}
|
||||
@@ -2943,8 +2895,7 @@ def _parse_session_key(session_key: str) -> "dict | None":
|
||||
"""Parse a session key (``agent:main:{platform}:{chat_type}:{chat_id}[:{extra}...]``).
|
||||
|
||||
For group/channel sessions the suffix may be a user_id (per-user isolation), not a thread_id,
|
||||
so ``thread_id`` is left out to avoid mis-routing.
|
||||
"""
|
||||
so ``thread_id`` is left out to avoid mis-routing."""
|
||||
parts = session_key.split(":")
|
||||
if len(parts) >= 5 and parts[0] == "agent" and parts[1] == "main":
|
||||
result = {"platform": parts[2], "chat_type": parts[3], "chat_id": parts[4]}
|
||||
@@ -3029,8 +2980,7 @@ def _drain_gateway_watch_events(completion_queue) -> "list[dict]":
|
||||
"""Drain gateway-owned watch events without spinning on requeued events.
|
||||
|
||||
Foreign events (process completions, async delegations) requeued inside ``while not
|
||||
queue.empty()`` never terminate, so detach the batch first and requeue afterwards.
|
||||
"""
|
||||
queue.empty()`` never terminate, so detach the batch first and requeue afterwards."""
|
||||
watch_events: list[dict] = []
|
||||
requeue: list[dict] = []
|
||||
while not completion_queue.empty():
|
||||
@@ -3060,8 +3010,7 @@ def _normalize_empty_agent_response(
|
||||
"""Normalize empty/None agent responses into user-facing messages.
|
||||
|
||||
Covers ``failed``, work done (api_calls > 0) with no text, and never-ran (api_calls == 0,
|
||||
the post-/stop silent-drop from a stale generation token) with a retry hint.
|
||||
"""
|
||||
the post-/stop silent-drop from a stale generation token) with a retry hint."""
|
||||
if response:
|
||||
return response
|
||||
|
||||
@@ -3144,8 +3093,7 @@ def _should_clear_resume_pending_after_turn(agent_result: dict) -> bool:
|
||||
"""True only when a gateway turn really completed successfully.
|
||||
|
||||
``resume_pending`` is a durable restart-recovery marker; a soft interrupt can look like a normal
|
||||
empty result, and clearing then loses the signal.
|
||||
"""
|
||||
empty result, and clearing then loses the signal."""
|
||||
if not isinstance(agent_result, dict):
|
||||
return False
|
||||
if agent_result.get("interrupted"):
|
||||
@@ -3160,8 +3108,7 @@ def _preserve_queued_followup_history_offset(
|
||||
"""Carry the outer history offset through queued follow-up drains.
|
||||
|
||||
Each recursive ``_run_agent()`` advances ``history_offset``; uncorrected, the outer persistence
|
||||
step sees only the *last* queued turn as "new" and drops earlier ones.
|
||||
"""
|
||||
step sees only the *last* queued turn as "new" and drops earlier ones."""
|
||||
if not isinstance(followup_result, dict) or not isinstance(current_result, dict):
|
||||
return followup_result
|
||||
|
||||
@@ -3182,8 +3129,7 @@ async def _dispose_unused_adapter(adapter: "BasePlatformAdapter | None") -> None
|
||||
|
||||
A failed connect leaves it uninstalled so nothing else calls ``disconnect()``; resources opened
|
||||
in ``__init__`` (e.g. SQLite fds) would leak until GC (not prompt for asyncio-bound objects) and
|
||||
exhaust the fd ulimit over a long retry loop. ``adapter`` may be ``None`` (half-constructed).
|
||||
"""
|
||||
exhaust the fd ulimit over a long retry loop. ``adapter`` may be ``None`` (half-constructed)."""
|
||||
if adapter is None:
|
||||
return
|
||||
try:
|
||||
@@ -3212,8 +3158,7 @@ def _reconnect_backoff(attempt: int) -> int:
|
||||
def _reconnect_needs_attention(info: dict, now: float) -> bool:
|
||||
"""True when a reconnect-queue entry has waited long enough for NEEDS_ATTENTION.
|
||||
|
||||
``queued_at`` is re-stamped on each (re)entry, so only *continuous* failure escalates.
|
||||
"""
|
||||
``queued_at`` is re-stamped on each (re)entry, so only *continuous* failure escalates."""
|
||||
if _RECONNECT_ATTENTION_AFTER_SECONDS <= 0:
|
||||
return False # escalation disabled
|
||||
queued_at = info.get("queued_at")
|
||||
@@ -3915,8 +3860,7 @@ class GatewayRunner(
|
||||
"""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``
|
||||
would clobber it). Best-effort: a failed write must never disrupt a turn.
|
||||
"""
|
||||
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:
|
||||
@@ -4276,8 +4220,7 @@ class GatewayRunner(
|
||||
|
||||
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.
|
||||
"""
|
||||
idle time accumulates to the 30-min timeout."""
|
||||
if interrupt_depth == 0:
|
||||
from agent.session_activity import ActivityProvenance
|
||||
|
||||
@@ -4561,8 +4504,7 @@ def _start_gateway_housekeeping(stop_event: threading.Event, adapters=None, loop
|
||||
|
||||
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.
|
||||
"""
|
||||
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),
|
||||
@@ -4635,8 +4577,7 @@ async def _shutdown_mcp_servers_nonblocking(timeout: float = 5.0) -> bool:
|
||||
|
||||
``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.
|
||||
On timeout shutdown proceeds and the daemon thread is left to finish or die.
|
||||
"""
|
||||
On timeout shutdown proceeds and the daemon thread is left to finish or die."""
|
||||
|
||||
def _do() -> None:
|
||||
try:
|
||||
@@ -4737,8 +4678,7 @@ 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.
|
||||
"""
|
||||
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)
|
||||
@@ -4808,8 +4748,7 @@ 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).
|
||||
"""
|
||||
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:
|
||||
hermes_home = str(get_hermes_home())
|
||||
@@ -5076,8 +5015,7 @@ 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)``.
|
||||
"""
|
||||
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)
|
||||
|
||||
Reference in New Issue
Block a user