diff --git a/gateway/run_goals.py b/gateway/run_goals.py index 61a202ad32..8d532e529a 100644 --- a/gateway/run_goals.py +++ b/gateway/run_goals.py @@ -209,7 +209,6 @@ class GatewayGoalsMixin: return except Exception as exc: logger.debug("goal continuation: post-delivery callback registration failed: %s", exc) - await _deliver() async def _post_turn_manager(self, session_entry: Any, label: str, module: str, load): @@ -258,7 +257,6 @@ class GatewayGoalsMixin: # Deferred until the visible final response is delivered, else "✓ Goal achieved" precedes it. if msg and source is not None: await self._defer_goal_status_notice_after_delivery(source, msg) - prompt = decision.get("continuation_prompt") or "" if not decision.get("should_continue") or not prompt or source is None: return diff --git a/gateway/run_inbound.py b/gateway/run_inbound.py index 83c8afec96..9633c98a28 100644 --- a/gateway/run_inbound.py +++ b/gateway/run_inbound.py @@ -40,11 +40,8 @@ class GatewayInboundMixin: self, event: "MessageEvent", source: SessionSource ) -> Optional["MessageEvent"]: """Run the ``pre_gateway_dispatch`` plugin hook; None = drop, else the (maybe rewritten) event. - - Results: ``{"action": "skip", "reason": ...}`` → drop; ``{"action": "rewrite", "text": ...}`` - → replace ``event.text``; ``{"action": "allow"}`` / None → normal dispatch. Runs BEFORE - auth so plugins can handle unauthorized senders without the pairing flow. - """ + Results: ``{"action": "skip"}`` → drop; ``{"action": "rewrite", "text"}`` → replace ``event.text``; + ``allow``/None → normal dispatch. Runs BEFORE auth so plugins can handle unauthorized senders.""" try: from hermes_cli.lifecycle import invoke_hook as _invoke_hook _hook_results = _invoke_hook( @@ -89,30 +86,26 @@ class GatewayInboundMixin: code = pairing_store.generate_code(platform_name, source.user_id, source.user_name or "") adapter = self._adapter_for_source(source) if code: - if adapter: - store_profile = getattr(pairing_store, "profile", None) - profile_arg = ( - f"-p {store_profile} " - if isinstance(store_profile, str) and store_profile and store_profile != "default" - else "" - ) - await adapter.send( - source.chat_id, - f"Hi~ I don't recognize you yet!\n\n" - f"Here's your pairing code: `{code}`\n\n" - f"Ask the bot owner to run:\n" - f"`hermes {profile_arg}pairing approve " - f"{platform_name} {code}`" - ) - return - if adapter: - await adapter.send( - source.chat_id, - "Too many pairing requests right now~ " - "Please try again later!" + store_profile = getattr(pairing_store, "profile", None) + profile_arg = ( + f"-p {store_profile} " + if isinstance(store_profile, str) and store_profile and store_profile != "default" + else "" ) - # Record rate limit so subsequent messages are silently ignored - pairing_store._record_rate_limit(platform_name, source.user_id) + reply = ( + f"Hi~ I don't recognize you yet!\n\n" + f"Here's your pairing code: `{code}`\n\n" + f"Ask the bot owner to run:\n" + f"`hermes {profile_arg}pairing approve " + f"{platform_name} {code}`" + ) + else: + reply = "Too many pairing requests right now~ Please try again later!" + if adapter: + await adapter.send(source.chat_id, reply) + if not code: + # Record rate limit so subsequent messages are silently ignored + pairing_store._record_rate_limit(platform_name, source.user_id) async def _hm_admit_event( self, event: "MessageEvent" @@ -121,6 +114,8 @@ class GatewayInboundMixin: (the ``pre_gateway_dispatch`` hook may have rewritten ``event``).""" from gateway.run import _is_slack_ignored_channel source = event.source + # getattr(self, ...) throughout: bare test runners build GatewayRunner via object.__new__. + _config = getattr(self, "config", None) # 🔴 Cross-session leak guard: this per-message task was create_task()'d with a copy of the # spawning context, which may carry ANOTHER message's HERMES_SESSION_* ContextVars; until @@ -133,8 +128,10 @@ class GatewayInboundMixin: # Most adapters resolve profile routes in build_source(); internal/voice paths construct # SessionSource directly, so resolve those here as the shared fail-closed ingress gate. + # Strict boolean marker: require the literal True so duck-typed test/internal sources with + # dynamic attributes are not mistaken for a rejection. if ( - getattr(getattr(self, "config", None), "multiplex_profiles", False) + getattr(_config, "multiplex_profiles", False) and not getattr(source, "profile", None) and getattr(source, "profile_route_rejected", False) is not True ): @@ -144,9 +141,6 @@ class GatewayInboundMixin: source.profile = self._profile_name_for_source(source) except ProfileRouteRejected: source.profile_route_rejected = True - - # Strict boolean marker: require the literal True so duck-typed test/internal sources with - # dynamic attributes are not mistaken for a rejection. if getattr(source, "profile_route_rejected", False) is True: logger.warning( "Dropping inbound message because its explicit profile route " @@ -154,17 +148,15 @@ class GatewayInboundMixin: ) return None - # Internal events (e.g. background-process notifications) skip user authorization. - is_internal = bool(getattr(event, "internal", False)) + is_internal = bool(getattr(event, "internal", False)) # e.g. background-process notifications # Ignored-channel guard runs FIRST — before startup-restore queueing, plugin hooks, auth, # and session setup — so an ignored channel can never reach pairing/auth/session state. - # getattr: bare test runners construct GatewayRunner via object.__new__ without config. _chat_id = getattr(source, "chat_id", None) if ( not is_internal and getattr(source, "platform", None) == Platform.SLACK - and _is_slack_ignored_channel(getattr(self, "config", None), _chat_id) + and _is_slack_ignored_channel(_config, _chat_id) ): logger.info("Dropping Slack message from configured ignored channel %s", _chat_id) return None @@ -202,17 +194,13 @@ class GatewayInboundMixin: ): await self._hm_offer_pairing_code(source) return None - return event, source, False def _hm_estop_turn_allowed(self, event: "MessageEvent", source: SessionSource) -> bool: - """Whether a turn may bypass the global emergency stop. - - Pause blocks NEW agent turns, never running work or control traffic: recognized slash - commands (/status, /approve, ... and /pause off as the in-band resume path) and replies - owned by in-flight work — pending update prompt, steering a running session, pending - slash-confirm, dangerous-command approval — all pass through. - """ + """Whether a turn may bypass the global emergency stop: pause blocks NEW agent turns, never + running work or control traffic — recognized slash commands (incl. /pause off, the in-band + resume) and replies owned by in-flight work (pending update prompt, running session, + pending slash-confirm, dangerous-command approval) all pass through.""" with suppress(Exception): _estop_cmd = event.get_command() if _estop_cmd: @@ -224,8 +212,7 @@ class GatewayInboundMixin: _estop_state = self._peek_session_state(_estop_key) if _estop_state is not None and _estop_state.persistent.update_prompt_pending: return True - # Steering / interrupting in-flight work (also covers pending clarify + tool - # approvals held by the running agent). + # A running session covers steering plus pending clarify / tool approvals it holds. if self._is_session_running(_estop_key): return True from tools import slash_confirm as _estop_confirm_mod @@ -245,9 +232,9 @@ class GatewayInboundMixin: return None try: from agent.estop import paused_reply as _estop_paused_reply - _paused_notice = _estop_paused_reply() except ImportError: return None + _paused_notice = _estop_paused_reply() if _paused_notice is None or self._hm_estop_turn_allowed(event, source): return None logger.info( @@ -262,32 +249,22 @@ class GatewayInboundMixin: """Atomically hand *response_text* to the detached update process; returns the OSError str.""" from gateway.run import _hermes_home response_path = _hermes_home / ".update_response" - prompt_path = _hermes_home / ".update_prompt.json" try: tmp = response_path.with_suffix(".tmp") tmp.write_text(response_text, encoding="utf-8") tmp.replace(response_path) - prompt_path.unlink(missing_ok=True) + (_hermes_home / ".update_prompt.json").unlink(missing_ok=True) except OSError as e: return str(e) return None - def _hm_update_prompt_reply( - self, event: "MessageEvent", _quick_key: str, allow_gateway_control: bool - ) -> Optional[str]: - """Consume a reply to a pending ``/update`` prompt; None when nothing was consumed. - - Routes the reply to the detached update process via ``.update_response``. Recognized slash - commands must bypass this or /new, /help etc. get silently consumed as update answers. - """ + def _hm_update_prompt_reply(self, event: "MessageEvent", _quick_key: str) -> Optional[str]: + """Consume a reply to a pending ``/update`` prompt (routed to the detached update process via + ``.update_response``); None when nothing was consumed. Recognized slash commands must bypass + this or /new, /help etc. get silently consumed as update answers.""" _up_state = self._peek_session_state(_quick_key) - if not ( - allow_gateway_control - and _up_state is not None - and _up_state.persistent.update_prompt_pending - ): + if _up_state is None or not _up_state.persistent.update_prompt_pending: return None - raw = (event.text or "").strip() # Accept /approve and /deny as shorthand for yes/no cmd = event.get_command() _recognized_cmd = None @@ -301,7 +278,7 @@ class GatewayInboundMixin: from hermes_cli.commands import resolve_command as _resolve_update_cmd _cmd_def = _resolve_update_cmd(cmd) _recognized_cmd = _cmd_def.name if _cmd_def else None - response_text = "" if _recognized_cmd else raw + response_text = "" if _recognized_cmd else (event.text or "").strip() if response_text: err = self._hm_write_update_response(response_text) if err is not None: @@ -327,40 +304,40 @@ class GatewayInboundMixin: return None async def _hm_clarify_reply( - self, event: "MessageEvent", source: SessionSource, _quick_key: str, - allow_gateway_control: bool, + self, event: "MessageEvent", source: SessionSource, _quick_key: str ) -> Optional[str]: """Intercept a reply to a pending clarify prompt; None when the message falls through. - - Open-ended prompts and "Other" responses are free text; direct replies to multi-choice - prompts are accepted too ("2" → second option). Resolved/retained replies return "" so - adapters that emit the agent's response don't double-post — the agent itself produces the - next user-facing message. - """ - _clarify_mod = None + Free text answers open-ended/"Other" prompts; "2" answers a multi-choice one. Resolved/retained + replies return "" so adapters don't double-post — the agent produces the next user-facing message.""" try: from tools import clarify_gateway as _clarify_mod _pending_clarify = _clarify_mod.get_pending_for_session(_quick_key, include_choice_prompts=True) except Exception: - _pending_clarify = None - if not (allow_gateway_control and _pending_clarify is not None and _clarify_mod is not None): + return None + if _pending_clarify is None: return None _clarify_has_audio = bool(self._pending_event_audio_paths(event)) _raw_clarify_reply = await self._prepare_clarify_reply_text(event) - _log_ids = (_quick_key, _pending_clarify.clarify_id) - if _clarify_has_audio and not _raw_clarify_reply: + + def _retain(why: str) -> str: logger.info( - "Gateway retained pending clarify after voice transcription " - "produced no usable text (session=%s, id=%s)", *_log_ids, + "Gateway retained pending clarify after %s (session=%s, id=%s)", + why, _quick_key, _pending_clarify.clarify_id, ) return "" + + if _clarify_has_audio and not _raw_clarify_reply: + return _retain("voice transcription produced no usable text") # Slash commands: the user wanted a command, not to answer the clarify. Leave it pending so # they can retry; on timeout the agent unblocks with an empty response. if not _raw_clarify_reply or _raw_clarify_reply.startswith("/"): return None _text_outcome = _clarify_mod.attempt_text_response_for_session(_quick_key, _raw_clarify_reply) if _text_outcome == _clarify_mod.TEXT_RESOLVED: - logger.info("Gateway intercepted clarify text response (session=%s, id=%s)", *_log_ids) + logger.info( + "Gateway intercepted clarify text response (session=%s, id=%s)", + _quick_key, _pending_clarify.clarify_id, + ) # The clarify callback pauses the platform typing/status indicator while waiting so # Slack users can type; the active agent resumes now, so re-enable its indicator. _clarify_adapter = self._adapter_for_source(source) @@ -373,11 +350,7 @@ class GatewayInboundMixin: if _text_outcome == _clarify_mod.TEXT_REJECTED_SELECTION: # Selection-shaped but invalid (out-of-range number, bad comma-list): keep the clarify # armed for retry — don't cancel, don't treat as an unrelated follow-up. - logger.info( - "Gateway retained pending clarify after invalid " - "selection attempt (session=%s, id=%s)", *_log_ids, - ) - return "" + return _retain("invalid selection attempt") if _text_outcome == _clarify_mod.TEXT_REJECTED_PROSE: # Native-choice prompts reject unmatched prose so it continues through normal busy # routing. Release this clarify first: redirect() degrades to steer() while tools @@ -398,23 +371,18 @@ class GatewayInboundMixin: "cancel": "cancel", "nevermind": "cancel", "no": "cancel", } - async def _hm_slash_confirm_reply( - self, event: "MessageEvent", _quick_key: str, allow_gateway_control: bool - ) -> Optional[str]: - """Resolve a reply to a pending slash-confirm prompt; None when the message falls through. - - Accepts /approve, /always, /cancel and short aliases. Anything else falls through — a stale - pending confirm does NOT block other commands. A pending dangerous-command approval takes - precedence: /approve there unblocks the waiting tool thread. - """ + async def _hm_slash_confirm_reply(self, event: "MessageEvent", _quick_key: str) -> Optional[str]: + """Resolve a reply (/approve, /always, /cancel + aliases) to a pending slash-confirm prompt; + None when it falls through — a stale pending confirm does NOT block other commands. A pending + dangerous-command approval takes precedence: /approve there unblocks the waiting tool thread.""" from tools import slash_confirm as _slash_confirm_mod _pending_confirm = _slash_confirm_mod.get_pending(_quick_key) - _tool_approval_live = False + if not _pending_confirm: + return None with suppress(Exception): from tools.approval import has_blocking_approval - _tool_approval_live = has_blocking_approval(_quick_key) - if not (allow_gateway_control and _pending_confirm and not _tool_approval_live): - return None + if has_blocking_approval(_quick_key): + return None # Accept bang-prefixed replies (`!always`, `!cancel`) verbatim: Slack/Matrix show the `!` # prefix (typed `/` is blocked in Slack threads) and adapters only rewrite # `!` — confirm keywords aren't commands, so the `!` survives to here. @@ -434,13 +402,10 @@ class GatewayInboundMixin: return None def _hm_evict_idle_stale_agent(self, _quick_key: str) -> None: - """Evict a leaked lock from a hung/crashed handler. - - With inactivity-based timeout active tasks can run for hours, so evict only when the agent - has been *idle* past the threshold (or has no activity tracker and its wall-clock age is - extreme). The pending sentinel is never evicted: it has no get_activity_summary(), so the - idle check would read inf and race the async setup path. - """ + """Evict a leaked lock from a hung/crashed handler: only when the agent has been *idle* past + the threshold (active tasks can run for hours), or has no activity tracker and an extreme + wall-clock age. The pending sentinel is never evicted (no get_activity_summary() → idle + reads inf and would race the async setup path).""" from gateway.run import _AGENT_PENDING_SENTINEL, _float_env _raw_stale_timeout = _float_env("HERMES_AGENT_TIMEOUT", 1800) _quick_state = self._peek_session_state(_quick_key) @@ -481,11 +446,9 @@ class GatewayInboundMixin: logger.warning( "Evicting stale _running_agents entry for %s " "(age: %.0fs, idle: %.0fs, timeout: %.0fs)%s", - _quick_key, _stale_age, _stale_idle, - _raw_stale_timeout, _stale_detail, + _quick_key, _stale_age, _stale_idle, _raw_stale_timeout, _stale_detail, ) - self._invalidate_session_run_generation(_quick_key, reason="stale_running_agent_eviction") - self._release_running_agent_state(_quick_key) + self._hm_evict_running_agent(_quick_key, "stale_running_agent_eviction") def _hm_evict_reaped_agent(self, _quick_key: str) -> None: """Evict the in-memory turn slot of a session whose durable row was ended while the gateway @@ -499,28 +462,19 @@ class GatewayInboundMixin: _reap_peek = getattr(_reap_store, "peek_session_id", None) _is_ended = getattr(_reap_store, "_is_session_ended_in_db", None) _reap_sid = _reap_peek(_quick_key) if callable(_reap_peek) else None - if ( - isinstance(_reap_sid, str) - and _reap_sid - and callable(_is_ended) - and _is_ended(_reap_sid) is True - ): + if isinstance(_reap_sid, str) and _reap_sid and callable(_is_ended) and _is_ended(_reap_sid) is True: logger.warning( "Evicting stale _running_agents entry for %s — " "durable session %s is ended (reaped) in state.db; " - "healing routing on next message (#99106)", - _quick_key, _reap_sid, + "healing routing on next message (#99106)", _quick_key, _reap_sid, ) - self._invalidate_session_run_generation(_quick_key, reason="reaped_session_eviction") - self._release_running_agent_state(_quick_key) + self._hm_evict_running_agent(_quick_key, "reaped_session_eviction") except Exception: logger.debug("reaped-session staleness check failed", exc_info=True) - def _hm_evict_stale_running_agent(self, _quick_key: str) -> None: - """Evict a leaked/reaped ``_running_agents`` slot before the busy-session fast-path.""" - self._hm_evict_idle_stale_agent(_quick_key) - if self._is_session_running(_quick_key): - self._hm_evict_reaped_agent(_quick_key) + def _hm_evict_running_agent(self, _quick_key: str, reason: str) -> None: + self._invalidate_session_run_generation(_quick_key, reason=reason) + self._release_running_agent_state(_quick_key) def _hm_merge_pending_for_source( self, source: SessionSource, _quick_key: str, event: "MessageEvent", *, merge_text: bool = False @@ -567,27 +521,24 @@ class GatewayInboundMixin: self, event: "MessageEvent", source: SessionSource, _quick_key: str, effective_busy_input_mode: str ) -> bool: """Queue a Telegram text follow-up that lands within the post-start grace window.""" - _telegram_followup_grace = float(os.getenv("HERMES_TELEGRAM_FOLLOWUP_GRACE_SECONDS", "3.0")) + _grace = float(os.getenv("HERMES_TELEGRAM_FOLLOWUP_GRACE_SECONDS", "3.0")) _grace_state = self._peek_session_state(_quick_key) _started_at = _grace_state.turn.started_ts if _grace_state else 0 if not ( - source.platform == Platform.TELEGRAM - and event.message_type == MessageType.TEXT - and _telegram_followup_grace > 0 - and _started_at - and (time.time() - _started_at) <= _telegram_followup_grace + source.platform == Platform.TELEGRAM and event.message_type == MessageType.TEXT + and _grace > 0 and _started_at and (time.time() - _started_at) <= _grace ): return False logger.debug( "Telegram follow-up arrived %.2fs after run start for %s — queueing without interrupt", time.time() - _started_at, _quick_key, ) - adapter = self._adapter_for_source(source) - if adapter: - if effective_busy_input_mode == "queue": + if effective_busy_input_mode != "queue": + self._hm_merge_pending_for_source(source, _quick_key, event, merge_text=True) + else: + adapter = self._adapter_for_source(source) + if adapter: self._enqueue_fifo(_quick_key, event, adapter) - else: - self._hm_merge_pending_for_source(source, _quick_key, event, merge_text=True) return True @staticmethod @@ -616,11 +567,8 @@ class GatewayInboundMixin: from gateway.run import _build_media_placeholder # Text-only corrections redirect the live turn (preserving displayed context) when the # runtime supports it; media/voice and older runtimes use the interrupt path below. - if ( - self._hm_text_only(event) - and getattr(running_agent, "_supports_active_turn_redirect", False) is True - and hasattr(running_agent, "redirect") - ): + _can_redirect = getattr(running_agent, "_supports_active_turn_redirect", False) is True + if self._hm_text_only(event) and _can_redirect and hasattr(running_agent, "redirect"): try: if running_agent.redirect((event.text or "").strip()): logger.debug("PRIORITY redirect for session %s", _quick_key) @@ -634,12 +582,11 @@ class GatewayInboundMixin: event, self._adapter_for_source(source), source, event.text or "", log_context="Voice-priority-interrupt", ) - elif not _interrupt_text and (getattr(event, "media_urls", None) or []): + elif not _interrupt_text and getattr(event, "media_urls", None): _interrupt_text = _build_media_placeholder(event) - # Delivered via adapter._pending_messages (read by _run_agent); never also buffered on self. + # Delivered via adapter._pending_messages (read by _run_agent); never also buffered on self + # — that copy was never consumed and grew unbounded. running_agent.interrupt(_interrupt_text) - # The interrupt message is delivered via adapter._pending_messages (read by _run_agent); - # don't also buffer it on self — that copy was never consumed and grew unbounded. async def _hm_handle_running_session_message( self, event: "MessageEvent", source: SessionSource, _quick_key: str @@ -657,15 +604,12 @@ class GatewayInboundMixin: _ra_state = self._peek_session_state(_quick_key) running_agent = _ra_state.turn.agent if _ra_state else None - if running_agent is _AGENT_PENDING_SENTINEL: - # Agent is being set up but not ready yet. - if event.get_command() == "stop": - # Force-clean the sentinel so the session is unlocked. + if running_agent is _AGENT_PENDING_SENTINEL: # agent still being set up + if event.get_command() == "stop": # force-clean the sentinel so the session is unlocked self._release_running_agent_state(_quick_key) logger.info("HARD STOP (pending) for session %s — sentinel cleared", _quick_key) return EphemeralReply("⚡ Force-stopped. The agent was still starting — session unlocked.") - # Queue the message so it will be picked up after the agent starts. - self._hm_merge_pending_for_source(source, _quick_key, event, merge_text=True) + self._hm_merge_pending_for_source(source, _quick_key, event, merge_text=True) # picked up after start return None if self._draining: queue_during_drain = self._queue_during_drain_enabled(effective_busy_input_mode) @@ -711,20 +655,17 @@ class GatewayInboundMixin: if not target: return None target = target if target.startswith("/") else f"/{target}" + event.text = f"{target} {event.get_command_args().strip()}".strip() target_command = target.lstrip("/") - user_args = event.get_command_args().strip() - event.text = f"{target} {user_args}".strip() return target_command.split()[0] if target_command else target_command async def _hm_command_hooks( self, event: "MessageEvent", source: SessionSource, _quick_key: str, command: str, canonical: str ) -> Tuple[bool, Optional[str], Optional[str]]: - """Fire ``pre_command`` (observer) and ``command:`` (interceptor) hooks. - - Returns ``(handled, result, new_command)``; ``new_command`` is set when a handler rewrote - the command. The running-agent intercept path deliberately does NOT fire these — a slow or - hostile plugin must not interfere with the operator's escape hatches for a live agent. - """ + """Fire ``pre_command`` (observer) and ``command:`` (interceptor) hooks → + ``(handled, result, new_command)`` (``new_command`` set when a handler rewrote the command). + The running-agent path deliberately does NOT fire these — a slow or hostile plugin must not + interfere with the operator's escape hatches for a live agent.""" raw_args = event.get_command_args().strip() platform = source.platform.value if source.platform else "" try: @@ -760,11 +701,9 @@ class GatewayInboundMixin: return True, message, None if decision == "rewrite": new_command = str(hook_result.get("command_name", "")).strip().lstrip("/") - if not new_command: - continue - new_args = str(hook_result.get("raw_args", "")).strip() - event.text = f"/{new_command} {new_args}".strip() - return False, None, event.get_command() + if new_command: + event.text = f"/{new_command} {str(hook_result.get('raw_args', '')).strip()}".strip() + return False, None, event.get_command() return False, None, None async def _hm_resolve_command( @@ -834,7 +773,6 @@ class GatewayInboundMixin: async def _hm_cmd_egress(self, event, source, _quick_key): from hermes_cli.proxy_cli import format_status_text - return True, format_status_text() async def _hm_rewrite_turn_to_prompt(self, event, source, name: str, ack: str, build) -> Tuple[bool, Optional[str]]: @@ -847,31 +785,20 @@ class GatewayInboundMixin: return True, f"Could not start /{name} — please try again." return False, None + # /learn and /plan: ack, rewrite the turn to a builder prompt, fall through to the agent. async def _hm_cmd_learn(self, event, source, _quick_key): from agent.learn_prompt import build_learn_prompt - _learn_req = event.get_command_args().strip() - _ack = ( - "Learning a skill from what you described…" - if _learn_req - else "Learning a skill from this conversation…" - ) - return await self._hm_rewrite_turn_to_prompt( - event, source, "learn", _ack, lambda: build_learn_prompt(_learn_req) - ) + req = event.get_command_args().strip() + _ack = f"Learning a skill from {'what you described' if req else 'this conversation'}…" + return await self._hm_rewrite_turn_to_prompt(event, source, "learn", _ack, lambda: build_learn_prompt(req)) async def _hm_cmd_plan(self, event, source, _quick_key): from agent.plan_prompt import build_plan_prompt - _plan_task = event.get_command_args().strip() - _ack = ( - f"Planning: {_plan_task[:80]}{'…' if len(_plan_task) > 80 else ''}" - if _plan_task - else "Planning from this conversation's context…" - ) - return await self._hm_rewrite_turn_to_prompt( - event, source, "plan", _ack, lambda: build_plan_prompt(_plan_task) - ) + task = event.get_command_args().strip() + _ack = f"Planning: {task[:80]}{'…' if len(task) > 80 else ''}" if task else "Planning from this conversation's context…" + return await self._hm_rewrite_turn_to_prompt(event, source, "plan", _ack, lambda: build_plan_prompt(task)) async def _hm_cmd_init(self, event, source, _quick_key): # /init builds the prompt first: the ack wording depends on whether AGENTS.md exists. @@ -958,11 +885,8 @@ class GatewayInboundMixin: _moa_state = self._session_state(_quick_key) event._moa_restore_override = _moa_state.conversation.model_override _moa_state.conversation.model_override = { - "provider": "moa", - "model": moa_cfg["default_preset"], - "base_url": "moa://local", - "api_key": "moa-virtual-provider", - "api_mode": "chat_completions", + "provider": "moa", "model": moa_cfg["default_preset"], "base_url": "moa://local", + "api_key": "moa-virtual-provider", "api_mode": "chat_completions", } self._evict_cached_agent(_quick_key) event._moa_disable_after_turn = True @@ -993,9 +917,9 @@ class GatewayInboundMixin: return False, None async def _hm_run_exec_quick_command(self, command: str, exec_cmd: str) -> str: - """Run a ``type: exec`` quick command in the gateway process (30 s cap, sanitized env).""" + """Run a ``type: exec`` quick command in the gateway process (30 s cap, sanitized env — the + gateway process has every API key in os.environ; output is redacted too).""" try: - # Sanitized env: the gateway process has every API key in os.environ. from tools.environments.local import build_subprocess_env proc = await asyncio.create_subprocess_shell( exec_cmd, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, @@ -1006,7 +930,7 @@ class GatewayInboundMixin: if output: from agent.redact import redact_sensitive_text output = redact_sensitive_text(output) - return output if output else "Command returned no output." + return output or "Command returned no output." except asyncio.TimeoutError: return "Quick command timed out (30s)." except Exception as e: @@ -1015,10 +939,8 @@ class GatewayInboundMixin: async def _hm_dispatch_quick_and_plugin_commands( self, event: "MessageEvent", source: SessionSource, command: Optional[str] ) -> Tuple[bool, Optional[str], Optional[str]]: - """Drain gate, user-defined quick commands (exec/alias) and plugin slash commands. - - Returns ``(handled, result, command)`` — an alias quick command rewrites ``command``. - """ + """Drain gate, user-defined quick commands (exec/alias) and plugin slash commands → + ``(handled, result, command)``; an alias quick command rewrites ``command``.""" if self._draining: return True, f"⏳ Gateway is {self._status_action_gerund()} and is not accepting new work right now.", command @@ -1044,12 +966,11 @@ class GatewayInboundMixin: return True, f"Quick command '/{command}' has no target defined.", command command = new_command # Fall through to normal command dispatch below - # Plugin-registered slash commands + # Plugin-registered slash commands. Underscores normalize to hyphens so Telegram's + # underscored autocomplete form matches plugin commands registered with hyphens. if command: try: from hermes_cli.plugins import get_plugin_command_handler - # Normalize underscores to hyphens so Telegram's underscored autocomplete form - # matches plugin commands registered with hyphens (see _build_telegram_menu). plugin_handler = get_plugin_command_handler(command.replace("_", "-")) if plugin_handler: result = plugin_handler(event.get_command_args().strip()) @@ -1058,7 +979,6 @@ class GatewayInboundMixin: return True, str(result) if result else None, command except Exception as e: logger.warning("Plugin command dispatch failed: %s", e) - return False, None, command def _hm_bundle_slash_rewrite( @@ -1118,12 +1038,9 @@ class GatewayInboundMixin: def _hm_skill_slash_rewrite( self, event: "MessageEvent", source: SessionSource, _quick_key: str, command: Optional[str] ) -> Optional[str]: - """Rewrite ``/`` / ``/`` invocations into the skill prompt on ``event.text``. - - Returns a reply string when the command is disabled/unknown/failed, else None. - resolve_skill_command_key() handles the Telegram underscore/hyphen round-trip so - /claude_code from Telegram autocomplete still resolves to the claude-code skill. - """ + """Rewrite ``/`` / ``/`` invocations into the skill prompt on ``event.text``; + returns a reply string when the command is disabled/unknown/failed, else None. + resolve_skill_command_key() handles the Telegram underscore/hyphen round-trip (/claude_code).""" if not command or self._hm_bundle_slash_rewrite(event, source, _quick_key, command): return None try: @@ -1189,13 +1106,15 @@ class GatewayInboundMixin: async def _hm_pending_reply_intercepts( self, event: "MessageEvent", source: SessionSource, _quick_key: str ) -> Optional[str]: - """Replies owned by in-flight work: pending /update prompt, clarify, slash-confirm.""" - allow = event.allow_gateway_control - _reply = self._hm_update_prompt_reply(event, _quick_key, allow) + """Replies owned by in-flight work: pending /update prompt, clarify, slash-confirm. + Only events that may control the gateway (``allow_gateway_control``) can answer them.""" + if not event.allow_gateway_control: + return None + _reply = self._hm_update_prompt_reply(event, _quick_key) if _reply is None: - _reply = await self._hm_clarify_reply(event, source, _quick_key, allow) + _reply = await self._hm_clarify_reply(event, source, _quick_key) if _reply is None: - _reply = await self._hm_slash_confirm_reply(event, _quick_key, allow) + _reply = await self._hm_slash_confirm_reply(event, _quick_key) return _reply async def _hm_dispatch_idle_commands( @@ -1203,29 +1122,22 @@ class GatewayInboundMixin: ) -> Tuple[bool, Optional[str]]: """Idle path: resolve + dispatch slash commands; rewriting commands fall through to the agent.""" _handled, _result, command, canonical = await self._hm_resolve_command(event, source, _quick_key) - if _handled: - return True, _result - _handled, _result = await self._hm_dispatch_canonical_command(event, source, _quick_key, canonical) - if _handled: - return True, _result - _handled, _result, command = await self._hm_dispatch_quick_and_plugin_commands(event, source, command) - if _handled: - return True, _result - _skill_reply = self._hm_skill_slash_rewrite(event, source, _quick_key, command) - if _skill_reply is not None: - return True, _skill_reply - return False, None + if not _handled: + _handled, _result = await self._hm_dispatch_canonical_command(event, source, _quick_key, canonical) + if not _handled: + _handled, _result, command = await self._hm_dispatch_quick_and_plugin_commands(event, source, command) + if not _handled: + _result = self._hm_skill_slash_rewrite(event, source, _quick_key, command) + _handled = _result is not None + return _handled, _result def _hm_rescue_orphaned_fifo( self, event: "MessageEvent", source: SessionSource, is_internal: bool, _quick_key: str ) -> Tuple["MessageEvent", SessionSource, bool]: - """FIFO orphan rescue: re-stage overflow events left behind by an idle session. - - A session that went idle with a populated overflow (post-turn drain never promoted, e.g. a - compression-demoted follow-up) silently orphaned those events. The oldest orphan runs as - THIS turn and the incoming event is parked behind the chain. Skipped for control commands - and internal events. - """ + """FIFO orphan rescue: a session that went idle with a populated overflow (post-turn drain + never promoted, e.g. a compression-demoted follow-up) silently orphaned those events. The + oldest orphan runs as THIS turn and the incoming event is parked behind the chain. Skipped + for control commands and internal events.""" try: _orphan_adapter = self._adapter_for_source(source) if _orphan_adapter is None or getattr(event, "internal", False) or event.get_command(): @@ -1239,21 +1151,15 @@ class GatewayInboundMixin: # Same session key by construction; carry the orphan's own source so reply anchors / # thread metadata point at the message actually being answered. _rescued_source = getattr(_rescued, "source", None) - return ( - _rescued, - _rescued_source if _rescued_source is not None else source, - bool(getattr(_rescued, "internal", False)), - ) + source = _rescued_source if _rescued_source is not None else source + return _rescued, source, bool(getattr(_rescued, "internal", False)) except Exception: logger.debug("FIFO orphan rescue pre-claim failed for %s", _quick_key, exc_info=True) return event, source, is_internal async def _handle_message(self, event: MessageEvent) -> Optional[str]: - """Handle an incoming message from any platform. - - Pipeline: auth → command check → running-agent interrupt → get/create session → build - context → run agent → return response. - """ + """Handle an incoming message from any platform: auth → command check → running-agent + interrupt → get/create session → build context → run agent → return response.""" from gateway.run import _AGENT_PENDING_SENTINEL _admitted = await self._hm_admit_event(event) if _admitted is None: @@ -1269,7 +1175,10 @@ class GatewayInboundMixin: if _reply is not None: return _reply - self._hm_evict_stale_running_agent(_quick_key) + # Evict a leaked/reaped ``_running_agents`` slot before the busy-session fast-path. + self._hm_evict_idle_stale_agent(_quick_key) + if self._is_session_running(_quick_key): + self._hm_evict_reaped_agent(_quick_key) if self._is_session_running(_quick_key): return await self._hm_handle_running_session_message(event, source, _quick_key) @@ -1279,24 +1188,22 @@ class GatewayInboundMixin: # Pending exec approvals go through /approve and /deny only — no bare-text matching, or a # conversational "yes" would execute a dangerous command. - - if not is_internal and await asyncio.to_thread(self._is_telegram_topic_root_lobby, source): - # Debounce the lobby reminder so a user who forgets about topic mode and fires ten - # prompts doesn't get ten copies. - if self._should_send_telegram_lobby_reminder(source): - return self._telegram_topic_root_lobby_message() - return None - - # External-drain new-turn gate: when NAS engaged an external drain (.drain_request.json, - # seen by _drain_control_watcher), refuse to START new turns so the in-flight set can only - # fall to zero. Internal/system events bypass; reversible. - if self._external_drain_active and not is_internal: - logger.info("Refusing new turn for session %s — external drain active.", _quick_key) - return ( - "⏳ This agent is draining for a maintenance action and isn't " - "accepting new turns right now. It'll be back in a moment — " - "please resend shortly." - ) + if not is_internal: + if await asyncio.to_thread(self._is_telegram_topic_root_lobby, source): + # Debounced so a user who forgets about topic mode doesn't get ten reminders. + if self._should_send_telegram_lobby_reminder(source): + return self._telegram_topic_root_lobby_message() + return None + # External-drain new-turn gate: when NAS engaged an external drain (.drain_request.json, + # seen by _drain_control_watcher), refuse to START new turns so the in-flight set can + # only fall to zero. Reversible. + if self._external_drain_active: + logger.info("Refusing new turn for session %s — external drain active.", _quick_key) + return ( + "⏳ This agent is draining for a maintenance action and isn't " + "accepting new turns right now. It'll be back in a moment — " + "please resend shortly." + ) # Claim this session before any await: many awaits sit between here and _run_agent # registering the real AIAgent; without this sentinel a second message during any of them @@ -1318,9 +1225,7 @@ class GatewayInboundMixin: try: try: - _agent_result = await self._handle_message_with_agent( - event, source, _quick_key, _run_generation - ) + _agent_result = await self._handle_message_with_agent(event, source, _quick_key, _run_generation) except TurnLeaseTimeoutError as exc: # A rejected message, not a completed turn: return before the /goal judge so it # cannot consume the resend notice and enqueue a synthetic continuation loop. @@ -1389,17 +1294,14 @@ class GatewayInboundMixin: thread_sessions_per_user=getattr(self.config, "thread_sessions_per_user", False), ) if _is_shared_multi_user and source.user_name: - # The display name is attacker-influenceable on platforms where participants set their - # own name: neutralize newlines/control chars or a hostile name can masquerade as a fake - # markdown section (mirrors build_session_context_prompt's treatment). + # Display names are attacker-influenceable: neutralize newlines/control chars or a + # hostile name masquerades as a fake markdown section (mirrors build_session_context_prompt). _safe_user_name = neutralize_untrusted_inline_text(source.user_name) - # Slack: expose the current author's verifiable user ID so "mention me again" requests - # have a trusted `<@U...>` target for the CURRENT speaker (display names are ambiguous). - # The user_id comes from the event envelope (not user-editable), so no neutralization. + # Slack: expose the CURRENT speaker's verifiable `<@U...>` id so "mention me again" has a + # trusted target (display names are ambiguous). user_id comes from the envelope, not user-editable. if source.platform == Platform.SLACK and source.user_id: _safe_user_name = f"{_safe_user_name} | Slack user <@{source.user_id}>" message_text = f"[{_safe_user_name}] {message_text}" - # After the sender-prefix so the prefix applies only to the trigger message, not the backfill. if getattr(event, "channel_context", None): message_text = f"{event.channel_context}\n\n[New message]\n{message_text}" @@ -1409,12 +1311,9 @@ class GatewayInboundMixin: def _classify_inbound_media( event: MessageEvent, pending_stt_prepared: bool ) -> Tuple[list, list, list, list]: - """Split ``event.media_urls`` into (image, STT-voice, audio-file, video) paths. - - Per-attachment MIME wins over the message-level type (a document sent alongside an image - must not be routed as an image). MessageType.AUDIO / mixed DOCUMENT audio is a file - attachment, never STT. - """ + """Split ``event.media_urls`` into (image, STT-voice, audio-file, video) paths. Per-attachment + MIME wins over the message-level type (a document sent alongside an image must not be routed + as an image). MessageType.AUDIO / mixed DOCUMENT audio is a file attachment, never STT.""" from gateway.run import _event_media_is_audio, _event_media_is_image, _event_media_is_stt_input image_paths, audio_paths, audio_file_paths, video_paths = [], [], [], [] for i, path in enumerate(event.media_urls or []): @@ -1499,8 +1398,7 @@ class GatewayInboundMixin: from tools.credential_files import to_agent_visible_cache_path basename = os.path.basename(path) parts = basename.split("_", 2) - display_name = parts[2] if len(parts) >= 3 else basename - return re.sub(r'[^\w.\- ]', '_', display_name), to_agent_visible_cache_path(path) + return re.sub(r'[^\w.\- ]', '_', parts[2] if len(parts) >= 3 else basename), to_agent_visible_cache_path(path) @classmethod def _prepend_inbound_media_file_notes(cls, message_text: str, audio_file_paths: list[str], video_paths: list[str]) -> str: @@ -1511,7 +1409,7 @@ class GatewayInboundMixin: ): for _path in paths: _display, _agent_path = cls._inbound_attachment_display_name(_path) - _note = ( + message_text = ( f"[The user sent {kind}: '{_display}'. " f"It is saved at: {_agent_path}. " f"Its content is not inlined here. If the user's request involves " @@ -1519,8 +1417,8 @@ class GatewayInboundMixin: f"example by passing the path to {tool} — " f"instead of asking the user to describe it. Only ask what to do " f"with it if their intent is genuinely unclear.]" + f"\n\n{message_text}" ) - message_text = f"{_note}\n\n{message_text}" return message_text @classmethod @@ -1543,10 +1441,8 @@ class GatewayInboundMixin: continue mtype = event.media_types[i] if i < len(event.media_types) else "" if mtype in {"", "application/octet-stream"}: - if os.path.splitext(path)[1].lower() in _TEXT_EXTENSIONS: - mtype = "text/plain" - else: - mtype = _mimetypes.guess_type(path)[0] or "application/octet-stream" + _is_text = os.path.splitext(path)[1].lower() in _TEXT_EXTENSIONS + mtype = "text/plain" if _is_text else (_mimetypes.guess_type(path)[0] or "application/octet-stream") # Every accepted file gets a note — a non-text/non-application MIME (font/*, model/*) # must still tell the agent the file exists. display_name, agent_path = cls._inbound_attachment_display_name(path) @@ -1611,13 +1507,13 @@ class GatewayInboundMixin: source=source, session_key=session_key, user_config=_msg_cfg, ) _msg_base_url = _msg_runtime.get("base_url") or "" - _is_dict_cfg = isinstance(_msg_model_cfg, dict) - _msg_configured_model = ( - _msg_model_cfg.get("default") or _msg_model_cfg.get("model") if _is_dict_cfg else _msg_model_cfg - ) + if isinstance(_msg_model_cfg, dict): + _msg_configured_model = _msg_model_cfg.get("default") or _msg_model_cfg.get("model") + else: + _msg_configured_model = _msg_model_cfg # (no dict → no pin was read; ctx is already None) if _msg_model != _msg_configured_model: _msg_config_ctx = None - if _msg_config_ctx is not None and _is_dict_cfg: + if _msg_config_ctx is not None: try: from hermes_cli.route_identity import should_clear_context_pin_async @@ -1677,17 +1573,12 @@ class GatewayInboundMixin: self, *, event: MessageEvent, source: SessionSource, history: List[Dict[str, Any]], session_key: Optional[str] = None, ) -> Optional[str]: - """Prepare inbound event text for the agent. - - Shared by the normal inbound and queued follow-up paths so attribution, image enrichment, - STT, document notes, reply context and @ references behave the same. Side effect: buffers - per-session native image paths when the model supports native vision; the caller consumes - that buffer at ``run_conversation``. Empty list means the text vision path already ran. - """ + """Prepare inbound event text for the agent. Shared by the normal inbound and queued + follow-up paths so attribution, image enrichment, STT, document notes, reply context and + @ references behave the same. Side effect: buffers per-session native image paths when the + model supports native vision; the caller consumes that buffer at ``run_conversation``.""" _pending_stt_prepared = hasattr(event, "_gateway_pending_stt_text") - message_text = ( - getattr(event, "_gateway_pending_stt_text", None) if _pending_stt_prepared else event.text - ) or "" + message_text = (event._gateway_pending_stt_text if _pending_stt_prepared else event.text) or "" # Prefer the caller's resolved session key so this write key matches the consume key at the # run_conversation site; derive it here only for tests and legacy standalone callers. session_key = session_key or self._session_key_for_source(source) @@ -1695,9 +1586,7 @@ class GatewayInboundMixin: self._consume_pending_native_image_paths(session_key) message_text = self._prefix_inbound_sender_context(event, source, message_text) - image_paths, audio_paths, audio_file_paths, video_paths = self._classify_inbound_media( - event, _pending_stt_prepared - ) + image_paths, audio_paths, audio_file_paths, video_paths = self._classify_inbound_media(event, _pending_stt_prepared) if image_paths: message_text = await self._enrich_inbound_images(source, session_key, message_text, image_paths) if audio_paths: @@ -1725,18 +1614,14 @@ class GatewayInboundMixin: """Return raw text or successful voice transcripts for a clarify reply.""" if not self._pending_event_audio_paths(event): return (event.text or "").strip() - _, successful_transcripts = await self._transcribe_pending_audio_event_once(event, "") - return "\n\n".join( - transcript.strip() for transcript in successful_transcripts if transcript.strip() - ) + return "\n\n".join(t.strip() for t in successful_transcripts if t.strip()) def _consume_pending_native_image_paths(self, session_key: str) -> List[str]: state = self._peek_session_state(session_key) - if state is None or not state.persistent.native_image_paths: - return [] - paths = list(state.persistent.native_image_paths) - state.persistent.native_image_paths = [] + paths = list(state.persistent.native_image_paths or []) if state is not None else [] + if paths: + state.persistent.native_image_paths = [] return paths async def _mark_durable_active_turn(self, event: "MessageEvent", session_key: str) -> bool: @@ -1748,10 +1633,10 @@ class GatewayInboundMixin: return False if not token: return False - # Private event attributes are process-local ownership state. Keep the - # token out of public metadata, transcripts, and platform payloads. - setattr(event, "_gateway_active_turn_session_key", session_key) - setattr(event, "_gateway_active_turn_token", token) + # Private event attributes are process-local ownership state: keep the token out of public + # metadata, transcripts, and platform payloads. + event._gateway_active_turn_session_key = session_key + event._gateway_active_turn_token = token return True async def _clear_durable_active_turn(self, event: "MessageEvent") -> bool: @@ -1832,19 +1717,16 @@ class GatewayInboundMixin: def _log_result(completed) -> None: try: - accepted = completed.result() + if completed.result(): + return + what, exc = "was not routed", None except (asyncio.CancelledError, concurrent.futures.CancelledError): return - except Exception: - logger.warning( - "Plugin message injection failed: plugin=%s session=%s", - plugin_id, session_key, exc_info=True, - ) - return - if not accepted: - logger.warning( - "Plugin message injection was not routed: plugin=%s session=%s", plugin_id, session_key, - ) + except Exception as err: + what, exc = "failed", err + logger.warning( + "Plugin message injection %s: plugin=%s session=%s", what, plugin_id, session_key, exc_info=exc, + ) future.add_done_callback(_log_result) return True @@ -1864,19 +1746,19 @@ class GatewayInboundMixin: source = dataclasses.replace(entry.origin) try: - if not self._is_user_authorized(source, allow_adapter_delegation=False): - logger.warning( - "Plugin message injection denied by current gateway authorization: " - "plugin=%s session=%s", - plugin_id, session_key, - ) - return False + authorized = self._is_user_authorized(source, allow_adapter_delegation=False) except Exception: logger.warning( "Plugin message injection authorization check failed: plugin=%s session=%s", plugin_id, session_key, exc_info=True, ) return False + if not authorized: + logger.warning( + "Plugin message injection denied by current gateway authorization: " + "plugin=%s session=%s", plugin_id, session_key, + ) + return False adapter = self._adapter_for_source(source) if adapter is None: @@ -1886,10 +1768,8 @@ class GatewayInboundMixin: text=content, message_type=MessageType.TEXT, source=source, internal=True, allow_gateway_control=False, metadata={ - "hermes_plugin_id": plugin_id, - "hermes_plugin_injection": True, - "gateway_session_key": session_key, - "gateway_session_id": entry.session_id, + "hermes_plugin_id": plugin_id, "hermes_plugin_injection": True, + "gateway_session_key": session_key, "gateway_session_id": entry.session_id, "gateway_session_strict": True, }, )) @@ -1904,12 +1784,10 @@ class GatewayInboundMixin: user_config: Optional[dict] = None, provider: Optional[str] = None, model: Optional[str] = None, ) -> str: - """Resolve image-input routing (``"native"`` / ``"text"``) for the effective model this turn. - - See agent/image_routing.py. Gateway sessions can carry /model overrides and image - preprocessing runs before AIAgent sets the auxiliary_client runtime globals, so resolve the - per-session runtime bundle the upcoming turn will use, not just the persisted default. - """ + """Resolve image-input routing (``"native"`` / ``"text"``) for the effective model this turn + (see agent/image_routing.py). Sessions can carry /model overrides and this runs before AIAgent + sets the auxiliary_client runtime globals, so resolve the per-session runtime bundle the + upcoming turn will use, not just the persisted default.""" try: from agent.image_routing import decide_image_input_mode from agent.auxiliary_client import _read_main_model, _read_main_provider @@ -1925,15 +1803,13 @@ class GatewayInboundMixin: turn_model, runtime_kwargs = self._resolve_session_agent_runtime( source=source, session_key=session_key, user_config=cfg, ) + rk = runtime_kwargs if isinstance(runtime_kwargs, dict) else {} if not resolved_model and isinstance(turn_model, str): resolved_model = turn_model.strip() - rk = runtime_kwargs if isinstance(runtime_kwargs, dict) else {} - runtime_provider = rk.get("provider") - runtime_requested_provider = rk.get("requested_provider") - if not resolved_provider and isinstance(runtime_provider, str): - resolved_provider = runtime_provider.strip() - if isinstance(runtime_requested_provider, str): - resolved_requested_provider = runtime_requested_provider.strip() + if not resolved_provider and isinstance(rk.get("provider"), str): + resolved_provider = rk["provider"].strip() + if isinstance(rk.get("requested_provider"), str): + resolved_requested_provider = rk["requested_provider"].strip() except Exception as exc: logger.debug( "image_routing: session runtime resolution failed, falling back to config — %s", @@ -1950,10 +1826,8 @@ class GatewayInboundMixin: async def _enrich_message_with_vision(self, user_text: str, image_paths: List[str]) -> str: """Auto-analyze user-attached images with the vision tool and prepend the descriptions. - Description *and* local cache path are injected so the model understands the image without - a tool call and can re-examine it with vision_analyze. - """ + a tool call and can re-examine it with vision_analyze.""" from tools.vision_tools import vision_analyze_tool from agent.memory_manager import sanitize_context @@ -1964,7 +1838,6 @@ class GatewayInboundMixin: "If it is a chart, diagram, or scientific figure, include the important " "labels, legend, and key values. Skip decorative details." ) - enriched_parts = [] for path in image_paths: try: @@ -1972,26 +1845,25 @@ class GatewayInboundMixin: result = json.loads(await vision_analyze_tool(image_url=path, user_prompt=analysis_prompt)) if result.get("success"): description = sanitize_context(result.get("analysis", "")) - enriched_parts.append( + note = ( f"[The user sent an image~ Here's what I can see:\n{description}]\n" f"[If you need a closer look, use vision_analyze with " f"image_url: {path} ~]" ) else: - enriched_parts.append( + note = ( "[The user sent an image but I couldn't quite see it " "this time (>_<) You can try looking at it yourself " f"with vision_analyze using image_url: {path}]" ) except Exception as e: logger.error("Vision auto-analysis error: %s", e) - enriched_parts.append( + note = ( f"[The user sent an image but something went wrong when I " f"tried to look at it~ You can try examining it yourself " f"with vision_analyze using image_url: {path}]" ) - - # Vision descriptions first, then the user's original text + enriched_parts.append(note) if not enriched_parts: return user_text prefix = "\n\n".join(enriched_parts) @@ -2011,12 +1883,8 @@ class GatewayInboundMixin: """One minimal neutral marker for every STT failure. Never mention "no STT provider" or setup steps — persisted in history they make the model keep volunteering STT-setup advice.""" from tools.credential_files import to_agent_visible_cache_path - agent_path = to_agent_visible_cache_path(os.path.abspath(path)) - return ( - "[voice message could not be transcribed automatically; " - f"the audio is available at: {agent_path}]" - ) + return f"[voice message could not be transcribed automatically; the audio is available at: {agent_path}]" async def _transcribe_one_clip(self, path: str, transcribe_audio, transcribe_audio_local_fallback) -> Tuple[Optional[str], str]: """``(transcript_or_None, note)`` for one clip via configured STT with local fallback.""" @@ -2046,12 +1914,9 @@ class GatewayInboundMixin: async def _enrich_message_with_transcription( self, user_text: str, audio_paths: List[str] ) -> tuple[str, List[str]]: - """Transcribe voice clips with the configured STT provider and prepend the transcripts. - - Returns ``(enriched_text, successful_transcripts)``: transcripts of successfully - transcribed clips in input order (empty if every clip failed or STT is disabled) so callers - can echo them back before the agent loop. - """ + """Transcribe voice clips with the configured STT provider and prepend the transcripts → + ``(enriched_text, successful_transcripts)``; the transcripts (input order; empty if every clip + failed or STT is disabled) let callers echo them back before the agent loop.""" from gateway.run import _probe_audio_duration audio_paths = list(dict.fromkeys(audio_paths)) if not getattr(self.config, "stt_enabled", True): @@ -2061,9 +1926,7 @@ class GatewayInboundMixin: duration_str = await _probe_audio_duration(abs_path) suffix = f" (duration: {duration_str})" if duration_str else "" notes.append(f"[The user sent a voice message: {abs_path}{suffix}]") - if not notes: - return user_text, [] - return self._prepend_media_prefix("\n\n".join(notes), user_text), [] + return (self._prepend_media_prefix("\n\n".join(notes), user_text) if notes else user_text), [] try: from tools.transcription_tools import ( @@ -2089,7 +1952,7 @@ class GatewayInboundMixin: enriched_parts.append(self._untranscribed_audio_note(path)) if enriched_parts: - return self._prepend_media_prefix("\n\n".join(enriched_parts), user_text), successful_transcripts + user_text = self._prepend_media_prefix("\n\n".join(enriched_parts), user_text) return user_text, successful_transcripts def _pending_event_audio_paths(self, event) -> List[str]: @@ -2103,23 +1966,15 @@ class GatewayInboundMixin: async def _transcribe_pending_audio_event_once( self, event, user_text: Optional[str] = None ) -> tuple[str | None, List[str]]: - """Transcribe a pending audio event once and cache the result on the event. - - The interrupt monitor and the pending-drain path both need the transcript; caching keeps - it to one STT call and one transcript echo per platform message. - """ + """Transcribe a pending audio event once and cache the result on the event: the interrupt + monitor and the pending-drain path both need it — one STT call and one echo per message.""" if hasattr(event, "_gateway_pending_stt_text"): - cached_transcripts = getattr(event, "_gateway_pending_stt_transcripts", []) or [] - return event._gateway_pending_stt_text, list(cached_transcripts) - + return event._gateway_pending_stt_text, list(getattr(event, "_gateway_pending_stt_transcripts", []) or []) audio_paths = self._pending_event_audio_paths(event) if not audio_paths: return user_text if user_text is not None else (getattr(event, "text", None) or None), [] - text = user_text if user_text is not None else (getattr(event, "text", "") or "") - enriched_text, successful_transcripts = await self._enrich_message_with_transcription( - text, audio_paths - ) + enriched_text, successful_transcripts = await self._enrich_message_with_transcription(text, audio_paths) event._gateway_pending_stt_text = enriched_text event._gateway_pending_stt_transcripts = list(successful_transcripts) return enriched_text, successful_transcripts @@ -2128,28 +1983,24 @@ class GatewayInboundMixin: self, event, adapter, source, transcripts: List[str], *, metadata=None, log_context: str = "Transcript", ) -> None: - """Echo pending-event STT transcripts to the chat at most once. - - Tracked as a COUNT (not a set — identical transcripts are distinct deliveries): - ``merge_pending_message_event`` can append a second voice note and invalidate the cache, - and the re-run returns earlier transcripts as a prefix, so only the unsent tail is echoed. - """ + """Echo pending-event STT transcripts to the chat at most once. Tracked as a COUNT (not a + set — identical transcripts are distinct deliveries): ``merge_pending_message_event`` can + append a second voice note and invalidate the cache; the re-run returns earlier transcripts + as a prefix, so only the unsent tail is echoed.""" if not transcripts or not self._should_echo_stt_transcripts() or adapter is None: return already_echoed = int(getattr(event, "_gateway_pending_stt_echoed", 0) or 0) - unsent = transcripts[already_echoed:] - setattr(event, "_gateway_pending_stt_echoed", already_echoed + len(unsent)) - await self._echo_stt_transcripts(adapter, source, unsent, metadata=metadata, log_context=log_context) + event._gateway_pending_stt_echoed = max(already_echoed, len(transcripts)) + await self._echo_stt_transcripts( + adapter, source, transcripts[already_echoed:], metadata=metadata, log_context=log_context, + ) async def _transcribe_and_echo_pending_voice( self, event, adapter, source, text: str, *, log_context: str, metadata=_UNSET ) -> tuple[str, List[str]]: - """Transcribe a pending voice event and echo transcripts once. - - Returns ``(enriched_text, transcripts)`` for ``agent.interrupt()`` or the pending-drain - flow; ``(text, [])`` unchanged when there is no STT-eligible media (caller owns the - ``_build_media_placeholder`` fallback for empty ``text`` with non-audio media). - """ + """Transcribe a pending voice event and echo transcripts once → ``(enriched_text, + transcripts)`` for ``agent.interrupt()`` or the pending-drain flow; ``(text, [])`` when there + is no STT-eligible media (caller owns the ``_build_media_placeholder`` fallback).""" if not self._pending_event_audio_paths(event): return text, [] try: