Skip to content

Commit 97c096d

Browse files
teknium1kshitijk4poormemosr
authored
fix(gateway): process /queue'd messages after agent completion (NousResearch#2469)
* fix: respect DashScope v1 runtime mode for alibaba Remove the hardcoded Alibaba branch from resolve_runtime_provider() that forced api_mode='anthropic_messages' regardless of the base URL. Alibaba now goes through the generic API-key provider path, which auto-detects the protocol from the URL: - /apps/anthropic → anthropic_messages (via endswith check) - /v1 → chat_completions (default) This fixes Alibaba setup with OpenAI-compatible DashScope endpoints (e.g. coding-intl.dashscope.aliyuncs.com/v1) that were broken because runtime always forced Anthropic mode even when setup saved a /v1 URL. Based on PR NousResearch#2024 by @kshitijk4poor. * docs(skill): add split, merge, search examples to ocr-and-documents skill Adds pymupdf examples for PDF splitting, merging, and text search to the existing ocr-and-documents skill. No new dependencies — pymupdf already covers all three operations natively. * fix: replace all production print() calls with logger in rl_training_tool Replace all bare print() calls in production code paths with proper logger calls. - Add `import logging` and module-level `logger = logging.getLogger(__name__)` - Replace print() in _start_training_run() with logger.info() - Replace print() in _stop_training_run() with logger.info() - Replace print(Warning/Note) calls with logger.warning() and logger.info() Using the logging framework allows log level filtering, proper formatting, and log routing instead of always printing to stdout. * fix(gateway): process /queue'd messages after agent completion /queue stored messages in adapter._pending_messages but never consumed them after normal (non-interrupted) completion. The consumption path at line 5219 only checked pending messages when result.get('interrupted') was True — since /queue deliberately doesn't interrupt, queued messages were silently dropped. Now checks adapter._pending_messages after both interrupted AND normal completion. For queued messages (non-interrupt), the first response is delivered before recursing to process the queued follow-up. Skips the direct send when streaming already delivered the response. Reported by GhostMode on Discord. --------- Co-authored-by: kshitijk4poor <kshitijk4poor@users.noreply.github.com> Co-authored-by: memosr.eth <96793918+memosr@users.noreply.github.com>
1 parent 445ca62 commit 97c096d

2 files changed

Lines changed: 202 additions & 14 deletions

File tree

gateway/run.py

Lines changed: 37 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -5226,22 +5226,31 @@ async def monitor_for_interrupt():
52265226
self._effective_model = None
52275227
self._effective_provider = None
52285228

5229-
# Check if we were interrupted and have a pending message
5229+
# Check if we were interrupted OR have a queued message (/queue).
52305230
result = result_holder[0]
52315231
adapter = self.adapters.get(source.platform)
52325232

5233-
# Get pending message from adapter if interrupted.
5233+
# Get pending message from adapter.
52345234
# Use session_key (not source.chat_id) to match adapter's storage keys.
52355235
pending = None
5236-
if result and result.get("interrupted") and adapter:
5237-
pending_event = adapter.get_pending_message(session_key) if session_key else None
5238-
if pending_event:
5239-
pending = pending_event.text
5240-
elif result.get("interrupt_message"):
5241-
pending = result.get("interrupt_message")
5236+
if result and adapter and session_key:
5237+
if result.get("interrupted"):
5238+
# Interrupted — consume the interrupt message
5239+
pending_event = adapter.get_pending_message(session_key)
5240+
if pending_event:
5241+
pending = pending_event.text
5242+
elif result.get("interrupt_message"):
5243+
pending = result.get("interrupt_message")
5244+
else:
5245+
# Normal completion — check for /queue'd messages that were
5246+
# stored without triggering an interrupt.
5247+
pending_event = adapter.get_pending_message(session_key)
5248+
if pending_event:
5249+
pending = pending_event.text
5250+
logger.debug("Processing queued message after agent completion: '%s...'", pending[:40])
52425251

52435252
if pending:
5244-
logger.debug("Processing interrupted message: '%s...'", pending[:40])
5253+
logger.debug("Processing pending message: '%s...'", pending[:40])
52455254

52465255
# Clear the adapter's interrupt event so the next _run_agent call
52475256
# doesn't immediately re-trigger the interrupt before the new agent
@@ -5263,11 +5272,25 @@ async def monitor_for_interrupt():
52635272
adapter.queue_message(session_key, pending)
52645273
return result_holder[0] or {"final_response": response, "messages": history}
52655274

5266-
# Don't send the interrupted response to the user — it's just noise
5267-
# like "Operation interrupted." They already know they sent a new
5268-
# message, so go straight to processing it.
5269-
5270-
# Now process the pending message with updated history
5275+
was_interrupted = result.get("interrupted")
5276+
if not was_interrupted:
5277+
# Queued message after normal completion — deliver the first
5278+
# response before processing the queued follow-up.
5279+
# Skip if streaming already delivered it.
5280+
_sc = stream_consumer_holder[0]
5281+
_already_streamed = _sc and getattr(_sc, "already_sent", False)
5282+
first_response = result.get("final_response", "")
5283+
if first_response and not _already_streamed:
5284+
try:
5285+
await adapter.send(source.chat_id, first_response,
5286+
metadata=getattr(event, "metadata", None))
5287+
except Exception as e:
5288+
logger.warning("Failed to send first response before queued message: %s", e)
5289+
# else: interrupted — discard the interrupted response ("Operation
5290+
# interrupted." is just noise; the user already knows they sent a
5291+
# new message).
5292+
5293+
# Process the pending message with updated history
52715294
updated_history = result.get("messages", history)
52725295
return await self._run_agent(
52735296
message=pending,
Lines changed: 165 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,165 @@
1+
"""Tests for /queue message consumption after normal agent completion.
2+
3+
Verifies that messages queued via /queue (which store in
4+
adapter._pending_messages WITHOUT triggering an interrupt) are consumed
5+
after the agent finishes its current task — not silently dropped.
6+
"""
7+
8+
import asyncio
9+
from unittest.mock import AsyncMock, MagicMock, patch
10+
11+
import pytest
12+
13+
from gateway.platforms.base import (
14+
BasePlatformAdapter,
15+
MessageEvent,
16+
MessageType,
17+
PlatformConfig,
18+
Platform,
19+
)
20+
21+
22+
# ---------------------------------------------------------------------------
23+
# Minimal adapter for testing pending message storage
24+
# ---------------------------------------------------------------------------
25+
26+
class _StubAdapter(BasePlatformAdapter):
27+
def __init__(self):
28+
super().__init__(PlatformConfig(enabled=True, token="test"), Platform.TELEGRAM)
29+
30+
async def connect(self) -> bool:
31+
return True
32+
33+
async def disconnect(self) -> None:
34+
self._mark_disconnected()
35+
36+
async def send(self, chat_id, content, reply_to=None, metadata=None):
37+
from gateway.platforms.base import SendResult
38+
return SendResult(success=True, message_id="msg-1")
39+
40+
async def get_chat_info(self, chat_id):
41+
return {"id": chat_id, "type": "dm"}
42+
43+
44+
# ---------------------------------------------------------------------------
45+
# Tests
46+
# ---------------------------------------------------------------------------
47+
48+
class TestQueueMessageStorage:
49+
"""Verify /queue stores messages correctly in adapter._pending_messages."""
50+
51+
def test_queue_stores_message_in_pending(self):
52+
adapter = _StubAdapter()
53+
session_key = "telegram:user:123"
54+
event = MessageEvent(
55+
text="do this next",
56+
message_type=MessageType.TEXT,
57+
source=MagicMock(chat_id="123", platform=Platform.TELEGRAM),
58+
message_id="q1",
59+
)
60+
adapter._pending_messages[session_key] = event
61+
62+
assert session_key in adapter._pending_messages
63+
assert adapter._pending_messages[session_key].text == "do this next"
64+
65+
def test_get_pending_message_consumes_and_clears(self):
66+
adapter = _StubAdapter()
67+
session_key = "telegram:user:123"
68+
event = MessageEvent(
69+
text="queued prompt",
70+
message_type=MessageType.TEXT,
71+
source=MagicMock(chat_id="123", platform=Platform.TELEGRAM),
72+
message_id="q2",
73+
)
74+
adapter._pending_messages[session_key] = event
75+
76+
retrieved = adapter.get_pending_message(session_key)
77+
assert retrieved is not None
78+
assert retrieved.text == "queued prompt"
79+
# Should be consumed (cleared)
80+
assert adapter.get_pending_message(session_key) is None
81+
82+
def test_queue_does_not_set_interrupt_event(self):
83+
"""The whole point of /queue — no interrupt signal."""
84+
adapter = _StubAdapter()
85+
session_key = "telegram:user:123"
86+
87+
# Simulate an active session (agent running)
88+
adapter._active_sessions[session_key] = asyncio.Event()
89+
90+
# Store a queued message (what /queue does)
91+
event = MessageEvent(
92+
text="queued",
93+
message_type=MessageType.TEXT,
94+
source=MagicMock(),
95+
message_id="q3",
96+
)
97+
adapter._pending_messages[session_key] = event
98+
99+
# The interrupt event should NOT be set
100+
assert not adapter._active_sessions[session_key].is_set()
101+
assert not adapter.has_pending_interrupt(session_key)
102+
103+
def test_regular_message_sets_interrupt_event(self):
104+
"""Contrast: regular messages DO trigger interrupt."""
105+
adapter = _StubAdapter()
106+
session_key = "telegram:user:123"
107+
108+
adapter._active_sessions[session_key] = asyncio.Event()
109+
110+
# Simulate regular message arrival (what handle_message does)
111+
event = MessageEvent(
112+
text="new message",
113+
message_type=MessageType.TEXT,
114+
source=MagicMock(),
115+
message_id="m1",
116+
)
117+
adapter._pending_messages[session_key] = event
118+
adapter._active_sessions[session_key].set() # this is what handle_message does
119+
120+
assert adapter.has_pending_interrupt(session_key)
121+
122+
123+
class TestQueueConsumptionAfterCompletion:
124+
"""Verify that pending messages are consumed after normal completion."""
125+
126+
def test_pending_message_available_after_normal_completion(self):
127+
"""After agent finishes without interrupt, pending message should
128+
still be retrievable from adapter._pending_messages."""
129+
adapter = _StubAdapter()
130+
session_key = "telegram:user:123"
131+
132+
# Simulate: agent starts, /queue stores a message, agent finishes
133+
adapter._active_sessions[session_key] = asyncio.Event()
134+
event = MessageEvent(
135+
text="process this after",
136+
message_type=MessageType.TEXT,
137+
source=MagicMock(),
138+
message_id="q4",
139+
)
140+
adapter._pending_messages[session_key] = event
141+
142+
# Agent finishes (no interrupt)
143+
del adapter._active_sessions[session_key]
144+
145+
# The queued message should still be retrievable
146+
retrieved = adapter.get_pending_message(session_key)
147+
assert retrieved is not None
148+
assert retrieved.text == "process this after"
149+
150+
def test_multiple_queues_last_one_wins(self):
151+
"""If user /queue's multiple times, last message overwrites."""
152+
adapter = _StubAdapter()
153+
session_key = "telegram:user:123"
154+
155+
for text in ["first", "second", "third"]:
156+
event = MessageEvent(
157+
text=text,
158+
message_type=MessageType.TEXT,
159+
source=MagicMock(),
160+
message_id=f"q-{text}",
161+
)
162+
adapter._pending_messages[session_key] = event
163+
164+
retrieved = adapter.get_pending_message(session_key)
165+
assert retrieved.text == "third"

0 commit comments

Comments
 (0)