diff --git a/gateway/run_adapters.py b/gateway/run_adapters.py index 06a56ee5bd..700f559388 100644 --- a/gateway/run_adapters.py +++ b/gateway/run_adapters.py @@ -511,9 +511,7 @@ class GatewayAdapterLifecycleMixin: # Does _process_handoff accept the profile argument? Test stand-ins bind a one-arg callable. try: import inspect as _inspect - _process_takes_profile = len( - _inspect.signature(self._process_handoff).parameters - ) >= 2 + _process_takes_profile = len(_inspect.signature(self._process_handoff).parameters) >= 2 except Exception: _process_takes_profile = False @@ -959,22 +957,7 @@ class GatewayAdapterLifecycleMixin: active = get_active_profile_name() or "default" connected = 0 - # Resource claim -> owning profile (credential: same account; listener: same bind+port). - claimed: Dict[tuple, str] = {} - for _plat, _ad in self.adapters.items(): - fp = self._adapter_credential_fingerprint(_ad) - if fp is not None: - claimed[(_plat, fp)] = active - listener_claim = self._adapter_listener_claim(_plat, _ad) - if listener_claim is not None: - claimed[listener_claim] = active - # A queued retryable primary still owns its credential and listener; reserve both. - for retry_info in getattr(self, "_failed_platforms", {}).values(): - for claim_name in ("credential_claim", "listener_claim"): - retry_claim = retry_info.get(claim_name) - if isinstance(retry_claim, tuple): - claimed[retry_claim] = active - + claimed = self._primary_resource_claims(active) profile_homes = _multiplex_profile_homes(self.config) for profile_name, profile_home in profile_homes: if profile_name == active: @@ -997,28 +980,49 @@ class GatewayAdapterLifecycleMixin: profile_name, e, exc_info=True, ) - # Record the served set for `hermes status`. "Served" = eligible for shared routing, HTTP - # prefixes, cron and runtime scope — broader than "has a connected secondary adapter". + self._record_served_profiles(active, profile_homes) + return connected + + def _primary_resource_claims(self, active: str) -> Dict[tuple, str]: + """Resource claim -> owning profile for every live or queued primary adapter. + + Credential claims stop two profiles polling one account; listener claims stop sidecars with + distinct credentials binding one endpoint. A queued retryable primary still owns both. + """ + claimed: Dict[tuple, str] = {} + for _plat, _ad in self.adapters.items(): + for claim in ( + self._adapter_credential_claim(_plat, _ad), + self._adapter_listener_claim(_plat, _ad), + ): + if claim is not None: + claimed[claim] = active + for retry_info in getattr(self, "_failed_platforms", {}).values(): + for claim_name in ("credential_claim", "listener_claim"): + retry_claim = retry_info.get(claim_name) + if isinstance(retry_claim, tuple): + claimed[retry_claim] = active + return claimed + + def _record_served_profiles(self, active: str, profile_homes) -> None: + """Record the served set for `hermes status` and seed per-profile PairingStores. + + "Served" = eligible for shared routing, HTTP prefixes, cron and runtime scope — broader + than "has a connected secondary adapter". + """ try: from gateway.status import write_runtime_status from gateway.pairing import PairingStore - served = [active] + sorted( - name for name, _home in profile_homes if name != active - ) - # Per-profile PairingStores so authz routes pairing checks to the right whitelist. + served = [active] + sorted(name for name, _home in profile_homes if name != active) for name in served: if name and name not in self.pairing_stores: self.pairing_stores[name] = ( - self.pairing_store - if name == active - else PairingStore(profile=name) + self.pairing_store if name == active else PairingStore(profile=name) ) write_runtime_status(served_profiles=served) except Exception: logger.debug("could not record served_profiles", exc_info=True) - return connected - async def _load_secondary_profile_config(self, profile_name: str, profile_home: "Path"): """Hydrate + enter ``profile_home``'s scope once; return its gateway config. @@ -1356,11 +1360,7 @@ class GatewayAdapterLifecycleMixin: if platform not in profile_map: profile_map[platform] = adapter self._sync_voice_mode_state_to_adapter(adapter) - logger.info( - "✓ %s reconnected (profile: %s)", - platform.value, - profile_name, - ) + logger.info("✓ %s reconnected (profile: %s)", platform.value, profile_name) await self._redeliver_failed_obligations_for_platform( platform, profile=profile_name ) @@ -1796,8 +1796,7 @@ class GatewayAdapterLifecycleMixin: # resolve the receiving adapter even once the routed profile is stamped below. registry = ( (getattr(self, "_profile_adapters", None) or {}).get(profile_name) - if profile_name - else getattr(self, "adapters", None) + if profile_name else getattr(self, "adapters", None) ) or {} adapter = registry.get(platform) if adapter is not None: diff --git a/gateway/run_busy.py b/gateway/run_busy.py index 938aff4326..18554cac86 100644 --- a/gateway/run_busy.py +++ b/gateway/run_busy.py @@ -374,26 +374,27 @@ class GatewayBusySessionMixin: else (None if event.source.platform == Platform.TELEGRAM and event.source.thread_id else event.message_id) ) + async def _send_busy_reply(self, event: MessageEvent, adapter, content: str, *, plain_anchor: bool = False) -> None: + """Send a busy-path reply anchored to the event (thread metadata included).""" + reply_anchor = self._reply_anchor_for_event(event) + await adapter._send_with_retry( + chat_id=event.source.chat_id, + content=content, + reply_to=reply_anchor if plain_anchor else self._busy_reply_to(event, reply_anchor), + metadata=self._thread_metadata_for_source(event.source, reply_anchor), + ) + async def _send_busy_drain_notice(self, event: MessageEvent, session_key: str, effective_mode: str) -> None: """Busy path while the gateway is restarting/stopping: queue (if allowed) and tell the user.""" adapter = self._adapter_for_source(event.source) if not adapter: return - - reply_anchor = self._reply_anchor_for_event(event) - thread_meta = self._thread_metadata_for_source(event.source, reply_anchor) if self._queue_during_drain_enabled(effective_mode): self._queue_or_replace_pending_event(session_key, event) message = f"⏳ Gateway {self._status_action_gerund()} — queued for the next turn after it comes back." else: message = f"⏳ Gateway is {self._status_action_gerund()} and is not accepting another turn right now." - - await adapter._send_with_retry( - chat_id=event.source.chat_id, - content=message, - reply_to=self._busy_reply_to(event, reply_anchor), - metadata=thread_meta, - ) + await self._send_busy_reply(event, adapter, message) # Bare-word approval replies → (verb, args) for the synthesized slash command. _PLAINTEXT_APPROVAL_WORDS: Dict[str, tuple] = { @@ -434,13 +435,7 @@ class GatewayBusySessionMixin: if _adapter and _reply: _text, _eph_ttl = _adapter._unwrap_ephemeral(_reply) if _text: - _anchor = self._reply_anchor_for_event(event) - await _adapter._send_with_retry( - chat_id=event.source.chat_id, - content=_text, - reply_to=_anchor, - metadata=self._thread_metadata_for_source(event.source, _anchor), - ) + await self._send_busy_reply(event, _adapter, _text, plain_anchor=True) return True except Exception: logger.warning( @@ -671,15 +666,8 @@ class GatewayBusySessionMixin: return message async def _send_busy_ack_reply(self, event: MessageEvent, adapter, message: str) -> None: - reply_anchor = self._reply_anchor_for_event(event) - thread_meta = self._thread_metadata_for_source(event.source, reply_anchor) try: - await adapter._send_with_retry( - chat_id=event.source.chat_id, - content=message, - reply_to=self._busy_reply_to(event, reply_anchor), - metadata=thread_meta, - ) + await self._send_busy_reply(event, adapter, message) except Exception as e: logger.debug("Failed to send busy-ack: %s", e) @@ -786,31 +774,35 @@ class GatewayBusySessionMixin: await self._send_busy_ack_reply(event, adapter, message) return True + # Slash name → handler method is ``_handle__command`` (``-`` → ``_``) except these. + _COMMAND_HANDLER_ALIASES = {"bg": "_handle_background_command", "sethome": "_handle_set_home_command"} + # Ordinary slash handlers shared by idle and busy dispatch. + _PLAIN_COMMANDS = ( + "status", "context", "restart", "approve", "deny", "pause", "agents", "bg", "btw", + "kanban", "subgoal", "heartbeat", "busy", "yolo", "verbose", "footer", "help", + "commands", "profile", "update", "version", + ) + # Dispatched only on the idle path (busy dispatch has its own allowlist). + _IDLE_COMMANDS = ( + "topic", "whoami", "platform", "stop", "reasoning", "memory", "skills", "fast", + "approvals", "model", "codex-runtime", "personality", "suggestions", "save", "retry", + "sethome", "compress", "usage", "topup", "insights", "reload-mcp", "reload-skills", + "bundles", "debug", "title", "resume", "sessions", "branch", "rollback", "diff", "goal", + "loop", "refine", "review", "voice", + ) + + def _command_handler_table(self, names) -> Dict[str, Any]: + return { + name: getattr( + self, + self._COMMAND_HANDLER_ALIASES.get(name, f"_handle_{name.replace('-', '_')}_command"), + ) + for name in names + } + def _gateway_plain_command_handlers(self): """Return ordinary slash handlers shared by idle and busy dispatch.""" - return { - "status": self._handle_status_command, - "context": self._handle_context_command, - "restart": self._handle_restart_command, - "approve": self._handle_approve_command, - "deny": self._handle_deny_command, - "pause": self._handle_pause_command, - "agents": self._handle_agents_command, - "bg": self._handle_background_command, - "btw": self._handle_btw_command, - "kanban": self._handle_kanban_command, - "subgoal": self._handle_subgoal_command, - "heartbeat": self._handle_heartbeat_command, - "busy": self._handle_busy_command, - "yolo": self._handle_yolo_command, - "verbose": self._handle_verbose_command, - "footer": self._handle_footer_command, - "help": self._handle_help_command, - "commands": self._handle_commands_command, - "profile": self._handle_profile_command, - "update": self._handle_update_command, - "version": self._handle_version_command, - } + return self._command_handler_table(self._PLAIN_COMMANDS) async def _send_command_ack(self, source, text: str, label: str) -> None: """Best-effort acknowledgment for a slash command that falls through to agent processing.""" @@ -825,43 +817,7 @@ class GatewayBusySessionMixin: def _gateway_idle_command_handlers(self): """Slash handlers dispatched only on the idle path (busy dispatch has its own allowlist).""" - return { - "topic": self._handle_topic_command, - "whoami": self._handle_whoami_command, - "platform": self._handle_platform_command, - "stop": self._handle_stop_command, - "reasoning": self._handle_reasoning_command, - "memory": self._handle_memory_command, - "skills": self._handle_skills_command, - "fast": self._handle_fast_command, - "approvals": self._handle_approvals_command, - "model": self._handle_model_command, - "codex-runtime": self._handle_codex_runtime_command, - "personality": self._handle_personality_command, - "suggestions": self._handle_suggestions_command, - "save": self._handle_save_command, - "retry": self._handle_retry_command, - "sethome": self._handle_set_home_command, - "compress": self._handle_compress_command, - "usage": self._handle_usage_command, - "topup": self._handle_topup_command, - "insights": self._handle_insights_command, - "reload-mcp": self._handle_reload_mcp_command, - "reload-skills": self._handle_reload_skills_command, - "bundles": self._handle_bundles_command, - "debug": self._handle_debug_command, - "title": self._handle_title_command, - "resume": self._handle_resume_command, - "sessions": self._handle_sessions_command, - "branch": self._handle_branch_command, - "rollback": self._handle_rollback_command, - "diff": self._handle_diff_command, - "goal": self._handle_goal_command, - "loop": self._handle_loop_command, - "refine": self._handle_refine_command, - "review": self._handle_review_command, - "voice": self._handle_voice_command, - } + return self._command_handler_table(self._IDLE_COMMANDS) # busy_handler key (hermes_cli/commands.py CommandDef) → mid-run variant method name. _BUSY_SPECIAL_HANDLERS: Dict[str, str] = {