diff --git a/gateway/run.py b/gateway/run.py index a20f0bee7e..d33527820a 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -94,8 +94,7 @@ _TELEGRAM_NOISY_STATUS_RE = re.compile( r"|stale\s+connections\s+from\s+a\s+previous\s+provider\s+issue" rf"|{re.escape(COMPACTION_DONE_STATUS)}" r")", - re.IGNORECASE | re.DOTALL, -) + re.IGNORECASE | re.DOTALL) _HYGIENE_COOLDOWN_LADDER_MULTIPLIERS = (1, 3, 9) # Ceiling on an escalated cooldown (cf. _RECONNECT_BACKOFF_CAP): an operator-raised base × ladder reaches @@ -138,8 +137,7 @@ def _hygiene_cooldown_for_failure(gateway, session_key: str, base_cooldown_secon state.hygiene_failure_streak += 1 streak = state.hygiene_failure_streak multiplier = _HYGIENE_COOLDOWN_LADDER_MULTIPLIERS[ - min(streak, len(_HYGIENE_COOLDOWN_LADDER_MULTIPLIERS)) - 1 - ] + min(streak, len(_HYGIENE_COOLDOWN_LADDER_MULTIPLIERS)) - 1] return min(base_cooldown_seconds * multiplier, _HYGIENE_COOLDOWN_MAX_SECONDS) @@ -164,8 +162,7 @@ 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, -) -> bool: + approx_tokens: int, new_tokens: int) -> bool: """True when a hygiene run actually recovered the session (extracted to be unit testable). Requires: no abort; transcript actually rewritten (the no-op path reuses pre-compression counts); @@ -178,22 +175,19 @@ def hygiene_compaction_recovered( def _hygiene_compression_timeout_message( - *, total_exhausted: bool, elapsed: float, idle_timeout: float, progress_observed: bool -) -> str: + *, total_exhausted: bool, elapsed: float, idle_timeout: float, progress_observed: bool) -> str: """Describe the host timeout that actually ended hygiene compression.""" if total_exhausted: progress = " after summary output was observed" if progress_observed else "" 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 " "output from the summary model. No messages were dropped — continuing " "without compression. Run /compress to retry, /reset for a clean " - "session, or check your auxiliary.compression model configuration." - ) + "session, or check your auxiliary.compression model configuration.") def _cached_agent_for_hygiene(gateway, session_key: str): @@ -212,8 +206,7 @@ def _cached_agent_for_hygiene(gateway, session_key: str): 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, -) -> str: + approx_tokens: int, timeout_seconds: float, failure_cooldown_seconds: float = 300.0) -> str: """Session hygiene for ``codex_app_server`` sessions. The real context is the server-side thread; the local transcript is a never-replayed mirror, so @@ -288,8 +281,7 @@ def hygiene_wait_should_extend( def _record_hygiene_cooldown( - gateway, session_id: str, cooldown_seconds: float, error: Optional[str] = None -) -> None: + gateway, session_id: str, cooldown_seconds: float, error: Optional[str] = None) -> None: """Persist a session-hygiene compression-failure cooldown to the state DB. Shares the in-conversation path's column/recorder so it survives restarts. ``error`` must be @@ -323,8 +315,7 @@ _COMPRESSION_PROGRESS_STATUS_RE = re.compile( 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) @@ -364,16 +355,14 @@ _GATEWAY_PROVIDER_POLICY_RE = re.compile( r"|disallowed" r"|moderation" r")", - re.IGNORECASE, -) + re.IGNORECASE) _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) _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. _CONNECTION_ERROR_MARKERS = ( @@ -446,8 +435,7 @@ def _gateway_platform_value(platform: Any) -> str: def _non_conversational_metadata( - metadata: Optional[Dict[str, Any]] = None, *, platform: Any = None -) -> Optional[Dict[str, Any]]: + 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": return metadata @@ -512,8 +500,7 @@ def _is_transient_network_error(exc: BaseException) -> bool: def _gateway_loop_exception_handler( - loop: "asyncio.AbstractEventLoop", context: Dict[str, Any] -) -> None: + loop: "asyncio.AbstractEventLoop", context: Dict[str, Any]) -> None: """Loop-level safety net for transient network errors (installed once by ``start_gateway``). Logs WARNING with traceback; non-transient errors go to the default handler so real bugs surface. @@ -567,8 +554,7 @@ 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, -) -> str: + allow_session: bool = True, smart_denied: bool = False) -> str: """Render the text fallback from approval capabilities, not platform names.""" cmd_preview = command[:200] + "..." if len(command) > 200 else command heading = "⚠️ **Dangerous command requires approval:**" @@ -583,8 +569,7 @@ def _format_exec_approval_fallback( choices.append(f"`{command_prefix}deny` to cancel") return ( f"{heading}\n```\n{cmd_preview}\n```\nReason: {description}\n\n" - + ", ".join(choices[:-1]) + f", or {choices[-1]}." - ) + + ", ".join(choices[:-1]) + f", or {choices[-1]}.") # Ordered: auth beats policy beats rate-limit beats connection; first match wins. _PROVIDER_ERROR_REPLIES = ( @@ -604,8 +589,7 @@ def _gateway_provider_error_reply(text: str) -> str: return reply return ( "⚠️ The model provider failed after retries. I kept raw provider details " - "out of chat; check gateway logs for diagnostics." - ) + "out of chat; check gateway logs for diagnostics.") # Provider/API failure envelope preambles (not ordinary assistant prose), anchored at line start. @@ -736,8 +720,7 @@ def _clarify_send_disposition(fut, *, session_key: str, clarify_mod) -> "str | N if outcome == "ambiguous": logger.warning( "Clarify prompt send timed out — treating as possibly-delivered " - "(no teardown; the registration stays armed for a late reply)" - ) + "(no teardown; the registration stays armed for a late reply)") return None @@ -789,8 +772,7 @@ 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, -) -> bool: + platform: Any = None, require_platform_override_for: set[Any] | None = None) -> bool: """Resolve a boolean display setting with optional platform-only opt-in. Scratch-text features are too noisy for threaded surfaces (Mattermost) under a global opt-in, @@ -800,8 +782,7 @@ def _resolve_gateway_display_bool( platform_only = {_gateway_platform_value(c) for c in (require_platform_override_for or set())} if ( current_platform in platform_only - and not _has_platform_display_override(user_config, platform_key, setting) - ): + and not _has_platform_display_override(user_config, platform_key, setting)): return False from gateway.display_config import resolve_display_setting @@ -956,8 +937,7 @@ def _stamp_hygiene_compression_provenance( def _is_fresh_gateway_interruption( - value: Any, *, now: Optional[float] = None, window_secs: Optional[float] = None -) -> bool: + 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). @@ -973,8 +953,7 @@ def _is_fresh_gateway_interruption( def build_resume_recovery_note( - reason: Optional[str], message: str = "", *, interactive: bool = True -) -> str: + reason: Optional[str], message: str = "", *, interactive: bool = True) -> str: """Build the resume-pending recovery system note for an interrupted turn. Empty ``message`` = startup auto-resume. Interactive platforms report the restore and ask what @@ -982,20 +961,17 @@ def build_resume_recovery_note( """ reason_phrase = ( "a gateway restart" if reason == "restart_timeout" - else "a gateway shutdown" if reason == "shutdown_timeout" else "a gateway interruption" - ) + else "a gateway shutdown" if reason == "shutdown_timeout" else "a gateway interruption") 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 unfinished work from the conversation history." ) elif interactive: resume_guidance = ( "Report to the user that the session was restored " - "successfully and ask what they would like to do next." - ) + "successfully and ask what they would like to do next.") tail_guidance = ( "Do NOT re-execute old tool calls — skip any unfinished work from the conversation history." ) @@ -1004,24 +980,20 @@ def build_resume_recovery_note( "No user is present on this non-interactive platform, " "so do NOT emit a 'session restored' acknowledgement " "or ask questions. Review the conversation history and " - "CONTINUE the interrupted task to completion." - ) + "CONTINUE the interrupted task to completion.") 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 " f"{reason_phrase}; the gateway is now back online. " f"Any restart/shutdown command in the history has already " f"run — do NOT re-execute or verify it. {resume_guidance} {tail_guidance}]" - + (f"\n\n{message}" if message else "") - ) + + (f"\n\n{message}" if message else "")) def _prepare_resume_pending_message( - reason: Optional[str], message: Optional[str], *, interactive: bool = True -) -> tuple[str, str]: + reason: Optional[str], message: Optional[str], *, interactive: bool = True) -> tuple[str, str]: """Return the recovery message and the user text to persist. Empty original (synthesized auto-resume): persist the note — a "" user row trips the pre-call @@ -1058,8 +1030,7 @@ def _build_replay_entry( role in ("user", "assistant") and isinstance(_sidecar, str) and _sidecar - and content == msg.get("content") - ): + and content == msg.get("content")): entry["api_content"] = _sidecar if role == "assistant": for _rkey in _ASSISTANT_REPLAY_FIELDS: @@ -1153,8 +1124,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, - inject_timestamps: bool = False, -) -> tuple[List[Dict[str, Any]], Optional[str]]: + inject_timestamps: bool = False) -> tuple[List[Dict[str, Any]], Optional[str]]: """Convert stored gateway transcript rows into agent replay messages. Keeping that context out of ``conversation_history`` stops consecutive-user repair merging it @@ -1215,8 +1185,7 @@ def _build_gateway_agent_history( def _select_cached_agent_history( - persisted_history: List[Dict[str, Any]], live_history: Any -) -> List[Dict[str, Any]]: + persisted_history: List[Dict[str, Any]], live_history: Any) -> List[Dict[str, Any]]: """Prefer the cached live transcript only when it is longer AND has a real, non-ephemeral unpersisted row; otherwise return ``persisted_history`` unchanged. @@ -1337,8 +1306,7 @@ from gateway.media_repair import tool_name_by_call_id as _tool_name_by_call_id def _collect_auto_append_media_tags( messages: List[Dict[str, Any]], history_offset: int = 0, - history_media_paths: Optional[set] = None, -) -> tuple[List[str], bool]: + history_media_paths: Optional[set] = None) -> tuple[List[str], bool]: """Collect real media tags from current-turn producer-tool results only. Producer allowlist: docs/logs/search results contain example MEDIA: strings that must never @@ -1446,8 +1414,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) import ssl @@ -1547,8 +1514,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) @@ -1765,8 +1731,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 @@ -1826,8 +1791,7 @@ _DOCKER_MEDIA_OUTPUT_CONTAINER_PATHS = {"/output", "/outputs"} from hermes_cli.config_defaults import DEFAULT_CONFIG as _DEFAULT_CONFIG os.environ["HERMES_TURN_LEASE_TIMEOUT"] = str( - _DEFAULT_CONFIG["agent"]["gateway_turn_lease_timeout"] -) + _DEFAULT_CONFIG["agent"]["gateway_turn_lease_timeout"]) # Bridge config.yaml values into env so os.getenv() picks them up. config.yaml unconditionally wins # over .env for these keys; a `not in os.environ` guard would let stale .env entries shadow config. @@ -1843,8 +1807,7 @@ _AGENT_ENV_BRIDGE = { "cron_drain_timeout": "HERMES_CRON_DRAIN_TIMEOUT", "gateway_auto_continue_freshness": "HERMES_AUTO_CONTINUE_FRESHNESS", "gateway_startup_restore_drain_timeout": "HERMES_STARTUP_RESTORE_DRAIN_TIMEOUT", - "gateway_startup_warmup_timeout": "HERMES_STARTUP_WARMUP_TIMEOUT", -} + "gateway_startup_warmup_timeout": "HERMES_STARTUP_WARMUP_TIMEOUT"} # config-authoritative knobs for the session-search index (env stays the cross-process carrier). _SESSIONS_ENV_BRIDGE = {"cjk_fts": "HERMES_CJK_FTS", "search_slow_ms": "HERMES_SEARCH_SLOW_MS"} _DISPLAY_ENV_BRIDGE = { @@ -1879,8 +1842,7 @@ def _bridge_max_turns_to_env(agent_cfg: Any) -> None: def _bridge_terminal_config_to_env(_terminal_cfg: dict) -> None: """Bridge nested ``terminal.*`` config to TERMINAL_* env vars (config.yaml overrides .env here).""" _terminal_backend = str( - _terminal_cfg.get("backend") or os.environ.get("TERMINAL_ENV") or "" - ).strip().lower() + _terminal_cfg.get("backend") or os.environ.get("TERMINAL_ENV") or "").strip().lower() _terminal_env_map = { "backend": "TERMINAL_ENV", "degraded_mode": "TERMINAL_DEGRADED_MODE", @@ -1975,8 +1937,7 @@ def _bridge_config_to_env(_cfg: dict) -> None: # Documented service-manager override: env wins when already set (other display bridges stay # config-authoritative for backwards compatibility). and "busy_steer_ack_enabled" in _display_cfg - and "HERMES_GATEWAY_BUSY_STEER_ACK_ENABLED" not in os.environ - ): + and "HERMES_GATEWAY_BUSY_STEER_ACK_ENABLED" not in os.environ): os.environ["HERMES_GATEWAY_BUSY_STEER_ACK_ENABLED"] = str(_display_cfg["busy_steer_ack_enabled"]) _tz_cfg = _cfg.get("timezone", "") if _tz_cfg and isinstance(_tz_cfg, str): @@ -1996,8 +1957,7 @@ def _bridge_config_to_env(_cfg: dict) -> None: # platform_connect_timeout is an escape hatch, unlike the bridges above: env WINS if already set. if ( "platform_connect_timeout" in _gateway_cfg - and not os.environ.get("HERMES_GATEWAY_PLATFORM_CONNECT_TIMEOUT", "").strip() - ): + and not os.environ.get("HERMES_GATEWAY_PLATFORM_CONNECT_TIMEOUT", "").strip()): os.environ["HERMES_GATEWAY_PLATFORM_CONNECT_TIMEOUT"] = str(_gateway_cfg["platform_connect_timeout"]) @@ -2071,8 +2031,7 @@ if not _configured_cwd or _configured_cwd in CWD_PLACEHOLDERS: terminal_backend=os.environ.get("TERMINAL_ENV", ""), messaging_cwd=os.getenv("MESSAGING_CWD"), docker_mount_cwd_to_workspace=os.getenv( - "TERMINAL_DOCKER_MOUNT_CWD_TO_WORKSPACE", "false" - ).lower() + "TERMINAL_DOCKER_MOUNT_CWD_TO_WORKSPACE", "false").lower() in {"true", "1", "yes"}, home_fallback=str(Path.home())) if _resolved_cwd is None: @@ -2081,11 +2040,9 @@ if not _configured_cwd or _configured_cwd in CWD_PLACEHOLDERS: 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.) @@ -2125,8 +2082,7 @@ from gateway.restart import ( DEFAULT_GATEWAY_POST_INTERRUPT_GRACE_TIMEOUT, # noqa: F401 (re-exported: run_* mixins + tests resolve gateway.run.) DEFAULT_GATEWAY_RESTART_AFTER_TURN_TIMEOUT, DEFAULT_GATEWAY_RESTART_DRAIN_TIMEOUT, - DEFAULT_GATEWAY_SIGNAL_INTERRUPT_GRACE_TIMEOUT, -) + DEFAULT_GATEWAY_SIGNAL_INTERRUPT_GRACE_TIMEOUT) logger = logging.getLogger(__name__) @@ -2176,8 +2132,7 @@ def _own_policy_open_startup_violation(config) -> Optional[str]: continue gateway_allow_all = _getenv("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 return f"{platform.value}: open policy without allow-all opt-in" @@ -2208,8 +2163,7 @@ _CONVERSATION_SCOPED_STATE: tuple = ( "_session_stall_notified", # Sidecar notes staged but never consumed (turn aborted before run_sync) must not leak into a # future conversation's first user message — session keys are source-derived and REUSED. - "_pending_turn_sidecar_notes", -) + "_pending_turn_sidecar_notes") from gateway.run_common import _UNSET # noqa: F401 (def-time sentinel shared with run_* mixins) @@ -2221,8 +2175,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 try: @@ -2263,8 +2216,7 @@ def _resolve_runtime_agent_kwargs() -> dict: capabilities = runtime.get("capabilities") capabilities = ( {k: v for k, v in capabilities.items() if isinstance(k, str) and isinstance(v, bool)} - if isinstance(capabilities, dict) else {} - ) + if isinstance(capabilities, dict) else {}) return {**_runtime_agent_kwargs(runtime), "max_tokens": max_tokens, "capabilities": capabilities} @@ -2495,8 +2447,7 @@ def _build_media_placeholder(event) -> str: def _build_document_context_note( - display_name: str, agent_path: str, mtype: str, *, content_inlined: bool = True -) -> str: + 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. ``content_inlined=False`` = cached without content, so tell the agent to read it. Binary docs @@ -2505,21 +2456,18 @@ def _build_document_context_note( if mtype.startswith("text/") and content_inlined: return ( f"[The user sent a text document: '{display_name}'. Its content has been included below. " - f"The file is also saved at: {agent_path}]" - ) + f"The file is also saved at: {agent_path}]") if mtype.startswith("text/"): return ( f"[The user sent a text document: '{display_name}'. It is saved at: {agent_path}. " f"Its content is not inlined here. Read the cached file yourself before answering " - f"when the user's request involves its contents.]" - ) + f"when the user's request involves its contents.]") return ( f"[The user sent a document: '{display_name}'. It is saved at: {agent_path}. " f"Its text is not inlined here (it's a binary format such as PDF or DOCX). " f"To read it, extract the document's text yourself — for example with the " f"terminal tool or the ocr-and-documents skill — before answering, instead " - f"of asking the user to paste the contents.]" - ) + f"of asking the user to paste the contents.]") def _format_duration(seconds: float) -> str: @@ -2590,8 +2538,7 @@ _INTERRUPT_REASON_GATEWAY_RESTART = "Gateway restarting" def _reap_gateway_turn_processes( task_id: str, process_baseline, *, source: str, - is_still_current: Optional[Callable[[], bool]] = None, -) -> int: + is_still_current: Optional[Callable[[], bool]] = None) -> int: """Reap only background processes created by one abandoned turn. ``task_id`` is session-scoped, so a *replacement* turn can spawn its own process mid-reap; @@ -2620,8 +2567,7 @@ def _reap_gateway_turn_processes( except Exception: # Detached daemon thread: an uncaught exception would only reach threading.excepthook. 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: logger.warning( @@ -2672,8 +2618,7 @@ def _dump_wedged_turn_stacks(task_id: str) -> None: 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, - is_still_current: Optional[Callable[[], bool]] = None, -) -> bool: + is_still_current: Optional[Callable[[], bool]] = None) -> bool: """Interrupt one timed-out turn and reap only processes it created.""" with cleanup_lock: if worker_done.is_set() or timeout_fired.is_set(): @@ -2696,16 +2641,14 @@ def _abandon_timed_out_gateway_turn( 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, - is_still_current: Optional[Callable[[], bool]] = None, -) -> None: + is_still_current: Optional[Callable[[], bool]] = None) -> None: """Thread watchdog that remains runnable when gateway asyncio is starved.""" while not worker_done.wait(max(0.01, poll_interval)): agent = agent_holder[0] if agent_holder else None @@ -2806,8 +2749,7 @@ def _check_unavailable_skill(command_name: str) -> str | None: if slug == normalized and declared_name in disabled: return ( f"The **{command_name}** skill is installed but disabled.\n" - f"Enable it with: `hermes skills config`" - ) + f"Enable it with: `hermes skills config`") # Check optional skills (shipped with repo but not installed) from hermes_constants import get_optional_skills_dir @@ -2825,8 +2767,7 @@ def _check_unavailable_skill(command_name: str) -> str | None: install_path = f"official/{'/'.join(rel.parts)}" return ( f"The **{command_name}** skill is available but not installed.\n" - f"Install it with: `hermes skills install {install_path}`" - ) + f"Install it with: `hermes skills install {install_path}`") except Exception: pass return None @@ -2945,8 +2886,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 -) -> list[str]: + chat_id: str, *, thread_id: Optional[str] = None, parent_id: Optional[str] = None) -> list[str]: """Ordered, de-duplicated ``channel_overrides`` lookup keys (matches ``resolve_channel_prompt``: exact id first, then parent — Discord threads inherit parent overrides).""" keys: list[str] = [] @@ -2964,8 +2904,7 @@ def _channel_override_lookup_keys( def _get_channel_override( config: GatewayConfig, platform: Platform, chat_id: str, *, thread_id: Optional[str] = None, - parent_id: Optional[str] = None, -) -> Optional[ChannelOverride]: + parent_id: Optional[str] = None) -> Optional[ChannelOverride]: """Per-channel override via chat_id, then thread_id, then parent_id; None if absent.""" platforms = getattr(config, "platforms", None) if not platforms: @@ -3024,8 +2963,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 -) -> str: + session_id: str, command: str, exit_code, output: str, duration_seconds=None) -> str: """One-line completion message for the ``concise`` display mode. Success is one status line; failure appends a short output tail (full output via process(log)). @@ -3074,8 +3012,7 @@ def _format_gateway_process_notification(evt: dict) -> "str | None": text = ( f"[IMPORTANT: Background process {_sid} matched " f"watch pattern \"{_pat}\".\n" - f"Command: {_cmd}\nMatched output:\n{_out}" - ) + f"Command: {_cmd}\nMatched output:\n{_out}") if _sup: text += f"\n({_sup} earlier matches were suppressed by rate limit)" text += "]" @@ -3103,8 +3040,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": requeue.append(evt) @@ -3120,8 +3056,7 @@ _gateway_runner_ref: _weakref.ref = lambda: None def _normalize_empty_agent_response( - agent_result: dict, response: str, *, history_len: int = 0 -) -> str: + agent_result: dict, response: str, *, history_len: int = 0) -> str: """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, @@ -3137,31 +3072,26 @@ def _normalize_empty_agent_response( # Persistence failures: suggesting /reset would destroy context without fixing storage. failure_reason = str(agent_result.get("failure_reason") or "") if failure_reason.startswith("session_persistence_failed") or ( - "session storage" in error_str - ): + "session storage" in error_str): if failure_reason.endswith(":disk") or "disk" in error_str: 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 for p in ("context", "token", "too large", "too long", "exceed", "payload") ) or ("400" in error_str and history_len > 50) 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" - "Try again or use /reset to start a fresh session." - ) + "Try again or use /reset to start a fresh session.") api_calls = int(agent_result.get("api_calls", 0) or 0) if agent_result.get("interrupted"): @@ -3170,8 +3100,7 @@ def _normalize_empty_agent_response( if api_calls == 0: return ( "⚠️ Your message was interrupted before processing started " - "(likely by a recent /stop). Please send it again." - ) + "(likely by a recent /stop). Please send it again.") return response if api_calls > 0: if _is_gateway_hidden_reasoning_incomplete_turn(agent_result): @@ -3181,15 +3110,13 @@ def _normalize_empty_agent_response( return f"⚠️ Processing stopped: {str(err)[:200]}. Try again." return ( "⚠️ Processing completed but no response was generated. " - "This may be a transient error — try sending your message again." - ) + "This may be a transient error — try sending your message again.") # api_calls == 0, not failed/interrupted: agent never ran (post-/stop race); don't drop silently. if api_calls == 0 and not agent_result.get("partial"): return ( "⚠️ Your message wasn't processed (the previous turn was still " - "being cleaned up). Please send it again." - ) + "being cleaned up). Please send it again.") return response @@ -3229,8 +3156,7 @@ def _should_clear_resume_pending_after_turn(agent_result: dict) -> bool: def _preserve_queued_followup_history_offset( - current_result: dict, followup_result: dict -) -> dict: + current_result: dict, followup_result: dict) -> dict: """Carry the outer history offset through queued follow-up drains. Each recursive ``_run_agent()`` advances ``history_offset``; uncorrected, the outer persistence @@ -3415,9 +3341,7 @@ class GatewayRunner( _restart_drain_timeout: float = DEFAULT_GATEWAY_RESTART_DRAIN_TIMEOUT _restart_after_turn_timeout: float = DEFAULT_GATEWAY_RESTART_AFTER_TURN_TIMEOUT _cron_drain_timeout: float = DEFAULT_GATEWAY_CRON_DRAIN_TIMEOUT - _signal_interrupt_grace_timeout: float = ( - DEFAULT_GATEWAY_SIGNAL_INTERRUPT_GRACE_TIMEOUT - ) + _signal_interrupt_grace_timeout: float = DEFAULT_GATEWAY_SIGNAL_INTERRUPT_GRACE_TIMEOUT _exit_code: Optional[int] = None _draining: bool = False _external_drain_active: bool = False @@ -3444,19 +3368,16 @@ class GatewayRunner( _session_run_generation = legacy_dict_property("_session_run_generation") _session_model_overrides = legacy_dict_property("_session_model_overrides") _pending_one_turn_model_restores = legacy_dict_property( - "_pending_one_turn_model_restores" - ) + "_pending_one_turn_model_restores") _session_reasoning_overrides = legacy_dict_property("_session_reasoning_overrides") _session_service_tier_overrides = legacy_dict_property( - "_session_service_tier_overrides" - ) + "_session_service_tier_overrides") _last_resolved_model = legacy_dict_property("_last_resolved_model") _queued_events = legacy_dict_property("_queued_events") _pending_turn_sidecar_notes = legacy_dict_property("_pending_turn_sidecar_notes") _pending_messages = legacy_dict_property("_pending_messages") _pending_native_image_paths_by_session = legacy_dict_property( - "_pending_native_image_paths_by_session" - ) + "_pending_native_image_paths_by_session") _session_ephemeral_pin = legacy_dict_property("_session_ephemeral_pin") _session_vc_last = legacy_dict_property("_session_vc_last") _pending_approvals = legacy_dict_property("_pending_approvals") @@ -3497,8 +3418,7 @@ class GatewayRunner( return [ (key, state.turn.agent) for key, state in self._sessions_map().items() - if state.turn.agent is not None - ] + if state.turn.agent is not None] # Loop-liveness / watchdog handles; class-level defaults so partially constructed test runners work. _loop_heartbeat_task: Optional["asyncio.Task"] = None _loop_floor_timer_handle: Optional[Any] = None @@ -3564,11 +3484,9 @@ class GatewayRunner( # default) is stale and no longer blocks idle/daily reset (NOT killed, only ignored). from tools.process_registry import process_registry _bg_max_age_hours = getattr( - self.config.default_reset_policy, "bg_process_max_age_hours", 24 - ) + self.config.default_reset_policy, "bg_process_max_age_hours", 24) _bg_max_age_seconds = ( - _bg_max_age_hours * 3600 if _bg_max_age_hours and _bg_max_age_hours > 0 else None - ) + _bg_max_age_hours * 3600 if _bg_max_age_hours and _bg_max_age_hours > 0 else None) self.session_store = SessionStore( self.config.sessions_dir, self.config, has_active_processes_fn=lambda key: process_registry.has_active_for_session( @@ -3686,8 +3604,7 @@ class GatewayRunner( "auxiliary.approval is unset): dangerous commands and " "execute_code scripts will BLOCK until a human approves " "them in chat. Enable security.tirith_enabled or configure " - "auxiliary.approval for unattended operation." - ) + "auxiliary.approval for unattended operation.") except Exception: logger.debug("approvals.mode startup check skipped", exc_info=True) @@ -3701,8 +3618,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) except Exception as e: @@ -3727,8 +3643,7 @@ class GatewayRunner( retention_days=int(_sess_cfg.get("retention_days", 90)), min_interval_hours=int(_sess_cfg.get("min_interval_hours", 24)), min_vacuum_interval_days=int( - _sess_cfg.get("min_vacuum_interval_days", 30) - ), + _sess_cfg.get("min_vacuum_interval_days", 30)), vacuum=bool(_sess_cfg.get("vacuum_after_prune", True)), sessions_dir=self.config.sessions_dir) except Exception as exc: @@ -3793,8 +3708,7 @@ class GatewayRunner( if cache is None: # Test runners built with object.__new__ skip __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 def _open(): @@ -3822,8 +3736,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 def _session_db(self) -> Any: @@ -3878,8 +3791,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: """Warn when Docker-backed gateways lack an explicit export mount: MEDIA delivery runs in the @@ -3919,8 +3831,7 @@ class GatewayRunner( "Docker backend is enabled for the messaging gateway but no explicit host-visible " "output mount (for example '/home/user/.hermes/cache/documents:/output') is configured. " "This is fine if the model already emits host-visible paths, but MEDIA file delivery can fail " - "for container-local paths like '/workspace/...' or '/output/...'." - ) + "for container-local paths like '/workspace/...' or '/output/...'.") _VOICE_MODE_PATH = _hermes_home / "gateway_voice_mode.json" @@ -3962,8 +3873,7 @@ class GatewayRunner( _TELEGRAM_LOBBY_REMINDER_COOLDOWN_S = 30.0 def _normalize_source_for_session_key( - self, source: SessionSource - ) -> SessionSource: + self, source: SessionSource) -> SessionSource: """Apply Telegram DM topic recovery to a source for session-key purposes. ``_handle_message_with_agent`` rewrites ``thread_id`` before deriving the session key, so @@ -4015,8 +3925,7 @@ class GatewayRunner( return { id(a) for _, a in self._running_agent_items() - if a is not None and a is not _AGENT_PENDING_SENTINEL - } + if a is not None and a is not _AGENT_PENDING_SENTINEL} def _snapshot_running_agents(self) -> Dict[str, Any]: return {k: a for k, a in self._running_agent_items() if a is not _AGENT_PENDING_SENTINEL} @@ -4079,8 +3988,7 @@ class GatewayRunner( return "default" def _is_user_authorized_for_source( - self, source: SessionSource, *, allow_adapter_delegation: bool = True - ) -> bool: + self, source: SessionSource, *, allow_adapter_delegation: bool = True) -> bool: """Authorize under the live transport's profile, not the routed runtime (which need not copy the shared bot token/allowlist); the transport home is stamped on the source for this read only.""" def _check() -> bool: @@ -4168,8 +4076,7 @@ class GatewayRunner( history: Any = None def _thread_metadata_for_source( - self, source, reply_to_message_id: Optional[str] = None - ) -> Optional[Dict[str, Any]]: + 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), @@ -4199,15 +4106,13 @@ class GatewayRunner( 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, - adapter: Optional[Any] = None, - ) -> Optional[Dict[str, Any]]: + adapter: Optional[Any] = None) -> Optional[Dict[str, Any]]: """Build thread metadata for synthetic sends that only have routing state.""" if thread_id is None: 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 # DM topic lanes need direct_messages_topic_id so synthetic sends reach the topic without a reply anchor. tid = str(thread_id) @@ -4223,8 +4128,7 @@ 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, - ) -> bool: + 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: return False @@ -4265,9 +4169,7 @@ class GatewayRunner( return set_session_vars( platform=context.source.platform.value, chat_id=context.source.chat_id, - chat_type=( - str(context.source.chat_type) if context.source.chat_type else "" - ), + chat_type=str(context.source.chat_type) if context.source.chat_type else "", chat_name=context.source.chat_name or "", thread_id=str(context.source.thread_id) if context.source.thread_id else "", user_id=str(context.source.user_id) if context.source.user_id else "", @@ -4290,8 +4192,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: """Return the gateway-owned executor for blocking agent work.""" @@ -4306,8 +4207,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 @@ -4440,8 +4340,7 @@ class GatewayRunner( ``_profile_name_for_source`` (sources bypassing ``build_source``), then the active profile.""" from gateway.profile_routing import ProfileRouteRejected from hermes_cli.profiles import ( - get_active_profile_name, get_profile_dir, profile_exists - ) + get_active_profile_name, get_profile_dir, profile_exists) from hermes_constants import get_hermes_home explicit_profile = None # explicitly requested (source or routing) vs. default fallback @@ -4515,8 +4414,7 @@ class GatewayRunner( def _run_planned_stop_watcher( stop_event: threading.Event, runner, loop: asyncio.AbstractEventLoop, shutdown_handler, *, - poll_interval: float = 0.5, -) -> None: + poll_interval: float = 0.5) -> None: """Poll for the planned-stop marker and trigger graceful shutdown. On Windows ``add_signal_handler`` is unavailable, so ``hermes gateway stop`` would never drain; @@ -4524,16 +4422,14 @@ def _run_planned_stop_watcher( would. On POSIX the signal handler consumes the marker first; ``_running``/``_draining`` guard re-triggers. """ 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(): try: if ( marker_path.exists() and not getattr(runner, "_draining", False) - and getattr(runner, "_running", False) - ): + and getattr(runner, "_running", False)): # A marker may target a PREVIOUS instance (different PID) that exited before stop() cleaned # up; firing on it means an "UNKNOWN" exit and a watchdog crash-loop. The probe unlinks # stale/malformed markers. @@ -4705,8 +4601,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) @@ -4721,8 +4616,7 @@ _HOUSEKEEPING_SHUTDOWN_DRAIN_TIMEOUT = 35.0 async def _await_thread_exit( - thread: Optional[threading.Thread], timeout: float, poll: float = 0.1 -) -> bool: + thread: Optional[threading.Thread], timeout: float, poll: float = 0.1) -> bool: """Wait for a daemon thread to exit WITHOUT blocking the event loop; True if it exited in time. A synchronous ``join()`` freezes the loop — fatal for the cron ticker, whose in-flight delivery is @@ -4927,8 +4821,7 @@ async def _start_gateway_replace_existing_instance(existing_pid: int, replace: b f"\n❌ Gateway already running (PID {existing_pid}).\n" f" Use 'hermes gateway restart' to replace it,\n" f" or 'hermes gateway stop' to kill it first.\n" - f" Or use 'hermes gateway run --replace' to auto-replace.\n" - ) + f" Or use 'hermes gateway run --replace' to auto-replace.\n") return False # Never signal a live process we cannot prove belongs to this HERMES_HOME: a poisoned PID record @@ -5074,8 +4967,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: _shutdown_ctx = None @@ -5102,8 +4994,7 @@ def _start_gateway_make_shutdown_signal_handler(runner, _signal_initiated_shutdo if _shutdown_ctx is not None: try: logger.warning( - "Shutdown context: %s", format_context_for_log(_shutdown_ctx) - ) + "Shutdown context: %s", format_context_for_log(_shutdown_ctx)) except Exception as _e: logger.debug("format_context_for_log failed: %s", _e) @@ -5112,8 +5003,7 @@ def _start_gateway_make_shutdown_signal_handler(runner, _signal_initiated_shutdo try: _diag_log = _hermes_home / "logs" / "gateway-shutdown-diag.log" spawn_async_diagnostic( - _diag_log, _shutdown_ctx["signal"], timeout_seconds=5.0 - ) + _diag_log, _shutdown_ctx["signal"], timeout_seconds=5.0) except Exception as _e: logger.debug("spawn_async_diagnostic failed: %s", _e) asyncio.create_task(runner.stop()) @@ -5130,21 +5020,18 @@ def _start_gateway_claim_pid_file() -> bool: if _current_pid is not None and _current_pid != os.getpid(): logger.error( "Another gateway instance (PID %d) started during our startup. " - "Exiting to avoid double-running.", _current_pid - ) + "Exiting to avoid double-running.", _current_pid) return False if not acquire_gateway_runtime_lock(): logger.error( - "Gateway runtime lock is already held by another instance. Exiting." - ) + "Gateway runtime lock is already held by another instance. Exiting.") return False try: write_pid_file() except FileExistsError: release_gateway_runtime_lock() logger.error( - "PID file race lost to another gateway instance. Exiting." - ) + "PID file race lost to another gateway instance. Exiting.") return False atexit.register(remove_pid_file) atexit.register(release_gateway_runtime_lock) @@ -5176,8 +5063,7 @@ async def _start_gateway_start_control_socket(runner): def _request() -> None: try: accepted_box.append( - runner.request_restart(detached=False, via_service=True) - ) + runner.request_restart(detached=False, via_service=True)) finally: _done.set() @@ -5191,8 +5077,7 @@ async def _start_gateway_start_control_socket(runner): "drain_timeout": _drain} _control_server = GatewayControlServer( - verb_handlers={"pause-for-update": _pause_for_update_handler} - ) + verb_handlers={"pause-for-update": _pause_for_update_handler}) if not await _control_server.start(): _control_server = None else: @@ -5210,13 +5095,11 @@ def _start_gateway_start_cron_and_housekeeping(runner): """ # 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 - ) + 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()} # Multiplex: tell the built-in ticker which profile homes to tick, else secondary profiles' cron jobs @@ -5237,15 +5120,13 @@ def _start_gateway_start_cron_and_housekeeping(runner): [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 # local due jobs, so only it receives the local external-drain dispatch gate. if isinstance(cron_provider, InProcessCronScheduler): cron_start_kwargs["can_dispatch"] = lambda: not ( - runner._draining or runner._external_drain_active - ) + 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") @@ -5296,8 +5177,7 @@ 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, - _signal_initiated_shutdown: list, -) -> bool: + _signal_initiated_shutdown: list) -> bool: """Post-``wait_for_shutdown`` teardown; returns the process exit verdict (True = exit 0).""" # Control socket first: once shutdown begins we are no longer a truthful "serving here" answer and a # successor must be able to bind. Early-exit paths rely on the atexit cleanup_files hook instead. @@ -5324,8 +5204,7 @@ async def _start_gateway_shutdown_tail( "Cron ticker did not exit within %.0fs of shutdown — an in-flight " "delivery may have been dropped.", _CRON_SHUTDOWN_DRAIN_TIMEOUT) await _await_thread_exit( - housekeeping_thread, timeout=_HOUSEKEEPING_SHUTDOWN_DRAIN_TIMEOUT - ) + housekeeping_thread, timeout=_HOUSEKEEPING_SHUTDOWN_DRAIN_TIMEOUT) # Stop the planned-stop watcher (daemon=True so this is belt-and-suspenders). _planned_stop_watcher_stop.set() @@ -5342,16 +5221,14 @@ async def _start_gateway_shutdown_tail( if _signal_initiated_shutdown[0] and not runner._restart_requested: logger.info( "Exiting with code 1 (signal-initiated shutdown without restart " - "request) so systemd Restart=on-failure can revive the gateway." - ) + "request) so systemd Restart=on-failure can revive the gateway.") return False # → sys.exit(1) in the caller # Older restart paths may reach here without ``runner.exit_code``; keep the non-zero fallback. if runner._restart_via_service: logger.info( "Exiting with code 75 (service-restart requested) so the service " - "manager relaunches the gateway." - ) + "manager relaunches the gateway.") raise SystemExit(75) return True @@ -5379,8 +5256,7 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = existing_pid = get_running_pid() if ( existing_pid is not None and existing_pid != os.getpid() - and not await _start_gateway_replace_existing_instance(existing_pid, replace) - ): + and not await _start_gateway_replace_existing_instance(existing_pid, replace)): return False _start_gateway_configure_logging(verbosity) @@ -5398,8 +5274,7 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = _signal_initiated_shutdown = [False] shutdown_signal_handler = _start_gateway_make_shutdown_signal_handler( - runner, _signal_initiated_shutdown - ) + runner, _signal_initiated_shutdown) def restart_signal_handler(): runner.request_restart(detached=False, via_service=True) @@ -5501,8 +5376,7 @@ async def start_gateway(config: Optional[GatewayConfig] = None, replace: bool = _shutdown_gateway_health_export(runner) cron_stop, cron_provider, cron_thread, housekeeping_thread = ( - _start_gateway_start_cron_and_housekeeping(runner) - ) + _start_gateway_start_cron_and_housekeeping(runner)) # READY is emitted only after adapters, cron and housekeeping reach their running boundary; # missing config/systemd runtime state leaves the watchdog disabled without changing behavior.