diff --git a/gateway/run_inbound.py b/gateway/run_inbound.py index 1eb873986c..695c419b5c 100644 --- a/gateway/run_inbound.py +++ b/gateway/run_inbound.py @@ -362,76 +362,64 @@ class GatewayInboundMixin: _quick_key: str, allow_gateway_control: bool, ) -> Optional[str]: - """Intercept a reply to a pending clarify prompt; None when the message falls through.""" - # Intercept replies to a pending clarify: open-ended prompts and "Other" responses are free - # text; direct replies to multi-choice prompts are accepted too ("2" → second option). + """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 try: from tools import clarify_gateway as _clarify_mod - _pending_clarify = _clarify_mod.get_pending_for_session( - _quick_key, include_choice_prompts=True, - ) + _pending_clarify = _clarify_mod.get_pending_for_session(_quick_key, include_choice_prompts=True) except Exception: _pending_clarify = None - if ( - allow_gateway_control - and _pending_clarify is not None - and _clarify_mod is not None - ): - _clarify_has_audio = bool(self._pending_event_audio_paths(event)) - _raw_clarify_reply = await self._prepare_clarify_reply_text(event) - if _clarify_has_audio and not _raw_clarify_reply: - logger.info( - "Gateway retained pending clarify after voice transcription " - "produced no usable text (session=%s, id=%s)", - _quick_key, - _pending_clarify.clarify_id, - ) - return "" - # Skip 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 _raw_clarify_reply and not _raw_clarify_reply.startswith("/"): - _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)", - _quick_key, _pending_clarify.clarify_id, - ) - # The clarify callback pauses the platform typing/status indicator while waiting - # so Slack users can type their answer. The active agent resumes as soon as this - # reply resolves the wait, so re-enable its indicator here too. - _clarify_adapter = self._adapter_for_source(source) - if _clarify_adapter: - try: - _clarify_adapter.resume_typing_for_chat(source.chat_id) - except Exception: - logger.debug( - "Failed to resume typing after clarify response", - exc_info=True, - ) - # Acknowledge with empty string so adapters that emit the agent's response don't - # double-post; the agent itself produces the next user-facing message. - return "" - 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)", - _quick_key, _pending_clarify.clarify_id, - ) - return "" - if _text_outcome == _clarify_mod.TEXT_REJECTED_PROSE: - # Native-choice prompts deliberately reject unmatched prose so it can continue - # through normal busy-message routing. Release this clarify first: redirect() - # degrades to steer() while tools execute, and that steer cannot drain until - # the clarify tool returns. - _clarify_mod.resolve_gateway_clarify( - _pending_clarify.clarify_id, - "", - ) + if not (allow_gateway_control and _pending_clarify is not None and _clarify_mod is not None): + return None + _clarify_has_audio = bool(self._pending_event_audio_paths(event)) + _raw_clarify_reply = await self._prepare_clarify_reply_text(event) + if _clarify_has_audio and not _raw_clarify_reply: + logger.info( + "Gateway retained pending clarify after voice transcription " + "produced no usable text (session=%s, id=%s)", + _quick_key, + _pending_clarify.clarify_id, + ) + return "" + # 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)", + _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) + if _clarify_adapter: + try: + _clarify_adapter.resume_typing_for_chat(source.chat_id) + except Exception: + logger.debug("Failed to resume typing after clarify response", exc_info=True) + return "" + 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)", + _quick_key, _pending_clarify.clarify_id, + ) + return "" + 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 + # execute, and that steer cannot drain until the clarify tool returns. + _clarify_mod.resolve_gateway_clarify(_pending_clarify.clarify_id, "") return None # Reply → choice for a pending slash-confirm prompt; the command spelling wins over the @@ -1334,6 +1322,70 @@ class GatewayInboundMixin: logger.debug("Skill command check failed (non-fatal): %s", e) return None + 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_gateway_control = event.allow_gateway_control + _reply = self._hm_update_prompt_reply(event, _quick_key, allow_gateway_control) + if _reply is not None: + return _reply + _reply = await self._hm_clarify_reply(event, source, _quick_key, allow_gateway_control) + if _reply is not None: + return _reply + return await self._hm_slash_confirm_reply(event, _quick_key, allow_gateway_control) + + async def _hm_dispatch_idle_commands( + self, event: "MessageEvent", source: SessionSource, _quick_key: str + ) -> 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 + + 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. + """ + try: + _orphan_adapter = self._adapter_for_source(source) + if ( + _orphan_adapter is not None + and not bool(getattr(event, "internal", False)) + and not event.get_command() + ): + _rescued = self._rescue_orphaned_overflow(_quick_key, _orphan_adapter) + if _rescued is not None: + # Into the slot when the chain was a single orphan (post-turn drain picks it + # up), otherwise into overflow behind the already-staged next orphan. + self._enqueue_fifo(_quick_key, event, _orphan_adapter) + event = _rescued + # 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) + if _rescued_source is not None: + source = _rescued_source + is_internal = 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. @@ -1350,121 +1402,49 @@ class GatewayInboundMixin: if _paused_notice is not None: return _paused_notice - # Replies owned by in-flight work: pending /update prompt, clarify, slash-confirm. _quick_key = self._session_key_for_source(source) - allow_gateway_control = event.allow_gateway_control - _update_reply = self._hm_update_prompt_reply(event, _quick_key, allow_gateway_control) - if _update_reply is not None: - return _update_reply - _clarify_reply = await self._hm_clarify_reply( - event, source, _quick_key, allow_gateway_control - ) - if _clarify_reply is not None: - return _clarify_reply - _confirm_reply = await self._hm_slash_confirm_reply( - event, _quick_key, allow_gateway_control - ) - if _confirm_reply is not None: - return _confirm_reply + _reply = await self._hm_pending_reply_intercepts(event, source, _quick_key) + if _reply is not None: + return _reply self._hm_evict_stale_running_agent(_quick_key) if self._is_session_running(_quick_key): return await self._hm_handle_running_session_message(event, source, _quick_key) - # 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 - ) + _handled, _result = await self._hm_dispatch_idle_commands(event, source, _quick_key) if _handled: return _result - _handled, _result = await self._hm_dispatch_canonical_command( - event, source, _quick_key, canonical - ) - if _handled: - return _result - _handled, _result, command = await self._hm_dispatch_quick_and_plugin_commands( - event, source, command - ) - if _handled: - return _result - _skill_reply = self._hm_skill_slash_rewrite(event, source, _quick_key, command) - if _skill_reply is not None: - return _skill_reply # 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 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 (stop accepting - # FIRST, then NAS polls active_agents==0). Internal/system events bypass; reversible. + # 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, - ) + 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 passes the "already running" guard and spins - # up a duplicate agent for the same session, corrupting the transcript. - _active_session_lease, _limit_message = self._claim_active_session_slot( - _quick_key, - source, - ) + # 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 + # passes the "already running" guard and spins up a duplicate agent for the same session. + _active_session_lease, _limit_message = self._claim_active_session_slot(_quick_key, source) if _limit_message is not None: - logger.info( - "Rejecting new active session %s: max_concurrent_sessions reached", - _quick_key, - ) + logger.info("Rejecting new active session %s: max_concurrent_sessions reached", _quick_key) return _limit_message - # ── 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. Re-stage them FIFO and - # enqueue this event behind them. Skipped for control commands and internal events. - try: - _orphan_adapter = self._adapter_for_source(source) - if ( - _orphan_adapter is not None - and not bool(getattr(event, "internal", False)) - and not event.get_command() - ): - _rescued = self._rescue_orphaned_overflow( - _quick_key, _orphan_adapter - ) - if _rescued is not None: - # The oldest orphan runs as THIS turn. Park the incoming event behind the rest of - # the chain: into the slot when the chain was a single orphan (post-turn drain - # picks it up), otherwise into overflow behind the already-staged next orphan. - self._enqueue_fifo(_quick_key, event, _orphan_adapter) - event = _rescued - # 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) - if _rescued_source is not None: - source = _rescued_source - is_internal = bool(getattr(_rescued, "internal", False)) - except Exception: - logger.debug( - "FIFO orphan rescue pre-claim failed for %s", - _quick_key, - exc_info=True, - ) + event, source, is_internal = self._hm_rescue_orphaned_fifo(event, source, is_internal, _quick_key) _claim_state = self._session_state(_quick_key) if _active_session_lease is not None: @@ -1480,8 +1460,8 @@ class GatewayInboundMixin: event, source, _quick_key, _run_generation ) except TurnLeaseTimeoutError as exc: - # A rejected message, not a completed turn: return before the /goal judge below so - # it cannot consume the resend notice and enqueue a synthetic continuation loop. + # 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. logger.error( "Rejecting turn for routing key %s on session %s after " "turn-lease timeout; transcript load was not started and " @@ -1496,30 +1476,25 @@ class GatewayInboundMixin: ) try: await self._run_post_turn_hooks( - agent_result=_agent_result, - source=source, - is_internal=is_internal, - event=event, + agent_result=_agent_result, source=source, is_internal=is_internal, event=event, ) except Exception as _goal_exc: logger.debug("post-turn hook failed: %s", _goal_exc) return _agent_result finally: - # MoA one-shot restore must run on EVERY exit path: the restore data lives on the - # per-turn event, so a restore in the try block is skipped when the handler raises and - # the override leaks permanently; finally covers success, exception and interrupt. + # MoA one-shot restore must run on EVERY exit path (success, exception, interrupt): + # the restore data lives on the per-turn event and would leak permanently otherwise. self._restore_moa_one_shot(event, _quick_key) self._restore_pending_one_turn_model_override(_quick_key) - # Normal completion/exception/interrupt clears this durable marker; SIGKILL/OOM skips - # finally, leaving it for the next unclean startup's recovery pass. + # SIGKILL/OOM skips finally, leaving the durable marker for the next unclean startup's + # recovery pass. await self._clear_durable_active_turn(event) - # Unconditional release covers every exit path: _release_running_agent_state is idempotent - # and, without a run_generation guard, clears the slot whichever generation holds it. This - # evicts the zombie left when session_reset bumps the generation mid-flight (gen-N's - # guarded release in _run_agent returns False; a sentinel-only check would lock forever). + # Unconditional, idempotent release without a run_generation guard: evicts the zombie + # left when session_reset bumps the generation mid-flight (gen-N's guarded release in + # _run_agent returns False; a sentinel-only check would lock forever). self._release_running_agent_state(_quick_key) - # Turn lease: release THIS turn's token — keyed by (routing key, run generation) so this - # unwind can only free the lease its own turn acquired, never a newer turn's. + # Turn lease is keyed by (routing key, run generation) so this unwind can only free + # the lease its own turn acquired, never a newer turn's. self._release_turn_lease(_quick_key, _run_generation) def _restore_moa_one_shot(self, event: "MessageEvent", quick_key: str) -> None: @@ -1555,31 +1530,24 @@ class GatewayInboundMixin: def _prefix_inbound_sender_context(self, event: MessageEvent, source: SessionSource, message_text: str) -> str: """Attribute the sender in shared multi-user sessions and prepend history-backfill channel context.""" - _group_sessions_per_user = getattr(self.config, "group_sessions_per_user", True) - _thread_sessions_per_user = getattr(self.config, "thread_sessions_per_user", False) _is_shared_multi_user = is_shared_multi_user_session( source, - group_sessions_per_user=_group_sessions_per_user, - thread_sessions_per_user=_thread_sessions_per_user, + group_sessions_per_user=getattr(self.config, "group_sessions_per_user", True), + thread_sessions_per_user=getattr(self.config, "thread_sessions_per_user", False), ) if _is_shared_multi_user and source.user_name: - # source.user_name is the platform display name — attacker-influenceable on any - # platform that lets participants set their own name. Neutralize newlines/control chars - # before interpolating it into every message, or a hostile name can masquerade as a - # fake markdown section (mirrors build_session_context_prompt's treatment). + # 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). _safe_user_name = neutralize_untrusted_inline_text(source.user_name) - # On Slack, expose the current author's verifiable user ID next to the display name: - # "mention me again" requests need a trusted `<@U...>` target for the CURRENT speaker — - # display names are ambiguous and historical mentions may point at someone else. The - # user_id comes from the Slack event envelope (not user-editable), so no neutralization. + # 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. if source.platform == Platform.SLACK and source.user_id: - _safe_user_name = ( - f"{_safe_user_name} | Slack user <@{source.user_id}>" - ) + _safe_user_name = f"{_safe_user_name} | Slack user <@{source.user_id}>" message_text = f"[{_safe_user_name}] {message_text}" - # Prepend history-backfill channel context after the sender-prefix so the prefix applies - # only to the trigger message, not the backfill block. + # 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}" return message_text @@ -1588,113 +1556,92 @@ 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.""" + """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: list[str] = [] audio_paths: list[str] = [] audio_file_paths: list[str] = [] video_paths: list[str] = [] - - if event.media_urls: - for i, path in enumerate(event.media_urls): - mtype = event.media_types[i] if i < len(event.media_types) else "" - # Classify images per-attachment: trust this attachment's own MIME, and only honour - # the message-level PHOTO type when the per-attachment MIME is unknown. Otherwise a - # document sent alongside an image gets mis-routed as an image and the provider 400s. - if _event_media_is_image(event, i): - image_paths.append(path) - # MessageType.AUDIO = audio file attachment (e.g. .mp3, .m4a) — never STT. - # Mixed DOCUMENT events also preserve audio as a file path instead of - # dropping it or treating it as a voice note. - if _event_media_is_audio(event, i): - if event.message_type in {MessageType.AUDIO, MessageType.DOCUMENT}: - audio_file_paths.append(path) - elif not pending_stt_prepared and _event_media_is_stt_input(event, i): - audio_paths.append(path) - if mtype.startswith("video/") or (not mtype and event.message_type == MessageType.VIDEO): - video_paths.append(path) + for i, path in enumerate(event.media_urls or []): + mtype = event.media_types[i] if i < len(event.media_types) else "" + if _event_media_is_image(event, i): + image_paths.append(path) + if _event_media_is_audio(event, i): + if event.message_type in {MessageType.AUDIO, MessageType.DOCUMENT}: + audio_file_paths.append(path) + elif not pending_stt_prepared and _event_media_is_stt_input(event, i): + audio_paths.append(path) + if mtype.startswith("video/") or (not mtype and event.message_type == MessageType.VIDEO): + video_paths.append(path) return image_paths, audio_paths, audio_file_paths, video_paths async def _enrich_inbound_images( self, source: SessionSource, session_key: str, message_text: str, image_paths: list[str] ) -> str: - # Decide routing: native (attach pixels) vs text (vision_analyze pre-run + prepend - # description). See agent/image_routing.py. Offloaded to a thread: the decision does - # blocking network I/O (models.dev fetch on cache miss, Ollama /api/show probe) whose - # timeout would otherwise stall the whole gateway event loop. + """Route images natively (attach pixels at run_conversation) or pre-analyze them into text.""" + # See agent/image_routing.py. Offloaded to a thread: the decision does blocking network I/O + # (models.dev fetch on cache miss, Ollama /api/show probe) that would stall the event loop. _img_mode = await asyncio.to_thread( - self._decide_image_input_mode, - source=source, - session_key=session_key, + self._decide_image_input_mode, source=source, session_key=session_key, ) if _img_mode == "native": - # Defer attachment to the run_conversation call site. - self._session_state( - session_key - ).persistent.native_image_paths = list(image_paths) + self._session_state(session_key).persistent.native_image_paths = list(image_paths) logger.info( "Image routing: native (model supports vision). %d image(s) will be attached inline.", len(image_paths), ) - else: - logger.info( - "Image routing: text (mode=%s). Pre-analyzing %d image(s) via vision_analyze.", - _img_mode, len(image_paths), + return message_text + logger.info( + "Image routing: text (mode=%s). Pre-analyzing %d image(s) via vision_analyze.", + _img_mode, len(image_paths), + ) + # Vision enrichment runs before AIAgent.run_conversation(), so bind this session's resolved + # runtime explicitly rather than consulting process-global compatibility mirrors. + vision_runtime = None + try: + turn_model, runtime_kwargs = self._resolve_session_agent_runtime( + source=source, session_key=session_key, ) - # Vision enrichment runs before AIAgent.run_conversation(), - # so bind this session's resolved runtime explicitly rather - # than consulting process-global compatibility mirrors. - vision_runtime = None + vision_runtime = dict(runtime_kwargs or {}) + vision_runtime["model"] = turn_model + except Exception: + logger.debug("vision enrichment: session runtime resolution failed", exc_info=True) + + from agent.auxiliary_client import scoped_runtime_main + + with scoped_runtime_main(vision_runtime): + return await self._enrich_message_with_vision(message_text, image_paths) + + async def _echo_stt_transcripts( + self, adapter, source: SessionSource, transcripts: List[str], *, metadata=None, log_context: str = "Transcript" + ) -> None: + """Send each transcript back as ``🎙️ "…"`` (best-effort; failures are logged, never raised).""" + for tx in transcripts: try: - turn_model, runtime_kwargs = self._resolve_session_agent_runtime( - source=source, - session_key=session_key, - ) - vision_runtime = dict(runtime_kwargs or {}) - vision_runtime["model"] = turn_model - except Exception: - logger.debug( - "vision enrichment: session runtime resolution failed", - exc_info=True, - ) - - from agent.auxiliary_client import scoped_runtime_main - - with scoped_runtime_main(vision_runtime): - message_text = await self._enrich_message_with_vision( - message_text, - image_paths, - ) - return message_text + await adapter.send(source.chat_id, f'🎙️ "{tx}"', metadata=metadata) + except Exception as echo_exc: + logger.debug("%s echo failed (non-fatal): %s", log_context, echo_exc) async def _enrich_inbound_voice( self, event: MessageEvent, source: SessionSource, message_text: str, audio_paths: list[str] ) -> str: message_text, _successful_transcripts = await self._enrich_message_with_transcription( - message_text, - audio_paths, + message_text, audio_paths, ) - # Echo each successful transcript back to the user immediately when configured. Lets - # users verify STT quality in real-time, while allowing quiet STT for users who only - # want the agent to receive the transcription. + # Echo each successful transcript back immediately when configured so users can verify STT + # quality in real time (quiet STT stays available for users who only want the agent to + # receive it). On transcription failure do NOT send a hardcoded notice: that bypassed the + # LLM and produced two replies; enrichment leaves one neutral marker for a localized reply. if _successful_transcripts and self._should_echo_stt_transcripts(): _echo_adapter = self._adapter_for_source(source) _echo_meta = self._thread_metadata_for_source(source, self._reply_anchor_for_event(event)) if _echo_adapter: - for _tx in _successful_transcripts: - try: - await _echo_adapter.send( - source.chat_id, - f'🎙️ "{_tx}"', - metadata=_echo_meta, - ) - except Exception as _echo_exc: - logger.debug( - "Transcript echo failed (non-fatal): %s", _echo_exc, - ) - # On transcription failure, do NOT send a hardcoded notice here: that bypassed the - # LLM and produced two replies (one pre-canned, TTS'd in the wrong language). - # Enrichment leaves a single neutral marker so the LLM gives one localized reply. + await self._echo_stt_transcripts(_echo_adapter, source, _successful_transcripts, metadata=_echo_meta) return message_text @staticmethod @@ -1775,10 +1722,10 @@ class GatewayInboundMixin: @staticmethod def _prepend_inbound_reply_context(event: MessageEvent, source: SessionSource, message_text: str) -> str: + """Prepend the Discord triggering-message id and the reply-to pointer.""" # Discord: surface the triggering message id per-turn on the user message rather than in the - # cached system prompt. message_id changes every turn, so baking it into - # build_session_context_prompt() would bust the agent-cache signature and rebuild the - # AIAgent every message (destroying prompt caching). + # cached system prompt — message_id changes every turn and would bust the agent-cache + # signature (rebuilding the AIAgent every message destroys prompt caching). if ( source is not None and getattr(source, "platform", None) == Platform.DISCORD @@ -1793,15 +1740,11 @@ class GatewayInboundMixin: ) if getattr(event, "reply_to_text", None) and event.reply_to_message_id: - # Always inject the reply-to pointer — even when the quoted text already appears in - # history. The prefix isn't deduplication, it's disambiguation: it tells the agent - # *which* prior message the user is referencing. Token overhead is minimal. + # Always inject the reply-to pointer even when the quoted text is already in history: + # it's disambiguation (*which* prior message), not deduplication. reply_snippet = event.reply_to_text[:500] if getattr(event, "reply_to_is_own_message", False): - message_text = ( - f'[Replying to your previous message: "{reply_snippet}"]\n\n' - f"{message_text}" - ) + message_text = f'[Replying to your previous message: "{reply_snippet}"]\n\n{message_text}' else: message_text = f'[Replying to: "{reply_snippet}"]\n\n{message_text}' return message_text @@ -1973,22 +1916,11 @@ class GatewayInboundMixin: ) -> Optional[str]: """Run inbound preprocessing under the routed profile when multiplexed.""" from gateway.run import _async_profile_runtime_scope + kwargs = dict(event=event, source=source, history=history, session_key=session_key) if getattr(getattr(self, "config", None), "multiplex_profiles", False): - async with _async_profile_runtime_scope( - self._resolve_profile_home_for_source(source) - ): - return await self._prepare_inbound_message_text( - event=event, - source=source, - history=history, - session_key=session_key, - ) - return await self._prepare_inbound_message_text( - event=event, - source=source, - history=history, - session_key=session_key, - ) + async with _async_profile_runtime_scope(self._resolve_profile_home_for_source(source)): + return await self._prepare_inbound_message_text(**kwargs) + return await self._prepare_inbound_message_text(**kwargs) async def _prepare_clarify_reply_text(self, event) -> str: """Return raw text or successful voice transcripts for a clarify reply.""" @@ -2036,7 +1968,11 @@ class GatewayInboundMixin: return True async def _clear_durable_active_turn(self, event: "MessageEvent") -> bool: - """Best-effort CAS clear of the marker owned by *event*.""" + """Best-effort CAS clear of the marker owned by *event* (3 attempts). + + Never lets marker cleanup block agent/lease release; a stale marker is bounded by the + agent timeout and the clean-start orphan-marker discard path. + """ session_key = getattr(event, "_gateway_active_turn_session_key", None) token = getattr(event, "_gateway_active_turn_token", None) try: @@ -2045,33 +1981,20 @@ class GatewayInboundMixin: last_error: Optional[Exception] = None for attempt in range(1, 4): try: - return bool( - await self.async_session_store.clear_turn_active( - session_key, token - ) - ) + return bool(await self.async_session_store.clear_turn_active(session_key, token)) except Exception as exc: last_error = exc if attempt < 3: logger.debug( "Retrying active-turn marker cleanup for %s (%d/3): %s", - session_key, - attempt, - exc, + session_key, attempt, exc, ) - # Never let marker cleanup block agent/lease release; a stale marker is bounded by the - # agent timeout and the clean-start orphan-marker discard path. logger.warning( - "Could not clear active-turn marker for %s after 3 attempts: %s", - session_key, - last_error, + "Could not clear active-turn marker for %s after 3 attempts: %s", session_key, last_error, ) return False finally: - for attr in ( - "_gateway_active_turn_session_key", - "_gateway_active_turn_token", - ): + for attr in ("_gateway_active_turn_session_key", "_gateway_active_turn_token"): with suppress(AttributeError): delattr(event, attr) @@ -2097,16 +2020,14 @@ class GatewayInboundMixin: content: str, plugin_id: str, ) -> bool: - """Schedule a plugin-triggered turn on the live gateway loop.""" + """Schedule a plugin-triggered turn on the live gateway loop (thread-safe).""" from gateway.run import safe_schedule_threadsafe loop = getattr(self, "_gateway_loop", None) if not getattr(self, "_running", False) or loop is None or loop.is_closed(): return False coro = self._dispatch_plugin_message_injection( - session_key=session_key, - content=content, - plugin_id=plugin_id, + session_key=session_key, content=content, plugin_id=plugin_id, ) try: current_loop = asyncio.get_running_loop() @@ -2118,10 +2039,7 @@ class GatewayInboundMixin: future = loop.create_task(coro) except Exception: coro.close() - logger.warning( - "Plugin message injection scheduling failed", - exc_info=True, - ) + logger.warning("Plugin message injection scheduling failed", exc_info=True) return False self._background_tasks.add(future) future.add_done_callback(self._background_tasks.discard) @@ -2144,16 +2062,12 @@ class GatewayInboundMixin: except Exception: logger.warning( "Plugin message injection failed: plugin=%s session=%s", - plugin_id, - session_key, - exc_info=True, + 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, + "Plugin message injection was not routed: plugin=%s session=%s", plugin_id, session_key, ) future.add_done_callback(_log_result) @@ -2167,35 +2081,28 @@ class GatewayInboundMixin: plugin_id: str, ) -> bool: """Route a plugin-triggered turn through the session's live adapter.""" - if not getattr(self, "_running", False) or getattr(self, "_draining", False): - return False + def _accepting() -> bool: + return getattr(self, "_running", False) and not getattr(self, "_draining", False) - entry = await self.async_session_store.lookup_by_session_key(session_key) - if entry is None or entry.origin is None: + if not _accepting(): return False - if not getattr(self, "_running", False) or getattr(self, "_draining", False): + entry = await self.async_session_store.lookup_by_session_key(session_key) + if entry is None or entry.origin is None or not _accepting(): return False source = dataclasses.replace(entry.origin) try: - if not self._is_user_authorized( - source, - allow_adapter_delegation=False, - ): + 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, + plugin_id, session_key, ) return False except Exception: logger.warning( - "Plugin message injection authorization check failed: " - "plugin=%s session=%s", - plugin_id, - session_key, - exc_info=True, + "Plugin message injection authorization check failed: plugin=%s session=%s", + plugin_id, session_key, exc_info=True, ) return False @@ -2220,9 +2127,7 @@ class GatewayInboundMixin: await adapter.handle_message(event) logger.info( "Plugin message injection dispatched: plugin=%s session=%s session_id=%s", - plugin_id, - session_key, - entry.session_id, + plugin_id, session_key, entry.session_id, ) return True @@ -2235,13 +2140,11 @@ class GatewayInboundMixin: provider: Optional[str] = None, model: Optional[str] = None, ) -> str: - """Resolve image-input routing for the effective model this turn. + """Resolve image-input routing (``"native"`` / ``"text"``) for the effective model this turn. - Returns ``"native"`` (attach pixels on the user turn) or ``"text"`` (pre-analyze with - vision_analyze and prepend the description); 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. + 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. """ try: from agent.image_routing import decide_image_input_mode @@ -2254,22 +2157,16 @@ class GatewayInboundMixin: resolved_requested_provider = "" needs_session_runtime = not resolved_provider or not resolved_model - has_session_identity = source is not None or session_key - if needs_session_runtime and has_session_identity: + if needs_session_runtime and (source is not None or session_key): try: turn_model, runtime_kwargs = self._resolve_session_agent_runtime( - source=source, - session_key=session_key, - user_config=cfg, + source=source, session_key=session_key, user_config=cfg, ) if not resolved_model and isinstance(turn_model, str): resolved_model = turn_model.strip() - runtime_provider = runtime_kwargs.get("provider") if isinstance(runtime_kwargs, dict) else None - runtime_requested_provider = ( - runtime_kwargs.get("requested_provider") - if isinstance(runtime_kwargs, dict) - else None - ) + 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): @@ -2286,10 +2183,7 @@ class GatewayInboundMixin: resolved_model = _read_main_model() return decide_image_input_mode( - resolved_provider, - resolved_model, - cfg, - requested_provider=resolved_requested_provider, + resolved_provider, resolved_model, cfg, requested_provider=resolved_requested_provider, ) except Exception as exc: logger.debug("image_routing: decision failed, falling back to text — %s", exc) @@ -2300,11 +2194,10 @@ class GatewayInboundMixin: user_text: str, image_paths: List[str], ) -> str: - """Auto-analyze user-attached images with the vision tool and prepend the descriptions to - the message text. + """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. Returns the enriched message string. + 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 @@ -2321,14 +2214,10 @@ class GatewayInboundMixin: for path in image_paths: try: logger.debug("Auto-analyzing user image: %s", path) - result_json = await vision_analyze_tool( - image_url=path, - user_prompt=analysis_prompt, - ) + result_json = await vision_analyze_tool(image_url=path, user_prompt=analysis_prompt) result = json.loads(result_json) if result.get("success"): - description = result.get("analysis", "") - description = sanitize_context(description) + description = sanitize_context(result.get("analysis", "")) enriched_parts.append( 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 " @@ -2348,13 +2237,11 @@ class GatewayInboundMixin: f"with vision_analyze using image_url: {path}]" ) - # Combine: vision descriptions first, then the user's original text - if enriched_parts: - prefix = "\n\n".join(enriched_parts) - if user_text: - return f"{prefix}\n\n{user_text}" - return prefix - return user_text + # Vision descriptions first, then the user's original text + if not enriched_parts: + return user_text + prefix = "\n\n".join(enriched_parts) + return f"{prefix}\n\n{user_text}" if user_text else prefix _EMPTY_TEXT_PLACEHOLDER = "(The user sent a message with no text content)" @@ -2512,24 +2399,12 @@ class GatewayInboundMixin: ``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. """ - if ( - not transcripts - or not self._should_echo_stt_transcripts() - or adapter is None - ): + 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)) - for tx in unsent: - try: - await adapter.send( - source.chat_id, - f'🎙️ "{tx}"', - metadata=metadata, - ) - except Exception as echo_exc: - logger.debug("%s echo failed (non-fatal): %s", log_context, echo_exc) + await self._echo_stt_transcripts(adapter, source, unsent, metadata=metadata, log_context=log_context) async def _transcribe_and_echo_pending_voice( self,