Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
50 changes: 46 additions & 4 deletions gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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")
Expand Down
69 changes: 68 additions & 1 deletion tests/gateway/test_queue_consumption.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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)


Loading