diff --git a/gateway/run_adapters.py b/gateway/run_adapters.py index 3d62e6e6db..ba6e146dca 100644 --- a/gateway/run_adapters.py +++ b/gateway/run_adapters.py @@ -38,11 +38,9 @@ class GatewayAdapterLifecycleMixin: @staticmethod async def _wait_or_detach(task: "asyncio.Future", timeout: float) -> bool: - """Wait up to ``timeout`` for ``task``; on deadline (or our own cancellation) detach it. - - Not ``asyncio.wait_for``: that WAITS for the cancelled child to exit, so a connect()/close() - that swallows ``CancelledError`` blocks recovery forever. Returns True if it finished in time. - """ + """Wait up to ``timeout`` for ``task``; on deadline (or our own cancellation) detach it. Not + ``asyncio.wait_for``: that WAITS for the cancelled child, so a connect()/close() swallowing + ``CancelledError`` blocks recovery forever. True if it finished in time.""" try: done, _pending = await asyncio.wait({task}, timeout=timeout) except asyncio.CancelledError: @@ -67,11 +65,8 @@ class GatewayAdapterLifecycleMixin: return True async def _safe_adapter_disconnect(self, adapter, platform) -> None: - """Call adapter.disconnect() defensively (bounded, never raises). - - After a failed connect() partial resources (ClientSession, poll tasks, subprocesses) would - otherwise leak; must tolerate partial-init state. - """ + """Call adapter.disconnect() defensively (bounded, never raises, tolerates partial-init state): + after a failed connect() partial resources (ClientSession, poll tasks, subprocesses) leak.""" timeout = self._adapter_disconnect_timeout_secs() label = platform.value if platform is not None else "adapter" with _log_suppressed(logging.DEBUG, "Defensive %s disconnect after failed connect raised: %s", label): @@ -82,18 +77,14 @@ class GatewayAdapterLifecycleMixin: ) async def _bounded_adapter_teardown(self, adapter, platform, *, profile: Optional[str] = None) -> None: - """Tear down one adapter on the shutdown path with bounded awaits (never raises). - - Unbounded, a half-dead transport stalls shutdown past systemd's ``TimeoutStopSec``; the - SIGKILL skips ``atexit`` PID-file cleanup and the next start dies with "PID file race lost". - """ + """Tear down one adapter on the shutdown path with bounded awaits (never raises). Unbounded, + a half-dead transport stalls past systemd's ``TimeoutStopSec``; the SIGKILL skips ``atexit`` + PID-file cleanup and the next start dies with "PID file race lost".""" timeout = self._adapter_disconnect_timeout_secs() suffix = f" (profile: {profile})" if profile else "" started_at = time.monotonic() try: - if not await self._await_adapter_cleanup_with_timeout( - adapter.cancel_background_tasks(), timeout - ): + if not await self._await_adapter_cleanup_with_timeout(adapter.cancel_background_tasks(), timeout): logger.warning( "✗ %s background-task cancel timed out after %.1fs - forcing continue%s", platform.value, timeout, suffix, @@ -177,11 +168,9 @@ class GatewayAdapterLifecycleMixin: await self.hooks.emit(event_name, ctx) async def _handle_adapter_fatal_error(self, adapter: BasePlatformAdapter) -> None: - """React to an adapter failure after startup (retryable → background reconnect queue). - - Detached task: the notification arrives on the failing adapter's own polling task, which - the handler's disconnect can cancel mid-flight, stranding the platform half-handled. - """ + """React to an adapter failure after startup (retryable → background reconnect queue). Runs + detached: the notification arrives on the failing adapter's own polling task, which the + handler's disconnect can cancel mid-flight, stranding the platform half-handled.""" tasks = getattr(self, "_fatal_handler_tasks", None) if tasks is None: tasks = self._fatal_handler_tasks = set() @@ -363,16 +352,12 @@ class GatewayAdapterLifecycleMixin: def _spawn_supervised( self, coro_factory, name, *, restart=True, _attempt=0, on_spawn=None, on_give_up=None ): - """Launch a long-lived background task with task-level supervision. - - Catches outer-loop/pre-try exceptions a bare ``create_task`` drops silently. Restarts with - capped backoff up to ``_MAX_SUPERVISED_RESTARTS`` rapid failures; the counter resets after - ``_SUPERVISED_HEALTHY_SECS`` of healthy running. Each spawn uses a fresh ``Context`` (an - inherited delegated-child marker would make the Kanban dispatcher reject its own writes). - ``on_spawn`` fires on EVERY spawn incl. respawns — callers tracking the handle elsewhere - MUST pass it or a respawn leaves a stale handle and a SECOND watcher. ``on_give_up(name)`` - fires when the restart budget is spent. - """ + """Launch a long-lived supervised background task: exceptions a bare ``create_task`` drops + are logged, and it respawns with capped backoff up to ``_MAX_SUPERVISED_RESTARTS`` rapid + failures (counter resets after ``_SUPERVISED_HEALTHY_SECS`` healthy). Fresh ``Context`` per + spawn (an inherited delegated-child marker would make the Kanban dispatcher reject its own + writes). ``on_spawn`` fires on EVERY spawn incl. respawns — handle trackers MUST pass it or a + respawn leaves a stale handle and a SECOND watcher; ``on_give_up(name)`` fires at budget end.""" # Spawn timestamp lets ``_done`` tell a rapid crash-loop from a healthy-run-then-crash. _started = time.monotonic() @@ -434,11 +419,9 @@ class GatewayAdapterLifecycleMixin: return task async def _handoff_watcher(self, interval: float = 2.0, drain_timeout: float = 30.0) -> None: - """Process pending CLI→gateway session handoffs from ``state.db``. - - Claim atomically (pending → running), re-bind the home channel to the CLI session_id, - dispatch a synthetic ``MessageEvent``, mark ``completed``/``failed``. - """ + """Process pending CLI→gateway session handoffs from ``state.db``: claim atomically (pending + → running), re-bind the home channel to the CLI session_id, dispatch a synthetic event, mark + ``completed``/``failed``.""" from gateway.run import _async_profile_runtime_scope, _handoff_watch_scopes, _reclaim_stale await asyncio.sleep(5) # let platforms connect before dispatching through them @@ -470,12 +453,9 @@ class GatewayAdapterLifecycleMixin: inflight.pop(session_id, None) async def _tick(profile_name: Optional[str] = None) -> None: - """One poll of the CURRENTLY-SCOPED store; ``profile_name`` (None = root) routes - delivery to that profile's OWN adapter. - - A closure, not a method: tests bind ``_handoff_watcher`` onto a ``SimpleNamespace`` - with only ``_session_db``/``_running``/``_process_handoff``. - """ + """One poll of the CURRENTLY-SCOPED store; ``profile_name`` (None = root) routes delivery + to that profile's OWN adapter. A closure, not a method: tests bind ``_handoff_watcher`` onto + a ``SimpleNamespace`` with only ``_session_db``/``_running``/``_process_handoff``.""" session_db = getattr(self, "_session_db", None) if session_db is None: return @@ -637,8 +617,7 @@ class GatewayAdapterLifecycleMixin: queued_for / 3600.0, info.get("attempts", 0), ) self._update_platform_runtime_status( - platform.value, platform_state="retrying", needs_attention=True, - retrying_since=retrying_since_iso, + platform.value, platform_state="retrying", needs_attention=True, retrying_since=retrying_since_iso ) def _mark_platform_fatal(self, status_key: str, adapter) -> None: @@ -763,9 +742,7 @@ class GatewayAdapterLifecycleMixin: try: self._schedule_resume_pending_sessions(platform=platform) except Exception: - logger.debug( - "resume-pending reschedule after %s reconnect failed", platform.value, exc_info=True - ) + logger.debug("resume-pending reschedule after %s reconnect failed", platform.value, exc_info=True) async def _cancel_secondary_profile_reconnect_tasks(self) -> None: """Cancel profile-scoped reconnects before tearing down their registry, so a reconnect @@ -794,10 +771,8 @@ class GatewayAdapterLifecycleMixin: async def _start_secondary_profile_adapters(self) -> int: """Bring up adapters for every non-active profile (multiplex only); returns connected count. - - Each profile connects under its own HERMES_HOME + secret scope; credential/listener - collisions are refused here — the only point seeing every profile's credentials together. - """ + Each profile connects under its own HERMES_HOME + secret scope; credential/listener collisions + are refused here — the only point seeing every profile's credentials together.""" from gateway.run import ( MultiplexConfigError, SecondaryPortBindingConfigError, _multiplex_profile_homes ) @@ -825,9 +800,7 @@ class GatewayAdapterLifecycleMixin: except MultiplexConfigError: raise except Exception as e: - logger.error( - "Failed to start adapters for profile '%s': %s", profile_name, e, exc_info=True - ) + logger.error("Failed to start adapters for profile '%s': %s", profile_name, e, exc_info=True) self._record_served_profiles(active, profile_homes) return connected @@ -863,11 +836,9 @@ class GatewayAdapterLifecycleMixin: write_runtime_status(served_profiles=served) async def _load_secondary_profile_config(self, profile_name: str, profile_home: "Path"): - """Hydrate + enter ``profile_home``'s scope once; return its gateway config. - - Raises ``MultiplexConfigError`` (open dm/group policy) or ``SecondaryPortBindingConfigError`` - (the default profile owns the single shared HTTP listener). - """ + """Hydrate + enter ``profile_home``'s scope once; return its gateway config. Raises + ``MultiplexConfigError`` (open dm/group policy) or ``SecondaryPortBindingConfigError`` (the + default profile owns the single shared HTTP listener).""" from gateway.run import ( MultiplexConfigError, SecondaryPortBindingConfigError, _load_gateway_runtime_config, _own_policy_open_startup_violation, _profile_runtime_scope, @@ -897,8 +868,7 @@ class GatewayAdapterLifecycleMixin: raise MultiplexConfigError( f"Profile '{profile_name}' enables {violation}. " "Enable GATEWAY_ALLOW_ALL_USERS or the platform allow-all flag " - "for that profile, or change dm_policy/group_policy away from " - "'open'." + "for that profile, or change dm_policy/group_policy away from 'open'." ) port_binding_platforms = sorted( @@ -1074,11 +1044,9 @@ class GatewayAdapterLifecycleMixin: adapter._hermes_profile_name = profile_name async def _secondary_reconnect_attempt(self, profile_name: str, platform: Platform): - """One scoped attempt to rebuild+connect a secondary adapter → ``(adapter, success)``. - - ``(None, None)`` = give up for good (disabled, credential removed, adapter unavailable). - Caller tears down a RETURNED adapter; one whose configure/connect raised is torn down here. - """ + """One scoped attempt to rebuild+connect a secondary adapter → ``(adapter, success)``; + ``(None, None)`` = give up for good (disabled, credential removed, adapter unavailable). Caller + tears down a RETURNED adapter; one whose configure/connect raised is torn down here.""" from gateway.run import _platform_has_bot_credential, _profile_runtime_scope # Lazy + per-attempt: keeps test monkeypatches on these modules live. from hermes_cli.profiles import get_profile_dir @@ -1162,8 +1130,7 @@ class GatewayAdapterLifecycleMixin: attempts += 1 backoff = _reconnect_backoff(attempts) logger.info( - "Secondary %s reconnect retry in %ds (profile: %s)", platform.value, backoff, - profile_name, + "Secondary %s reconnect retry in %ds (profile: %s)", platform.value, backoff, profile_name ) await asyncio.sleep(backoff) finally: diff --git a/gateway/run_notifications.py b/gateway/run_notifications.py index caaeed1463..df42a624c3 100644 --- a/gateway/run_notifications.py +++ b/gateway/run_notifications.py @@ -245,31 +245,24 @@ class GatewayNotificationsMixin: if switched is None: logger.warning( "Async-delegation completion could not bind routing key %s to " - "owning session %s; dropping injection.", - session_entry.session_key, target_session_id, + "owning session %s; dropping injection.", session_entry.session_key, target_session_id, ) return None logger.info( - "Pinned async-delegation completion to owning session %s " - "(was %s) for routing key %s (#57498)", + "Pinned async-delegation completion to owning session %s (was %s) for routing key %s (#57498)", target_session_id, prior_session_id, session_entry.session_key, ) return switched async def _deliver_media_from_response( - self, response: str, event: MessageEvent, adapter, - thread_metadata: Optional[Dict[str, Any]] = None, + self, response: str, event: MessageEvent, adapter, thread_metadata: Optional[Dict[str, Any]] = None ) -> None: - """Extract explicit MEDIA: tags from an already-streamed response and deliver them. - - The text is already delivered; this only handles file attachments. Unlike the non-streaming - path in ``gateway/platforms/base.py`` this rescan is EXPLICIT-ONLY: a bare local path in a - streamed reply was either shown as text or is stale inspected content, and promoting it sent - files the model never asked to deliver. MEDIA tags are NOT deduped against prior turns: a - MEDIA: directive in the final reply is a deliberate attach (incl. user-requested resends); - stale auto-appended tags are deduped upstream (_collect_auto_append_media_tags). - """ + """Deliver explicit MEDIA: tags from an already-streamed response (text already delivered). + EXPLICIT-ONLY, unlike the non-streaming path in ``gateway/platforms/base.py``: a bare local + path in a streamed reply is shown text or stale inspected content, and promoting it sent + files the model never asked for. MEDIA tags are NOT deduped against prior turns (a final-reply + directive is a deliberate attach); stale auto-appended tags are deduped upstream.""" from urllib.parse import quote as _quote with _log_suppressed(logging.WARNING, "Post-stream media extraction failed: %s"): # Capture [[as_document]] before extract_media strips it: image files then go through @@ -378,8 +371,7 @@ class GatewayNotificationsMixin: from gateway.run import _hermes_home return cls._UpdatePaths( pending=_hermes_home / ".update_pending.json", - claimed=_hermes_home / ".update_pending.claimed.json", - output=_hermes_home / ".update_output.txt", + claimed=_hermes_home / ".update_pending.claimed.json", output=_hermes_home / ".update_output.txt", exit_code=_hermes_home / ".update_exit_code", prompt=_hermes_home / ".update_prompt.json", response=_hermes_home / ".update_response", ) @@ -664,8 +656,7 @@ class GatewayNotificationsMixin: platform_cfg = self.config.platforms.get(platform) if platform_cfg is not None and not platform_cfg.gateway_restart_notification: logger.info( - "Restart notification suppressed: %s has gateway_restart_notification=false", - platform_str, + "Restart notification suppressed: %s has gateway_restart_notification=false", platform_str ) return None @@ -833,8 +824,7 @@ class GatewayNotificationsMixin: try: platform = Platform(platform_name) - # Reject arbitrary strings (dynamic pseudo-members): built-ins are always valid, plugin - # platforms must be registered in the platform registry. + # Reject dynamic pseudo-members: plugin platforms must be registered. if platform.value not in _BUILTIN_PLATFORM_VALUES: try: from gateway.platform_registry import platform_registry @@ -851,15 +841,12 @@ class GatewayNotificationsMixin: scope_id = _opt("scope_id") if scope_id is None and chat_type not in ("dm", "thread"): - # Reconstructed (non-persisted) source for a scoped chat with no scope discriminator: a - # relay connector's fail-closed tenant guard may decline the reply unless user_id resolves - # it. Don't fail — DMs/author-bound chats still route and native adapters need no - # scope_id — but warn so a post-restart egress decline isn't silent. + # Reconstructed scoped-chat source without scope_id: a relay connector's tenant guard may + # decline the reply. Warn, don't fail (native adapters need no scope_id). logger.warning( "Synthetic event source for %s chat=%s (%s) reconstructed " "without scope_id; scoped relay egress may be declined by " - "the connector's tenant guard (user_id fallback only).", - platform_name, chat_id, chat_type, + "the connector's tenant guard (user_id fallback only).", platform_name, chat_id, chat_type, ) return SessionSource( platform=platform, chat_id=chat_id, chat_type=chat_type, thread_id=_opt("thread_id"), @@ -916,15 +903,9 @@ class GatewayNotificationsMixin: return False def _resolve_injection_adapter(self, platform_name: str): - """Adapter for a synthetic-event platform: transport resolver first, literal scan as fallback. - - Alias-aware resolution (relay-plane): one adapter under Platform.RELAY fronts N logical - platforms, so a literal ``p.value == platform_name`` scan misses "slack" and drops the - completion as "no gateway route". Native adapter wins; relay is eligible only when it - advertises fronting the logical platform. The legacy literal scan is still correct for - native adapters and keeps minimal runner stubs (tests) / exotic platform strings working - when the resolver can't run. - """ + """Adapter for a synthetic-event platform: alias-aware transport resolver first (one + Platform.RELAY adapter fronts N logical platforms; native wins), literal ``p.value`` scan as + fallback for minimal runner stubs / exotic platform strings when the resolver can't run.""" from gateway.run import resolve_delivery_transport try: _transport = resolve_delivery_transport(Platform(platform_name), self.config, self.adapters) @@ -1042,15 +1023,10 @@ class GatewayNotificationsMixin: return seen async def _classify_completion_target(self, parent_session_id: str) -> str: - """Classify an async-completion delivery target before adapter acceptance. - - - ``"deliver"``: spawning session live (or compression-rotated with a live continuation); - proves deliverability only, the resolver still retargets. - - ``"terminal"``: parent gone for good (unknown / explicit user boundary like /new); drop - the durable row rather than falsely ack or replay forever. - - ``"retry"``: transient uncertainty (DB unavailable, rotation mid-flight); release the - claim for a later consumer; the attempt cap bounds churn. - """ + """Classify an async-completion target before adapter acceptance: ``"deliver"`` (spawning + session live or compression-rotated with a live continuation; the resolver still retargets), + ``"terminal"`` (parent gone for good — unknown / user boundary like /new; drop the durable row + rather than falsely ack), ``"retry"`` (DB unavailable / rotation mid-flight; release the claim).""" from gateway.run import _USER_BOUNDARY_END_REASONS session_db = getattr(self, "_session_db", None) if session_db is None: @@ -1139,8 +1115,7 @@ class GatewayNotificationsMixin: logger.warning( "Background process %s completion targets " "permanently-gone session %s (user boundary such as " - "/new); dropping notification (output remains " - "available via process(action='log')).", + "/new); dropping notification (output remains available via process(action='log')).", evt.get("session_id") or "", parent_session_id, ) claim.proceed = False @@ -1309,12 +1284,10 @@ class GatewayNotificationsMixin: # Some unit tests construct GatewayRunner with object.__new__. Keep the # batching seam lazy so those focused lifecycle tests remain valid. for attr, default in ( - ("_completion_notification_batches", dict), - ("_completion_notification_batch_tasks", dict), + ("_completion_notification_batches", dict), ("_completion_notification_batch_tasks", dict), ("_completion_notification_batch_flush_tasks", set), ("_completion_notification_batch_window", lambda: 0.1), - ("_completion_notification_batches_stopping", lambda: False), - ("_background_tasks", set), + ("_completion_notification_batches_stopping", lambda: False), ("_background_tasks", set), ): if not hasattr(self, attr): setattr(self, attr, default()) @@ -1360,14 +1333,10 @@ class GatewayNotificationsMixin: logger.debug(fail_msg, exc_info=True) async def _deliver_async_delegation_group(self, group: list[dict]) -> Optional[bool]: - """Deliver a same-session batch of async completions as ONE turn. - - Single-event groups ride the per-event path. Multi-event groups deliver the primary via - ``_deliver_completion_notification`` with consolidated text of every sibling THIS runner - claimed; sibling claims are acked only after adapter acceptance, and siblings claimed by - another consumer are excluded (no double delivery). Returns True after acceptance, False - to requeue the group, None when nothing is deliverable here (retry siblings requeued). - """ + """Deliver a same-session batch of async completions as ONE turn: the primary carries the + consolidated text of every sibling THIS runner claimed (siblings owned elsewhere are excluded; + their claims are acked only after adapter acceptance). True after acceptance, False to requeue + the group, None when nothing is deliverable here (retry siblings requeued).""" from gateway.run import _format_gateway_process_notification from tools.process_registry import process_registry as _pr deliverable: list[tuple[dict, str]] = [] @@ -1547,13 +1516,10 @@ class GatewayNotificationsMixin: ) async def _run_process_watcher(self, watcher: dict) -> None: - """Periodically check a background process and push updates to the user. - - Runs as an asyncio task. Stays silent when nothing changed. Auto-removes when the process - exits or is killed. Notification mode (``display.background_process_notifications``): - concise (default, one-line; failures append output tail) / all (running updates + final - raw output) / result (final raw only) / error (final raw only if exit != 0) / off. - """ + """Poll a background process and push updates until it exits. Mode + (``display.background_process_notifications``): concise (default one-liner; failures append + the output tail) / all (running updates + final raw) / result (final raw) / error (final raw + if exit != 0) / off.""" from tools.process_registry import format_process_notification, process_registry session_id = watcher["session_id"] interval = watcher["check_interval"] diff --git a/gateway/run_shutdown.py b/gateway/run_shutdown.py index 6a7263d174..90f65723f2 100644 --- a/gateway/run_shutdown.py +++ b/gateway/run_shutdown.py @@ -385,13 +385,10 @@ class GatewayShutdownMixin: return self.adapters.get(Platform.RELAY) async def _scale_to_zero_watcher(self, interval: float = 30.0) -> None: - """Watch for idle, drive the relay dormant, then self-suspend the machine. - - On sustained idle: status `draining` (NOT _running=False), relay go_dormant() (socket close, - NOT disconnect()), no mark_resume_pending (suspend preserves RAM), THEN suspend via the flaps - socket. The gateway owns the suspend because Fly autostop sees only INBOUND connections and - would freeze mid-job or before the relay flip. Off-Fly the watcher does not quiesce at all. - """ + """Watch for idle, drive the relay dormant, then self-suspend. On sustained idle: status + `draining` (NOT _running=False), relay go_dormant() (socket close, NOT disconnect()), no + mark_resume_pending (suspend preserves RAM), THEN suspend via the flaps socket — Fly autostop + sees only INBOUND connections and would freeze mid-job. Off-Fly: no quiesce at all.""" await asyncio.sleep(min(interval, 30.0)) # let startup settle while self._running: try: @@ -473,8 +470,7 @@ class GatewayShutdownMixin: self._external_drain_active = True logger.info( "External drain ENGAGED (.drain_request.json present) — refusing " - "new turns; %d in-flight turn(s) will finish. Process stays up.", - self._active_work_count(), + "new turns; %d in-flight turn(s) will finish. Process stays up.", self._active_work_count(), ) # Flip persisted lifecycle state so /api/status.gateway_busy / gateway_drainable track the # drain; active_agents is preserved (read-merge keeps the live count), only state changes. @@ -502,9 +498,8 @@ class GatewayShutdownMixin: from gateway.drain_control import drain_requested while self._running: try: - # drain_requested() does a synchronous read_text() on the marker file: at 1s cadence - # that is a blocking disk read on the event loop ~86k times/day, and under host I/O - # pressure one read can stall 30s+ and take every platform heartbeat down. Off-thread it. + # Off-thread: a synchronous marker read at 1s cadence can stall 30s+ under host I/O + # pressure and take every platform heartbeat down. if await asyncio.to_thread(drain_requested): self._enter_external_drain() # API and cron work live outside messaging's _running_agents map; refresh the @@ -1085,9 +1080,8 @@ class GatewayShutdownMixin: ) watcher_env = GatewayShutdownMixin._restart_watcher_env() project_root = Path(__file__).resolve().parent.parent - # Console python under CREATE_NO_WINDOW owns one hidden console inherited by the restart - # child, so nothing flashes. Do NOT swap in pythonw.exe — a console-less watcher forces - # every console-subsystem descendant to allocate a visible conhost. + # Console python under CREATE_NO_WINDOW: nothing flashes. NOT pythonw.exe — a console-less + # watcher makes every console-subsystem descendant allocate a visible conhost. watcher_python = sys.executable venv_dir = Path(watcher_env.get("VIRTUAL_ENV") or project_root / "venv") site_packages = venv_dir / "Lib" / "site-packages" @@ -1102,10 +1096,8 @@ class GatewayShutdownMixin: str(current_pid), str(restart_after_s), *hermes_cmd, "gateway", "restart", ] popen_kwargs = dict(stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, env=watcher_env) - # The watcher must break away from the parent CLI's job object (Desktop wrappers, Windows - # Terminal, schtasks) or it is reaped when the CLI exits. A job without - # JOB_OBJECT_LIMIT_BREAKAWAY_OK rejects CREATE_BREAKAWAY_FROM_JOB (OSError); retry once - # without the bit, preserving argv and the scrubbed env. + # Break away from the parent CLI's job object or be reaped when the CLI exits; a job without + # BREAKAWAY_OK rejects CREATE_BREAKAWAY_FROM_JOB (OSError) — retry once without the bit. try: subprocess.Popen(watcher_argv, **popen_kwargs, **windows_detach_popen_kwargs()) except OSError: @@ -1406,8 +1398,7 @@ class GatewayShutdownMixin: if _cron_at_start and _cron_timeout > timeout: logger.info( "Shutdown drain: %d in-flight cron job(s) — waiting up to " - "%.0fs for them (cron_drain_timeout=%.0fs, " - "restart_drain_timeout=%.0fs)", + "%.0fs for them (cron_drain_timeout=%.0fs, restart_drain_timeout=%.0fs)", _cron_at_start, _cron_timeout, _cron_drain_cfg, timeout, ) _drain_started_at = time.monotonic() @@ -1415,8 +1406,7 @@ class GatewayShutdownMixin: ctx.drain_elapsed = time.monotonic() - _drain_started_at logger.info( "Shutdown phase: drain done at +%.2fs (drain took %.2fs, " - "timed_out=%s, active_at_start=%d, active_now=%d, " - "cron_at_start=%d, cron_now=%d, " + "timed_out=%s, active_at_start=%d, active_now=%d, cron_at_start=%d, cron_now=%d, " "api_at_start=%d, api_now=%d, " "deferred_at_start=%d, deferred_now=%d)", ctx.elapsed(), ctx.drain_elapsed, ctx.timed_out, len(ctx.active_agents), self._running_agent_count(), _cron_at_start, @@ -1439,15 +1429,12 @@ class GatewayShutdownMixin: from gateway.run import GatewayRunner logger.warning( "Gateway drain timed out after %.1fs with %d active agent(s), " - "%d in-flight cron job(s), %d api_server run(s), and " - "%d deferred agent worker(s); " + "%d in-flight cron job(s), %d api_server run(s), and %d deferred agent worker(s); " "interrupting remaining work.", ctx.drain_elapsed, self._running_agent_count(), self._active_cron_job_count(), self._active_api_run_count(), ctx.deferred_count(), ) - # Mark forcibly-interrupted sessions resume_pending BEFORE interrupting so the next message - # auto-resumes instead of becoming a fresh session via suspend_recently_active(); stuck - # sessions still escalate via ``.restart_failure_counts`` (suspended=True wins). Uses the - # CURRENT _running_agents, not the drain-start snapshot, so cleanly finished sessions get no note. + # Mark resume_pending BEFORE interrupting so the next message auto-resumes (stuck sessions + # still escalate via .restart_failure_counts). CURRENT _running_agents, not the drain snapshot. await GatewayRunner._mark_running_sessions_resume_pending(self, "mark_resume_pending") reason = GatewayRunner._shutdown_interrupt_reason(self) self._interrupt_running_agents(reason) @@ -1568,11 +1555,9 @@ class GatewayShutdownMixin: def _stop_quiesce_and_close_session_dbs(self, timeout: float, ctx: "GatewayShutdownMixin._StopContext") -> None: """Quiesce the executor, then close SessionDB handles only if no worker is still live.""" from gateway.run import GatewayRunner, _EXECUTOR_QUIESCE_TIMEOUT - # Quiesce the thread pool BEFORE closing session DBs: otherwise a coroutine can mint a fresh - # pool (``_executor_closing`` still False) or an already-started run_in_executor future keeps - # writing after ``SessionDB.close()`` checkpointed the WAL — the late write reopens the handle - # and mints a split WAL generation (close-time corruption). Bounded and clamped to what is - # left of the watchdog leash minus a second for the close itself. + # Quiesce the thread pool BEFORE closing session DBs: a late executor write after + # SessionDB.close() checkpointed the WAL reopens the handle and splits the WAL generation + # (close-time corruption). Clamped to the remaining watchdog leash minus 1s for the close. _exec_quiesce_budget = max( 0.0, min(_EXECUTOR_QUIESCE_TIMEOUT, resolve_shutdown_watchdog_delay(timeout) - ctx.elapsed() - 1.0), ) @@ -1637,8 +1622,7 @@ class GatewayShutdownMixin: else: logger.info( "Skipping .clean_shutdown marker — drain timed out with " - "interrupted agents; next startup will suspend recently " - "active sessions." + "interrupted agents; next startup will suspend recently active sessions." ) # Stuck-loop detection: the counter increments for sessions active at each restart; at # the threshold (3 consecutive) the next startup auto-suspends the session. @@ -1661,14 +1645,12 @@ class GatewayShutdownMixin: self._exit_code = GATEWAY_SERVICE_RESTART_EXIT_CODE self._exit_reason = self._exit_reason or "Gateway restart requested" self._draining = False - # Terminal gateway_state: "stopped" by default, "running" on an UNEXPECTED signal (s6 SIGTERM - # on docker restart, OOM-kill) — container_boot.py only auto-starts gateways last seen - # "running". Operator stops write a planned-stop marker BEFORE signalling; restarts persist "stopped". + # Terminal gateway_state: "stopped", or "running" on an UNEXPECTED signal (docker restart, + # OOM) — container_boot.py only auto-starts gateways last seen "running". if getattr(self, "_signal_initiated_shutdown", False) and not self._restart_requested: logger.info( "Gateway stopped by an unexpected signal — persisting " - "gateway_state=running so container_boot auto-starts on " - "the next boot (issue #42675)" + "gateway_state=running so container_boot auto-starts on the next boot (issue #42675)" ) self._update_runtime_status("running", self._exit_reason) else: @@ -1694,9 +1676,8 @@ class GatewayShutdownMixin: async def _stop_impl(self) -> None: """Run every ``_stop_*`` phase under the thread-based shutdown watchdog.""" from gateway.run import GatewayRunner - # Thread-based watchdog (asyncio timeouts cannot recover a frozen loop): if teardown never - # finishes within drain+grace it dumps faulthandler stacks and os._exit so KeepAlive/systemd - # can revive. Skipped under pytest (no delayed hard-exit in the worker). + # Thread-based watchdog (asyncio timeouts cannot recover a frozen loop): dumps stacks and + # os._exit past drain+grace so the service manager revives us. Skipped under pytest. _watchdog_done = threading.Event() self._shutdown_watchdog_done = _watchdog_done # Shutdown-path tests and third-party runner doubles may only implement the older @@ -1706,10 +1687,8 @@ class GatewayShutdownMixin: ) if not os.environ.get("PYTEST_CURRENT_TEST"): arm_shutdown_watchdog( - resolve_shutdown_watchdog_delay(self._restart_drain_timeout), - done_event=_watchdog_done, - snapshot_fn=lambda: GatewayRunner._shutdown_watchdog_snapshot(self, ctx), - exit_code=1, + resolve_shutdown_watchdog_delay(self._restart_drain_timeout), done_event=_watchdog_done, + snapshot_fn=lambda: GatewayRunner._shutdown_watchdog_snapshot(self, ctx), exit_code=1, ) try: await GatewayRunner._stop_begin_teardown(self, ctx) @@ -1725,8 +1704,7 @@ class GatewayShutdownMixin: _watchdog_done.set() async def stop( - self, *, restart: bool = False, detached_restart: bool = False, - service_restart: bool = False, + self, *, restart: bool = False, detached_restart: bool = False, service_restart: bool = False ) -> None: """Stop the gateway and disconnect all adapters.""" from gateway.run import GatewayRunner diff --git a/gateway/run_startup.py b/gateway/run_startup.py index d8d3a1d2d0..086b87dcd3 100644 --- a/gateway/run_startup.py +++ b/gateway/run_startup.py @@ -130,8 +130,7 @@ class GatewayStartupMixin: await self._wait_bounded_or_release( {task}, timeout, "Turn-machinery warm-up still running after %.0fs; opening inbound gate anyway — the " - "first turn may see lazily initialized machinery (#99373). Warm-up continues in the " - "background.", + "first turn may see lazily initialized machinery (#99373). Warm-up continues in the background.", "boot turn-machinery warm-up failed after gate release", level=logging.DEBUG, ) @@ -158,13 +157,10 @@ class GatewayStartupMixin: return done async def _finish_startup_restore(self) -> None: - """Wait (BOUNDED) for startup auto-resume, then release + drain inbound. - - Bounded by ``_startup_restore_drain_timeout_secs`` so one pathological boot-resume turn - cannot hold the gate shut for every channel; on timeout the gate opens and resume turns - finish in the background (NOT cancelled). Safe because ``_schedule_resume_pending_sessions`` - claims each ``_running_agents`` slot SYNCHRONOUSLY first, so drained inbound queues behind. - """ + """Wait (BOUNDED by ``_startup_restore_drain_timeout_secs``) for startup auto-resume, then + release + drain inbound. On timeout the gate opens and resume turns finish in the background + (NOT cancelled) — safe because ``_schedule_resume_pending_sessions`` claims each + ``_running_agents`` slot SYNCHRONOUSLY first, so drained inbound queues behind.""" from gateway.run import _startup_restore_drain_timeout_secs tasks = list(getattr(self, "_startup_restore_tasks", []) or []) if tasks: @@ -203,15 +199,11 @@ class GatewayStartupMixin: return _report async def _await_startup_boot_sends(self, *, planned_restart_notification_pending: bool) -> None: - """Run boot-path sends without letting them pin the inbound restore gate. - - Awaiting the sends inline before the gate releases lets one Telegram flood-control sleep - freeze inbound on every platform. Same bounded ``asyncio.wait`` as the resume gate: on - timeout return and let the sends finish in the background (not cancelled). The ledger claim - + ``resume_pending`` clear run INLINE before the send task exists (bounded DB work): deferring - it let a hung notification expire the gate with zero rows claimed, so answered turns were - replayed AND redelivered. - """ + """Run boot-path sends without letting them pin the inbound restore gate (one Telegram + flood-control sleep must not freeze inbound on every platform): same bounded wait as the + resume gate, sends finish in the background on timeout. The ledger claim + ``resume_pending`` + clear run INLINE before the send task exists — deferring it let a hung notification expire the + gate with zero rows claimed, so answered turns were replayed AND redelivered.""" from gateway.run import _clear_planned_restart_notification, _startup_restore_drain_timeout_secs claimed = await self._claim_pending_obligations() @@ -259,15 +251,11 @@ class GatewayStartupMixin: return sendable async def _claim_pending_obligations(self) -> list: - """Claim recoverable delivery-ledger rows and clear their ``resume_pending`` flags. - - Pure DB work, no sends. Must run INLINE at startup BEFORE ``_schedule_resume_pending_sessions`` - and before the abandonable boot-send task exists: these sessions already produced their - answer, so clearing ``resume_pending`` here stops the resume path from re-running (and - re-paying for) the turn however long the sends take. Rows that were mid-send or previously - rejected carry a visible recovered-reply marker so a possible duplicate is labeled, never - silent (gateway/delivery_ledger.py). Returns the claimed rows for redelivery. - """ + """Claim recoverable delivery-ledger rows and clear their ``resume_pending`` flags (pure DB + work, no sends). Must run INLINE BEFORE ``_schedule_resume_pending_sessions`` and the + abandonable boot-send task: these sessions already produced their answer, so the resume path + must not re-run (re-pay for) the turn however long the sends take. Mid-send / rejected rows + carry a visible recovered-reply marker (gateway/delivery_ledger.py). Returns the rows.""" try: from gateway.delivery_ledger import ledger_enabled, sweep_recoverable if not await asyncio.to_thread(ledger_enabled): @@ -333,9 +321,7 @@ class GatewayStartupMixin: try: result = await adapter.send(chat_id=row["chat_id"], content=content, metadata=metadata) except Exception as send_err: - logger.warning( - "obligation %s: redelivery send raised: %s", row["obligation_id"], send_err - ) + logger.warning("obligation %s: redelivery send raised: %s", row["obligation_id"], send_err) result = None with _log_suppressed(logging.DEBUG, "delivery ledger update failed", exc_info=True): if result is not None and getattr(result, "success", False): @@ -348,8 +334,7 @@ class GatewayStartupMixin: ) else: await asyncio.to_thread( - mark_failed, row["obligation_id"], - str(getattr(result, "error", "") or "send failed"), + mark_failed, row["obligation_id"], str(getattr(result, "error", "") or "send failed") ) return redelivered @@ -358,9 +343,7 @@ class GatewayStartupMixin: try: platform = Platform(row["platform"]) except Exception: - logger.debug( - "obligation %s: unknown platform %r", row["obligation_id"], row.get("platform"), - ) + logger.debug("obligation %s: unknown platform %r", row["obligation_id"], row.get("platform")) return None if "profile" in row: adapter = self._authorization_adapter(platform, row.get("profile")) @@ -384,12 +367,9 @@ class GatewayStartupMixin: async def _redeliver_failed_obligations_for_platform( self, platform: Platform, *, profile: Optional[str] = None, ) -> int: - """Replay one adapter identity's transient failures after reconnect. - - The startup sweep cannot claim live-owner rows, and an adapter can reconnect without the - process exiting, so ``send_path_degraded`` responses would otherwise stay failed until the - next restart. Claim/clear/send are best-effort and reuse the startup redelivery contract. - """ + """Replay one adapter identity's transient failures after reconnect: the startup sweep cannot + claim live-owner rows, so ``send_path_degraded`` responses would otherwise stay failed until + the next restart. Best-effort; reuses the startup redelivery contract.""" try: from gateway.delivery_ledger import ledger_enabled, sweep_failed_for_runtime if not await asyncio.to_thread(ledger_enabled): @@ -405,9 +385,7 @@ class GatewayStartupMixin: # Clear before any send so the reconnect path cannot both redeliver an # already-produced answer and schedule the same agent turn for resume. - sendable = await self._clear_resume_pending_for_claimed_obligations( - claimed, require_success=True - ) + sendable = await self._clear_resume_pending_for_claimed_obligations(claimed, require_success=True) sendable_ids = {row["obligation_id"] for row in sendable} for row in claimed: if row["obligation_id"] not in sendable_ids: @@ -434,10 +412,8 @@ class GatewayStartupMixin: logger.warning("Failed to enumerate resume-pending sessions: %s", exc) return None - # Defense-3: break the SIGTERM-respawn loop. Only count this boot when there are restart- - # interrupted sessions to resume — a clean boot must not accrue toward the breaker. If too - # many such boots hit the window, skip auto-resume for THIS boot only: the gateway still - # serves inbound; the session stays resume_pending so a real user message can continue it. + # Restart-loop breaker: only boots WITH restart-interrupted sessions count; when tripped, skip + # auto-resume for THIS boot only (inbound still served; sessions stay resume_pending). if candidates: try: from gateway import restart_loop_guard as _rlg @@ -460,20 +436,14 @@ class GatewayStartupMixin: "longer authorized under the current allowlist", session_key, ) except Exception as exc: - logger.warning( - "Skipping auto-resume for %s: authorization check failed: %s", session_key, exc, - ) + logger.warning("Skipping auto-resume for %s: authorization check failed: %s", session_key, exc) return False def _schedule_resume_pending_sessions(self, platform=None) -> int: - """Auto-continue fresh restart-interrupted sessions after startup. - - Synthesizes the next turn once adapters are back online; the event text is empty so the - existing ``_is_resume_pending`` injection path owns the recovery wording. Sessions whose - adapter is not in ``self.adapters`` stay ``resume_pending`` for the reconnect watcher, which - re-calls this scoped to that ``platform`` (a reconnecting platform never touches another's - recoveries); sessions with a running agent are skipped so none is resumed twice. - """ + """Auto-continue fresh restart-interrupted sessions: synthesize an empty-text turn (the + ``_is_resume_pending`` injection path owns the wording). Sessions whose adapter is offline stay + ``resume_pending`` for the reconnect watcher, which re-calls this scoped to that ``platform``; + sessions with a running agent are skipped so none is resumed twice.""" from gateway.run import _AGENT_PENDING_SENTINEL, _auto_continue_freshness_window window = _auto_continue_freshness_window() candidates = self._resume_pending_candidates(platform) @@ -678,24 +648,19 @@ class GatewayStartupMixin: def _open_faulthandler_log(self): """Open (append) ``/gateway_faulthandler.log``, creating the directory.""" from gateway.run import get_hermes_home - log_dir = getattr(self.config, "log_dir", None) or os.path.join( - str(get_hermes_home()), "logs", - ) + log_dir = getattr(self.config, "log_dir", None) or os.path.join(str(get_hermes_home()), "logs") os.makedirs(log_dir, exist_ok=True) return open(os.path.join(log_dir, "gateway_faulthandler.log"), "a", encoding="utf-8") def _start_install_faulthandler(self) -> None: """Enable faulthandler (stderr or a log file) plus the SIGUSR2 stack-dump hook.""" - # Falls back to a log file when sys.stderr is None (Windows VBS / pythonw / detached - # service) — otherwise the gateway would die here and take every adapter offline. + # sys.stderr may be None (Windows VBS / pythonw / detached service): fall back to a log file. try: faulthandler.enable() except (RuntimeError, ValueError, OSError): with _log_suppressed(logging.DEBUG, "faulthandler.enable() unavailable", exc_info=True): faulthandler.enable(file=self._open_faulthandler_log(), all_threads=True) - # Also dump stacks to a file on SIGUSR2 for off-line analysis under a service manager that - # doesn't capture stderr. faulthandler.register()/SIGUSR2 are POSIX-only: skip on Windows - # (faulthandler.enable() above still covers fatal errors). + # SIGUSR2 stack dump to file for service managers that drop stderr; POSIX-only. _sigusr2 = getattr(signal, "SIGUSR2", None) if _sigusr2 is not None and hasattr(faulthandler, "register"): with _log_suppressed(logging.DEBUG, "Could not set up faulthandler file logging", exc_info=True): @@ -711,10 +676,8 @@ class GatewayStartupMixin: self._gateway_loop = None if self._gateway_loop is not None: self._start_loop_liveness_guards(self._gateway_loop) - # Loop confirmed live: the startup-liveness watchdog is done and the loop-liveness - # watchdog (armed above) takes over. Disarm even when loop guards are config-disabled — - # the startup watchdog covers only the pre-loop window. Deliberately inside this branch: - # if the loop isn't live, startup has NOT reached the milestone and it must stay armed. + # Loop live: the loop-liveness watchdog takes over from the startup watchdog. Disarm even + # when loop guards are config-disabled; only inside this branch (no live loop = stay armed). with _log_suppressed(logging.DEBUG, "Startup watchdog disarm failed", exc_info=True): from gateway.startup_watchdog import disarm_startup_watchdog disarm_startup_watchdog() @@ -752,9 +715,7 @@ class GatewayStartupMixin: logger.info("Active profile: %s", _profile) with suppress(Exception): from gateway.status import write_runtime_status - write_runtime_status( - gateway_state="starting", exit_reason=None, clear_profile_platforms=True - ) + write_runtime_status(gateway_state="starting", exit_reason=None, clear_profile_platforms=True) with _log_suppressed(logging.DEBUG, "gateway health OTLP export startup failed", exc_info=True): from hermes_cli.config import load_config from agent.monitoring.gateway_health_export import start_gateway_health_export @@ -910,9 +871,8 @@ class GatewayStartupMixin: if recovered: logger.info("Recovered %s background process(es) from previous run", recovered) - # Recover sessions active when the gateway last exited. Exact durable turn markers cover - # long-running work; the 120s recency heuristic remains as a fallback for turns from older - # versions without markers. SKIP after a clean shutdown — the previous process already drained. + # Recover sessions active at last exit (exact turn markers + 120s recency fallback for + # marker-less older turns). SKIP after a clean exit — the previous process already drained. _clean_marker = _hermes_home / ".clean_shutdown" if _clean_marker.exists(): logger.info("Previous gateway exited cleanly — skipping session suspension") @@ -954,10 +914,8 @@ class GatewayStartupMixin: return True, enabled_platform_count, _multiplex_skipped_platforms, _pending_connects if not platform_config.enabled: continue - # Under multiplexing, a platform may be enabled on the default profile's config.yaml - # while its bot token lives only in a secondary profile's .env. Starting the primary with - # an empty token fails at once and queues a reconnect loop that can never heal; the - # secondary starts its own adapter with the real token, so skip the empty primary. + # Multiplex: a platform enabled in the shared config.yaml may hold its token only in a + # secondary profile's .env; an empty primary would queue a reconnect loop that never heals. if _multiplex_on and not _platform_has_bot_credential(platform, platform_config): logger.info( "Skipping %s on default profile: no bot credential in this " @@ -976,8 +934,7 @@ class GatewayStartupMixin: else: logger.warning( "No adapter for '%s' -- is the plugin installed? " - "(platform is enabled in config.yaml but no plugin registered it)", - platform.value, + "(platform is enabled in config.yaml but no plugin registered it)", platform.value, ) continue # Under multiplexing the default profile needs the same whole-handler runtime scope as @@ -1010,8 +967,7 @@ class GatewayStartupMixin: # restart/shutdown requested mid-connect must cancel still-pending connects, clean up the # ones already completed, and abort startup. _task_map = { - asyncio.ensure_future(_connect_one_startup(p, c, a)): (p, c, a) - for (p, c, a) in _pending_connects + asyncio.ensure_future(_connect_one_startup(p, c, a)): (p, c, a) for (p, c, a) in _pending_connects } _pending_tasks = set(_task_map) while _pending_tasks: @@ -1064,9 +1020,8 @@ class GatewayStartupMixin: continue if outcome == "exception": logger.error("\u2717 %s error: %s", platform.value, exc) - # Same defensive cleanup path for exceptions -- an adapter that raised mid-connect - # may still have a live aiohttp.ClientSession or child subprocess. Unexpected - # exceptions are typically transient -- queue for retry. + # An adapter that raised mid-connect may still hold a ClientSession/subprocess; treat + # unexpected exceptions as transient and queue for retry. await self._safe_adapter_disconnect(adapter, platform) self._startup_queue_transient_failure( platform, adapter, platform_config, str(exc), startup_retryable_errors @@ -1091,10 +1046,8 @@ class GatewayStartupMixin: platform, adapter, platform_config, "failed to connect", startup_retryable_errors ) continue - # A live foreign holder of this bot token is a single-writer ownership conflict, not a - # blip — ``_acquire_platform_lock`` emits it retryable only so a MID-RUN reconnect can - # recover. At startup route it non-retryable: with nothing connected the gateway exits - # 78 instead of sitting alive and deaf in the retry queue. + # A live foreign token holder is an ownership conflict, not a blip: retryable only for + # MID-RUN reconnects; at startup route it non-retryable so the gateway exits 78, not deaf. _retryable = adapter.fatal_error_retryable and not is_global_startup_conflict(adapter.fatal_error_code) self._update_platform_runtime_status( platform.value, platform_state="retrying" if _retryable else "fatal", @@ -1164,10 +1117,8 @@ class GatewayStartupMixin: self._startup_fail_fatal_config(reason) return True if startup_nonretryable_errors: - # Mixed failure mode (some fatal, some transient). Exiting 78 here would let exit-78 - # supervisors take the gateway PERMANENTLY down over a network blip and deny the - # retryable ones their retry. Log the fatal side loudly, then fall through to the - # degraded/retry path: the watcher recovers the retryable; the rest stay parked. + # Mixed (some fatal, some transient): exiting 78 would take the gateway PERMANENTLY down + # over a blip. Log the fatal side loudly and fall through to the degraded/retry path. logger.error( "%d platform(s) fatally misconfigured and parked: %s. " "Staying alive so retryable platforms can recover.", @@ -1238,27 +1189,22 @@ class GatewayStartupMixin: from gateway.run import _planned_restart_notification_pending, _restart_notification_pending await self._start_post_connect_services(connected_count) - # Give freshly connected adapters a brief moment to settle before sending restart/startup - # lifecycle messages; in practice this helps Discord thread deliveries after reconnect. + # Let fresh adapters settle before lifecycle sends (helps Discord thread deliveries). if connected_count > 0: await asyncio.sleep(1.0) - # Capture, before _send_restart_notification() unlinks the marker, whether this process - # booted from a chat-originated /restart. One-shot signal for the /restart redelivery - # guard (_is_stale_restart_redelivery): a missing dedup marker only suppresses a /restart - # when we KNOW we just came out of a restart cycle. + # Before _send_restart_notification() unlinks the marker: did we boot from a chat /restart? + # One-shot signal for _is_stale_restart_redelivery. if _restart_notification_pending(): self._booted_from_restart = True - # Restart notification, home-channel startup notice, and obligation redelivery all call - # adapter.send(). Those sends must not pin the inbound restore gate — a Telegram flood- - # control sleep on this path froze every platform for the full penalty. + # Boot-path adapter.send() calls must not pin the inbound restore gate (a Telegram flood- + # control sleep here once froze every platform). await self._await_startup_boot_sends( planned_restart_notification_pending=_planned_restart_notification_pending(), ) - # Auto-continue fresh sessions interrupted by the previous restart/shutdown. resume_pending - # is cleared by the normal successful-turn path, so a failed auto-resume stays visible on the - # next user message. _await_startup_boot_sends already cleared sessions answered in the ledger. + # Auto-resume restart-interrupted sessions (ledger-answered ones were cleared above); a failed + # auto-resume stays visible on the next user message. self._schedule_resume_pending_sessions() await self._finish_startup_restore() @@ -1266,10 +1212,8 @@ class GatewayStartupMixin: # is broken before losing data. await self._send_session_db_warning_notifications() - # Drain recovered process watchers (crash-recovery checkpoint). Detach the batch atomically: - # reassigning a fresh list takes ownership of exactly the watchers present now, so one - # appended concurrently during the yield below isn't dropped by a clear() on the shared - # list. Yield every 100 to avoid O(n^2) loop blocking when recovering thousands. + # Resume recovered process watchers. Detach the batch atomically (fresh list, not clear(): a + # concurrent append during the yield must not be lost); yield every 100 to keep the loop live. with _log_suppressed(logging.ERROR, "Recovered watcher setup error: %s"): from tools.process_registry import process_registry watchers = process_registry.pending_watchers @@ -1382,8 +1326,7 @@ class GatewayStartupMixin: if _aborted: return True if self._start_handle_no_connections( - connected_count, enabled_platform_count, startup_retryable_errors, - startup_nonretryable_errors, + connected_count, enabled_platform_count, startup_retryable_errors, startup_nonretryable_errors ): return True @@ -1447,9 +1390,8 @@ class GatewayStartupMixin: handoff_config, handoff_adapters = self._handoff_resolve_scope(profile_name) - # Adapter must be live. A relay-fronted gateway registers ONE adapter under Platform.RELAY - # fronting N logical platforms, so a literal adapters.get(discord) misses a deliverable - # platform; resolve_delivery_transport is the alias-aware resolver (native adapter wins). + # Alias-aware transport: a relay-fronted gateway registers ONE Platform.RELAY adapter fronting + # N logical platforms, so a literal adapters.get() would miss a deliverable one. transport = resolve_delivery_transport(platform, handoff_config, handoff_adapters) if not transport: raise RuntimeError(f"platform '{platform_name}' is not active in this gateway") @@ -1468,16 +1410,12 @@ class GatewayStartupMixin: home_chat_id, f"Hermes — {cli_title}", ) except Exception as exc: - logger.debug( - "Handoff: create_handoff_thread raised on %s: %s", platform_name, exc, exc_info=True - ) + logger.debug("Handoff: create_handoff_thread raised on %s: %s", platform_name, exc, exc_info=True) new_thread_id = None effective_thread_id = new_thread_id or (str(home.thread_id) if home.thread_id else None) - # Telegram private-chat DM topics are shaped differently from group/forum threads by the - # inbound adapter: a handoff-created topic in a positive chat_id must use the DM-topic source - # shape (real user id == chat_id, so topic-mode checks and binding persistence match later - # inbound turns), or the synthetic turn binds a `thread` key while replies arrive on `dm`. + # Telegram private-chat DM topics use the DM-topic source shape (user_id == chat_id) so the + # synthetic turn binds the same key later inbound turns arrive on (`dm`, not `thread`). is_telegram_private_chat = ( platform == Platform.TELEGRAM and looks_like_telegram_private_chat_id(home_chat_id) ) @@ -1501,14 +1439,11 @@ class GatewayStartupMixin: ) def _handoff_session_key(self, dest, profile_name: Optional[str]) -> str: - """Build the destination session_key with the adapters' own rules. - - Thread keys omit user_id (thread_sessions_per_user default) so the next message shares it. - The key is namespaced to the queuing profile: a multiplexed gateway would otherwise build - ``agent:main:...`` while the profile's adapter routes inbound on ``agent::...``. - The store resolver is only the root fallback (None when multiplexing is off; old key - unchanged). The isinstance check is load-bearing: a Mock store returns a truthy MagicMock. - """ + """Destination session_key by the adapters' own rules. Thread keys omit user_id so the next + message shares it. Namespaced to the queuing profile (else a multiplexed gateway builds + ``agent:main:...`` while the profile's adapter routes on ``agent::...``); the store + resolver is only the root fallback. The isinstance check is load-bearing: a Mock store returns + a truthy MagicMock.""" platform_cfg = dest.handoff_config.platforms.get(dest.platform) extra = platform_cfg.extra if platform_cfg else {} handoff_profile = profile_name if (profile_name and profile_name != "default") else None @@ -1524,20 +1459,13 @@ class GatewayStartupMixin: logger.debug("Handoff: could not resolve profile namespace", exc_info=True) return build_session_key( dest.source, group_sessions_per_user=extra.get("group_sessions_per_user", True), - thread_sessions_per_user=extra.get("thread_sessions_per_user", False), - profile=handoff_profile, + thread_sessions_per_user=extra.get("thread_sessions_per_user", False), profile=handoff_profile, ) - async def _process_handoff( - self, row: Dict[str, Any], profile_name: Optional[str] = None, - ) -> None: - """Execute one handoff row. Raises on failure (caller marks failed). - - ``profile_name`` (``None`` = root) is the profile whose store queued this handoff. Under - multiplex it is load-bearing: ``self.adapters``/``self.config`` are the primary's (secondaries - live in ``_profile_adapters``), and the session key must be namespaced ``agent::...`` - or it binds a key nobody reads. Passing the name beats re-deriving it from the contextvar. - """ + async def _process_handoff(self, row: Dict[str, Any], profile_name: Optional[str] = None) -> None: + """Execute one handoff row; raises on failure (caller marks failed). ``profile_name`` (None = + root) is the profile whose store queued it — load-bearing under multiplex: secondaries live in + ``_profile_adapters`` and the key must be namespaced ``agent::...`` or nobody reads it.""" cli_session_id = row["id"] dest = await self._handoff_resolve_destination(row, profile_name) session_key = self._handoff_session_key(dest, profile_name) @@ -1570,8 +1498,7 @@ class GatewayStartupMixin: logger.info( "Handoff: dispatching synthetic turn for CLI session %s → %s " "(home=%s, thread=%s, session_key=%s)", - cli_session_id, dest.platform_name, dest.home.chat_id, dest.effective_thread_id, - session_key, + cli_session_id, dest.platform_name, dest.home.chat_id, dest.effective_thread_id, session_key, ) # Dispatch through the runner directly: adapter.handle_message would spawn a background task # and lose error visibility; inline _handle_message keeps success/failure observable.