diff --git a/gateway/run.py b/gateway/run.py index 9ddaa75c74..eb00badd6b 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -35,16 +35,10 @@ from typing import Callable, Dict, Optional, Any, List, Tuple, cast from agent.async_utils import safe_schedule_threadsafe from agent.conversation_compression import ( - COMPACTION_DONE_STATUS, - COMPACTION_STATUS, - COMPRESSION_RETRY_CONTEXT_REDUCED_STATUS_TEMPLATE, - COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE, - COMPRESSION_RETRY_TOKENS_STATUS_TEMPLATE, - COMPRESSION_RETRY_TOO_LARGE_STATUS_TEMPLATE, - IDLE_COMPACTION_STATUS_TEMPLATE, - PRE_API_COMPRESSION_STATUS_TEMPLATE, - PREFLIGHT_COMPRESSION_STATUS_TEMPLATE, -) + COMPACTION_DONE_STATUS, COMPACTION_STATUS, COMPRESSION_RETRY_CONTEXT_REDUCED_STATUS_TEMPLATE, + COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE, COMPRESSION_RETRY_TOKENS_STATUS_TEMPLATE, + COMPRESSION_RETRY_TOO_LARGE_STATUS_TEMPLATE, IDLE_COMPACTION_STATUS_TEMPLATE, + PRE_API_COMPRESSION_STATUS_TEMPLATE, PREFLIGHT_COMPRESSION_STATUS_TEMPLATE) from agent.conversation_loop import INTERRUPT_WAITING_FOR_MODEL_PREFIX from agent.interrupt_compat import request_hard_interrupt from agent.turn_context import compression_made_progress @@ -67,10 +61,7 @@ _ADAPTER_DISCONNECT_TIMEOUT_SECS_DEFAULT = 5.0 # by _classify_completion_target and _resolve_async_delegation_session so they can never disagree: # a reason the classifier "delivers" but the resolver drops would be acked and then silently lost. _USER_BOUNDARY_END_REASONS = ( - "session_reset", - "user_exit", - "session_switch", - "new_session", + "session_reset", "user_exit", "session_switch", "new_session" ) # Bound on a single stall-notify adapter.send so a wedged transport cannot block the stall watcher # pass; on timeout the latch stays clear and the next tick retries. @@ -131,9 +122,7 @@ def _gateway_session_db_inner(gateway): def _hygiene_cooldown_for_failure( - gateway, - session_key: str, - base_cooldown_seconds: float, + gateway, session_key: str, base_cooldown_seconds: float ) -> float: """Bump the hygiene failure streak and return the escalated cooldown. @@ -187,14 +176,8 @@ def _reset_hygiene_failure_streak(gateway, session_key: str) -> None: def hygiene_compaction_recovered( - *, - aborted: bool, - rotated: bool, - in_place: bool, - msg_count: int, - new_count: int, - approx_tokens: int, - new_tokens: int, + *, aborted: bool, rotated: bool, in_place: bool, msg_count: int, new_count: int, + approx_tokens: int, new_tokens: int, ) -> bool: """True when a hygiene run actually recovered the session (extracted to be unit testable). @@ -213,11 +196,7 @@ def hygiene_compaction_recovered( def _hygiene_compression_timeout_message( - *, - total_exhausted: bool, - elapsed: float, - idle_timeout: float, - progress_observed: bool, + *, total_exhausted: bool, elapsed: float, idle_timeout: float, progress_observed: bool ) -> str: """Describe the host timeout that actually ended hygiene compression.""" if total_exhausted: @@ -227,8 +206,7 @@ def _hygiene_compression_timeout_message( return ( "⚠️ Context compression reached its total ceiling after " f"{elapsed:.1f}s{progress}. No messages were dropped — continuing " - "without compression. Run /compress to retry or /reset for a clean " - "session." + "without compression. Run /compress to retry or /reset for a clean session." ) return ( f"⚠️ Context compression timed out after {idle_timeout:.1f}s with no " @@ -239,15 +217,8 @@ def _hygiene_compression_timeout_message( async def run_codex_hygiene_compaction( - gateway, - session_key: str, - session_id: str, - *, - auto_mode: str, - history: list, - approx_tokens: int, - timeout_seconds: float, - failure_cooldown_seconds: float = 300.0, + gateway, session_key: str, session_id: str, *, auto_mode: str, history: list, + approx_tokens: int, timeout_seconds: float, failure_cooldown_seconds: float = 300.0, ) -> str: """Session hygiene for ``codex_app_server`` sessions. @@ -295,9 +266,7 @@ async def run_codex_hygiene_compaction( # ContextVars on the Python runtimes Hermes currently ships. copy_context().run, lambda: agent._compress_context( - history, - "", - approx_tokens=approx_tokens, + history, "", approx_tokens=approx_tokens ), ) track_worker = getattr(gateway, "_track_deferred_agent_worker", None) @@ -308,8 +277,7 @@ async def run_codex_hygiene_compaction( track_worker(worker_future, agent) try: await asyncio.wait_for( - asyncio.shield(worker_future), - timeout=max(float(timeout_seconds), 1.0), + 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 @@ -319,22 +287,18 @@ async def run_codex_hygiene_compaction( gateway, session_id, failure_cooldown_seconds, - "codex app-server thread compaction timed out", - ) + "codex app-server thread compaction timed out") logger.warning( "Session hygiene: codex app-server thread compaction for " "session %s timed out after %.1fs; continuing without compaction", session_id, - timeout_seconds, - ) + timeout_seconds) return "failed:timeout" except Exception as exc: logger.warning( - "Session hygiene: codex app-server thread compaction for " - "session %s failed: %s", + "Session hygiene: codex app-server thread compaction for session %s failed: %s", session_id, - exc, - ) + exc) return f"failed:{exc}" count_after = getattr(compressor, "compression_count", 0) @@ -348,12 +312,7 @@ async def run_codex_hygiene_compaction( return "failed:no-boundary" def hygiene_wait_should_extend( - *, - idle: float, - timeout: float, - waited: float, - ceiling: float, - fence_cancelled: bool = False, + *, idle: float, timeout: float, waited: float, ceiling: float, fence_cancelled: bool = False ) -> bool: """Whether the hygiene host should keep waiting for a slow summary. @@ -366,10 +325,7 @@ def hygiene_wait_should_extend( def _record_hygiene_cooldown( - gateway, - session_id: str, - cooldown_seconds: float, - error: Optional[str] = None, + gateway, session_id: str, cooldown_seconds: float, error: Optional[str] = None ) -> None: """Persist a session-hygiene compression-failure cooldown to the state DB. @@ -401,19 +357,13 @@ _COMPRESSION_PROGRESS_STATUS_RE = re.compile( "|".join( _status_template_to_regex(_template) for _template in ( - COMPACTION_STATUS, - COMPACTION_DONE_STATUS, - PRE_API_COMPRESSION_STATUS_TEMPLATE, - PREFLIGHT_COMPRESSION_STATUS_TEMPLATE, - IDLE_COMPACTION_STATUS_TEMPLATE, - COMPRESSION_RETRY_TOO_LARGE_STATUS_TEMPLATE, - COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE, + COMPACTION_STATUS, COMPACTION_DONE_STATUS, PRE_API_COMPRESSION_STATUS_TEMPLATE, + PREFLIGHT_COMPRESSION_STATUS_TEMPLATE, IDLE_COMPACTION_STATUS_TEMPLATE, + COMPRESSION_RETRY_TOO_LARGE_STATUS_TEMPLATE, COMPRESSION_RETRY_MESSAGES_STATUS_TEMPLATE, COMPRESSION_RETRY_TOKENS_STATUS_TEMPLATE, - COMPRESSION_RETRY_CONTEXT_REDUCED_STATUS_TEMPLATE, - ) + COMPRESSION_RETRY_CONTEXT_REDUCED_STATUS_TEMPLATE) ), - re.IGNORECASE, -) + re.IGNORECASE) def _gateway_compression_progress_notices_enabled() -> bool: @@ -426,10 +376,7 @@ def _gateway_compression_progress_notices_enabled() -> bool: compression_cfg = config.get("compression") if isinstance(config, dict) else None if isinstance(compression_cfg, dict): return str(compression_cfg.get("progress_notices", False)).strip().lower() in { - "true", - "1", - "yes", - "on", + "true", "1", "yes", "on" } except Exception: pass @@ -464,12 +411,10 @@ _GATEWAY_PROVIDER_POLICY_RE = re.compile( _GATEWAY_AUTH_ERROR_RE = re.compile( r"(provider\s+authentication\s+failed|incorrect\s+api\s+key|invalid\s+api\s+key|\b401\b)", - re.IGNORECASE, -) + re.IGNORECASE) _GATEWAY_RATE_LIMIT_RE = re.compile( - r"(rate\s+limit|rate-limited|\b429\b|quota|usage\s+limit)", - re.IGNORECASE, + r"(rate\s+limit|rate-limited|\b429\b|quota|usage\s+limit)", re.IGNORECASE ) # Connection-failure markers: the first 8 also anchor the provider-failure envelope shape below. @@ -477,19 +422,15 @@ _CONNECTION_ERROR_MARKERS = ( r"(?:\w+\.)?(?:api\s*)?connection\s*(?:error|timeout)", r"(?:\w+\.)?connect\s*(?:error|timeout)", r"connection\s+refused", r"connection\s+reset", r"connection\s+aborted", r"actively\s+refused", r"winerror\s+10061", r"errno\s+111", r"no\s+route\s+to\s+host", r"network\s+is\s+unreachable", - r"cannot\s+connect", r"failed\s+to\s+establish", r"could\s+not\s+connect", -) + r"cannot\s+connect", r"failed\s+to\s+establish", r"could\s+not\s+connect") _GATEWAY_CONNECTION_ERROR_RE = re.compile("(" + "|".join(_CONNECTION_ERROR_MARKERS) + ")", re.IGNORECASE) _GATEWAY_SECRET_PATTERNS = ( re.compile(r"\bsk-[A-Za-z0-9][A-Za-z0-9_\-]{12,}\b"), - re.compile(r"\bgh[pousr]_[A-Za-z0-9_]{20,}\b"), - re.compile(r"\bxapp-\d+-[A-Za-z0-9\-]{20,}\b"), - re.compile(r"\bxox[baprs]-[A-Za-z0-9\-]{20,}\b"), - re.compile(r"\bhf_[A-Za-z0-9]{20,}\b"), + re.compile(r"\bgh[pousr]_[A-Za-z0-9_]{20,}\b"), re.compile(r"\bxapp-\d+-[A-Za-z0-9\-]{20,}\b"), + re.compile(r"\bxox[baprs]-[A-Za-z0-9\-]{20,}\b"), re.compile(r"\bhf_[A-Za-z0-9]{20,}\b"), re.compile(r"\bglpat-[A-Za-z0-9_\-]{20,}\b"), - re.compile(r"(?i)\b(Bearer\s+)[A-Za-z0-9._\-]{20,}\b"), -) + re.compile(r"(?i)\b(Bearer\s+)[A-Za-z0-9._\-]{20,}\b")) def _ensure_windows_gateway_venv_imports() -> None: @@ -547,9 +488,7 @@ def _gateway_platform_value(platform: Any) -> str: def _non_conversational_metadata( - metadata: Optional[Dict[str, Any]] = None, - *, - platform: Any = None, + metadata: Optional[Dict[str, Any]] = None, *, platform: Any = None ) -> Optional[Dict[str, Any]]: """Mark Discord lifecycle/status sends without changing other platforms.""" if _gateway_platform_value(platform) != "discord": @@ -574,8 +513,7 @@ def _interim_metadata( def _seed_hygiene_system_prompt( - agent: Any, - session_row: Optional[Dict[str, Any]], + agent: Any, session_row: Optional[Dict[str, Any]] ) -> bool: """Keep gateway hygiene from rebuilding a live session's system prompt. @@ -597,8 +535,7 @@ def _seed_hygiene_system_prompt( _TRANSIENT_NETWORK_ERROR_CLASS_NAMES = frozenset({ "TimedOut", "NetworkError", "ReadError", "WriteError", "ConnectError", "ConnectTimeout", "ReadTimeout", "WriteTimeout", "PoolTimeout", "RemoteProtocolError", "ServerDisconnectedError", - "ClientConnectorError", "ClientOSError", -}) + "ClientConnectorError", "ClientOSError"}) def _is_transient_network_error(exc: BaseException) -> bool: @@ -642,8 +579,7 @@ def _gateway_loop_exception_handler( task_name or "", type(exc).__name__, exc, - exc_info=(type(exc), exc, exc.__traceback__), - ) + exc_info=(type(exc), exc, exc.__traceback__)) return # Fall back to the default handler for anything we don't recognise. loop.default_exception_handler(context) @@ -682,13 +618,8 @@ def _redact_approval_command(cmd: "str | None") -> str: def _format_exec_approval_fallback( - command: str, - description: str, - command_prefix: str, - *, - allow_permanent: bool = True, - allow_session: bool = True, - smart_denied: bool = False, + command: str, description: str, command_prefix: str, *, allow_permanent: bool = True, + 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 @@ -717,8 +648,7 @@ _PROVIDER_ERROR_REPLIES = ( "error out of chat; check gateway logs for details or try rephrasing."), (_GATEWAY_RATE_LIMIT_RE, "⏱️ The model provider is rate-limiting requests. Please wait a moment and try again."), (_GATEWAY_CONNECTION_ERROR_RE, "⚠️ The model server is not responding — it looks like the configured " - "model endpoint is not running or is unreachable."), -) + "model endpoint is not running or is unreachable.")) def _gateway_provider_error_reply(text: str) -> str: @@ -736,14 +666,12 @@ def _gateway_provider_error_reply(text: str) -> str: _PROVIDER_ERROR_MARKERS = ( r"api\s+(?:call\s+)?failed", r"provider\s+authentication\s+failed", r"non-retryable\s+error", r"rate\s+limited\s+after\s+\d+\s+retries", r"error\s+code\s*:", r"http\s*\d{3}\b", - r"incorrect\s+api\s+key", r"invalid\s+api\s+key", -) + r"incorrect\s+api\s+key", r"invalid\s+api\s+key") _GATEWAY_PROVIDER_ERROR_SHAPE_RE = re.compile( r"^\s*(\W*\s*)?(" + "|".join(_PROVIDER_ERROR_MARKERS + _CONNECTION_ERROR_MARKERS[:8] + (r"all\s+connection\s+attempts\s+failed",)) + ")", - re.IGNORECASE, -) + re.IGNORECASE) def _looks_like_gateway_provider_error(text: str) -> bool: @@ -896,11 +824,7 @@ def _clarify_send_then_wait(fut, *, clarify_id: str, session_key: str, clarify_m def _resolve_progress_thread_id( - platform: Any, - source_thread_id: Any, - event_message_id: Any, - *, - reply_in_thread: bool = True, + platform: Any, source_thread_id: Any, event_message_id: Any, *, reply_in_thread: bool = True ) -> Optional[str]: """Return thread/root ID that progress/status bubbles should target. @@ -912,9 +836,7 @@ def _resolve_progress_thread_id( platform_key = str(platform_value or "").lower() if not reply_in_thread: if ( - source_thread_id - and event_message_id - and str(source_thread_id) == str(event_message_id) + source_thread_id and event_message_id and str(source_thread_id) == str(event_message_id) ): return None return str(source_thread_id) if source_thread_id else None @@ -938,13 +860,8 @@ def _has_platform_display_override(user_config: dict, platform_key: str, setting def _resolve_gateway_display_bool( - user_config: dict, - platform_key: str, - setting: str, - *, - default: bool = False, - platform: Any = None, - require_platform_override_for: set[Any] | None = None, + user_config: dict, platform_key: str, setting: str, *, default: bool = False, + platform: Any = None, require_platform_override_for: set[Any] | None = None, ) -> bool: """Resolve a boolean display setting with optional platform-only opt-in. @@ -953,8 +870,7 @@ def _resolve_gateway_display_bool( """ current_platform = _gateway_platform_value(platform or platform_key) platform_only = { - _gateway_platform_value(candidate) - for candidate in (require_platform_override_for or set()) + _gateway_platform_value(candidate) for candidate in (require_platform_override_for or set()) } if ( current_platform in platform_only @@ -1111,10 +1027,7 @@ def _float_env(name: str, default: float) -> float: def _stamp_hygiene_compression_provenance( - agent: Any, - desc: str, - provenance: "ActivityProvenance", - debug_label: str, + agent: Any, desc: str, provenance: "ActivityProvenance", debug_label: str ) -> None: """Best-effort activity provenance stamp for hygiene compression transitions.""" try: @@ -1124,10 +1037,7 @@ def _stamp_hygiene_compression_provenance( def _is_fresh_gateway_interruption( - value: Any, - *, - now: Optional[float] = None, - window_secs: Optional[float] = None, + value: Any, *, now: Optional[float] = None, window_secs: Optional[float] = None ) -> bool: """True when an interruption marker is fresh enough to auto-continue. @@ -1148,10 +1058,7 @@ def _is_fresh_gateway_interruption( def build_resume_recovery_note( - reason: Optional[str], - message: str = "", - *, - interactive: bool = True, + reason: Optional[str], message: str = "", *, interactive: bool = True ) -> str: """Build the resume-pending recovery system note for an interrupted turn. @@ -1167,8 +1074,7 @@ def build_resume_recovery_note( ) if message: resume_guidance = ( - "Address the user's NEW message below FIRST and focus " - "on what the user is asking now." + "Address the user's NEW message below FIRST and focus on what the user is asking now." ) tail_guidance = ( "Do NOT re-execute old tool calls — skip any " @@ -1192,8 +1098,7 @@ def build_resume_recovery_note( ) tail_guidance = ( "Do NOT re-run tool calls whose results already " - "appear in the history — resume from the first step " - "that has no recorded result." + "appear in the history — resume from the first step that has no recorded result." ) return ( f"[System note: The previous turn was interrupted by " @@ -1206,10 +1111,7 @@ def build_resume_recovery_note( def _prepare_resume_pending_message( - reason: Optional[str], - message: Optional[str], - *, - interactive: bool = True, + reason: Optional[str], message: Optional[str], *, interactive: bool = True ) -> tuple[str, str]: """Return the recovery message and the user text to persist. @@ -1217,8 +1119,7 @@ def _prepare_resume_pending_message( sanitizer every call. Real user text: persist clean words so the transcript stays scaffold-free. """ recovery_message = build_resume_recovery_note( - reason, message or "", interactive=interactive, - ) + reason, message or "", interactive=interactive) persist_message = ( message if isinstance(message, str) and message.strip() else recovery_message ) @@ -1235,15 +1136,11 @@ _ASSISTANT_REPLAY_FIELDS: tuple[str, ...] = ( "reasoning_details", "codex_reasoning_items", "codex_message_items", - "finish_reason", -) + "finish_reason") def _build_replay_entry( - role: str, - content: Any, - msg: Dict[str, Any], - preserve_timestamp: bool = False, + role: str, content: Any, msg: Dict[str, Any], preserve_timestamp: bool = False ) -> Dict[str, Any]: """Build a replay entry for a non-tool-calling message, preserving ``_ASSISTANT_REPLAY_FIELDS``. @@ -1357,9 +1254,7 @@ def _message_timestamps_enabled(user_config: Optional[dict]) -> bool: def _build_gateway_agent_history( - history: List[Dict[str, Any]], - *, - channel_prompt: Optional[str] = None, + history: List[Dict[str, Any]], *, channel_prompt: Optional[str] = None, inject_timestamps: bool = False, ) -> tuple[List[Dict[str, Any]], Optional[str]]: """Convert stored gateway transcript rows into agent replay messages. @@ -1370,8 +1265,7 @@ def _build_gateway_agent_history( from hermes_time import get_timezone as _get_msg_tz from gateway.message_timestamps import ( - render_user_content_with_timestamp as _render_msg_ts, - ) + render_user_content_with_timestamp as _render_msg_ts) _msg_tz = _get_msg_tz() agent_history: List[Dict[str, Any]] = [] @@ -1444,8 +1338,7 @@ def _build_gateway_agent_history( def _select_cached_agent_history( - persisted_history: List[Dict[str, Any]], - live_history: Any, + persisted_history: List[Dict[str, Any]], live_history: Any ) -> List[Dict[str, Any]]: """Prefer a cached live transcript only when it is longer and has at least one real, non-ephemeral unpersisted row; otherwise return ``persisted_history`` unchanged. @@ -1521,9 +1414,7 @@ def _last_transcript_timestamp(history: Optional[List[Dict[str, Any]]]) -> Any: # 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. _AUTO_APPEND_MEDIA_TOOL_NAMES = { - "text_to_speech", - "text_to_speech_tool", - "image_generate", + "text_to_speech", "text_to_speech_tool", "image_generate" } # ---- helpers: detect interrupted tool tails & auto-continue noise ---------- @@ -1579,8 +1470,7 @@ _TOOL_MEDIA_RE = re.compile( r'mp4|mov|avi|mkv|webm|ogg|opus|mp3|wav|m4a|' r'flac|epub|pdf|zip|rar|7z|docx?|xlsx?|pptx?|' r'txt|csv|apk|ipa))', - re.IGNORECASE, -) + re.IGNORECASE) # Shared with cron delivery and gateway background tasks — the repair must run on every surface @@ -1591,8 +1481,7 @@ from gateway.media_repair import ( # noqa: E402 def _collect_auto_append_media_tags( - messages: List[Dict[str, Any]], - history_offset: int = 0, + messages: List[Dict[str, Any]], history_offset: int = 0, history_media_paths: Optional[set] = None, ) -> tuple[List[str], bool]: """Collect real media tags from current-turn producer-tool results only. @@ -1711,8 +1600,7 @@ def _ensure_ssl_certs() -> None: if os.path.exists(configured_cert): return # user already configured it to a real file logging.getLogger(__name__).warning( - "Ignoring stale SSL_CERT_FILE=%r because the path does not exist", - configured_cert, + "Ignoring stale SSL_CERT_FILE=%r because the path does not exist", configured_cert ) os.environ.pop("SSL_CERT_FILE", None) @@ -1823,8 +1711,7 @@ def _reload_runtime_env_preserving_config_authority() -> None: return load_hermes_dotenv( - hermes_home=_hermes_home, - project_env=Path(__file__).resolve().parents[1] / '.env', + hermes_home=_hermes_home, project_env=Path(__file__).resolve().parents[1] / '.env' ) _bridge_max_turns_from_config(_hermes_home) @@ -1892,8 +1779,7 @@ def _multiplex_profile_homes(config: object) -> list[tuple[str, "Path"]]: return list( profiles_to_serve( - multiplex=True, - profile_allowlist=getattr(config, "multiplex_profile_allowlist", None), + multiplex=True, profile_allowlist=getattr(config, "multiplex_profile_allowlist", None) ) ) @@ -1955,8 +1841,7 @@ async def _reclaim_stale(runner: object) -> None: return try: ids = await reclaim( - "gateway stopped mid-handoff; state reclaimed at startup. " - "Re-run /handoff to try again." + "gateway stopped mid-handoff; state reclaimed at startup. Re-run /handoff to try again." ) except Exception: logger.debug("Stale-handoff reclaim raised", exc_info=True) @@ -1964,8 +1849,7 @@ async def _reclaim_stale(runner: object) -> None: if ids: logger.warning( "Reclaimed %d handoff(s) stranded in 'running' by a previous " - "gateway: %s", len(ids), ", ".join(str(i) for i in ids), - ) + "gateway: %s", len(ids), ", ".join(str(i) for i in ids)) def _terminal_scope_cwd(default: str = "") -> str: @@ -1996,11 +1880,8 @@ def _load_profile_secret_scope(profile_home: "Path") -> dict: @_contextmanager def _profile_runtime_scope( - profile_home: "Path", - prepared_secret_scope: Optional[dict] = None, - *, - hydrate_secrets: bool = True, -): + profile_home: "Path", prepared_secret_scope: Optional[dict] = None, *, + hydrate_secrets: bool = True): """Scope config/skills/memory AND credentials to a profile for one turn (multiplexed path only). (1) ``set_hermes_home_override`` redirects ``get_hermes_home()`` — a contextvar, so it reaches @@ -2010,8 +1891,7 @@ def _profile_runtime_scope( """ from hermes_constants import set_hermes_home_override, reset_hermes_home_override from agent.secret_scope import ( - set_secret_scope, - reset_secret_scope, + set_secret_scope, reset_secret_scope ) home_token = set_hermes_home_override(str(profile_home)) @@ -2066,8 +1946,7 @@ def load_gateway_config_for_runner() -> "GatewayConfig": return load_gateway_config() except Exception: logger.debug( - "multiplex default-scope config reload failed; using unscoped load", - exc_info=True, + "multiplex default-scope config reload failed; using unscoped load", exc_info=True ) return cfg @@ -2090,8 +1969,7 @@ async def _discover_gateway_mcp_tools(config: object) -> None: await loop.run_in_executor(None, copy_context().run, discover_mcp_tools) except Exception: logger.warning( - "MCP tool discovery failed for profile '%s'", profile_name, exc_info=True, - ) + "MCP tool discovery failed for profile '%s'", profile_name, exc_info=True) def _platform_has_bot_credential(platform: "Platform", platform_config: "PlatformConfig") -> bool: @@ -2119,8 +1997,7 @@ def _platform_has_bot_credential(platform: "Platform", platform_config: "Platfor if platform is Platform.MATRIX: extra = getattr(platform_config, "extra", None) or {} if all( - str(extra.get(key) or "").strip() - for key in ("homeserver", "user_id", "password") + str(extra.get(key) or "").strip() for key in ("homeserver", "user_id", "password") ): return True return False @@ -2159,8 +2036,7 @@ _SESSIONS_ENV_BRIDGE = {"cjk_fts": "HERMES_CJK_FTS", "search_slow_ms": "HERMES_S _DISPLAY_ENV_BRIDGE = { "busy_input_mode": "HERMES_GATEWAY_BUSY_INPUT_MODE", "busy_text_mode": "HERMES_GATEWAY_BUSY_TEXT_MODE", - "busy_ack_enabled": "HERMES_GATEWAY_BUSY_ACK_ENABLED", -} + "busy_ack_enabled": "HERMES_GATEWAY_BUSY_ACK_ENABLED"} def _bridge_section_to_env(section: Any, mapping: Dict[str, str]) -> None: @@ -2224,8 +2100,7 @@ def _bridge_terminal_config_to_env(_terminal_cfg: dict) -> None: "docker_shared_container_key": "TERMINAL_DOCKER_SHARED_CONTAINER_KEY", "docker_orphan_reaper": "TERMINAL_DOCKER_ORPHAN_REAPER", "sandbox_dir": "TERMINAL_SANDBOX_DIR", - "persistent_shell": "TERMINAL_PERSISTENT_SHELL", - } + "persistent_shell": "TERMINAL_PERSISTENT_SHELL"} for _cfg_key, _env_var in _terminal_env_map.items(): if _cfg_key not in _terminal_cfg: continue @@ -2348,13 +2223,11 @@ if _config_path.exists(): print( f" Warning: config.yaml → env bridge failed: " f"{type(_bridge_err).__name__}: {_bridge_err}", - file=sys.stderr, - ) + file=sys.stderr) print( " Gateway will fall back to .env values, which may not match " "your current config.yaml. Run `hermes doctor` to investigate.", - file=sys.stderr, - ) + file=sys.stderr) # Apply IPv4 preference if configured (before any HTTP clients are created). try: @@ -2400,39 +2273,26 @@ if not _configured_cwd or _configured_cwd in CWD_PLACEHOLDERS: "TERMINAL_DOCKER_MOUNT_CWD_TO_WORKSPACE", "false" ).lower() in {"true", "1", "yes"}, - home_fallback=str(Path.home()), - ) + home_fallback=str(Path.home())) if _resolved_cwd is None: os.environ.pop("TERMINAL_CWD", None) else: os.environ["TERMINAL_CWD"] = _resolved_cwd from gateway.config import ( - ChannelOverride, - Platform, - GatewayConfig, - PlatformConfig, - _getenv, - load_gateway_config, + ChannelOverride, Platform, GatewayConfig, PlatformConfig, _getenv, load_gateway_config ) from gateway.session import ( - AsyncSessionStore, - SessionStore, - SessionSource, - SessionContext, - build_session_key, + AsyncSessionStore, SessionStore, SessionSource, SessionContext, build_session_key ) from gateway.delivery import ( DeliveryRouter, resolve_delivery_transport, # noqa: F401 (re-exported: run_* mixins + tests resolve gateway.run.) ) from gateway.turn_lease import ( - SessionTurnLeaseRegistry, -) + SessionTurnLeaseRegistry) from gateway.session_state import ( - SessionState, - legacy_dict_property, - legacy_lease_token_property, + SessionState, legacy_dict_property, legacy_lease_token_property ) from gateway.authz_mixin import GatewayAuthorizationMixin from gateway.kanban_watchers import GatewayKanbanWatchersMixin @@ -2486,8 +2346,7 @@ _OWN_POLICY_OPEN_ENV = { Platform.WEIXIN: ("WEIXIN_DM_POLICY", "WEIXIN_GROUP_POLICY", "WEIXIN_ALLOW_ALL_USERS"), Platform.YUANBAO: ("YUANBAO_DM_POLICY", "YUANBAO_GROUP_POLICY", "YUANBAO_ALLOW_ALL_USERS"), Platform.QQBOT: (None, None, "QQ_ALLOW_ALL_USERS"), - Platform.WHATSAPP: ("WHATSAPP_DM_POLICY", "WHATSAPP_GROUP_POLICY", "WHATSAPP_ALLOW_ALL_USERS"), -} + Platform.WHATSAPP: ("WHATSAPP_DM_POLICY", "WHATSAPP_GROUP_POLICY", "WHATSAPP_ALLOW_ALL_USERS")} def _own_policy_open_startup_violation(config) -> Optional[str]: @@ -2501,12 +2360,10 @@ def _own_policy_open_startup_violation(config) -> Optional[str]: dm_env, group_env, allow_all_env = open_env extra = getattr(platform_config, "extra", None) or {} dm_policy = str( - extra.get("dm_policy") - or (_getenv(dm_env, "pairing") if dm_env else "pairing") + extra.get("dm_policy") or (_getenv(dm_env, "pairing") if dm_env else "pairing") ).strip().lower() group_policy = str( - extra.get("group_policy") - or (_getenv(group_env, "pairing") if group_env else "pairing") + extra.get("group_policy") or (_getenv(group_env, "pairing") if group_env else "pairing") ).strip().lower() if dm_policy != "open" and group_policy != "open": continue @@ -2514,8 +2371,7 @@ def _own_policy_open_startup_violation(config) -> Optional[str]: "GATEWAY_ALLOW_ALL_USERS", "" ).lower() in {"true", "1", "yes"} platform_opted_in = gateway_allow_all or ( - allow_all_env - and _getenv(allow_all_env, "").lower() in {"true", "1", "yes"} + allow_all_env and _getenv(allow_all_env, "").lower() in {"true", "1", "yes"} ) if platform_opted_in: continue @@ -2562,9 +2418,7 @@ def _resolve_runtime_agent_kwargs() -> dict: 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, + resolve_runtime_provider, format_runtime_provider_error, _get_model_config ) from hermes_cli.auth import AuthError, is_rate_limited_auth_error @@ -2632,8 +2486,7 @@ def _runtime_agent_kwargs(runtime: dict) -> dict: "command": runtime.get("command"), "args": list(runtime.get("args") or []), "credential_pool": runtime.get("credential_pool"), - "request_overrides": runtime.get("request_overrides"), - } + "request_overrides": runtime.get("request_overrides")} @dataclasses.dataclass(frozen=True) @@ -2701,13 +2554,8 @@ def _resolve_gateway_model_context(model: Optional[str] = None) -> _GatewayModel from hermes_cli.route_identity import should_clear_context_pin if should_clear_context_pin( - configured_model, - resolved_model, - configured_base_url, - base_url, - configured_provider, - provider, - ): + configured_model, resolved_model, configured_base_url, base_url, + configured_provider, provider): config_context_length = None except Exception: config_context_length = None @@ -2717,9 +2565,7 @@ def _resolve_gateway_model_context(model: Optional[str] = None) -> _GatewayModel from hermes_cli.config import get_custom_provider_context_length custom_ctx = get_custom_provider_context_length( - model=resolved_model, - base_url=base_url, - custom_providers=custom_providers, + model=resolved_model, base_url=base_url, custom_providers=custom_providers ) if custom_ctx: config_context_length = custom_ctx @@ -2727,13 +2573,9 @@ def _resolve_gateway_model_context(model: Optional[str] = None) -> _GatewayModel pass context_length = get_model_context_length( - resolved_model, - base_url=base_url or "", - api_key=api_key or "", - config_context_length=config_context_length, - provider=provider or "", - custom_providers=custom_providers, - ) + resolved_model, base_url=base_url or "", api_key=api_key or "", + config_context_length=config_context_length, provider=provider or "", + custom_providers=custom_providers) if config_context_length is not None: context_source = "config" elif context_length == DEFAULT_FALLBACK_CONTEXT: @@ -2742,19 +2584,14 @@ def _resolve_gateway_model_context(model: Optional[str] = None) -> _GatewayModel context_source = "detected" return _GatewayModelContext( - model=resolved_model, - provider=provider or "", - base_url=base_url or "", - context_length=context_length, - context_source=context_source, - ) + model=resolved_model, provider=provider or "", base_url=base_url or "", + context_length=context_length, context_source=context_source) def _resolve_runtime_agent_kwargs_for_provider(provider: str) -> dict: """Resolve runtime credentials for a specific provider (e.g. from channel override).""" from hermes_cli.runtime_provider import ( - resolve_runtime_provider, - format_runtime_provider_error, + resolve_runtime_provider, format_runtime_provider_error ) try: runtime = resolve_runtime_provider(requested=provider) @@ -2764,8 +2601,7 @@ def _resolve_runtime_agent_kwargs_for_provider(provider: str) -> dict: **_runtime_agent_kwargs(runtime), "request_overrides": dict(runtime.get("request_overrides") or {}), "capabilities": dict(runtime.get("capabilities") or {}), - "max_tokens": runtime.get("max_output_tokens"), - } + "max_tokens": runtime.get("max_output_tokens")} def _deep_merge_request_overrides(base: Optional[dict], override: Optional[dict]) -> dict: @@ -2791,9 +2627,7 @@ def _credential_pool_for_provider(provider: Optional[str]): ) except Exception: logger.debug( - "Failed to resolve credential pool for provider=%s", - provider, - exc_info=True, + "Failed to resolve credential pool for provider=%s", provider, exc_info=True ) return None @@ -2813,17 +2647,14 @@ def _try_resolve_fallback_provider() -> dict | None: from hermes_cli.fallback_config import resolve_entry_api_key runtime = resolve_runtime_provider( - requested=entry.get("provider"), - explicit_base_url=entry.get("base_url"), - explicit_api_key=resolve_entry_api_key(entry), - ) + requested=entry.get("provider"), explicit_base_url=entry.get("base_url"), + explicit_api_key=resolve_entry_api_key(entry)) # Log the literal config `provider`, not the resolved runtime category: an Ollama # fallback resolves via the OpenAI-compatible path and would log as "openrouter". logger.info( "Fallback provider resolved: %s model=%s", entry.get("provider") or runtime.get("provider"), - entry.get("model"), - ) + entry.get("model")) return {**_runtime_agent_kwargs(runtime), "model": entry.get("model")} except Exception as fb_exc: logger.debug("Fallback entry %s failed: %s", entry.get("provider"), fb_exc) @@ -2864,8 +2695,7 @@ def _event_media_is_stt_input(event, index: int) -> bool: if message_type in {MessageType.AUDIO, MessageType.DOCUMENT}: return False return ( - message_type == MessageType.VOICE - or _event_media_type_at(event, index).startswith("audio/") + message_type == MessageType.VOICE or _event_media_type_at(event, index).startswith("audio/") ) @@ -2894,11 +2724,7 @@ def _build_media_placeholder(event) -> str: def _build_document_context_note( - display_name: str, - agent_path: str, - mtype: str, - *, - content_inlined: bool = True, + display_name: str, agent_path: str, mtype: str, *, content_inlined: bool = True ) -> str: """Context note prepended to a user turn when they attach a document. @@ -2968,8 +2794,7 @@ async def _probe_audio_duration(path: str) -> Optional[str]: proc = await asyncio.create_subprocess_exec( "ffprobe", "-v", "error", "-show_entries", "format=duration", "-of", "default=noprint_wrappers=1:nokey=1", path, - stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, - ) + stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE) stdout, _ = await asyncio.wait_for(proc.communicate(), timeout=5.0) if proc.returncode == 0: return _format_duration(float(stdout.decode().strip())) @@ -2997,10 +2822,7 @@ _INTERRUPT_REASON_GATEWAY_RESTART = "Gateway restarting" def _reap_gateway_turn_processes( - task_id: str, - process_baseline, - *, - source: str, + task_id: str, process_baseline, *, source: str, is_still_current: Optional[Callable[[], bool]] = None, ) -> int: """Reap only background processes created by one abandoned turn. @@ -3020,33 +2842,26 @@ def _reap_gateway_turn_processes( "Skipping reap for turn %s (%s): a newer turn already " "claimed this session; it owns its own baseline.", task_id, - source, - ) + source) return 0 except Exception: logger.debug( "is_still_current check failed for turn %s (%s); reaping anyway", task_id, source, - exc_info=True, - ) + exc_info=True) from tools.process_registry import process_registry try: killed = process_registry.kill_started_since( - task_id, - process_baseline, - source=source, + task_id, process_baseline, source=source ) except Exception: # Runs on a detached daemon thread (fire-and-forget from interrupt and timeout paths); an # uncaught exception would only reach threading.excepthook. Swallow and log normally. logger.warning( - "Failed to reap background processes for turn %s (%s)", - task_id, - source, - exc_info=True, + "Failed to reap background processes for turn %s (%s)", task_id, source, exc_info=True ) return 0 if killed: @@ -3054,8 +2869,7 @@ def _reap_gateway_turn_processes( "Reaped %d background process(es) created by abandoned turn %s (%s)", killed, task_id, - source, - ) + source) return killed @@ -3065,8 +2879,7 @@ _TURN_STACK_DUMP_FRAME_MARKERS = ( "_run_sync_with_timeout_lifecycle", "finalize_turn", "end_turn", - "run_in_session", -) + "run_in_session") def _dump_wedged_turn_stacks(task_id: str) -> None: @@ -3089,36 +2902,27 @@ def _dump_wedged_turn_stacks(task_id: str) -> None: dumped += 1 if dumped > 8: logger.error( - "Wedged-turn stack dump for task %s truncated: more than " - "8 candidate threads", - task_id, - ) + "Wedged-turn stack dump for task %s truncated: more than 8 candidate threads", + task_id) break logger.error( "Wedged-turn stack dump (task=%s thread=%s ident=%s):\n%s", task_id, names.get(ident, "?"), ident, - "".join(stack[-25:]), - ) + "".join(stack[-25:])) if dumped == 0: logger.error( "Wedged-turn stack dump for task %s: no thread with " "turn-machinery frames found (worker may have already exited)", - task_id, - ) + task_id) except Exception: logger.debug("Wedged-turn stack dump failed", exc_info=True) def _abandon_timed_out_gateway_turn( - *, - agent_holder, - task_id: str, - process_baseline, - worker_done: threading.Event, - timeout_fired: threading.Event, - cleanup_lock: threading.Lock, + *, agent_holder, task_id: str, process_baseline, worker_done: threading.Event, + timeout_fired: threading.Event, cleanup_lock: threading.Lock, is_still_current: Optional[Callable[[], bool]] = None, ) -> bool: """Interrupt one timed-out turn and reap only processes it created.""" @@ -3140,30 +2944,18 @@ def _abandon_timed_out_gateway_turn( try: _reap_gateway_turn_processes( - task_id, - process_baseline, - source="gateway_turn_timeout", - is_still_current=is_still_current, - ) + task_id, process_baseline, source="gateway_turn_timeout", + is_still_current=is_still_current) except Exception: logger.warning( - "Failed to reap background processes for timed-out turn %s", - task_id, - exc_info=True, + "Failed to reap background processes for timed-out turn %s", task_id, exc_info=True ) return True def _watch_gateway_turn_inactivity( - *, - agent_holder, - task_id: str, - process_baseline, - timeout: float, - worker_done: threading.Event, - timeout_fired: threading.Event, - cleanup_lock: threading.Lock, - poll_interval: float = 5.0, + *, agent_holder, task_id: str, process_baseline, timeout: float, worker_done: threading.Event, + timeout_fired: threading.Event, cleanup_lock: threading.Lock, poll_interval: float = 5.0, is_still_current: Optional[Callable[[], bool]] = None, ) -> None: """Thread watchdog that remains runnable when gateway asyncio is starved.""" @@ -3180,26 +2972,17 @@ def _watch_gateway_turn_inactivity( if idle_seconds < timeout: continue _abandon_timed_out_gateway_turn( - agent_holder=agent_holder, - task_id=task_id, - process_baseline=process_baseline, - worker_done=worker_done, - timeout_fired=timeout_fired, - cleanup_lock=cleanup_lock, - is_still_current=is_still_current, - ) + agent_holder=agent_holder, task_id=task_id, process_baseline=process_baseline, + worker_done=worker_done, timeout_fired=timeout_fired, cleanup_lock=cleanup_lock, + is_still_current=is_still_current) return _CONTROL_INTERRUPT_MESSAGES = frozenset( { - _INTERRUPT_REASON_STOP.lower(), - _INTERRUPT_REASON_RESET.lower(), - _INTERRUPT_REASON_TIMEOUT.lower(), - _INTERRUPT_REASON_SSE_DISCONNECT.lower(), - _INTERRUPT_REASON_GATEWAY_SHUTDOWN.lower(), - _INTERRUPT_REASON_GATEWAY_RESTART.lower(), - } + _INTERRUPT_REASON_STOP.lower(), _INTERRUPT_REASON_RESET.lower(), + _INTERRUPT_REASON_TIMEOUT.lower(), _INTERRUPT_REASON_SSE_DISCONNECT.lower(), + _INTERRUPT_REASON_GATEWAY_SHUTDOWN.lower(), _INTERRUPT_REASON_GATEWAY_RESTART.lower()} ) @@ -3406,15 +3189,11 @@ def _checkpoint_agent_kwargs(config: dict | None) -> dict: return { "checkpoints_enabled": cp_cfg.get("enabled", defaults["enabled"]), "checkpoint_max_snapshots": cp_cfg.get( - "max_snapshots", defaults["max_snapshots"], - ), + "max_snapshots", defaults["max_snapshots"]), "checkpoint_max_total_size_mb": cp_cfg.get( - "max_total_size_mb", defaults["max_total_size_mb"], - ), + "max_total_size_mb", defaults["max_total_size_mb"]), "checkpoint_max_file_size_mb": cp_cfg.get( - "max_file_size_mb", defaults["max_file_size_mb"], - ), - } + "max_file_size_mb", defaults["max_file_size_mb"])} def _load_gateway_runtime_config() -> dict: @@ -3448,10 +3227,7 @@ def _resolve_gateway_model(config: dict | None = None) -> str: def _channel_override_lookup_keys( - chat_id: str, - *, - thread_id: Optional[str] = None, - parent_id: Optional[str] = None, + chat_id: str, *, thread_id: Optional[str] = None, parent_id: Optional[str] = None ) -> list[str]: """Ordered, de-duplicated ``channel_overrides`` lookup keys. @@ -3472,11 +3248,7 @@ def _channel_override_lookup_keys( def _get_channel_override( - config: GatewayConfig, - platform: Platform, - chat_id: str, - *, - thread_id: Optional[str] = None, + config: GatewayConfig, platform: Platform, chat_id: str, *, thread_id: Optional[str] = None, parent_id: Optional[str] = None, ) -> Optional[ChannelOverride]: """Per-channel override for this platform/chat_id, or None. @@ -3530,9 +3302,7 @@ def _parse_session_key(session_key: str) -> "dict | None": 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], + "platform": parts[2], "chat_type": parts[3], "chat_id": parts[4] } if len(parts) > 5 and parts[3] in {"dm", "thread"}: result["thread_id"] = parts[5] @@ -3549,11 +3319,7 @@ def _shorten_command_for_display(command: str, limit: int = 80) -> str: def _format_concise_process_notification( - session_id: str, - command: str, - exit_code, - output: str, - duration_seconds=None, + session_id: str, command: str, exit_code, output: str, duration_seconds=None ) -> str: """One-line completion message for the ``concise`` display mode. @@ -3636,10 +3402,7 @@ def _drain_gateway_watch_events(completion_queue) -> "list[dict]": break evt_type = evt.get("type", "completion") if evt_type in { - "watch_match", - "watch_disabled", - "watch_overflow_tripped", - "watch_overflow_released", + "watch_match", "watch_disabled", "watch_overflow_tripped", "watch_overflow_released" }: watch_events.append(evt) elif evt_type == "async_delegation": @@ -3657,10 +3420,7 @@ _gateway_runner_ref: _weakref.ref = lambda: None def _normalize_empty_agent_response( - agent_result: dict, - response: str, - *, - history_len: int = 0, + agent_result: dict, response: str, *, history_len: int = 0 ) -> str: """Normalize empty/None agent responses into user-facing messages. @@ -3686,14 +3446,12 @@ def _normalize_empty_agent_response( return ( "⚠️ Session storage was temporarily unavailable, so this " "turn was stopped to protect your conversation history. " - "Please check available disk space, then send your " - "message again." + "Please check available disk space, then send your message again." ) return ( "⚠️ Session storage was temporarily unavailable, so this " "turn was stopped to protect your conversation history. " - "Your message should already be saved — please send it " - "again in a moment." + "Your message should already be saved — please send it again in a moment." ) is_context_failure = any( p in error_str @@ -3702,8 +3460,7 @@ def _normalize_empty_agent_response( if is_context_failure: return ( "⚠️ Session too large for the model's context window.\n" - "Use /compact to compress the conversation, or " - "/reset to start fresh." + "Use /compact to compress the conversation, or /reset to start fresh." ) return ( f"The request failed: {str(error_detail)[:300]}\n" @@ -3779,8 +3536,7 @@ def _should_clear_resume_pending_after_turn(agent_result: dict) -> bool: def _preserve_queued_followup_history_offset( - current_result: dict, - followup_result: dict, + current_result: dict, followup_result: dict ) -> dict: """Carry the outer history offset through queued follow-up drains. @@ -3823,8 +3579,7 @@ async def _dispose_unused_adapter(adapter: "BasePlatformAdapter | None") -> None logger.debug( "Adapter dispose raised on unowned adapter %r", getattr(adapter, "name", type(adapter).__name__), - exc_info=True, - ) + exc_info=True) # Max seconds between platform reconnect retries (primary watcher and @@ -3907,8 +3662,7 @@ def _command_origin_for_source(source: Any) -> Optional[dict]: "platform": platform, "chat_id": str(chat_id), "chat_name": getattr(source, "chat_name", None), - "thread_id": getattr(source, "thread_id", None), - } + "thread_id": getattr(source, "thread_id", None)} except Exception: pass return None @@ -3942,8 +3696,7 @@ _BUILTIN_ADAPTERS: dict[Platform, tuple[str, str, str, str]] = { Platform.QQBOT: ("qqbot", "QQAdapter", "check_qq_requirements", "QQBot: aiohttp/httpx missing or QQ_APP_ID/QQ_CLIENT_SECRET not configured"), Platform.YUANBAO: ("yuanbao", "YuanbaoAdapter", "WEBSOCKETS_AVAILABLE", - "Yuanbao: websockets not installed. Run: pip install websockets"), -} + "Yuanbao: websockets not installed. Run: pip install websockets")} def _instantiate_builtin_adapter(platform: Platform, config: Any) -> Optional[BasePlatformAdapter]: @@ -3966,23 +3719,11 @@ def _instantiate_builtin_adapter(platform: Platform, config: Any) -> Optional[Ba class GatewayRunner( - GatewayAuthorizationMixin, - GatewayKanbanWatchersMixin, - GatewaySlashCommandsMixin, - GatewayVoiceMixin, - GatewayAdapterLifecycleMixin, - GatewayTopicThreadsMixin, - GatewayTurnMixin, - GatewayShutdownMixin, - GatewayBusySessionMixin, - GatewayConfigLoadersMixin, - GatewayStartupMixin, - GatewaySessionWatchersMixin, - GatewayNotificationsMixin, - GatewayInboundMixin, - GatewayGoalsMixin, - GatewayAgentCacheMixin, -): + GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, GatewaySlashCommandsMixin, + GatewayVoiceMixin, GatewayAdapterLifecycleMixin, GatewayTopicThreadsMixin, GatewayTurnMixin, + GatewayShutdownMixin, GatewayBusySessionMixin, GatewayConfigLoadersMixin, GatewayStartupMixin, + GatewaySessionWatchersMixin, GatewayNotificationsMixin, GatewayInboundMixin, GatewayGoalsMixin, + GatewayAgentCacheMixin): """Main gateway controller: manages adapter lifecycles, routes messages to/from the agent.""" # Class-level defaults so partial construction in tests doesn't @@ -4162,9 +3903,7 @@ class GatewayRunner( self.session_store = SessionStore( self.config.sessions_dir, self.config, has_active_processes_fn=lambda key: process_registry.has_active_for_session( - key, max_active_age=_bg_max_age_seconds, - ), - ) + key, max_active_age=_bg_max_age_seconds)) # One enforced loop-side boundary for the synchronous SessionStore: sync helpers keep using # ``session_store`` directly; async gateway handlers call this facade and await every op. self._async_session_store = AsyncSessionStore(self.session_store) @@ -4342,8 +4081,7 @@ class GatewayRunner( from gateway.session_db_recovery import RecoverableHandleCache self._session_db_handle_cache = RecoverableHandleCache( - handles=self._session_db_handles, - lock=self._session_db_handles_lock, + handles=self._session_db_handles, lock=self._session_db_handles_lock ) try: self._open_session_db_for_active_scope(raise_on_error=True) @@ -4366,8 +4104,7 @@ class GatewayRunner( if _sess_cfg.get("auto_archive", False): self._session_db._db.maybe_auto_archive( idle_days=float(_sess_cfg.get("auto_archive_days", 3)), - min_interval_hours=int(_sess_cfg.get("min_interval_hours", 24)), - ) + min_interval_hours=int(_sess_cfg.get("min_interval_hours", 24))) if _sess_cfg.get("auto_prune", False): # Construction-time, before the loop serves traffic; sync DB is fine. self._session_db._db.maybe_auto_prune_and_vacuum( @@ -4377,8 +4114,7 @@ class GatewayRunner( _sess_cfg.get("min_vacuum_interval_days", 30) ), vacuum=bool(_sess_cfg.get("vacuum_after_prune", True)), - sessions_dir=self.config.sessions_dir, - ) + sessions_dir=self.config.sessions_dir) except Exception as exc: logger.debug("state.db auto-maintenance skipped: %s", exc) @@ -4396,8 +4132,7 @@ class GatewayRunner( retention_days=int(_ckpt_cfg.get("retention_days", 7)), min_interval_hours=int(_ckpt_cfg.get("min_interval_hours", 24)), delete_orphans=False, - max_total_size_mb=int(_ckpt_cfg.get("max_total_size_mb", 500)), - ) + max_total_size_mb=int(_ckpt_cfg.get("max_total_size_mb", 500))) except Exception as exc: logger.debug("checkpoint auto-maintenance skipped: %s", exc) @@ -4458,8 +4193,7 @@ class GatewayRunner( # Compatibility for lightweight test runners built with # object.__new__ rather than GatewayRunner.__init__. cache = RecoverableHandleCache( - handles=self._session_db_handles, - lock=self._session_db_handles_lock, + handles=self._session_db_handles, lock=self._session_db_handles_lock ) self._session_db_handle_cache = cache @@ -4492,10 +4226,7 @@ class GatewayRunner( logger.info("SQLite session store recovered") return cache.get( - path, - _open, - raise_on_error=raise_on_error, - on_recovered=_recovered, + path, _open, raise_on_error=raise_on_error, on_recovered=_recovered ) @property @@ -4558,8 +4289,7 @@ class GatewayRunner( logger.info("Teams pipeline runtime bound to msgraph webhook ingress") elif self._teams_pipeline_runtime_error: logger.warning( - "Teams pipeline runtime unavailable: %s", - self._teams_pipeline_runtime_error, + "Teams pipeline runtime unavailable: %s", self._teams_pipeline_runtime_error ) def _warn_if_docker_media_delivery_is_risky(self) -> None: @@ -4639,11 +4369,9 @@ class GatewayRunner( except Exception: _profile = None return build_session_key( - source, - group_sessions_per_user=getattr(config, "group_sessions_per_user", True), + source, group_sessions_per_user=getattr(config, "group_sessions_per_user", True), thread_sessions_per_user=getattr(config, "thread_sessions_per_user", False), - profile=_profile, - ) + profile=_profile) # Telegram's General (pinned top) topic in forum-enabled private chats: clients variously omit @@ -4655,8 +4383,7 @@ class GatewayRunner( def _normalize_source_for_session_key( - self, - source: SessionSource, + self, source: SessionSource ) -> SessionSource: """Apply Telegram DM topic recovery to a source for session-key purposes. @@ -4710,11 +4437,8 @@ class GatewayRunner( def _update_runtime_status(self, gateway_state: Optional[str] = None, exit_reason: Optional[str] = None) -> None: _write_runtime_status_quiet( - gateway_state=gateway_state, - exit_reason=exit_reason, - restart_requested=self._restart_requested, - active_agents=self._active_work_count(), - ) + gateway_state=gateway_state, exit_reason=exit_reason, + restart_requested=self._restart_requested, active_agents=self._active_work_count()) def _persist_active_agents(self) -> None: """Persist the live in-flight agent count to ``gateway_state.json``. @@ -4813,10 +4537,7 @@ class GatewayRunner( def _is_user_authorized_for_source( - self, - source: SessionSource, - *, - allow_adapter_delegation: bool = True, + self, source: SessionSource, *, allow_adapter_delegation: bool = True ) -> bool: """Authorize under the live transport's profile, not the routed runtime. @@ -4830,8 +4551,7 @@ class GatewayRunner( if allow_adapter_delegation: return self._is_user_authorized(source) return self._is_user_authorized( - source, - allow_adapter_delegation=False, + source, allow_adapter_delegation=False ) authorization_home = getattr(source, "_authorization_profile_home", None) @@ -4853,8 +4573,7 @@ class GatewayRunner( "model": "Agent is running — wait or /stop first, then switch models.", "codex-runtime": ("Agent is running — wait or /stop first, then " "change runtime."), - "moa": "Agent is running — wait or /stop first, then run /moa.", - } + "moa": "Agent is running — wait or /stop first, then run /moa."} def _cache_session_source(self, session_key: str, source) -> None: @@ -4947,18 +4666,13 @@ class GatewayRunner( def _thread_metadata_for_source( - self, - source, - reply_to_message_id: Optional[str] = None, + self, source, reply_to_message_id: Optional[str] = None ) -> Optional[Dict[str, Any]]: """Build the metadata dict platforms need for thread-aware replies.""" metadata = self._thread_metadata_for_target( - getattr(source, "platform", None), - getattr(source, "chat_id", None), - 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), - ) + getattr(source, "platform", None), getattr(source, "chat_id", None), + 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's chat.startStream needs recipient_user_id/team_id, # which the relay connector fills from metadata.user_id/scope_id; the relay adapter's @@ -4985,13 +4699,8 @@ class GatewayRunner( return metadata def _thread_metadata_for_target( - self, - platform: Optional[Platform], - chat_id: Optional[str], - thread_id: Optional[str], - *, - chat_type: Optional[str] = None, - reply_to_message_id: Optional[str] = None, + self, platform: Optional[Platform], chat_id: Optional[str], thread_id: Optional[str], *, + chat_type: Optional[str] = None, reply_to_message_id: Optional[str] = None, adapter: Optional[Any] = None, ) -> Optional[Dict[str, Any]]: """Build thread metadata for synthetic sends that only have routing state.""" @@ -4999,11 +4708,7 @@ class GatewayRunner( return None metadata: Dict[str, Any] = {"thread_id": thread_id} if self._is_telegram_dm_topic_target( - platform, - chat_id, - thread_id, - chat_type=chat_type, - adapter=adapter, + platform, chat_id, thread_id, chat_type=chat_type, adapter=adapter ): metadata["telegram_dm_topic_reply_fallback"] = True # Telegram DM topic lanes need direct_messages_topic_id in metadata so synthetic/queued @@ -5021,12 +4726,8 @@ class GatewayRunner( @staticmethod def _is_telegram_dm_topic_target( - platform: Optional[Platform], - chat_id: Optional[str], - thread_id: Optional[str], - *, - chat_type: Optional[str] = None, - adapter: Optional[Any] = None, + platform: Optional[Platform], chat_id: Optional[str], thread_id: Optional[str], *, + chat_type: Optional[str] = None, adapter: Optional[Any] = None, ) -> bool: """Return True when a target is a Telegram private DM topic lane.""" if platform != Platform.TELEGRAM or thread_id is None: @@ -5064,8 +4765,7 @@ class GatewayRunner( # platforms are NOT listed here — they declare ``allow_update_command=True`` on their # ``PlatformEntry`` and are honored via the registry fallback in ``_handle_update_command``. _UPDATE_ALLOWED_PLATFORMS = frozenset({ - Platform.TELEGRAM, Platform.SLACK, Platform.WHATSAPP, - Platform.SIGNAL, Platform.MATRIX, + Platform.TELEGRAM, Platform.SLACK, Platform.WHATSAPP, Platform.SIGNAL, Platform.MATRIX, Platform.EMAIL, Platform.SMS, Platform.DINGTALK, Platform.FEISHU, Platform.WECOM, Platform.WECOM_CALLBACK, Platform.WEIXIN, Platform.BLUEBUBBLES, Platform.QQBOT, Platform.LOCAL, }) @@ -5103,8 +4803,7 @@ class GatewayRunner( message_id=str(context.source.message_id) if context.source.message_id else "", profile=getattr(context.source, "profile", "") or "", async_delivery=_async_delivery, - cron_session="", - ) + cron_session="") def _clear_session_env(self, tokens: list) -> None: """Restore session context variables to their pre-handler values.""" @@ -5116,10 +4815,7 @@ class GatewayRunner( loop = asyncio.get_running_loop() ctx = copy_context() return await loop.run_in_executor( - self._get_executor(), - ctx.run, - func, - *args, + self._get_executor(), ctx.run, func, *args ) def _get_executor(self) -> concurrent.futures.ThreadPoolExecutor: @@ -5135,8 +4831,7 @@ class GatewayRunner( executor = getattr(self, "_executor", None) if executor is None or getattr(executor, "_shutdown", False): executor = concurrent.futures.ThreadPoolExecutor( - max_workers=10, - thread_name_prefix="hermes-gateway", + max_workers=10, thread_name_prefix="hermes-gateway" ) self._executor = executor return executor @@ -5189,44 +4884,29 @@ class GatewayRunner( # them in at construction, so a mid-gateway edit would otherwise be silently ignored until some # other eviction. (section, key) tuples from the raw config dict; add new baked-in settings here. _CACHE_BUSTING_CONFIG_KEYS: tuple = ( - ("model", "context_length"), - ("model", "max_tokens"), - ("compression", "enabled"), - ("compression", "progress_notices"), - ("compression", "threshold"), - ("compression", "model_thresholds"), - ("compression", "threshold_tokens"), - ("compression", "codex_gpt55_autoraise"), - ("compression", "codex_app_server_auto"), + ("model", "context_length"), ("model", "max_tokens"), ("compression", "enabled"), + ("compression", "progress_notices"), ("compression", "threshold"), + ("compression", "model_thresholds"), ("compression", "threshold_tokens"), + ("compression", "codex_gpt55_autoraise"), ("compression", "codex_app_server_auto"), ("compression", "codex_responses_native"), - ("compression", "codex_responses_compact_threshold"), - ("compression", "in_place"), - ("compression", "checkpoint_required"), - ("compression", "micro_compact"), + ("compression", "codex_responses_compact_threshold"), ("compression", "in_place"), + ("compression", "checkpoint_required"), ("compression", "micro_compact"), ("compression", "micro_compact_every_n_turns"), - ("compression", "micro_compact_defrag_threshold_tokens"), - ("compression", "target_ratio"), - ("compression", "tail_mode"), - ("compression", "protect_last_n"), + ("compression", "micro_compact_defrag_threshold_tokens"), ("compression", "target_ratio"), + ("compression", "tail_mode"), ("compression", "protect_last_n"), ("compression", "proactive_prune_tokens"), ("compression", "proactive_prune_min_result_chars"), ("compression", "proactive_prune_min_reclaim_tokens"), - ("compression", "min_tail_user_messages"), - ("agent", "disabled_toolsets"), - ("memory", "provider"), - ("checkpoints", "enabled"), - ("checkpoints", "max_snapshots"), - ("checkpoints", "max_total_size_mb"), - ("checkpoints", "max_file_size_mb"), - ) + ("compression", "min_tail_user_messages"), ("agent", "disabled_toolsets"), + ("memory", "provider"), ("checkpoints", "enabled"), ("checkpoints", "max_snapshots"), + ("checkpoints", "max_total_size_mb"), ("checkpoints", "max_file_size_mb")) _HONCHO_CACHE_BUSTING_KEYS = ( "honcho.peer_name", "honcho.ai_peer", "honcho.pin_peer_name", "honcho.runtime_peer_prefix", - "honcho.user_peer_aliases", - ) + "honcho.user_peer_aliases") _HONCHO_CACHE_BUSTING_MEMO: dict[tuple[str, int | None], dict[str, Any]] = {} @@ -5276,18 +4956,13 @@ class GatewayRunner( from gateway.profile_routing import ProfileRouteRejected, match_profile_route try: matched = match_profile_route( - routes, - platform=source.platform.value, - guild_id=getattr(source, "guild_id", None), - chat_id=source.chat_id, - thread_id=getattr(source, "thread_id", None), - parent_chat_id=getattr(source, "parent_chat_id", None), - ) + routes, platform=source.platform.value, guild_id=getattr(source, "guild_id", None), + chat_id=source.chat_id, thread_id=getattr(source, "thread_id", None), + parent_chat_id=getattr(source, "parent_chat_id", None)) except Exception: logger.warning( "Profile route matching failed for %s/%s, falling back to default", - source.platform, source.chat_id, exc_info=True, - ) + source.platform, source.chat_id, exc_info=True) return None if matched: try: @@ -5297,22 +4972,19 @@ class GatewayRunner( "Rejecting profile route %r because the served-profile set " "could not be resolved", matched.name, - exc_info=True, - ) + 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.profile) raise ProfileRouteRejected(matched.name) return matched.profile logger.debug( "No profile route matched: platform=%s chat_id=%s thread_id=%s parent_chat_id=%s", source.platform.value, source.chat_id, - getattr(source, "thread_id", None), getattr(source, "parent_chat_id", None), - ) + getattr(source, "thread_id", None), getattr(source, "parent_chat_id", None)) return None def _resolve_profile_home_for_source(self, source: SessionSource) -> "Path": @@ -5324,9 +4996,7 @@ class GatewayRunner( """ from gateway.profile_routing import ProfileRouteRejected from hermes_cli.profiles import ( - get_active_profile_name, - get_profile_dir, - profile_exists, + get_active_profile_name, get_profile_dir, profile_exists ) from hermes_constants import get_hermes_home @@ -5352,8 +5022,7 @@ class GatewayRunner( explicit_profile, source.platform.value, source.chat_id, - getattr(source, "guild_id", None), - ) + getattr(source, "guild_id", None)) return get_hermes_home() return profile_dir except ProfileRouteRejected: @@ -5367,8 +5036,7 @@ class GatewayRunner( source.chat_id, getattr(source, "guild_id", None), explicit_profile or "(no profile)", - exc_info=True, - ) + exc_info=True) return get_hermes_home() @dataclasses.dataclass @@ -5410,11 +5078,7 @@ class GatewayRunner( def _run_planned_stop_watcher( - stop_event: threading.Event, - runner, - loop: asyncio.AbstractEventLoop, - shutdown_handler, - *, + 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. @@ -5427,8 +5091,7 @@ def _run_planned_stop_watcher( guard against re-triggering; the handler tolerates ``signal=None``. """ from gateway.status import ( - _get_planned_stop_marker_path, - planned_stop_marker_targets_self, + _get_planned_stop_marker_path, planned_stop_marker_targets_self ) marker_path = _get_planned_stop_marker_path() while not stop_event.is_set(): @@ -5470,10 +5133,8 @@ def _housekeeping_channel_directory(adapters, loop) -> None: # build_channel_directory is async (Slack web calls) and this is a background thread: # schedule onto the gateway loop and wait briefly so refresh failures still log. fut = safe_schedule_threadsafe( - build_channel_directory(adapters), loop, - logger=logger, - log_message="Channel directory refresh scheduling error", - ) + build_channel_directory(adapters), loop, logger=logger, + log_message="Channel directory refresh scheduling error") if fut is not None: fut.result(timeout=30) @@ -5481,28 +5142,19 @@ def _housekeeping_channel_directory(adapters, loop) -> None: def _housekeeping_media_caches() -> None: """Every platform media cache prunes on the same hourly cadence (24h max age).""" from gateway.platforms.base import ( - cleanup_audio_cache, - cleanup_document_cache, - cleanup_image_cache, - cleanup_screenshot_cache, - cleanup_video_cache, - ) + cleanup_audio_cache, cleanup_document_cache, cleanup_image_cache, cleanup_screenshot_cache, + cleanup_video_cache) from tools.tool_result_storage import cleanup_spillover_cache from tools.environments.local import cleanup_terminal_temp_cache from tools.bot_mode_dm import cleanup_bot_dm_cache from tools.bot_relay import cleanup_bot_relay_artifacts for cache_name, cleanup_fn in ( - ("Image", cleanup_image_cache), - ("Document", cleanup_document_cache), - ("Audio", cleanup_audio_cache), - ("Video", cleanup_video_cache), - ("Screenshot", cleanup_screenshot_cache), - ("Spillover", cleanup_spillover_cache), - ("Terminal temp", cleanup_terminal_temp_cache), - ("Bot DM", cleanup_bot_dm_cache), - ("Bot relay", cleanup_bot_relay_artifacts), - ): + ("Image", cleanup_image_cache), ("Document", cleanup_document_cache), + ("Audio", cleanup_audio_cache), ("Video", cleanup_video_cache), + ("Screenshot", cleanup_screenshot_cache), ("Spillover", cleanup_spillover_cache), + ("Terminal temp", cleanup_terminal_temp_cache), ("Bot DM", cleanup_bot_dm_cache), + ("Bot relay", cleanup_bot_relay_artifacts)): def _one(name=cache_name, fn=cleanup_fn): removed = fn(max_age_hours=24) if removed: @@ -5557,8 +5209,7 @@ def _housekeeping_auto_archive() -> None: try: _adb.maybe_auto_archive( idle_days=float(_sess_cfg.get("auto_archive_days", 3)), - min_interval_hours=int(_sess_cfg.get("min_interval_hours", 24)), - ) + min_interval_hours=int(_sess_cfg.get("min_interval_hours", 24))) finally: release_or_close(_adb) @@ -5573,8 +5224,7 @@ def _housekeeping_deferred_fts_retry() -> None: if callable(_retry) and _retry(): logger.info( "Deferred state.db FTS rebuild completed in-process for %s; full-text search restored.", - getattr(_sdb, "db_path", "state.db"), - ) + getattr(_sdb, "db_path", "state.db")) def _housekeeping_memory_trim() -> None: @@ -5595,8 +5245,7 @@ def _start_gateway_housekeeping(stop_event: threading.Event, adapters=None, loop 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), - (60, "Paste sweep", _housekeeping_paste_sweep), - ] + (60, "Paste sweep", _housekeeping_paste_sweep)] if cron_provider is not None: chores.append((5, "Misfire catch-up sweep", lambda: _housekeeping_misfire_catch_up(cron_provider, adapters, loop))) chores += [ @@ -5605,8 +5254,7 @@ def _start_gateway_housekeeping(stop_event: threading.Event, adapters=None, loop (60, "Org sync pull tick", _housekeeping_org_skill_sync), (60, "Auto-archive tick", _housekeeping_auto_archive), (1, "Deferred FTS retry tick", _housekeeping_deferred_fts_retry), - (1, "gateway housekeeping memory trim", _housekeeping_memory_trim), - ] + (1, "gateway housekeeping memory trim", _housekeeping_memory_trim)] logger.info("Gateway housekeeping started (interval=%ds)", interval) tick_count = 0 @@ -5635,8 +5283,7 @@ def _stop_cron_provider(provider) -> None: provider.stop() except SystemExit as exc: logger.warning( - "Cron provider stop() attempted to exit the gateway with code %s; ignoring", - exc.code, + "Cron provider stop() attempted to exit the gateway with code %s; ignoring", exc.code ) except Exception as exc: logger.debug("Cron provider stop() error: %s", exc) @@ -5694,8 +5341,7 @@ async def _shutdown_mcp_servers_nonblocking(timeout: float = 5.0) -> bool: logger.warning( "MCP shutdown did not finish within %.1fs; continuing gateway " "teardown (background thread will be reaped at process exit)", - timeout, - ) + timeout) return done @@ -5730,15 +5376,8 @@ def _replace_target_belongs_to_other_profile(existing_pid: int) -> bool: """ 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, - ) + _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() @@ -5747,17 +5386,14 @@ def _replace_target_belongs_to_other_profile(existing_pid: int) -> bool: record = _read_pid_record(_get_pid_path()) if not isinstance(record, dict) or not _record_looks_like_gateway(record): logger.warning( - "Refusing --replace: no valid gateway pid record to prove " - "ownership of PID %s.", - existing_pid, - ) + "Refusing --replace: no valid gateway pid record to prove ownership of PID %s.", + existing_pid) return True record_pid = _pid_from_record(record) if record_pid != existing_pid: logger.warning( - "Refusing --replace: pid record names %s, not target %s.", - record_pid, existing_pid, + "Refusing --replace: pid record names %s, not target %s.", record_pid, existing_pid ) return True @@ -5768,8 +5404,7 @@ def _replace_target_belongs_to_other_profile(existing_pid: int) -> bool: logger.warning( "Refusing --replace: pid record start-time does not match " "the live process %s (stale/PID-reuse record).", - existing_pid, - ) + existing_pid) return True recorded_home = record.get("hermes_home") @@ -5778,8 +5413,7 @@ def _replace_target_belongs_to_other_profile(existing_pid: int) -> bool: logger.warning( "Refusing --replace: pid record predates hermes_home " "stampings; ownership of PID %s unprovable.", - existing_pid, - ) + existing_pid) return True if not _same_hermes_home(recorded_home, our_home): @@ -5788,8 +5422,7 @@ def _replace_target_belongs_to_other_profile(existing_pid: int) -> bool: "HERMES_HOME (%s, ours %s). Remove the stale PID record or " "stop the owning profile explicitly.", recorded_home, - our_home, - ) + our_home) return True # Readable-argv consistency check (never authority): an explicit profile flag / HERMES_HOME= that @@ -5805,19 +5438,16 @@ def _replace_target_belongs_to_other_profile(existing_pid: int) -> bool: "Refusing --replace: target PID %s command line explicitly " "advertises a different profile than HERMES_HOME %s.", existing_pid, - our_home, - ) + our_home) return True return False except Exception: # Destructive action + unknown ownership => fail closed (#89315). logger.warning( - "cross-profile --replace ownership probe failed for PID %s; " - "refusing to signal", + "cross-profile --replace ownership probe failed for PID %s; refusing to signal", existing_pid, - exc_info=True, - ) + exc_info=True) return True @@ -5902,17 +5532,14 @@ async def _start_gateway_replace_existing_instance(existing_pid: int, replace: b 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, + get_process_start_time, remove_pid_file, terminate_pid ) if not replace: hermes_home = str(get_hermes_home()) logger.error( "Another gateway instance is already running (PID %d, HERMES_HOME=%s). " "Use 'hermes gateway restart' to replace it, or 'hermes gateway stop' first.", - existing_pid, hermes_home, - ) + existing_pid, hermes_home) print( f"\n❌ Gateway already running (PID {existing_pid}).\n" f" Use 'hermes gateway restart' to replace it,\n" @@ -5932,8 +5559,7 @@ async def _start_gateway_replace_existing_instance(existing_pid: int, replace: b "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(), - ) + _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) @@ -5977,8 +5603,7 @@ async def _start_gateway_replace_existing_instance(existing_pid: int, replace: b logger.error( "Old gateway (PID %d) still appears alive after SIGKILL; " "aborting replacement to avoid a duplicate gateway.", - existing_pid, - ) + existing_pid) _clear_takeover_marker_quiet() return False # Old gateway confirmed dead — reap orphaned children (POSIX; mirrors Windows taskkill /T) so @@ -6080,9 +5705,7 @@ def _start_gateway_make_shutdown_signal_handler(runner, _signal_initiated_shutdo # signal handler: stdlib + /proc only, no subprocesses (a sync `ps aux` here once blocked ~3s). try: from gateway.shutdown_forensics import ( - format_context_for_log, - snapshot_shutdown_context, - spawn_async_diagnostic, + format_context_for_log, snapshot_shutdown_context, spawn_async_diagnostic ) _shutdown_ctx = snapshot_shutdown_context(received_signal) except Exception as _e: @@ -6092,13 +5715,11 @@ def _start_gateway_make_shutdown_signal_handler(runner, _signal_initiated_shutdo if planned_takeover: logger.info( "Received %s as a planned --replace takeover — exiting cleanly", - _shutdown_ctx["signal"] if _shutdown_ctx else "SIGTERM", - ) + _shutdown_ctx["signal"] if _shutdown_ctx else "SIGTERM") elif planned_stop: logger.info( "Received %s as a planned gateway stop — exiting cleanly", - _shutdown_ctx["signal"] if _shutdown_ctx else "SIGTERM/SIGINT", - ) + _shutdown_ctx["signal"] if _shutdown_ctx else "SIGTERM/SIGINT") else: _signal_initiated_shutdown[0] = True # Mirror onto the runner so _stop_impl can suppress the gateway_state=stopped persist for @@ -6107,8 +5728,7 @@ def _start_gateway_make_shutdown_signal_handler(runner, _signal_initiated_shutdo runner._signal_initiated_shutdown = True logger.info( "Received %s — initiating shutdown", - _shutdown_ctx["signal"] if _shutdown_ctx else "SIGTERM/SIGINT", - ) + _shutdown_ctx["signal"] if _shutdown_ctx else "SIGTERM/SIGINT") # Always log who/what triggered the signal — the most useful line for "gateway keeps dying" # tickets. One line, key=value, parent_cmdline last (often long). @@ -6137,12 +5757,8 @@ def _start_gateway_claim_pid_file() -> bool: """Claim the runtime lock + PID file (O_EXCL winner is the authoritative gateway). False = lost.""" import atexit from gateway.status import ( - acquire_gateway_runtime_lock, - get_running_pid, - release_gateway_runtime_lock, - remove_pid_file, - write_pid_file, - ) + acquire_gateway_runtime_lock, get_running_pid, release_gateway_runtime_lock, + remove_pid_file, write_pid_file) _current_pid = get_running_pid() if _current_pid is not None and _current_pid != os.getpid(): logger.error( @@ -6206,8 +5822,7 @@ async def _start_gateway_start_control_socket(runner): "pausing": accepted, "already_stopping": not accepted, "pid": os.getpid(), - "drain_timeout": _drain, - } + "drain_timeout": _drain} _control_server = GatewayControlServer( verb_handlers={"pause-for-update": _pause_for_update_handler} @@ -6230,15 +5845,12 @@ def _start_gateway_start_cron_and_housekeeping(runner): # Start the background cron scheduler via the resolved provider so scheduled jobs fire # automatically. Pass the event loop so cron delivery can use live adapters (E2EE support). from cron.scheduler_provider import ( - InProcessCronScheduler, - resolve_cron_scheduler, - scheduler_for_profile_mode, + InProcessCronScheduler, resolve_cron_scheduler, scheduler_for_profile_mode ) cron_stop = threading.Event() multiplex_cron = bool(getattr(runner.config, "multiplex_profiles", False)) cron_provider = scheduler_for_profile_mode( - resolve_cron_scheduler(), - multiplex_profiles=multiplex_cron, + resolve_cron_scheduler(), multiplex_profiles=multiplex_cron ) cron_start_kwargs: Dict[str, Any] = {"adapters": runner.adapters, "loop": asyncio.get_running_loop()} @@ -6246,8 +5858,7 @@ def _start_gateway_start_cron_and_housekeeping(runner): # process-global HERMES_HOME is iterated and secondary profiles' cron jobs show as "scheduled" # with a valid next_run_at but never execute because no ticker owns that store. if ( - isinstance(cron_provider, InProcessCronScheduler) - and multiplex_cron + isinstance(cron_provider, InProcessCronScheduler) and multiplex_cron ): try: profile_homes = _multiplex_profile_homes(runner.config) @@ -6266,12 +5877,10 @@ def _start_gateway_start_cron_and_housekeeping(runner): logger.info( "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], - ) + [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, + "Could not resolve profile homes for multiplex cron: %s", exc ) # External cron providers own their remote scheduling contract; only the in-process ticker polls @@ -6281,12 +5890,8 @@ def _start_gateway_start_cron_and_housekeeping(runner): runner._draining or runner._external_drain_active ) cron_thread = threading.Thread( - target=cron_provider.start, - args=(cron_stop,), - kwargs=cron_start_kwargs, - daemon=True, - name="cron-scheduler", - ) + target=cron_provider.start, args=(cron_stop,), kwargs=cron_start_kwargs, daemon=True, + name="cron-scheduler") cron_thread.start() # Preflight tell for the hosted fire path: an external cron provider fires over HTTP to THIS @@ -6306,8 +5911,7 @@ def _start_gateway_start_cron_and_housekeeping(runner): "missing from this gateway process's environment. Restart " "the gateway through its supervisor (`hermes gateway " "restart`) so the profile env loads.", - getattr(cron_provider, "name", "external"), - ) + getattr(cron_provider, "name", "external")) # Gateway-only periodic housekeeping (channel dir, cache cleanup, paste sweep, curator) — runs # independently of the active cron provider; shares cron_stop as the shutdown signal. @@ -6317,24 +5921,17 @@ def _start_gateway_start_cron_and_housekeeping(runner): kwargs={ "adapters": runner.adapters, "loop": asyncio.get_running_loop(), - "cron_provider": cron_provider, - }, + "cron_provider": cron_provider}, daemon=True, - name="gateway-housekeeping", - ) + name="gateway-housekeeping") housekeeping_thread.start() return cron_stop, cron_provider, cron_thread, housekeeping_thread async def _start_gateway_shutdown_tail( - runner, - _control_server, - cron_stop: threading.Event, - cron_provider, - cron_thread: threading.Thread, - housekeeping_thread: threading.Thread, - _planned_stop_watcher_stop: threading.Event, - _planned_stop_watcher_thread: threading.Thread, + runner, _control_server, cron_stop: threading.Event, cron_provider, + cron_thread: threading.Thread, housekeeping_thread: threading.Thread, + _planned_stop_watcher_stop: threading.Event, _planned_stop_watcher_thread: threading.Thread, _signal_initiated_shutdown: list, ) -> bool: """Post-``wait_for_shutdown`` teardown; returns the process exit verdict (True = exit 0).""" @@ -6368,8 +5965,7 @@ async def _start_gateway_shutdown_tail( 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, - ) + "delivery may have been dropped.", _CRON_SHUTDOWN_DRAIN_TIMEOUT) await _await_thread_exit( housekeeping_thread, timeout=_HOUSEKEEPING_SHUTDOWN_DRAIN_TIMEOUT ) @@ -6487,10 +6083,8 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = _planned_stop_watcher_stop = threading.Event() _planned_stop_watcher_thread = threading.Thread( target=_run_planned_stop_watcher, - args=(_planned_stop_watcher_stop, runner, loop, shutdown_signal_handler), - daemon=True, - name="planned-stop-watcher", - ) + args=(_planned_stop_watcher_stop, runner, loop, shutdown_signal_handler), daemon=True, + name="planned-stop-watcher") _planned_stop_watcher_thread.start() # Claim the PID file BEFORE bringing up any platform adapters: two concurrent `gateway run @@ -6545,8 +6139,7 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = recovered = recover_pending_to_db() if recovered: logger.info( - "Recovered %d pending message(s) from shutdown flush", recovered, - ) + "Recovered %d pending message(s) from shutdown flush", recovered) except Exception: pass if runner.should_exit_cleanly: @@ -6590,16 +6183,8 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = await runner.wait_for_shutdown() return await _start_gateway_shutdown_tail( - runner, - _control_server, - cron_stop, - cron_provider, - cron_thread, - housekeeping_thread, - _planned_stop_watcher_stop, - _planned_stop_watcher_thread, - _signal_initiated_shutdown, - ) + runner, _control_server, cron_stop, cron_provider, cron_thread, housekeeping_thread, + _planned_stop_watcher_stop, _planned_stop_watcher_thread, _signal_initiated_shutdown) def _guard_corrupt_user_config() -> None: @@ -6610,8 +6195,7 @@ def _guard_corrupt_user_config() -> None: Same policy and escape hatch (``HERMES_IGNORE_USER_CONFIG=1``) as ``hermes_cli/main.py``. """ from hermes_cli.config import ( - InvalidUserConfigError, - require_parseable_user_config, + InvalidUserConfigError, require_parseable_user_config ) try: @@ -6637,8 +6221,7 @@ def main(): # reapers can identify this gateway (and its child tree dies with it on Windows). Best-effort. try: from hermes_cli.process_identity import ( - attach_self_to_kill_on_close_job, - register_self, + attach_self_to_kill_on_close_job, register_self ) register_self("gateway")