fix(gateway): filter finalized first response before queued follow-up send
A successful turn returning exactly NO_REPLY (or another exact silence marker) leaked the literal control token when a second message was queued before the first turn finished. The queued-follow-up recovery branch sends the first final_response directly through adapter.send() and predates the silence filter added to the normal completed-turn path. Use the finalized task result for that recovery delivery rather than the raw result_holder copy, then apply the existing is_intentional_silence_agent_result() predicate before sending. This keeps the established contract: only successful exact-marker turns suppress; substantive prose and failed results still send; stream-confirmed responses still skip the resend; and persisted history is untouched. Using finalized output also preserves normal empty/failure normalization on this direct path. Integration regressions run the real Slack _run_agent queued-follow-up flow: one proves NO_REPLY never reaches the adapter while the second turn still runs; another proves an empty failed first turn sends its normalized error before the queued follow-up. Existing filter tests cover every supported marker, prose mentions, and failed-result semantics.
This commit is contained in:
committed by
Teknium
parent
3d7e1c5f43
commit
59fdd41f5a
+24
-3
@@ -21819,14 +21819,35 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
|
||||
pass
|
||||
except Exception as e:
|
||||
logger.debug("Stream consumer wait before queued message failed: %s", e)
|
||||
_previewed = bool(result.get("response_previewed"))
|
||||
first_response = result.get("final_response", "")
|
||||
# The queued branch needs raw ``result`` for interruption,
|
||||
# history, and recursion state, but delivery must use the
|
||||
# finalized task result. The latter contains empty/failure
|
||||
# normalization and any final response processing applied by
|
||||
# _run_agent_task; sending the raw copy bypasses those steps.
|
||||
_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,
|
||||
)
|
||||
if first_response and not _already_streamed:
|
||||
# Apply the same predicate as the normal completed-turn path.
|
||||
# This direct queued-send branch predates intentional-silence
|
||||
# filtering, so without this check it leaks 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 and not _already_streamed:
|
||||
try:
|
||||
logger.info(
|
||||
"Queued follow-up for session %s: final stream delivery not confirmed; sending first response before continuing.",
|
||||
|
||||
@@ -697,6 +697,48 @@ class QueuedCommentaryAgent:
|
||||
}
|
||||
|
||||
|
||||
class QueuedSilenceAgent:
|
||||
"""First turn is intentionally silent; queued follow-up still runs."""
|
||||
|
||||
calls = 0
|
||||
|
||||
def __init__(self, **kwargs):
|
||||
self.tools = []
|
||||
|
||||
def run_conversation(self, message, conversation_history=None, task_id=None):
|
||||
type(self).calls += 1
|
||||
return {
|
||||
"final_response": "NO_REPLY" if type(self).calls == 1 else "follow-up processed",
|
||||
"messages": [],
|
||||
"api_calls": 1,
|
||||
}
|
||||
|
||||
|
||||
class QueuedFailedEmptyAgent:
|
||||
"""First turn fails empty; its normalized error must send before follow-up."""
|
||||
|
||||
calls = 0
|
||||
|
||||
def __init__(self, **kwargs):
|
||||
self.tools = []
|
||||
|
||||
def run_conversation(self, message, conversation_history=None, task_id=None):
|
||||
type(self).calls += 1
|
||||
if type(self).calls == 1:
|
||||
return {
|
||||
"final_response": "",
|
||||
"messages": [],
|
||||
"api_calls": 1,
|
||||
"failed": True,
|
||||
"error": "provider exploded",
|
||||
}
|
||||
return {
|
||||
"final_response": "follow-up processed",
|
||||
"messages": [],
|
||||
"api_calls": 1,
|
||||
}
|
||||
|
||||
|
||||
class BackgroundReviewAgent:
|
||||
def __init__(self, **kwargs):
|
||||
self.background_review_callback = kwargs.get("background_review_callback")
|
||||
@@ -1090,6 +1132,52 @@ async def test_run_agent_queued_message_does_not_treat_commentary_as_final(monke
|
||||
assert "final response 1" in sent_texts
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_run_agent_suppresses_silent_first_turn_and_processes_queued_followup(
|
||||
monkeypatch, tmp_path,
|
||||
):
|
||||
"""Regression: queued direct-send must not leak NO_REPLY to the channel."""
|
||||
QueuedSilenceAgent.calls = 0
|
||||
adapter, result = await _run_with_agent(
|
||||
monkeypatch,
|
||||
tmp_path,
|
||||
QueuedSilenceAgent,
|
||||
session_id="sess-queued-silence",
|
||||
pending_text="queued follow-up",
|
||||
platform=Platform.SLACK,
|
||||
chat_id="C123",
|
||||
thread_id="1712345678.000100",
|
||||
)
|
||||
|
||||
sent_texts = [call["content"] for call in adapter.sent]
|
||||
assert QueuedSilenceAgent.calls == 2
|
||||
assert result["final_response"] == "follow-up processed"
|
||||
assert "NO_REPLY" not in sent_texts
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_run_agent_sends_normalized_failure_before_queued_followup(
|
||||
monkeypatch, tmp_path,
|
||||
):
|
||||
"""Queued delivery uses finalized output, not the raw empty agent result."""
|
||||
QueuedFailedEmptyAgent.calls = 0
|
||||
adapter, result = await _run_with_agent(
|
||||
monkeypatch,
|
||||
tmp_path,
|
||||
QueuedFailedEmptyAgent,
|
||||
session_id="sess-queued-failed-empty",
|
||||
pending_text="queued follow-up",
|
||||
platform=Platform.SLACK,
|
||||
chat_id="C123",
|
||||
thread_id="1712345678.000100",
|
||||
)
|
||||
|
||||
sent_texts = [call["content"] for call in adapter.sent]
|
||||
assert QueuedFailedEmptyAgent.calls == 2
|
||||
assert result["final_response"] == "follow-up processed"
|
||||
assert any("The request failed: provider exploded" in text for text in sent_texts)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_run_agent_defers_background_review_notification_until_release(monkeypatch, tmp_path):
|
||||
adapter, result = await _run_with_agent(
|
||||
|
||||
Reference in New Issue
Block a user