diff --git a/gateway/run.py b/gateway/run.py index d05b1d49ec..9888b39907 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -27050,50 +27050,45 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew ) return get_hermes_home() - async def _run_agent_inner( - self, - message: str, - context_prompt: str, - history: List[Dict[str, Any]], - source: SessionSource, - session_id: str, - session_key: str = None, - run_generation: Optional[int] = None, - _interrupt_depth: int = 0, - event_message_id: Optional[str] = None, - inbound_message_id: Optional[str] = None, - channel_prompt: Optional[str] = None, - moa_config: Optional[dict] = None, - persist_user_message: Optional[Any] = None, - persist_user_timestamp: Optional[float] = None, - persist_user_display_kind: Optional[str] = None, - message_type: Optional[str] = None, - ) -> Dict[str, Any]: - """Run the agent; returns the full run_conversation result dict. + @dataclasses.dataclass + class _RunAgentDisplay: + """Per-turn display / progress settings resolved by ``_run_agent_display_settings``.""" - Keys: "final_response", "messages", "api_calls", "completed". - """ - # ---- Proxy mode: delegate to remote API server ---- - if self._get_proxy_url(): - return await self._run_agent_via_proxy( - message=message, - context_prompt=context_prompt, - history=history, - source=source, - session_id=session_id, - session_key=session_key, - run_generation=run_generation, - event_message_id=event_message_id, - ) + user_config: Any = None + platform_key: Any = None + enabled_toolsets: Any = None + disabled_toolsets: Any = None + resolve_display_setting: Any = None + progress_mode: Any = None + progress_grouping: Any = None + _display_surface_mode: Any = None + tool_progress_enabled: Any = None + _live_status_mode: Any = None + _live_status_adapter: Any = None + log_mode_enabled: Any = None + log_queue: Any = None + interim_assistant_messages_enabled: Any = None + _thinking_enabled: Any = None + _native_slack_task_cards: Any = None + needs_progress_queue: Any = None + _generic_status_phrase: Any = None - from run_agent import AIAgent - import queue + @dataclasses.dataclass + class _RunAgentWorker: + """Executor future + inactivity-watchdog handles for one ``_run_agent_inner`` turn.""" - def _run_still_current() -> bool: - if run_generation is None or not session_key: - return True - return self._is_session_run_current(session_key, run_generation) + executor_task: Any = None + agent_timeout: Optional[float] = None + agent_warning: Optional[float] = None + task_id: str = "" + process_baseline: Any = None + worker_done: Any = None + timeout_fired: Any = None + cleanup_lock: Any = None + is_current: Any = None + def _run_agent_display_settings(self, source: SessionSource) -> "GatewayRunner._RunAgentDisplay": + """Resolve per-platform display, progress, status and streaming-surface settings for a turn.""" user_config = _load_gateway_config() platform_key = _platform_config_key(source.platform) @@ -27253,9 +27248,60 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew needs_progress_queue = ( tool_progress_enabled or _thinking_enabled or _native_slack_task_cards ) + return self._RunAgentDisplay( + user_config=user_config, + platform_key=platform_key, + enabled_toolsets=enabled_toolsets, + disabled_toolsets=disabled_toolsets, + resolve_display_setting=resolve_display_setting, + progress_mode=progress_mode, + progress_grouping=progress_grouping, + _display_surface_mode=_display_surface_mode, + tool_progress_enabled=tool_progress_enabled, + _live_status_mode=_live_status_mode, + _live_status_adapter=_live_status_adapter, + log_mode_enabled=log_mode_enabled, + log_queue=log_queue, + interim_assistant_messages_enabled=interim_assistant_messages_enabled, + _thinking_enabled=_thinking_enabled, + _native_slack_task_cards=_native_slack_task_cards, + needs_progress_queue=needs_progress_queue, + _generic_status_phrase=_generic_status_phrase, + ) + + def _run_agent_build_turn_context( + self, + disp: "GatewayRunner._RunAgentDisplay", + AIAgent: Any, + *, + message: str, + context_prompt: str, + history: List[Dict[str, Any]], + source: SessionSource, + session_id: str, + session_key: Optional[str], + run_generation: Optional[int], + _interrupt_depth: int, + event_message_id: Optional[str], + inbound_message_id: Optional[str], + channel_prompt: Optional[str], + moa_config: Optional[dict], + persist_user_message: Optional[Any], + persist_user_timestamp: Optional[float], + persist_user_display_kind: Optional[str], + ) -> Tuple[TurnContext, TurnRunner, Any]: + """Build the progress queues / holders, the ``TurnContext`` and its ``TurnRunner``. + + Returns ``(turn_ctx, turn_runner, cleanup_adapter)``; the progress-bubble cleanup flags travel + on ``turn_ctx._cleanup_progress`` / ``turn_ctx._cleanup_msg_ids``. + """ + def _run_still_current() -> bool: + if run_generation is None or not session_key: + return True + return self._is_session_run_current(session_key, run_generation) # Queue for progress messages (thread-safe) - progress_queue = queue.Queue() if needs_progress_queue else None + progress_queue = queue.Queue() if disp.needs_progress_queue else None last_tool = [None] # Mutable container for tracking in closure last_progress_msg = [None] # Track last message for dedup repeat_count = [0] # How many times the same message repeated @@ -27286,7 +27332,7 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew # Auto-cleanup of temporary progress bubbles (Telegram + any adapter that implements # ``delete_message``). Failed runs skip cleanup so the bubbles remain as breadcrumbs. _cleanup_progress = bool( - resolve_display_setting(user_config, platform_key, "cleanup_progress") + disp.resolve_display_setting(disp.user_config, disp.platform_key, "cleanup_progress") ) _cleanup_adapter = self._adapter_for_source(source) if _cleanup_progress else None # getattr, not attribute access — same duck-typed-adapter guard as the edit_message check in @@ -27308,14 +27354,14 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew turn_ctx = TurnContext( source=source, _run_still_current=_run_still_current, - _live_status_adapter=_live_status_adapter, - _live_status_mode=_live_status_mode, - _thinking_enabled=_thinking_enabled, - progress_mode=progress_mode, - progress_grouping=progress_grouping, - tool_progress_enabled=tool_progress_enabled, + _live_status_adapter=disp._live_status_adapter, + _live_status_mode=disp._live_status_mode, + _thinking_enabled=disp._thinking_enabled, + progress_mode=disp.progress_mode, + progress_grouping=disp.progress_grouping, + tool_progress_enabled=disp.tool_progress_enabled, progress_queue=progress_queue, - log_queue=log_queue, + log_queue=disp.log_queue, last_progress_msg=last_progress_msg, last_tool=last_tool, last_was_terminal_block=last_was_terminal_block, @@ -27326,14 +27372,14 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew _cleanup_msg_ids=_cleanup_msg_ids, message=message, AIAgent=AIAgent, - resolve_display_setting=resolve_display_setting, - user_config=user_config, - enabled_toolsets=enabled_toolsets, - disabled_toolsets=disabled_toolsets, - log_mode_enabled=log_mode_enabled, - interim_assistant_messages_enabled=interim_assistant_messages_enabled, - needs_progress_queue=needs_progress_queue, - _native_slack_task_cards=_native_slack_task_cards, + resolve_display_setting=disp.resolve_display_setting, + user_config=disp.user_config, + enabled_toolsets=disp.enabled_toolsets, + disabled_toolsets=disp.disabled_toolsets, + log_mode_enabled=disp.log_mode_enabled, + interim_assistant_messages_enabled=disp.interim_assistant_messages_enabled, + needs_progress_queue=disp.needs_progress_queue, + _native_slack_task_cards=disp._native_slack_task_cards, _voice_ack_fired=_voice_ack_fired, _voice_ack_guild=_voice_ack_guild, _voice_ack_loop=_voice_ack_loop, @@ -27360,7 +27406,18 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew turn_ctx.native_tool_complete_callback = ( turn_runner.native_tool_complete_callback ) + return turn_ctx, turn_runner, _cleanup_adapter + def _run_agent_progress_threading( + self, + source: SessionSource, + event_message_id: Optional[str], + _native_slack_task_cards: bool, + ) -> Tuple[Optional[dict], Optional[str], Any, Optional[str]]: + """Resolve where progress bubbles are threaded. + + Returns ``(progress_metadata, progress_reply_to, progress_thread_id, relay_prospective_thread_id)``. + """ # Background task accumulating tool lines into one edited progress message. Threading metadata # is platform-specific: Slack DM threading needs the event_message_id fallback; Telegram forum # topics use message_thread_id and Hermes-created private DM topic lanes need thread metadata @@ -27459,96 +27516,68 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew or _relay_prospective_thread_id else None ) + return _progress_metadata, _progress_reply_to, _progress_thread_id, _relay_prospective_thread_id - async def write_tool_log(): - """Drain log_queue and append tool-call lines to tool_calls.log (tool_progress=log). + async def _run_agent_write_tool_log(self, log_queue: Any) -> None: + """Drain log_queue and append tool-call lines to tool_calls.log (tool_progress=log). - RotatingFileHandler (5MB × 3) bounds the log; RedactingFormatter keeps secrets off disk. - """ - if log_queue is None: - return - from logging.handlers import RotatingFileHandler + RotatingFileHandler (5MB × 3) bounds the log; RedactingFormatter keeps secrets off disk. + """ + if log_queue is None: + return + from logging.handlers import RotatingFileHandler - from agent.redact import RedactingFormatter + from agent.redact import RedactingFormatter - log_dir = _hermes_home / "logs" - log_dir.mkdir(parents=True, exist_ok=True) - file_handler = RotatingFileHandler( - log_dir / "tool_calls.log", - maxBytes=5 * 1024 * 1024, - backupCount=3, - encoding="utf-8", - ) - file_handler.setFormatter(RedactingFormatter("%(message)s")) - tool_logger = logging.getLogger(f"hermes.tool_calls.{id(log_queue)}") - tool_logger.setLevel(logging.INFO) - tool_logger.propagate = False - tool_logger.addHandler(file_handler) - try: - while True: - try: - tool_logger.info("%s", log_queue.get_nowait()) - except queue.Empty: - await asyncio.sleep(0.3) - except Exception as e: - logger.error("write_tool_log error: %s", e) - await asyncio.sleep(1) - except asyncio.CancelledError: - pass - finally: - # Drain remaining entries before closing so late tool calls - # from the final iteration aren't lost. - while True: - try: - tool_logger.info("%s", log_queue.get_nowait()) - except queue.Empty: - break - except Exception: - break - tool_logger.removeHandler(file_handler) + log_dir = _hermes_home / "logs" + log_dir.mkdir(parents=True, exist_ok=True) + file_handler = RotatingFileHandler( + log_dir / "tool_calls.log", + maxBytes=5 * 1024 * 1024, + backupCount=3, + encoding="utf-8", + ) + file_handler.setFormatter(RedactingFormatter("%(message)s")) + tool_logger = logging.getLogger(f"hermes.tool_calls.{id(log_queue)}") + tool_logger.setLevel(logging.INFO) + tool_logger.propagate = False + tool_logger.addHandler(file_handler) + try: + while True: try: - file_handler.flush() - file_handler.close() + tool_logger.info("%s", log_queue.get_nowait()) + except queue.Empty: + await asyncio.sleep(0.3) + except Exception as e: + logger.error("write_tool_log error: %s", e) + await asyncio.sleep(1) + except asyncio.CancelledError: + pass + finally: + # Drain remaining entries before closing so late tool calls + # from the final iteration aren't lost. + while True: + try: + tool_logger.info("%s", log_queue.get_nowait()) + except queue.Empty: + break except Exception: - pass + break + tool_logger.removeHandler(file_handler) + try: + file_handler.flush() + file_handler.close() + except Exception: + pass - # Extracted to TurnRunner.send_progress_messages; the threading metadata above is published - # onto the shared TurnContext where the original closure's captured locals were bound. - turn_ctx._progress_metadata = _progress_metadata - turn_ctx._progress_reply_to = _progress_reply_to - send_progress_messages = turn_runner.send_progress_messages - - # We need to share the agent instance for interrupt support - agent_holder = [None] # Mutable container for the agent instance - turn_ctx.agent_holder = agent_holder - result_holder = [None] # Mutable container for the result - tools_holder = [None] # Mutable container for the tool definitions - stream_consumer_holder = [None] # Mutable container for stream consumer - # streaming PCM audio consumer. Created on the gateway event-loop thread (NOT in run_sync's - # executor worker) so outer finalisation / interrupt paths can reference it without a NameError. - streaming_tts_consumer_holder: list = [None] - turn_ctx.result_holder = result_holder - turn_ctx.tools_holder = tools_holder - turn_ctx.stream_consumer_holder = stream_consumer_holder - turn_ctx.streaming_tts_consumer_holder = streaming_tts_consumer_holder - - # Bridge sync step_callback → async hooks.emit for agent:step events - _loop_for_step = asyncio.get_running_loop() - _hooks_ref = self.hooks - - # Bridge extracted to TurnRunner._step_callback_sync; the loop and - # hooks refs bound just above are published at their original site. - turn_ctx._loop_for_step = _loop_for_step - turn_ctx._hooks_ref = _hooks_ref - turn_ctx._step_callback_sync = turn_runner._step_callback_sync - - # Bridge sync event_callback → async hooks.emit for lifecycle events (e.g. session:compress - # after a compression split); extracted to TurnRunner._event_callback_sync. - turn_ctx._event_callback_sync = turn_runner._event_callback_sync - - # Bridge sync status_callback → async adapter.send for context pressure - _status_adapter = self._adapter_for_source(source) - _status_chat_id = source.chat_id + def _run_agent_status_thread_metadata( + self, + source: SessionSource, + event_message_id: Optional[str], + _progress_thread_id: Any, + _relay_prospective_thread_id: Optional[str], + ) -> Optional[Dict[str, Any]]: + """Thread metadata for status / approval / stream sends (Feishu carries the reply anchor).""" if source.platform == Platform.FEISHU and source.thread_id and event_message_id: # Feishu topics only keep messages inside the topic when they are sent via the reply API # with reply_in_thread=true. Status/approval/stream paths usually only get metadata, so @@ -27575,14 +27604,15 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew _status_thread_metadata = { "reply_to_message_id": event_message_id } + return _status_thread_metadata - # Bridge extracted to TurnRunner._status_callback_sync; publish the status wiring computed - # above onto the shared TurnContext at the exact original binding site. - turn_ctx._status_adapter = _status_adapter - turn_ctx._status_chat_id = _status_chat_id - turn_ctx._status_thread_metadata = _status_thread_metadata - turn_ctx._status_callback_sync = turn_runner._status_callback_sync - + def _run_agent_start_streaming_tts( + self, + source: SessionSource, + message_type: Optional[str], + _status_thread_metadata: Optional[Dict[str, Any]], + streaming_tts_consumer_holder: list, + ) -> None: # Streaming TTS consumer setup. Created on the gateway event-loop thread (here), NOT inside # run_sync's executor worker: the outer interrupt / finalisation paths reference the consumer # via ``streaming_tts_consumer_holder[0]`` and would hit a cross-scope NameError. @@ -27616,911 +27646,919 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew except Exception as _stts_err: logger.debug("Could not set up streaming TTS consumer: %s", _stts_err) - # run_sync extracted to TurnRunner.run_sync (bound method; executor call unchanged). Its - # closed-over locals travel on turn_ctx; `nonlocal message` rebinds became ctx.message writes. - run_sync = turn_runner.run_sync - - # Start the progress sender if enabled. Gate on needs_progress_queue (tool_progress OR - # thinking_progress), not tool_progress alone: the sender drains BOTH tool-progress lines and - # _thinking scratch bubbles — a tool_progress-only gate left thinking-only queues never drained. - progress_task = None - if needs_progress_queue: - progress_task = asyncio.create_task(send_progress_messages()) - - # Start the tool-call log writer when tool_progress == "log". - log_task = None - if log_mode_enabled: - log_task = asyncio.create_task(write_tool_log()) - - # Start stream consumer task — polls for consumer creation since it - # happens inside run_sync (thread pool) after the agent is constructed. - stream_task = None - - async def _start_stream_consumer(): - """Wait for the stream consumer to be created, then run it.""" - for _ in range(200): # Up to 10s wait - if stream_consumer_holder[0] is not None: - await stream_consumer_holder[0].run() - return - await asyncio.sleep(0.05) - - stream_task = asyncio.create_task(_start_stream_consumer()) - - # Track this agent as running for this session (for interrupt support) - # We do this in a callback after the agent is created - async def track_agent(): - # Wait for agent to be created - while agent_holder[0] is None: - await asyncio.sleep(0.05) - if not session_key: + async def _run_agent_stream_consumer_task(self, stream_consumer_holder: list) -> None: + """Wait for the stream consumer to be created, then run it.""" + for _ in range(200): # Up to 10s wait + if stream_consumer_holder[0] is not None: + await stream_consumer_holder[0].run() return - # Only promote the sentinel to the real agent if this run is still current. If /stop or - # /new bumped the generation while we were spinning up, leave the newer run's slot alone - # — we'll be discarded by the stale-result check in _handle_message_with_agent. - if run_generation is not None and not self._is_session_run_current( - session_key, run_generation - ): - logger.info( - "Skipping stale agent promotion for %s — generation %s is no longer current", - session_key or "", - run_generation, - ) - return - self._session_state(session_key).turn.agent = agent_holder[0] - if self._draining: - self._update_runtime_status("draining") + await asyncio.sleep(0.05) - tracking_task = asyncio.create_task(track_agent()) + async def _run_agent_track_agent( + self, + session_key: Optional[str], + run_generation: Optional[int], + agent_holder: list, + ) -> None: + """Track this agent as running for the session (interrupt support) once it is created.""" + # Wait for agent to be created + while agent_holder[0] is None: + await asyncio.sleep(0.05) + if not session_key: + return + # Only promote the sentinel to the real agent if this run is still current. If /stop or + # /new bumped the generation while we were spinning up, leave the newer run's slot alone + # — we'll be discarded by the stale-result check in _handle_message_with_agent. + if run_generation is not None and not self._is_session_run_current( + session_key, run_generation + ): + logger.info( + "Skipping stale agent promotion for %s — generation %s is no longer current", + session_key or "", + run_generation, + ) + return + self._session_state(session_key).turn.agent = agent_holder[0] + if self._draining: + self._update_runtime_status("draining") + async def _run_agent_monitor_for_interrupt( + self, + source: SessionSource, + session_key: Optional[str], + agent_holder: list, + _interrupt_detected: "asyncio.Event", + streaming_tts_consumer_holder: list, + ) -> None: # Monitor adapter interrupts (new messages). PRIMARY interrupt path for regular text: Level 1 # (base.py) catches them before _handle_message(), so the Level 2 running_agent.interrupt() path # never fires. The inactivity poll loop has a BACKUP check in case this task dies silently. - _interrupt_detected = asyncio.Event() # shared with backup check + if not session_key: + return - async def monitor_for_interrupt(): - if not session_key: - return - - while True: - await asyncio.sleep(0.2) # Check every 200ms - try: - # Re-resolve adapter each iteration so reconnects don't - # leave us holding a stale reference. - _adapter = self._adapter_for_source(source) - if not _adapter: - continue - # Must use session_key (build_session_key output), NOT source.chat_id: the adapter - # stores interrupt events under the full session key. - if hasattr(_adapter, 'has_pending_interrupt') and _adapter.has_pending_interrupt(session_key): - agent = agent_holder[0] - if agent: - # Peek WITHOUT consuming: the message must stay in _pending_messages for the - # post-run _dequeue_pending_event() (full MessageEvent + media). Popping here - # races: the agent may finish before checking _interrupt_requested, losing it. - _peek_event = _adapter._pending_messages.get(session_key) - pending_text = None - if _peek_event is not None: - pending_text = _peek_event.text or "" - # Transcribe audio BEFORE signaling the agent, so voice messages interrupt - # with the real transcript, not an empty string / file-path placeholder. - _media_urls = getattr(_peek_event, "media_urls", None) or [] - if self._pending_event_audio_paths(_peek_event): - pending_text, _ = await self._transcribe_and_echo_pending_voice( - _peek_event, - _adapter, - source, - pending_text, - log_context="Voice-interrupt", - metadata={"thread_id": source.thread_id} if source.thread_id else None, - ) - elif not pending_text and _media_urls: - pending_text = _build_media_placeholder(_peek_event) - logger.debug("Interrupt detected from adapter, signaling agent...") - agent.interrupt(pending_text) - _interrupt_detected.set() - # Abort streaming TTS on barge-in (#60671). - _stts = streaming_tts_consumer_holder[0] - if _stts is not None: - _stts.abort("barge-in") - break - except asyncio.CancelledError: - raise - except Exception as _mon_err: - logger.debug("monitor_for_interrupt error (will retry): %s", _mon_err) - - interrupt_monitor = asyncio.create_task(monitor_for_interrupt()) - - # Periodic "still working" notifications so the user knows the agent hasn't died. Config: - # agent.gateway_notify_interval or HERMES_AGENT_NOTIFY_INTERVAL env; default 180s. - _NOTIFY_INTERVAL_RAW = _float_env("HERMES_AGENT_NOTIFY_INTERVAL", 180) - _NOTIFY_INTERVAL = _NOTIFY_INTERVAL_RAW if _NOTIFY_INTERVAL_RAW > 0 else None - _long_running_mode = _display_surface_mode( - "long_running_notifications", - default=True, - allow_generic=True, - ) - if _long_running_mode == "off": - _NOTIFY_INTERVAL = None - _notify_start = time.time() - - async def _notify_long_running(): - if _NOTIFY_INTERVAL is None: - return # Notifications disabled (gateway_notify_interval: 0) - _notify_adapter = self._adapter_for_source(source) - if not _notify_adapter: - return - # Track the heartbeat message id to edit in place where supported (Telegram, Discord, - # Slack, ...) instead of a new "Still working" bubble every interval. - _heartbeat_msg_id: Optional[str] = None - while True: - await asyncio.sleep(_NOTIFY_INTERVAL) - # Stop heartbeating once this run no longer owns the session slot or the executor has - # finished, else a stale "running: delegate_task" bubble outlives its run. _executor_task - # is bound just after this task is scheduled; tolerate the brief window before then. - try: - _exec_ref = _executor_task - except NameError: - _exec_ref = None - if not self._should_emit_long_running_notification( - session_key, agent_holder[0], _exec_ref - ): - break - _elapsed_mins = int((time.time() - _notify_start) // 60) - # Default heartbeat is terse (elapsed + current tool); the verbose iteration counter is - # gated on busy_ack_detail so users can opt in per platform. - _agent_ref = agent_holder[0] - _status_detail = "" - _want_iteration_detail = bool( - resolve_display_setting( - user_config, - platform_key, - "busy_ack_detail", - True, - ) - ) - if _agent_ref and hasattr(_agent_ref, "get_activity_summary"): - try: - _a = _agent_ref.get_activity_summary() - _parts = [] - if _want_iteration_detail: - _parts.append( - f"iteration {_a['api_call_count']}/{_a['max_iterations']}" - ) - _action = _a.get("current_tool") or _a.get("last_activity_desc") - if _action: - _parts.append(str(_action)) - if _parts: - _status_detail = " — " + ", ".join(_parts) - except Exception: - pass - _heartbeat_text = ( - _generic_status_phrase("status") - if _long_running_mode == "generic" - else f"⏳ Working — {_elapsed_mins} min{_status_detail}" - ) - try: - _notify_res = None - if _heartbeat_msg_id: - try: - _notify_res = await _notify_adapter.edit_message( - source.chat_id, - _heartbeat_msg_id, - _heartbeat_text, - ) - except Exception as _ee: - logger.debug("Heartbeat edit failed: %s", _ee) - _notify_res = None - if not (_notify_res and getattr(_notify_res, "success", False)): - _notify_res = await _notify_adapter.send( - source.chat_id, - _heartbeat_text, - metadata=_interim_metadata(_non_conversational_metadata(_status_thread_metadata, platform=source.platform)), - ) - if getattr(_notify_res, "success", False) and getattr( - _notify_res, "message_id", None - ): - _heartbeat_msg_id = str(_notify_res.message_id) - if _cleanup_progress: - _cleanup_msg_ids.append(_heartbeat_msg_id) - except Exception as _ne: - logger.debug("Long-running notification error: %s", _ne) - - _notify_task = asyncio.create_task(_notify_long_running()) - - def _stream_confirmed_final_delivery( - consumer, - final_text: str, - *, - previewed: bool = False, - ) -> bool: - """Return True only when the actual final reply reached the user.""" - if consumer is None: - return False - if getattr(consumer, "final_response_sent", False): - # A successful finalize call is not proof the *content* was final: the edit may carry - # only the last preview snapshot. Reconcile against the recorded turn-final payload: - # only a demonstrable mismatch (False, incl. payload-less split delivery) overrides - # the flag; None keeps legacy trust so timeout dedup isn't regressed. - matcher = getattr(consumer, "delivered_final_matches", None) - if callable(matcher): - try: - if matcher(final_text) is False: - return False - except Exception: - pass - return True - if previewed: - has_delivered_text = getattr(consumer, "has_delivered_text", None) - if callable(has_delivered_text): - try: - return bool(has_delivered_text(final_text)) - except Exception: - return False - return False - - try: - # Thread pool so we don't block. *Inactivity* timeout, not wall-clock: the agent may run for - # hours while actively calling tools / streaming, but a hung API call or stuck tool is killed. - # agent.gateway_timeout / HERMES_AGENT_TIMEOUT (env wins); default 1800s; 0 = unlimited. - _agent_timeout_raw = _float_env("HERMES_AGENT_TIMEOUT", 1800) - _agent_timeout = _agent_timeout_raw if _agent_timeout_raw > 0 else None - _agent_warning_raw = _float_env("HERMES_AGENT_TIMEOUT_WARNING", 900) - _agent_warning = _agent_warning_raw if _agent_warning_raw > 0 else None - _warning_fired = False - - # A background=true process intentionally survives a successful turn, so capture - # existing IDs and reap only children created by THIS turn if it times out. The daemon - # watchdog is independent of asyncio: cgroup memory reclaim can starve the loop that - # runs the normal timeout poll, and cleanup must not wait for the loop to recover. - from tools.process_registry import process_registry - - _turn_task_id = session_id or "" - _turn_process_baseline = process_registry.snapshot_running_ids(_turn_task_id) - turn_ctx.process_task_id = _turn_task_id - turn_ctx.process_baseline = _turn_process_baseline - _turn_worker_done = threading.Event() - _turn_timeout_fired = threading.Event() - _turn_cleanup_lock = threading.Lock() - # task_id is session-scoped, not turn-scoped: gate the eventual reap on this exact claim still - # being current, so a replacement turn on the same session that starts before the watchdog - # fires doesn't get its own fresh process killed by this turn's stale baseline. - _turn_run_generation = run_generation - _turn_is_current = ( - (lambda: self._is_session_run_current(session_key, _turn_run_generation)) - if _turn_run_generation is not None - else (lambda: True) - ) - - def _run_sync_with_timeout_lifecycle(): - try: - return run_sync() - finally: - _turn_worker_done.set() - # `.turn.agent` is only reset to _AGENT_PENDING_SENTINEL when the *next* turn is - # claimed, so this agent stays reachable from _interrupt_and_clear_session() - # until then. Clearing ownership markers the instant our worker finishes means a - # /stop on the finished turn no longer reaps background work it left running. - _finished_agent = agent_holder[0] if agent_holder else None - if _finished_agent is not None: - _finished_agent._gateway_turn_process_task_id = "" - _finished_agent._gateway_turn_process_baseline = frozenset() - - if _agent_timeout is not None: - threading.Thread( - target=_watch_gateway_turn_inactivity, - kwargs={ - "agent_holder": agent_holder, - "task_id": _turn_task_id, - "process_baseline": _turn_process_baseline, - "timeout": _agent_timeout, - "worker_done": _turn_worker_done, - "timeout_fired": _turn_timeout_fired, - "cleanup_lock": _turn_cleanup_lock, - "poll_interval": 5.0, - "is_still_current": _turn_is_current, - }, - name=f"gateway-turn-watchdog-{_turn_task_id[:12]}", - daemon=True, - ).start() - _executor_task = asyncio.ensure_future( - self._run_in_executor_with_context(_run_sync_with_timeout_lifecycle) - ) - - _inactivity_timeout = False - _POLL_INTERVAL = 5.0 - - if _agent_timeout is None: - # Unlimited — still poll periodically for backup interrupt - # detection in case monitor_for_interrupt() silently died. - response = None - while True: - done, _ = await asyncio.wait( - {_executor_task}, timeout=_POLL_INTERVAL - ) - if done: - response = _executor_task.result() - break - # Backup interrupt check: if the monitor task died or - # missed the interrupt, catch it here. - if not _interrupt_detected.is_set() and session_key: - _backup_adapter = self._adapter_for_source(source) - _backup_agent = agent_holder[0] - if (_backup_adapter and _backup_agent - and hasattr(_backup_adapter, 'has_pending_interrupt') - and _backup_adapter.has_pending_interrupt(session_key)): - _bp_event = _backup_adapter._pending_messages.get(session_key) - _bp_text = _bp_event.text if _bp_event else None - if _bp_event is not None: - _bp_media_urls = getattr(_bp_event, "media_urls", None) or [] - if self._pending_event_audio_paths(_bp_event): - _bp_text, _ = await self._transcribe_and_echo_pending_voice( - _bp_event, - _backup_adapter, - source, - _bp_text or "", - log_context="Voice-backup-interrupt", - metadata={"thread_id": source.thread_id} if source.thread_id else None, - ) - elif not _bp_text and _bp_media_urls: - _bp_text = _build_media_placeholder(_bp_event) - logger.info( - "Backup interrupt detected for session %s " - "(monitor task state: %s)", - session_key, - "done" if interrupt_monitor.done() else "running", - ) - _backup_agent.interrupt(_bp_text) - _interrupt_detected.set() - # Abort streaming TTS on barge-in (#60671). - _stts = streaming_tts_consumer_holder[0] - if _stts is not None: - _stts.abort("barge-in") - - else: - # Poll the agent's built-in activity tracker (updated by _touch_activity() on every tool - # call, API call, and stream delta) every few seconds. - response = None - while True: - done, _ = await asyncio.wait( - {_executor_task}, timeout=_POLL_INTERVAL - ) - if done: - # Prefer the real result when the worker finished even if the watchdog fired in - # the same window: the completed run already persisted its reply, so the "agent - # inactive" diagnostic would contradict the stored transcript. - response = _executor_task.result() - break - if _turn_timeout_fired.is_set(): - _inactivity_timeout = True - break - # Agent still running — check inactivity. - _agent_ref = agent_holder[0] - _idle_secs = 0.0 - if _agent_ref and hasattr(_agent_ref, "get_activity_summary"): - try: - _act = _agent_ref.get_activity_summary() - _idle_secs = _act.get("seconds_since_activity", 0.0) - except Exception: - pass - # Staged warning: fire once before escalating to full timeout. - if (not _warning_fired and _agent_warning is not None - and _idle_secs >= _agent_warning): - _warning_fired = True - _warn_adapter = self._adapter_for_source(source) - if _warn_adapter: - _elapsed_warn = int(_agent_warning // 60) or 1 - _remaining_mins = int((_agent_timeout - _agent_warning) // 60) or 1 - try: - await _warn_adapter.send( - source.chat_id, - f"⚠️ No activity for {_elapsed_warn} min. " - f"If the agent does not respond soon, it will " - f"be timed out in {_remaining_mins} min. " - f"You can continue waiting or use /reset.", - metadata=_interim_metadata(_status_thread_metadata), + while True: + await asyncio.sleep(0.2) # Check every 200ms + try: + # Re-resolve adapter each iteration so reconnects don't + # leave us holding a stale reference. + _adapter = self._adapter_for_source(source) + if not _adapter: + continue + # Must use session_key (build_session_key output), NOT source.chat_id: the adapter + # stores interrupt events under the full session key. + if hasattr(_adapter, 'has_pending_interrupt') and _adapter.has_pending_interrupt(session_key): + agent = agent_holder[0] + if agent: + # Peek WITHOUT consuming: the message must stay in _pending_messages for the + # post-run _dequeue_pending_event() (full MessageEvent + media). Popping here + # races: the agent may finish before checking _interrupt_requested, losing it. + _peek_event = _adapter._pending_messages.get(session_key) + pending_text = None + if _peek_event is not None: + pending_text = _peek_event.text or "" + # Transcribe audio BEFORE signaling the agent, so voice messages interrupt + # with the real transcript, not an empty string / file-path placeholder. + _media_urls = getattr(_peek_event, "media_urls", None) or [] + if self._pending_event_audio_paths(_peek_event): + pending_text, _ = await self._transcribe_and_echo_pending_voice( + _peek_event, + _adapter, + source, + pending_text, + log_context="Voice-interrupt", + metadata={"thread_id": source.thread_id} if source.thread_id else None, ) - except Exception as _warn_err: - logger.debug("Inactivity warning send error: %s", _warn_err) - if _idle_secs >= _agent_timeout: - _inactivity_timeout = True - threading.Thread( - target=_abandon_timed_out_gateway_turn, - kwargs={ - "agent_holder": agent_holder, - "task_id": _turn_task_id, - "process_baseline": _turn_process_baseline, - "worker_done": _turn_worker_done, - "timeout_fired": _turn_timeout_fired, - "cleanup_lock": _turn_cleanup_lock, - "is_still_current": _turn_is_current, - }, - name=f"gateway-turn-reaper-{_turn_task_id[:12]}", - daemon=True, - ).start() + elif not pending_text and _media_urls: + pending_text = _build_media_placeholder(_peek_event) + logger.debug("Interrupt detected from adapter, signaling agent...") + agent.interrupt(pending_text) + _interrupt_detected.set() + # Abort streaming TTS on barge-in (#60671). + _stts = streaming_tts_consumer_holder[0] + if _stts is not None: + _stts.abort("barge-in") break - # Backup interrupt check (same as unlimited path). - if not _interrupt_detected.is_set() and session_key: - _backup_adapter = self._adapter_for_source(source) - _backup_agent = agent_holder[0] - if (_backup_adapter and _backup_agent - and hasattr(_backup_adapter, 'has_pending_interrupt') - and _backup_adapter.has_pending_interrupt(session_key)): - _bp_event = _backup_adapter._pending_messages.get(session_key) - _bp_text = _bp_event.text if _bp_event else None - if _bp_event is not None: - _bp_media_urls = getattr(_bp_event, "media_urls", None) or [] - if self._pending_event_audio_paths(_bp_event): - _bp_text, _ = await self._transcribe_and_echo_pending_voice( - _bp_event, - _backup_adapter, - source, - _bp_text or "", - log_context="Voice-backup-interrupt", - metadata={"thread_id": source.thread_id} if source.thread_id else None, - ) - elif not _bp_text and _bp_media_urls: - _bp_text = _build_media_placeholder(_bp_event) - logger.info( - "Backup interrupt detected for session %s " - "(monitor task state: %s)", - session_key, - "done" if interrupt_monitor.done() else "running", - ) - _backup_agent.interrupt(_bp_text) - _interrupt_detected.set() - # Abort streaming TTS on barge-in (#60671). - _stts = streaming_tts_consumer_holder[0] - if _stts is not None: - _stts.abort("barge-in") + except asyncio.CancelledError: + raise + except Exception as _mon_err: + logger.debug("monitor_for_interrupt error (will retry): %s", _mon_err) - if _inactivity_timeout: - # Build a diagnostic summary from the agent's activity tracker. - _timed_out_agent = agent_holder[0] - _activity = {} - if _timed_out_agent and hasattr(_timed_out_agent, "get_activity_summary"): - with suppress(Exception): - _activity = _timed_out_agent.get_activity_summary() - - _last_desc = _activity.get("last_activity_desc", "unknown") - _secs_ago = _activity.get("seconds_since_activity", 0) - _cur_tool = _activity.get("current_tool") - _iter_n = _activity.get("api_call_count", 0) - _iter_max = _activity.get("max_iterations", 0) - - logger.error( - "Agent idle for %.0fs (timeout %.0fs) in session %s " - "| last_activity=%s | iteration=%s/%s | tool=%s", - _secs_ago, _agent_timeout, session_key, - _last_desc, _iter_n, _iter_max, - _cur_tool or "none", - ) - - # Interrupt the agent if it's still running so the thread - # pool worker is freed. - if _timed_out_agent: - request_hard_interrupt(_timed_out_agent, _INTERRUPT_REASON_TIMEOUT) - - _timeout_mins = int(_agent_timeout // 60) or 1 - - # Construct a user-facing message with diagnostic context. - _diag_lines = [ - f"⏱️ Agent inactive for {_timeout_mins} min — no tool calls " - f"or API responses." - ] - if _cur_tool: - _diag_lines.append( - f"The agent appears stuck on tool `{_cur_tool}` " - f"({_secs_ago:.0f}s since last activity, " - f"iteration {_iter_n}/{_iter_max})." - ) - else: - _diag_lines.append( - f"Last activity: {_last_desc} ({_secs_ago:.0f}s ago, " - f"iteration {_iter_n}/{_iter_max}). " - "The agent may have been waiting on an API response." - ) - _diag_lines.append( - "To increase the limit, set agent.gateway_timeout in config.yaml " - "(value in seconds, 0 = no limit) and restart the gateway.\n" - "Try again, or use /reset to start fresh." - ) - - response = { - "final_response": "\n".join(_diag_lines), - "messages": result_holder[0].get("messages", []) if result_holder[0] else [], - "api_calls": _iter_n, - "tools": tools_holder[0] or [], - "history_offset": 0, - "failed": True, - } - - # Persist fallback-model switches so /model shows the actually-active model. Skip - # eviction when the run failed — evicting forces MCP reinit on the next message for no - # benefit (bad model → fallback → evict → recreate → same 400 loop burning CPU). - _agent = agent_holder[0] - _result_for_fb = result_holder[0] - _run_failed = _result_for_fb.get("failed") if _result_for_fb else False - if _agent is not None and hasattr(_agent, 'model') and not _run_failed: - _cfg_model = _resolve_gateway_model() - # Normalize _cfg_model as AIAgent.__init__ does so a vendor-prefixed config value - # matches the agent's stripped model on native providers — otherwise the cached agent - # is evicted every turn, destroying prompt caching. Aggregators keep the vendor slug. + @staticmethod + def _run_agent_stream_confirmed_final_delivery( + consumer, + final_text: str, + *, + previewed: bool = False, + ) -> bool: + """Return True only when the actual final reply reached the user.""" + if consumer is None: + return False + if getattr(consumer, "final_response_sent", False): + # A successful finalize call is not proof the *content* was final: the edit may carry + # only the last preview snapshot. Reconcile against the recorded turn-final payload: + # only a demonstrable mismatch (False, incl. payload-less split delivery) overrides + # the flag; None keeps legacy trust so timeout dedup isn't regressed. + matcher = getattr(consumer, "delivered_final_matches", None) + if callable(matcher): try: - from hermes_cli.model_normalize import ( - _AGGREGATOR_PROVIDERS, - normalize_model_for_provider, - ) - _agent_provider = getattr(_agent, 'provider', '') or '' - if _agent_provider and _agent_provider not in _AGGREGATOR_PROVIDERS: - _cfg_model = normalize_model_for_provider(_cfg_model, _agent_provider) + if matcher(final_text) is False: + return False except Exception: pass - if _agent.model != _cfg_model and not self._is_intentional_model_switch(session_key, _agent.model): - # Fallback activated on a successful run — evict cached - # agent so the next message retries the primary model. - self._evict_cached_agent(session_key) - - # Check if we were interrupted OR have a queued message (/queue). - result = result_holder[0] - adapter = self._adapter_for_source(source) - - # Finalize the streaming-TTS consumer. finish() runs on the outer event-loop thread so - # early returns from run_sync are also finalised. wait_complete() drains queued audio; - # on timeout abort unconditionally — if audio was audible keep suppression (no replay - # from the start); if not, the whole-file fallback is permitted. - _stts = streaming_tts_consumer_holder[0] - if _stts is not None: - _stts.finish() + return True + if previewed: + has_delivered_text = getattr(consumer, "has_delivered_text", None) + if callable(has_delivered_text): try: - await _stts.wait_complete(timeout=10.0) - except Exception as _stts_done_err: - logger.debug("streaming TTS wait_complete error: %s", _stts_done_err) - if not _stts.done: - # Timeout before or after audible audio: abort to free the consumer task. Audible - # streams retain suppression; silent streams stay eligible for whole-file fallback. - _stts.abort("streaming TTS finalisation timeout") - await _stts.wait_complete(timeout=2.0) - if _stts.suppress_whole_file and adapter is not None: - _mark_turn = getattr(adapter, "_mark_streaming_tts_completed_turn", None) - if callable(_mark_turn): - _mark_turn(session_key, run_generation) + return bool(has_delivered_text(final_text)) + except Exception: + return False + return False - # Get pending message from adapter. - # Use session_key (not source.chat_id) to match adapter's storage keys. - pending_event = None - pending = None - if result and adapter and session_key: - pending_event = _dequeue_pending_event(adapter, session_key) - # /queue overflow: after consuming the adapter's "next-up" slot, promote the next - # queued event into it so the recursive run's drain will see it. Keeping the slot - # occupied for the whole FIFO chain preserves order and makes a mid-chain /queue - # route to overflow instead of jumping the queue. - pending_event = self._promote_queued_event(session_key, adapter, pending_event) - if result.get("interrupted") and not pending_event and result.get("interrupt_message"): - interrupt_message = result.get("interrupt_message") - if _is_control_interrupt_message(interrupt_message): - logger.info( - "Ignoring control interrupt message for session %s: %s", - session_key or "?", - interrupt_message, - ) - else: - pending = interrupt_message - elif pending_event: - # Transcribe audio on the dequeued event BEFORE it becomes the next user turn, so - # queued/interrupting voice messages drain with the real transcript, not a file path. - _pending_text = pending_event.text or "" - _media_urls = getattr(pending_event, "media_urls", None) or [] - if self._pending_event_audio_paths(pending_event): - pending, _ = await self._transcribe_and_echo_pending_voice( - pending_event, - adapter, + def _run_agent_start_turn_worker( + self, + turn_ctx: TurnContext, + run_sync: Callable[[], Any], + agent_holder: list, + session_id: str, + session_key: Optional[str], + run_generation: Optional[int], + ) -> "GatewayRunner._RunAgentWorker": + """Schedule ``run_sync`` on the executor plus the inactivity watchdog thread.""" + # Thread pool so we don't block. *Inactivity* timeout, not wall-clock: the agent may run for + # hours while actively calling tools / streaming, but a hung API call or stuck tool is killed. + # agent.gateway_timeout / HERMES_AGENT_TIMEOUT (env wins); default 1800s; 0 = unlimited. + _agent_timeout_raw = _float_env("HERMES_AGENT_TIMEOUT", 1800) + _agent_timeout = _agent_timeout_raw if _agent_timeout_raw > 0 else None + _agent_warning_raw = _float_env("HERMES_AGENT_TIMEOUT_WARNING", 900) + _agent_warning = _agent_warning_raw if _agent_warning_raw > 0 else None + + # A background=true process intentionally survives a successful turn, so capture + # existing IDs and reap only children created by THIS turn if it times out. The daemon + # watchdog is independent of asyncio: cgroup memory reclaim can starve the loop that + # runs the normal timeout poll, and cleanup must not wait for the loop to recover. + from tools.process_registry import process_registry + + _turn_task_id = session_id or "" + _turn_process_baseline = process_registry.snapshot_running_ids(_turn_task_id) + turn_ctx.process_task_id = _turn_task_id + turn_ctx.process_baseline = _turn_process_baseline + _turn_worker_done = threading.Event() + _turn_timeout_fired = threading.Event() + _turn_cleanup_lock = threading.Lock() + # task_id is session-scoped, not turn-scoped: gate the eventual reap on this exact claim still + # being current, so a replacement turn on the same session that starts before the watchdog + # fires doesn't get its own fresh process killed by this turn's stale baseline. + _turn_run_generation = run_generation + _turn_is_current = ( + (lambda: self._is_session_run_current(session_key, _turn_run_generation)) + if _turn_run_generation is not None + else (lambda: True) + ) + + def _run_sync_with_timeout_lifecycle(): + try: + return run_sync() + finally: + _turn_worker_done.set() + # `.turn.agent` is only reset to _AGENT_PENDING_SENTINEL when the *next* turn is + # claimed, so this agent stays reachable from _interrupt_and_clear_session() + # until then. Clearing ownership markers the instant our worker finishes means a + # /stop on the finished turn no longer reaps background work it left running. + _finished_agent = agent_holder[0] if agent_holder else None + if _finished_agent is not None: + _finished_agent._gateway_turn_process_task_id = "" + _finished_agent._gateway_turn_process_baseline = frozenset() + + if _agent_timeout is not None: + threading.Thread( + target=_watch_gateway_turn_inactivity, + kwargs={ + "agent_holder": agent_holder, + "task_id": _turn_task_id, + "process_baseline": _turn_process_baseline, + "timeout": _agent_timeout, + "worker_done": _turn_worker_done, + "timeout_fired": _turn_timeout_fired, + "cleanup_lock": _turn_cleanup_lock, + "poll_interval": 5.0, + "is_still_current": _turn_is_current, + }, + name=f"gateway-turn-watchdog-{_turn_task_id[:12]}", + daemon=True, + ).start() + _executor_task = asyncio.ensure_future( + self._run_in_executor_with_context(_run_sync_with_timeout_lifecycle) + ) + return self._RunAgentWorker( + executor_task=_executor_task, + agent_timeout=_agent_timeout, + agent_warning=_agent_warning, + task_id=_turn_task_id, + process_baseline=_turn_process_baseline, + worker_done=_turn_worker_done, + timeout_fired=_turn_timeout_fired, + cleanup_lock=_turn_cleanup_lock, + is_current=_turn_is_current, + ) + + async def _run_agent_backup_interrupt_check( + self, + source: SessionSource, + session_key: Optional[str], + agent_holder: list, + _interrupt_detected: "asyncio.Event", + interrupt_monitor: "asyncio.Task", + streaming_tts_consumer_holder: list, + ) -> None: + """Backup interrupt check: if the monitor task died or missed the interrupt, catch it here.""" + if not _interrupt_detected.is_set() and session_key: + _backup_adapter = self._adapter_for_source(source) + _backup_agent = agent_holder[0] + if (_backup_adapter and _backup_agent + and hasattr(_backup_adapter, 'has_pending_interrupt') + and _backup_adapter.has_pending_interrupt(session_key)): + _bp_event = _backup_adapter._pending_messages.get(session_key) + _bp_text = _bp_event.text if _bp_event else None + if _bp_event is not None: + _bp_media_urls = getattr(_bp_event, "media_urls", None) or [] + if self._pending_event_audio_paths(_bp_event): + _bp_text, _ = await self._transcribe_and_echo_pending_voice( + _bp_event, + _backup_adapter, source, - _pending_text, - log_context="Voice-drain", + _bp_text or "", + log_context="Voice-backup-interrupt", metadata={"thread_id": source.thread_id} if source.thread_id else None, ) - if not pending: - pending = _build_media_placeholder(pending_event) - else: - pending = _pending_text or _build_media_placeholder(pending_event) - if pending: - logger.debug("Processing queued message after agent completion: '%s...'", pending[:40]) + elif not _bp_text and _bp_media_urls: + _bp_text = _build_media_placeholder(_bp_event) + logger.info( + "Backup interrupt detected for session %s " + "(monitor task state: %s)", + session_key, + "done" if interrupt_monitor.done() else "running", + ) + _backup_agent.interrupt(_bp_text) + _interrupt_detected.set() + # Abort streaming TTS on barge-in (#60671). + _stts = streaming_tts_consumer_holder[0] + if _stts is not None: + _stts.abort("barge-in") - # Leftover /steer: a steer arriving after the last tool batch (e.g. during the final API - # call) comes back in result["pending_steer"]; deliver it as the next user turn, not drop it. - if result and not pending and not pending_event: - _leftover_steer = result.get("pending_steer") - if _leftover_steer: - pending = _leftover_steer - logger.debug("Delivering leftover /steer as next turn: '%s...'", pending[:40]) + async def _run_agent_await_turn_worker( + self, + worker: "GatewayRunner._RunAgentWorker", + *, + source: SessionSource, + session_key: Optional[str], + agent_holder: list, + result_holder: list, + tools_holder: list, + _interrupt_detected: "asyncio.Event", + interrupt_monitor: "asyncio.Task", + streaming_tts_consumer_holder: list, + _status_thread_metadata: Optional[Dict[str, Any]], + ) -> Any: + """Poll the executor future (inactivity timeout + backup interrupt checks); return its result. - # Safety net: if the pending text is a slash command (e.g. "/stop", "/new"), discard it - # — commands should never be passed to the agent as user input. - if pending and pending.strip().startswith("/"): - _pending_parts = pending.strip().split(None, 1) - _pending_cmd_word = _pending_parts[0][1:].lower() if _pending_parts else "" - if _pending_cmd_word: + On inactivity timeout the result is a synthetic failed run dict carrying the diagnostic. + """ + _warning_fired = False + _inactivity_timeout = False + _POLL_INTERVAL = 5.0 + + if worker.agent_timeout is None: + # Unlimited — still poll periodically for backup interrupt + # detection in case monitor_for_interrupt() silently died. + response = None + while True: + done, _ = await asyncio.wait( + {worker.executor_task}, timeout=_POLL_INTERVAL + ) + if done: + response = worker.executor_task.result() + break + # Backup interrupt check: if the monitor task died or + # missed the interrupt, catch it here. + await self._run_agent_backup_interrupt_check( + source, + session_key, + agent_holder, + _interrupt_detected, + interrupt_monitor, + streaming_tts_consumer_holder, + ) + + else: + # Poll the agent's built-in activity tracker (updated by _touch_activity() on every tool + # call, API call, and stream delta) every few seconds. + response = None + while True: + done, _ = await asyncio.wait( + {worker.executor_task}, timeout=_POLL_INTERVAL + ) + if done: + # Prefer the real result when the worker finished even if the watchdog fired in + # the same window: the completed run already persisted its reply, so the "agent + # inactive" diagnostic would contradict the stored transcript. + response = worker.executor_task.result() + break + if worker.timeout_fired.is_set(): + _inactivity_timeout = True + break + # Agent still running — check inactivity. + _agent_ref = agent_holder[0] + _idle_secs = 0.0 + if _agent_ref and hasattr(_agent_ref, "get_activity_summary"): try: - from hermes_cli.commands import resolve_command as _rc_pending - if _rc_pending(_pending_cmd_word): - logger.info( - "Discarding command '/%s' from pending queue — " - "commands must not be passed as agent input", - _pending_cmd_word, - ) - pending_event = None - pending = None + _act = _agent_ref.get_activity_summary() + _idle_secs = _act.get("seconds_since_activity", 0.0) except Exception: pass - - if self._draining and (pending_event or pending): - logger.info( - "Discarding pending follow-up for session %s during gateway %s", - session_key or "?", - self._status_action_label(), - ) - pending_event = None - pending = None - - if pending_event or pending: - logger.debug("Processing pending message: '%s...'", pending[:40]) - - # Clear the adapter's interrupt event so the next _run_agent call doesn't re-trigger the - # interrupt before the new agent's first API call (infinite loop otherwise). - if adapter and hasattr(adapter, '_active_sessions') and session_key and session_key in adapter._active_sessions: - adapter._active_sessions[session_key].clear() - - # Cap recursion depth to prevent resource exhaustion when the - # user sends multiple messages while the agent keeps failing. (#816) - if _interrupt_depth >= self._MAX_INTERRUPT_DEPTH: - logger.warning( - "Interrupt recursion depth %d reached for session %s — " - "queueing message instead of recursing.", - _interrupt_depth, session_key, - ) - adapter = self._adapter_for_source(source) - if adapter and pending_event: - merge_pending_message_event(adapter._pending_messages, session_key, pending_event) - elif adapter and hasattr(adapter, 'queue_message'): - adapter.queue_message(session_key, pending) - return result_holder[0] or {"final_response": response, "messages": history} - - was_interrupted = result.get("interrupted") - if not was_interrupted: - # Queued message after normal completion: deliver the first response before the - # queued follow-up, unless streaming already delivered it. - _sc = stream_consumer_holder[0] - if _sc and stream_task: + # Staged warning: fire once before escalating to full timeout. + if (not _warning_fired and worker.agent_warning is not None + and _idle_secs >= worker.agent_warning): + _warning_fired = True + _warn_adapter = self._adapter_for_source(source) + if _warn_adapter: + _elapsed_warn = int(worker.agent_warning // 60) or 1 + _remaining_mins = int((worker.agent_timeout - worker.agent_warning) // 60) or 1 try: - await asyncio.wait_for(stream_task, timeout=5.0) - except (asyncio.TimeoutError, asyncio.CancelledError): - stream_task.cancel() - with suppress(asyncio.CancelledError): - await stream_task - except Exception as e: - logger.debug("Stream consumer wait before queued message failed: %s", e) - # The queued branch needs raw ``result`` for interruption, history, and - # recursion state, but delivery must use the finalized task result — it carries - # empty/failure normalization and final-response processing from _run_agent_task. - _delivery_result = response if isinstance(response, dict) else (result or {}) - _previewed = bool(_delivery_result.get("response_previewed")) - first_response = _delivery_result.get("final_response", "") - _already_streamed = _stream_confirmed_final_delivery( - _sc, - first_response, - previewed=_previewed, - ) - # Same predicate as the normal completed-turn path: this direct queued-send branch - # predates intentional-silence filtering and would leak the literal marker. - try: - from gateway.response_filters import is_intentional_silence_agent_result - _intentional_silence = is_intentional_silence_agent_result( - _delivery_result, first_response, - ) - except Exception: - _intentional_silence = False - if _intentional_silence: - logger.info( - "Queued follow-up for session %s: suppressing intentional silence marker before continuing.", - session_key or "?", - ) - elif first_response: - try: - if _already_streamed: - logger.info( - "Queued follow-up for session %s: final text delivery confirmed; delivering explicit media before continuing.", - session_key or "?", - ) - else: - logger.info( - "Queued follow-up for session %s: final stream delivery not confirmed; sending first response before continuing.", - session_key or "?", - ) - await self._deliver_queued_first_response( - first_response, - source=source, - adapter=adapter, - metadata=_status_thread_metadata, - event_message_id=event_message_id, - text_already_delivered=_already_streamed, - deliver_media=not _delivery_result.get("failed"), - stream_consumer=_sc, + await _warn_adapter.send( + source.chat_id, + f"⚠️ No activity for {_elapsed_warn} min. " + f"If the agent does not respond soon, it will " + f"be timed out in {_remaining_mins} min. " + f"You can continue waiting or use /reset.", + metadata=_interim_metadata(_status_thread_metadata), ) - except Exception as e: - logger.warning("Failed to send first response before queued message: %s", e) - # Release deferred bg-review notifications now that the first response is delivered: - # pop from the adapter's callback dict (no double-fire in base.py's finally) and call. - if getattr(type(adapter), "pop_post_delivery_callback", None) is not None: - _bg_cb = adapter.pop_post_delivery_callback( - session_key, - generation=run_generation, - ) - if callable(_bg_cb): - try: - _bg_result = _bg_cb() - if inspect.isawaitable(_bg_result): - await _bg_result - except Exception: - pass - elif adapter and hasattr(adapter, "_post_delivery_callbacks"): - _bg_cb = adapter._post_delivery_callbacks.pop(session_key, None) - if callable(_bg_cb): - try: - _bg_result = _bg_cb() - if inspect.isawaitable(_bg_result): - await _bg_result - except Exception: - pass - # else: interrupted — discard the response ("Operation interrupted." is noise; the user - # knows they sent a new message). + except Exception as _warn_err: + logger.debug("Inactivity warning send error: %s", _warn_err) + if _idle_secs >= worker.agent_timeout: + _inactivity_timeout = True + threading.Thread( + target=_abandon_timed_out_gateway_turn, + kwargs={ + "agent_holder": agent_holder, + "task_id": worker.task_id, + "process_baseline": worker.process_baseline, + "worker_done": worker.worker_done, + "timeout_fired": worker.timeout_fired, + "cleanup_lock": worker.cleanup_lock, + "is_still_current": worker.is_current, + }, + name=f"gateway-turn-reaper-{worker.task_id[:12]}", + daemon=True, + ).start() + break + # Backup interrupt check (same as unlimited path). + await self._run_agent_backup_interrupt_check( + source, + session_key, + agent_holder, + _interrupt_detected, + interrupt_monitor, + streaming_tts_consumer_holder, + ) - updated_history = result.get("messages", history) - next_source = source - next_message = pending - next_message_id = None - next_channel_prompt = None - next_session_key = session_key - # Carry the pending event's message_type into the recursive call so queued voice turns - # can stream TTS and re-mark the generation for the final delivered turn. - next_message_type = None - if pending_event is not None: - next_source = getattr(pending_event, "source", None) or source - if self._is_goal_continuation_event(pending_event) and not self._goal_still_active_for_session(session_id): - logger.info( - "Discarding stale goal continuation for session %s — goal is no longer active", - session_key or "?", - ) - return result - # Resolve the follow-up's session key BEFORE preparing the inbound text: - # _prepare_inbound_message_text buffers native image paths under the key given, and - # the recursive _run_agent consumes them under next_session_key — mismatch drops them. - try: - next_session_key = self._session_key_for_source(next_source) - except Exception: - logger.debug( - "Queued follow-up session-key resolution failed; reusing %s", - session_key or "?", - exc_info=True, - ) - next_message = await self._prepare_profile_scoped_inbound_message_text( - event=pending_event, - source=next_source, - history=updated_history, - session_key=next_session_key, + if _inactivity_timeout: + # Build a diagnostic summary from the agent's activity tracker. + _timed_out_agent = agent_holder[0] + _activity = {} + if _timed_out_agent and hasattr(_timed_out_agent, "get_activity_summary"): + with suppress(Exception): + _activity = _timed_out_agent.get_activity_summary() + + _last_desc = _activity.get("last_activity_desc", "unknown") + _secs_ago = _activity.get("seconds_since_activity", 0) + _cur_tool = _activity.get("current_tool") + _iter_n = _activity.get("api_call_count", 0) + _iter_max = _activity.get("max_iterations", 0) + + logger.error( + "Agent idle for %.0fs (timeout %.0fs) in session %s " + "| last_activity=%s | iteration=%s/%s | tool=%s", + _secs_ago, worker.agent_timeout, session_key, + _last_desc, _iter_n, _iter_max, + _cur_tool or "none", + ) + + # Interrupt the agent if it's still running so the thread + # pool worker is freed. + if _timed_out_agent: + request_hard_interrupt(_timed_out_agent, _INTERRUPT_REASON_TIMEOUT) + + _timeout_mins = int(worker.agent_timeout // 60) or 1 + + # Construct a user-facing message with diagnostic context. + _diag_lines = [ + f"⏱️ Agent inactive for {_timeout_mins} min — no tool calls " + f"or API responses." + ] + if _cur_tool: + _diag_lines.append( + f"The agent appears stuck on tool `{_cur_tool}` " + f"({_secs_ago:.0f}s since last activity, " + f"iteration {_iter_n}/{_iter_max})." + ) + else: + _diag_lines.append( + f"Last activity: {_last_desc} ({_secs_ago:.0f}s ago, " + f"iteration {_iter_n}/{_iter_max}). " + "The agent may have been waiting on an API response." + ) + _diag_lines.append( + "To increase the limit, set agent.gateway_timeout in config.yaml " + "(value in seconds, 0 = no limit) and restart the gateway.\n" + "Try again, or use /reset to start fresh." + ) + + response = { + "final_response": "\n".join(_diag_lines), + "messages": result_holder[0].get("messages", []) if result_holder[0] else [], + "api_calls": _iter_n, + "tools": tools_holder[0] or [], + "history_offset": 0, + "failed": True, + } + return response + + def _run_agent_evict_on_fallback( + self, session_key: Optional[str], agent_holder: list, result_holder: list, + ) -> None: + # Persist fallback-model switches so /model shows the actually-active model. Skip + # eviction when the run failed — evicting forces MCP reinit on the next message for no + # benefit (bad model → fallback → evict → recreate → same 400 loop burning CPU). + _agent = agent_holder[0] + _result_for_fb = result_holder[0] + _run_failed = _result_for_fb.get("failed") if _result_for_fb else False + if _agent is not None and hasattr(_agent, 'model') and not _run_failed: + _cfg_model = _resolve_gateway_model() + # Normalize _cfg_model as AIAgent.__init__ does so a vendor-prefixed config value + # matches the agent's stripped model on native providers — otherwise the cached agent + # is evicted every turn, destroying prompt caching. Aggregators keep the vendor slug. + try: + from hermes_cli.model_normalize import ( + _AGGREGATOR_PROVIDERS, + normalize_model_for_provider, + ) + _agent_provider = getattr(_agent, 'provider', '') or '' + if _agent_provider and _agent_provider not in _AGGREGATOR_PROVIDERS: + _cfg_model = normalize_model_for_provider(_cfg_model, _agent_provider) + except Exception: + pass + if _agent.model != _cfg_model and not self._is_intentional_model_switch(session_key, _agent.model): + # Fallback activated on a successful run — evict cached + # agent so the next message retries the primary model. + self._evict_cached_agent(session_key) + + async def _run_agent_finalize_streaming_tts( + self, + streaming_tts_consumer_holder: list, + adapter: Any, + session_key: Optional[str], + run_generation: Optional[int], + ) -> None: + # Finalize the streaming-TTS consumer. finish() runs on the outer event-loop thread so + # early returns from run_sync are also finalised. wait_complete() drains queued audio; + # on timeout abort unconditionally — if audio was audible keep suppression (no replay + # from the start); if not, the whole-file fallback is permitted. + _stts = streaming_tts_consumer_holder[0] + if _stts is not None: + _stts.finish() + try: + await _stts.wait_complete(timeout=10.0) + except Exception as _stts_done_err: + logger.debug("streaming TTS wait_complete error: %s", _stts_done_err) + if not _stts.done: + # Timeout before or after audible audio: abort to free the consumer task. Audible + # streams retain suppression; silent streams stay eligible for whole-file fallback. + _stts.abort("streaming TTS finalisation timeout") + await _stts.wait_complete(timeout=2.0) + if _stts.suppress_whole_file and adapter is not None: + _mark_turn = getattr(adapter, "_mark_streaming_tts_completed_turn", None) + if callable(_mark_turn): + _mark_turn(session_key, run_generation) + + async def _run_agent_drain_pending( + self, + result: Any, + adapter: Any, + source: SessionSource, + session_key: Optional[str], + ) -> Tuple[Any, Optional[str]]: + """Dequeue the adapter's pending / interrupt / leftover-steer follow-up. + + Returns ``(pending_event, pending)``. + """ + # Get pending message from adapter. + # Use session_key (not source.chat_id) to match adapter's storage keys. + pending_event = None + pending = None + if result and adapter and session_key: + pending_event = _dequeue_pending_event(adapter, session_key) + # /queue overflow: after consuming the adapter's "next-up" slot, promote the next + # queued event into it so the recursive run's drain will see it. Keeping the slot + # occupied for the whole FIFO chain preserves order and makes a mid-chain /queue + # route to overflow instead of jumping the queue. + pending_event = self._promote_queued_event(session_key, adapter, pending_event) + if result.get("interrupted") and not pending_event and result.get("interrupt_message"): + interrupt_message = result.get("interrupt_message") + if _is_control_interrupt_message(interrupt_message): + logger.info( + "Ignoring control interrupt message for session %s: %s", + session_key or "?", + interrupt_message, ) - if next_message is None: - return result - next_message_id = self._reply_anchor_for_event(pending_event) - next_channel_prompt = getattr(pending_event, "channel_prompt", None) - next_message_type = getattr(pending_event, "message_type", None) + else: + pending = interrupt_message + elif pending_event: + # Transcribe audio on the dequeued event BEFORE it becomes the next user turn, so + # queued/interrupting voice messages drain with the real transcript, not a file path. + _pending_text = pending_event.text or "" + _media_urls = getattr(pending_event, "media_urls", None) or [] + if self._pending_event_audio_paths(pending_event): + pending, _ = await self._transcribe_and_echo_pending_voice( + pending_event, + adapter, + source, + _pending_text, + log_context="Voice-drain", + metadata={"thread_id": source.thread_id} if source.thread_id else None, + ) + if not pending: + pending = _build_media_placeholder(pending_event) + else: + pending = _pending_text or _build_media_placeholder(pending_event) + if pending: + logger.debug("Processing queued message after agent completion: '%s...'", pending[:40]) - # Clear the prior logical turn's completed streaming marker so the recursive turn's - # streaming TTS isn't suppressed by that completion. - _clear_adapter = self._adapter_for_source(source) - if _clear_adapter is not None and session_key and run_generation is not None: - _completed_turns = getattr(_clear_adapter, "_streaming_tts_completed_turns", None) - if _completed_turns is not None: - _prior_key = getattr(_clear_adapter, "_streaming_tts_turn_key", None) - if callable(_prior_key): - _pk = _prior_key(session_key, run_generation) - if _pk: - _completed_turns.discard(_pk) + # Leftover /steer: a steer arriving after the last tool batch (e.g. during the final API + # call) comes back in result["pending_steer"]; deliver it as the next user turn, not drop it. + if result and not pending and not pending_event: + _leftover_steer = result.get("pending_steer") + if _leftover_steer: + pending = _leftover_steer + logger.debug("Delivering leftover /steer as next turn: '%s...'", pending[:40]) - # Restart the typing indicator for the follow-up turn; the outer - # _process_message_background typing task is alive but may be stale. - _followup_adapter = self._adapter_for_source(source) - if _followup_adapter: - with suppress(Exception): - await _followup_adapter.send_typing( - source.chat_id, - metadata=_status_thread_metadata, + # Safety net: if the pending text is a slash command (e.g. "/stop", "/new"), discard it + # — commands should never be passed to the agent as user input. + if pending and pending.strip().startswith("/"): + _pending_parts = pending.strip().split(None, 1) + _pending_cmd_word = _pending_parts[0][1:].lower() if _pending_parts else "" + if _pending_cmd_word: + try: + from hermes_cli.commands import resolve_command as _rc_pending + if _rc_pending(_pending_cmd_word): + logger.info( + "Discarding command '/%s' from pending queue — " + "commands must not be passed as agent input", + _pending_cmd_word, ) + pending_event = None + pending = None + except Exception: + pass - # Re-baseline the cached agent's message_count before recursing into the /queue follow-up: - # the coherence guard would otherwise rebuild on OUR OWN flushed rows and destroy the - # prompt-cache prefix; _handle_message_with_agent re-baselines only after the chain ends. - await self._refresh_agent_cache_message_count(session_key, session_id) + if self._draining and (pending_event or pending): + logger.info( + "Discarding pending follow-up for session %s during gateway %s", + session_key or "?", + self._status_action_label(), + ) + pending_event = None + pending = None + return pending_event, pending - followup_result = await self._run_agent( - message=next_message, - context_prompt=context_prompt, - history=updated_history, - source=next_source, - session_id=session_id, - session_key=next_session_key, - run_generation=run_generation, - _interrupt_depth=_interrupt_depth + 1, - event_message_id=next_message_id, - channel_prompt=next_channel_prompt, - message_type=next_message_type, + async def _run_agent_deliver_first_response( + self, + *, + source: SessionSource, + adapter: Any, + session_key: Optional[str], + run_generation: Optional[int], + event_message_id: Optional[str], + response: Any, + result: Any, + stream_consumer_holder: list, + stream_task: Any, + _status_thread_metadata: Optional[Dict[str, Any]], + ) -> None: + # Queued message after normal completion: deliver the first response before the + # queued follow-up, unless streaming already delivered it. + _sc = stream_consumer_holder[0] + if _sc and stream_task: + try: + await asyncio.wait_for(stream_task, timeout=5.0) + except (asyncio.TimeoutError, asyncio.CancelledError): + stream_task.cancel() + with suppress(asyncio.CancelledError): + await stream_task + except Exception as e: + logger.debug("Stream consumer wait before queued message failed: %s", e) + # The queued branch needs raw ``result`` for interruption, history, and + # recursion state, but delivery must use the finalized task result — it carries + # empty/failure normalization and final-response processing from _run_agent_task. + _delivery_result = response if isinstance(response, dict) else (result or {}) + _previewed = bool(_delivery_result.get("response_previewed")) + first_response = _delivery_result.get("final_response", "") + _already_streamed = self._run_agent_stream_confirmed_final_delivery( + _sc, + first_response, + previewed=_previewed, + ) + # Same predicate as the normal completed-turn path: this direct queued-send branch + # predates intentional-silence filtering and would leak the literal marker. + try: + from gateway.response_filters import is_intentional_silence_agent_result + _intentional_silence = is_intentional_silence_agent_result( + _delivery_result, first_response, + ) + except Exception: + _intentional_silence = False + if _intentional_silence: + logger.info( + "Queued follow-up for session %s: suppressing intentional silence marker before continuing.", + session_key or "?", + ) + elif first_response: + try: + if _already_streamed: + logger.info( + "Queued follow-up for session %s: final text delivery confirmed; delivering explicit media before continuing.", + session_key or "?", + ) + else: + logger.info( + "Queued follow-up for session %s: final stream delivery not confirmed; sending first response before continuing.", + session_key or "?", + ) + await self._deliver_queued_first_response( + first_response, + source=source, + adapter=adapter, + metadata=_status_thread_metadata, + event_message_id=event_message_id, + text_already_delivered=_already_streamed, + deliver_media=not _delivery_result.get("failed"), + stream_consumer=_sc, ) - return _preserve_queued_followup_history_offset(result, followup_result) - finally: - # Stop progress sender, interrupt monitor, and notification task - if progress_task: - progress_task.cancel() - if log_task: - log_task.cancel() - interrupt_monitor.cancel() - _notify_task.cancel() + except Exception as e: + logger.warning("Failed to send first response before queued message: %s", e) + # Release deferred bg-review notifications now that the first response is delivered: + # pop from the adapter's callback dict (no double-fire in base.py's finally) and call. + if getattr(type(adapter), "pop_post_delivery_callback", None) is not None: + _bg_cb = adapter.pop_post_delivery_callback( + session_key, + generation=run_generation, + ) + if callable(_bg_cb): + try: + _bg_result = _bg_cb() + if inspect.isawaitable(_bg_result): + await _bg_result + except Exception: + pass + elif adapter and hasattr(adapter, "_post_delivery_callbacks"): + _bg_cb = adapter._post_delivery_callbacks.pop(session_key, None) + if callable(_bg_cb): + try: + _bg_result = _bg_cb() + if inspect.isawaitable(_bg_result): + await _bg_result + except Exception: + pass - # Wait for stream consumer to finish its final edit - if stream_task: - # If the agent never created a stream consumer (non-streaming path, or a test stub - # returning synchronously) there is nothing to flush — cancel now instead of waiting - # out the 5s timeout polling for a consumer that will never arrive. - _has_stream_consumer = ( - stream_consumer_holder - and stream_consumer_holder[0] is not None + async def _run_agent_queued_followup( + self, + *, + source: SessionSource, + adapter: Any, + session_id: str, + session_key: Optional[str], + run_generation: Optional[int], + _interrupt_depth: int, + event_message_id: Optional[str], + context_prompt: str, + history: List[Dict[str, Any]], + pending: Optional[str], + pending_event: Any, + response: Any, + result: Any, + result_holder: list, + stream_consumer_holder: list, + stream_task: Any, + _status_thread_metadata: Optional[Dict[str, Any]], + ) -> Any: + """Run the queued / interrupting follow-up as the next turn (recursive ``_run_agent``).""" + logger.debug("Processing pending message: '%s...'", pending[:40]) + + # Clear the adapter's interrupt event so the next _run_agent call doesn't re-trigger the + # interrupt before the new agent's first API call (infinite loop otherwise). + if adapter and hasattr(adapter, '_active_sessions') and session_key and session_key in adapter._active_sessions: + adapter._active_sessions[session_key].clear() + + # Cap recursion depth to prevent resource exhaustion when the + # user sends multiple messages while the agent keeps failing. (#816) + if _interrupt_depth >= self._MAX_INTERRUPT_DEPTH: + logger.warning( + "Interrupt recursion depth %d reached for session %s — " + "queueing message instead of recursing.", + _interrupt_depth, session_key, + ) + adapter = self._adapter_for_source(source) + if adapter and pending_event: + merge_pending_message_event(adapter._pending_messages, session_key, pending_event) + elif adapter and hasattr(adapter, 'queue_message'): + adapter.queue_message(session_key, pending) + return result_holder[0] or {"final_response": response, "messages": history} + + was_interrupted = result.get("interrupted") + if not was_interrupted: + await self._run_agent_deliver_first_response( + source=source, + adapter=adapter, + session_key=session_key, + run_generation=run_generation, + event_message_id=event_message_id, + response=response, + result=result, + stream_consumer_holder=stream_consumer_holder, + stream_task=stream_task, + _status_thread_metadata=_status_thread_metadata, + ) + # else: interrupted — discard the response ("Operation interrupted." is noise; the user + # knows they sent a new message). + + updated_history = result.get("messages", history) + next_source = source + next_message = pending + next_message_id = None + next_channel_prompt = None + next_session_key = session_key + # Carry the pending event's message_type into the recursive call so queued voice turns + # can stream TTS and re-mark the generation for the final delivered turn. + next_message_type = None + if pending_event is not None: + next_source = getattr(pending_event, "source", None) or source + if self._is_goal_continuation_event(pending_event) and not self._goal_still_active_for_session(session_id): + logger.info( + "Discarding stale goal continuation for session %s — goal is no longer active", + session_key or "?", ) - if not _has_stream_consumer: + return result + # Resolve the follow-up's session key BEFORE preparing the inbound text: + # _prepare_inbound_message_text buffers native image paths under the key given, and + # the recursive _run_agent consumes them under next_session_key — mismatch drops them. + try: + next_session_key = self._session_key_for_source(next_source) + except Exception: + logger.debug( + "Queued follow-up session-key resolution failed; reusing %s", + session_key or "?", + exc_info=True, + ) + next_message = await self._prepare_profile_scoped_inbound_message_text( + event=pending_event, + source=next_source, + history=updated_history, + session_key=next_session_key, + ) + if next_message is None: + return result + next_message_id = self._reply_anchor_for_event(pending_event) + next_channel_prompt = getattr(pending_event, "channel_prompt", None) + next_message_type = getattr(pending_event, "message_type", None) + + # Clear the prior logical turn's completed streaming marker so the recursive turn's + # streaming TTS isn't suppressed by that completion. + _clear_adapter = self._adapter_for_source(source) + if _clear_adapter is not None and session_key and run_generation is not None: + _completed_turns = getattr(_clear_adapter, "_streaming_tts_completed_turns", None) + if _completed_turns is not None: + _prior_key = getattr(_clear_adapter, "_streaming_tts_turn_key", None) + if callable(_prior_key): + _pk = _prior_key(session_key, run_generation) + if _pk: + _completed_turns.discard(_pk) + + # Restart the typing indicator for the follow-up turn; the outer + # _process_message_background typing task is alive but may be stale. + _followup_adapter = self._adapter_for_source(source) + if _followup_adapter: + with suppress(Exception): + await _followup_adapter.send_typing( + source.chat_id, + metadata=_status_thread_metadata, + ) + + # Re-baseline the cached agent's message_count before recursing into the /queue follow-up: + # the coherence guard would otherwise rebuild on OUR OWN flushed rows and destroy the + # prompt-cache prefix; _handle_message_with_agent re-baselines only after the chain ends. + await self._refresh_agent_cache_message_count(session_key, session_id) + + followup_result = await self._run_agent( + message=next_message, + context_prompt=context_prompt, + history=updated_history, + source=next_source, + session_id=session_id, + session_key=next_session_key, + run_generation=run_generation, + _interrupt_depth=_interrupt_depth + 1, + event_message_id=next_message_id, + channel_prompt=next_channel_prompt, + message_type=next_message_type, + ) + return _preserve_queued_followup_history_offset(result, followup_result) + + async def _run_agent_cleanup_turn_tasks( + self, + *, + progress_task: Any, + log_task: Any, + interrupt_monitor: "asyncio.Task", + _notify_task: "asyncio.Task", + tracking_task: "asyncio.Task", + stream_task: Any, + stream_consumer_holder: list, + streaming_tts_consumer_holder: list, + session_key: Optional[str], + run_generation: Optional[int], + ) -> None: + """``finally`` half of a turn: cancel background tasks, flush stream, release the session slot.""" + # Stop progress sender, interrupt monitor, and notification task + if progress_task: + progress_task.cancel() + if log_task: + log_task.cancel() + interrupt_monitor.cancel() + _notify_task.cancel() + + # Wait for stream consumer to finish its final edit + if stream_task: + # If the agent never created a stream consumer (non-streaming path, or a test stub + # returning synchronously) there is nothing to flush — cancel now instead of waiting + # out the 5s timeout polling for a consumer that will never arrive. + _has_stream_consumer = ( + stream_consumer_holder + and stream_consumer_holder[0] is not None + ) + if not _has_stream_consumer: + stream_task.cancel() + with suppress(asyncio.CancelledError): + await stream_task + else: + try: + await asyncio.wait_for(stream_task, timeout=5.0) + except (asyncio.TimeoutError, asyncio.CancelledError): stream_task.cancel() with suppress(asyncio.CancelledError): await stream_task - else: - try: - await asyncio.wait_for(stream_task, timeout=5.0) - except (asyncio.TimeoutError, asyncio.CancelledError): - stream_task.cancel() - with suppress(asyncio.CancelledError): - await stream_task - # Unconditional abort + bounded wait for the streaming-TTS consumer: covers cancellation / - # exception paths where the normal finalisation block was skipped. - _stts_finally = streaming_tts_consumer_holder[0] - if _stts_finally is not None and not _stts_finally.done: - _stts_finally.abort("cleanup") - with suppress(Exception): - await _stts_finally.wait_complete(timeout=2.0) + # Unconditional abort + bounded wait for the streaming-TTS consumer: covers cancellation / + # exception paths where the normal finalisation block was skipped. + _stts_finally = streaming_tts_consumer_holder[0] + if _stts_finally is not None and not _stts_finally.done: + _stts_finally.abort("cleanup") + with suppress(Exception): + await _stts_finally.wait_complete(timeout=2.0) - # Clean up tracking - tracking_task.cancel() - if session_key: - # Release the slot only if this run's generation still owns it: a /stop or /new that - # bumped the generation while we unwound already installed its own state; keep it. - self._release_running_agent_state( - session_key, run_generation=run_generation - ) - if self._draining: - self._update_runtime_status("draining") + # Clean up tracking + tracking_task.cancel() + if session_key: + # Release the slot only if this run's generation still owns it: a /stop or /new that + # bumped the generation while we unwound already installed its own state; keep it. + self._release_running_agent_state( + session_key, run_generation=run_generation + ) + if self._draining: + self._update_runtime_status("draining") - # Wait for cancelled tasks - for task in [progress_task, log_task, interrupt_monitor, tracking_task, _notify_task]: - if task: - try: - await task - except asyncio.CancelledError: - pass - except Exception: - # A background task that died of a non-cancellation error (transport drop in - # a progress/card publish) must not abort the cleanup path — everything - # after this loop (final-delivery bookkeeping) still runs (review B7). - logger.debug( - "background turn task failed during cleanup", - exc_info=True, - ) + # Wait for cancelled tasks + for task in [progress_task, log_task, interrupt_monitor, tracking_task, _notify_task]: + if task: + try: + await task + except asyncio.CancelledError: + pass + except Exception: + # A background task that died of a non-cancellation error (transport drop in + # a progress/card publish) must not abort the cleanup path — everything + # after this loop (final-delivery bookkeeping) still runs (review B7). + logger.debug( + "background turn task failed during cleanup", + exc_info=True, + ) + async def _run_agent_mark_streamed_delivery( + self, + response: Any, + stream_consumer_holder: list, + source: SessionSource, + session_key: Optional[str], + ) -> None: # If streaming already delivered the response, skip the caller's send() — but never when the # agent failed (the error is unseen content) or on "(empty)": interim text ("Let me search…") # set already_sent but is NOT the final answer; suppressing would leave the user with silence. @@ -28553,7 +28591,7 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew _transformed = bool(response.get("response_transformed")) # Suppress the normal send only when the actual final reply reached the user (streamed, or # interim preview of that *exact* text); commentary shown during a compression/split isn't it. - _streamed = _stream_confirmed_final_delivery( + _streamed = self._run_agent_stream_confirmed_final_delivery( _sc, _final, previewed=_previewed, @@ -28648,6 +28686,16 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew len(_final), ) + def _run_agent_schedule_bubble_cleanup( + self, + response: Any, + _cleanup_progress: bool, + _cleanup_adapter: Any, + _cleanup_msg_ids: List[str], + source: SessionSource, + session_key: Optional[str], + run_generation: Optional[int], + ) -> None: # Schedule deletion of tracked temporary progress bubbles after the final response lands; failed # runs keep them as breadcrumbs. Only on adapters with ``delete_message``; failures swallowed. if ( @@ -28686,6 +28734,406 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew ) except Exception as _rpe: logger.debug("Post-delivery cleanup registration failed: %s", _rpe) + def _run_agent_bind_turn_wiring( + self, + turn_ctx: TurnContext, + turn_runner: TurnRunner, + source: SessionSource, + event_message_id: Optional[str], + _progress_metadata: Optional[dict], + _progress_reply_to: Optional[str], + _progress_thread_id: Any, + _relay_prospective_thread_id: Optional[str], + ) -> Optional[Dict[str, Any]]: + """Publish progress metadata, result holders and the sync→async bridges onto ``turn_ctx``. + + Returns ``_status_thread_metadata``; the holders are read back via ``turn_ctx.*_holder``. + """ + # Extracted to TurnRunner.send_progress_messages; the threading metadata above is published + # onto the shared TurnContext where the original closure's captured locals were bound. + turn_ctx._progress_metadata = _progress_metadata + turn_ctx._progress_reply_to = _progress_reply_to + + # We need to share the agent instance for interrupt support + agent_holder = [None] # Mutable container for the agent instance + turn_ctx.agent_holder = agent_holder + result_holder = [None] # Mutable container for the result + tools_holder = [None] # Mutable container for the tool definitions + stream_consumer_holder = [None] # Mutable container for stream consumer + # streaming PCM audio consumer. Created on the gateway event-loop thread (NOT in run_sync's + # executor worker) so outer finalisation / interrupt paths can reference it without a NameError. + streaming_tts_consumer_holder: list = [None] + turn_ctx.result_holder = result_holder + turn_ctx.tools_holder = tools_holder + turn_ctx.stream_consumer_holder = stream_consumer_holder + turn_ctx.streaming_tts_consumer_holder = streaming_tts_consumer_holder + + # Bridge sync step_callback → async hooks.emit for agent:step events + _loop_for_step = asyncio.get_running_loop() + _hooks_ref = self.hooks + + # Bridge extracted to TurnRunner._step_callback_sync; the loop and + # hooks refs bound just above are published at their original site. + turn_ctx._loop_for_step = _loop_for_step + turn_ctx._hooks_ref = _hooks_ref + turn_ctx._step_callback_sync = turn_runner._step_callback_sync + + # Bridge sync event_callback → async hooks.emit for lifecycle events (e.g. session:compress + # after a compression split); extracted to TurnRunner._event_callback_sync. + turn_ctx._event_callback_sync = turn_runner._event_callback_sync + + # Bridge sync status_callback → async adapter.send for context pressure + _status_adapter = self._adapter_for_source(source) + _status_chat_id = source.chat_id + _status_thread_metadata = self._run_agent_status_thread_metadata( + source, event_message_id, _progress_thread_id, _relay_prospective_thread_id, + ) + + # Bridge extracted to TurnRunner._status_callback_sync; publish the status wiring computed + # above onto the shared TurnContext at the exact original binding site. + turn_ctx._status_adapter = _status_adapter + turn_ctx._status_chat_id = _status_chat_id + turn_ctx._status_thread_metadata = _status_thread_metadata + turn_ctx._status_callback_sync = turn_runner._status_callback_sync + return _status_thread_metadata + + async def _run_agent_notify_long_running( + self, + disp: "GatewayRunner._RunAgentDisplay", + *, + source: SessionSource, + session_key: Optional[str], + agent_holder: list, + _executor_task_holder: list, + _NOTIFY_INTERVAL: Optional[float], + _long_running_mode: str, + _notify_start: float, + _status_thread_metadata: Optional[Dict[str, Any]], + _cleanup_progress: bool, + _cleanup_msg_ids: List[str], + ) -> None: + """Periodic \"still working\" heartbeat (edited in place where the adapter supports it). + + ``_executor_task_holder[0]`` is populated once the executor future exists; tolerate the + brief window before then (it reads as None). + """ + if _NOTIFY_INTERVAL is None: + return # Notifications disabled (gateway_notify_interval: 0) + _notify_adapter = self._adapter_for_source(source) + if not _notify_adapter: + return + # Track the heartbeat message id to edit in place where supported (Telegram, Discord, + # Slack, ...) instead of a new "Still working" bubble every interval. + _heartbeat_msg_id: Optional[str] = None + while True: + await asyncio.sleep(_NOTIFY_INTERVAL) + # Stop heartbeating once this run no longer owns the session slot or the executor has + # finished, else a stale "running: delegate_task" bubble outlives its run. _executor_task + # is bound just after this task is scheduled; tolerate the brief window before then. + _exec_ref = _executor_task_holder[0] + if not self._should_emit_long_running_notification( + session_key, agent_holder[0], _exec_ref + ): + break + _elapsed_mins = int((time.time() - _notify_start) // 60) + # Default heartbeat is terse (elapsed + current tool); the verbose iteration counter is + # gated on busy_ack_detail so users can opt in per platform. + _agent_ref = agent_holder[0] + _status_detail = "" + _want_iteration_detail = bool( + disp.resolve_display_setting( + disp.user_config, + disp.platform_key, + "busy_ack_detail", + True, + ) + ) + if _agent_ref and hasattr(_agent_ref, "get_activity_summary"): + try: + _a = _agent_ref.get_activity_summary() + _parts = [] + if _want_iteration_detail: + _parts.append( + f"iteration {_a['api_call_count']}/{_a['max_iterations']}" + ) + _action = _a.get("current_tool") or _a.get("last_activity_desc") + if _action: + _parts.append(str(_action)) + if _parts: + _status_detail = " — " + ", ".join(_parts) + except Exception: + pass + _heartbeat_text = ( + disp._generic_status_phrase("status") + if _long_running_mode == "generic" + else f"⏳ Working — {_elapsed_mins} min{_status_detail}" + ) + try: + _notify_res = None + if _heartbeat_msg_id: + try: + _notify_res = await _notify_adapter.edit_message( + source.chat_id, + _heartbeat_msg_id, + _heartbeat_text, + ) + except Exception as _ee: + logger.debug("Heartbeat edit failed: %s", _ee) + _notify_res = None + if not (_notify_res and getattr(_notify_res, "success", False)): + _notify_res = await _notify_adapter.send( + source.chat_id, + _heartbeat_text, + metadata=_interim_metadata(_non_conversational_metadata(_status_thread_metadata, platform=source.platform)), + ) + if getattr(_notify_res, "success", False) and getattr( + _notify_res, "message_id", None + ): + _heartbeat_msg_id = str(_notify_res.message_id) + if _cleanup_progress: + _cleanup_msg_ids.append(_heartbeat_msg_id) + except Exception as _ne: + logger.debug("Long-running notification error: %s", _ne) + + async def _run_agent_inner( + self, + message: str, + context_prompt: str, + history: List[Dict[str, Any]], + source: SessionSource, + session_id: str, + session_key: str = None, + run_generation: Optional[int] = None, + _interrupt_depth: int = 0, + event_message_id: Optional[str] = None, + inbound_message_id: Optional[str] = None, + channel_prompt: Optional[str] = None, + moa_config: Optional[dict] = None, + persist_user_message: Optional[Any] = None, + persist_user_timestamp: Optional[float] = None, + persist_user_display_kind: Optional[str] = None, + message_type: Optional[str] = None, + ) -> Dict[str, Any]: + """Run the agent; returns the full run_conversation result dict. + + Keys: "final_response", "messages", "api_calls", "completed". + """ + # ---- Proxy mode: delegate to remote API server ---- + if self._get_proxy_url(): + return await self._run_agent_via_proxy( + message=message, + context_prompt=context_prompt, + history=history, + source=source, + session_id=session_id, + session_key=session_key, + run_generation=run_generation, + event_message_id=event_message_id, + ) + + from run_agent import AIAgent + + disp = self._run_agent_display_settings(source) + _display_surface_mode = disp._display_surface_mode + needs_progress_queue = disp.needs_progress_queue + log_mode_enabled = disp.log_mode_enabled + log_queue = disp.log_queue + + turn_ctx, turn_runner, _cleanup_adapter = self._run_agent_build_turn_context( + disp, + AIAgent, + message=message, + context_prompt=context_prompt, + history=history, + source=source, + session_id=session_id, + session_key=session_key, + run_generation=run_generation, + _interrupt_depth=_interrupt_depth, + event_message_id=event_message_id, + inbound_message_id=inbound_message_id, + channel_prompt=channel_prompt, + moa_config=moa_config, + persist_user_message=persist_user_message, + persist_user_timestamp=persist_user_timestamp, + persist_user_display_kind=persist_user_display_kind, + ) + _cleanup_progress = turn_ctx._cleanup_progress + _cleanup_msg_ids = turn_ctx._cleanup_msg_ids + + ( + _progress_metadata, + _progress_reply_to, + _progress_thread_id, + _relay_prospective_thread_id, + ) = self._run_agent_progress_threading(source, event_message_id, disp._native_slack_task_cards) + + _status_thread_metadata = self._run_agent_bind_turn_wiring( + turn_ctx, + turn_runner, + source, + event_message_id, + _progress_metadata, + _progress_reply_to, + _progress_thread_id, + _relay_prospective_thread_id, + ) + send_progress_messages = turn_runner.send_progress_messages + agent_holder = turn_ctx.agent_holder + result_holder = turn_ctx.result_holder + tools_holder = turn_ctx.tools_holder + stream_consumer_holder = turn_ctx.stream_consumer_holder + streaming_tts_consumer_holder = turn_ctx.streaming_tts_consumer_holder + + self._run_agent_start_streaming_tts( + source, message_type, _status_thread_metadata, streaming_tts_consumer_holder, + ) + + # run_sync extracted to TurnRunner.run_sync (bound method; executor call unchanged). Its + # closed-over locals travel on turn_ctx; `nonlocal message` rebinds became ctx.message writes. + run_sync = turn_runner.run_sync + + # Start the progress sender if enabled. Gate on needs_progress_queue (tool_progress OR + # thinking_progress), not tool_progress alone: the sender drains BOTH tool-progress lines and + # _thinking scratch bubbles — a tool_progress-only gate left thinking-only queues never drained. + progress_task = None + if needs_progress_queue: + progress_task = asyncio.create_task(send_progress_messages()) + + # Start the tool-call log writer when tool_progress == "log". + log_task = None + if log_mode_enabled: + log_task = asyncio.create_task(self._run_agent_write_tool_log(log_queue)) + + # Start stream consumer task — polls for consumer creation since it + # happens inside run_sync (thread pool) after the agent is constructed. + stream_task = None + stream_task = asyncio.create_task(self._run_agent_stream_consumer_task(stream_consumer_holder)) + + # Track this agent as running for this session (for interrupt support) + # We do this in a callback after the agent is created + tracking_task = asyncio.create_task( + self._run_agent_track_agent(session_key, run_generation, agent_holder) + ) + + _interrupt_detected = asyncio.Event() # shared with backup check + interrupt_monitor = asyncio.create_task( + self._run_agent_monitor_for_interrupt( + source, + session_key, + agent_holder, + _interrupt_detected, + streaming_tts_consumer_holder, + ) + ) + + # Periodic "still working" notifications so the user knows the agent hasn't died. Config: + # agent.gateway_notify_interval or HERMES_AGENT_NOTIFY_INTERVAL env; default 180s. + _NOTIFY_INTERVAL_RAW = _float_env("HERMES_AGENT_NOTIFY_INTERVAL", 180) + _NOTIFY_INTERVAL = _NOTIFY_INTERVAL_RAW if _NOTIFY_INTERVAL_RAW > 0 else None + _long_running_mode = _display_surface_mode( + "long_running_notifications", + default=True, + allow_generic=True, + ) + if _long_running_mode == "off": + _NOTIFY_INTERVAL = None + _notify_start = time.time() + _executor_task_holder: list = [None] # bound once the executor future exists (see below) + _notify_task = asyncio.create_task( + self._run_agent_notify_long_running( + disp, + source=source, + session_key=session_key, + agent_holder=agent_holder, + _executor_task_holder=_executor_task_holder, + _NOTIFY_INTERVAL=_NOTIFY_INTERVAL, + _long_running_mode=_long_running_mode, + _notify_start=_notify_start, + _status_thread_metadata=_status_thread_metadata, + _cleanup_progress=_cleanup_progress, + _cleanup_msg_ids=_cleanup_msg_ids, + ) + ) + + try: + worker = self._run_agent_start_turn_worker( + turn_ctx, run_sync, agent_holder, session_id, session_key, run_generation, + ) + _executor_task_holder[0] = worker.executor_task # read late by _notify_long_running + response = await self._run_agent_await_turn_worker( + worker, + source=source, + session_key=session_key, + agent_holder=agent_holder, + result_holder=result_holder, + tools_holder=tools_holder, + _interrupt_detected=_interrupt_detected, + interrupt_monitor=interrupt_monitor, + streaming_tts_consumer_holder=streaming_tts_consumer_holder, + _status_thread_metadata=_status_thread_metadata, + ) + + self._run_agent_evict_on_fallback(session_key, agent_holder, result_holder) + + # Check if we were interrupted OR have a queued message (/queue). + result = result_holder[0] + adapter = self._adapter_for_source(source) + + await self._run_agent_finalize_streaming_tts( + streaming_tts_consumer_holder, adapter, session_key, run_generation, + ) + + pending_event, pending = await self._run_agent_drain_pending( + result, adapter, source, session_key, + ) + + if pending_event or pending: + return await self._run_agent_queued_followup( + source=source, + adapter=adapter, + session_id=session_id, + session_key=session_key, + run_generation=run_generation, + _interrupt_depth=_interrupt_depth, + event_message_id=event_message_id, + context_prompt=context_prompt, + history=history, + pending=pending, + pending_event=pending_event, + response=response, + result=result, + result_holder=result_holder, + stream_consumer_holder=stream_consumer_holder, + stream_task=stream_task, + _status_thread_metadata=_status_thread_metadata, + ) + finally: + await self._run_agent_cleanup_turn_tasks( + progress_task=progress_task, + log_task=log_task, + interrupt_monitor=interrupt_monitor, + _notify_task=_notify_task, + tracking_task=tracking_task, + stream_task=stream_task, + stream_consumer_holder=stream_consumer_holder, + streaming_tts_consumer_holder=streaming_tts_consumer_holder, + session_key=session_key, + run_generation=run_generation, + ) + + await self._run_agent_mark_streamed_delivery( + response, stream_consumer_holder, source, session_key, + ) + self._run_agent_schedule_bubble_cleanup( + response, + _cleanup_progress, + _cleanup_adapter, + _cleanup_msg_ids, + source, + session_key, + run_generation, + ) return response