refactor(gateway): run_inbound — _handle_message phase helpers, shared STT echo, plugin-injection/vision/reply-context compaction

This commit is contained in:
Teknium
2026-09-02 18:50:20 -07:00
parent 8a323df49f
commit cca4b89724
+268 -393
View File
@@ -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,