diff --git a/gateway/run.py b/gateway/run.py index d7f44c6a4b945..27e4a397d64cd 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -8347,6 +8347,43 @@ async def _prepare_busy_steer_text(self, event: MessageEvent) -> str: return text return (enriched_text or text).strip() + def _requeue_interrupt_depth_pending( + self, + session_key: str, + source, + *, + adapter: Any, + pending_event: Optional[MessageEvent], + pending_text: Optional[str], + event_message_id: Optional[str], + channel_prompt: Optional[str], + ) -> bool: + """Park the next interrupt follow-up when recursive drain hits its cap.""" + if not adapter or not session_key: + return False + + queued_event = pending_event + if queued_event is None and pending_text: + queued_event = MessageEvent( + text=pending_text, + message_type=MessageType.TEXT, + source=source, + message_id=event_message_id, + channel_prompt=channel_prompt, + ) + + pending_slot = getattr(adapter, "_pending_messages", None) + if queued_event is not None and pending_slot is not None: + if getattr(queued_event, "source", None) is None: + queued_event.source = source + self._enqueue_fifo(session_key, queued_event, adapter) + return True + + if pending_text and hasattr(adapter, "queue_message"): + adapter.queue_message(session_key, pending_text) + return True + return False + async def _handle_active_session_busy_message(self, event: MessageEvent, session_key: str) -> bool: # --- Authorization gate (#17775) --- # The cold path (_handle_message) checks _is_user_authorized before @@ -24292,10 +24329,15 @@ def _stream_confirmed_final_delivery( _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) + self._requeue_interrupt_depth_pending( + session_key, + source, + adapter=adapter, + pending_event=pending_event, + pending_text=pending, + event_message_id=event_message_id, + channel_prompt=channel_prompt, + ) return result_holder[0] or {"final_response": response, "messages": history} was_interrupted = result.get("interrupted") diff --git a/tests/gateway/test_queue_consumption.py b/tests/gateway/test_queue_consumption.py index ad258b00233e4..84fa0ae5434db 100644 --- a/tests/gateway/test_queue_consumption.py +++ b/tests/gateway/test_queue_consumption.py @@ -169,6 +169,74 @@ def test_promote_stages_overflow_when_slot_already_populated(self): # gets the next-in-line item. assert adapter._pending_messages[session_key].text == "Q2" + def test_interrupt_depth_requeues_plain_interrupt_text_as_event(self): + """Depth-capped interrupt strings must not require adapter.queue_message().""" + from gateway.run import GatewayRunner + + runner = GatewayRunner.__new__(GatewayRunner) + runner._queued_events = {} + adapter = _StubAdapter() + session_key = "telegram:user:voice-depth" + source = MagicMock(chat_id="123", platform=Platform.TELEGRAM, profile=None) + + stored = runner._requeue_interrupt_depth_pending( + session_key, + source, + adapter=adapter, + pending_event=None, + pending_text='"transcribed voice follow-up"', + event_message_id="voice-1", + channel_prompt="voice-channel", + ) + + assert stored is True + queued = adapter._pending_messages[session_key] + assert queued.text == '"transcribed voice follow-up"' + assert queued.message_type == MessageType.TEXT + assert queued.source is source + assert queued.message_id == "voice-1" + assert queued.channel_prompt == "voice-channel" + assert runner._queued_events == {} + + def test_interrupt_depth_requeues_event_behind_existing_slot(self): + """Depth-capped full events should preserve FIFO ordering.""" + from gateway.run import GatewayRunner + + runner = GatewayRunner.__new__(GatewayRunner) + runner._queued_events = {} + adapter = _StubAdapter() + session_key = "telegram:user:voice-depth-fifo" + source = MagicMock(chat_id="123", platform=Platform.TELEGRAM, profile=None) + + adapter._pending_messages[session_key] = MessageEvent( + text="already queued", + message_type=MessageType.TEXT, + source=source, + message_id="text-1", + ) + voice_event = MessageEvent( + text="", + message_type=MessageType.VOICE, + source=source, + message_id="voice-2", + media_urls=["/tmp/voice-2.ogg"], + media_types=["audio/ogg"], + ) + + stored = runner._requeue_interrupt_depth_pending( + session_key, + source, + adapter=adapter, + pending_event=voice_event, + pending_text=None, + event_message_id=None, + channel_prompt=None, + ) + + assert stored is True + assert adapter._pending_messages[session_key].text == "already queued" + assert runner._queued_events[session_key] == [voice_event] + class TestBusyInputModeQueueFifo: """Regression coverage for issue #28503. @@ -219,4 +287,3 @@ def test_rapid_text_followups_are_queued_in_fifo_order(self): ] assert runner._queue_depth(session_key, adapter=adapter) == len(texts) -