fix(gateway): queued terminal turns carry their own display kind into silence shaping

The recursive _run_agent for a queued (/queue) follow-up passed no
persist_user_display_kind, so an internal (self-injected) follow-up ending
in a bare silence marker got the visible "silence marker rejected" fallback
meant for humans, and a human follow-up behind an internal opener could be
swallowed. Pass the same rule the top-level turn uses, stamp the terminal
turn's kind next to queued_terminal_inbound_id, and let the shaper prefer
it over the opener's kind.

Co-authored-by: KoNit-K <124019182+KoNit-K@users.noreply.github.com>
This commit is contained in:
kshitijk4poor
2026-09-15 11:59:36 +05:30
committed by kshitij
parent e1bdb026c8
commit a75628bc27
2 changed files with 60 additions and 2 deletions

View File

@@ -1391,8 +1391,13 @@ class GatewayTurnMixin:
if _is_gateway_hidden_reasoning_incomplete_turn(agent_result):
response = ""
_intentional_silence = self._is_intentional_silence(agent_result, response)
# A queued (/queue) chain's TERMINAL turn owns the silence verdict, not the event that
# opened the chain: an internal follow-up may go silent, a human one must not.
_silence_kind = persist_user_display_kind
if isinstance(agent_result, dict) and "queued_terminal_display_kind" in agent_result:
_silence_kind = agent_result["queued_terminal_display_kind"]
if _intentional_silence and not self._should_swallow_silence(
agent_result, response, display_kind=persist_user_display_kind,
agent_result, response, display_kind=_silence_kind,
):
# the current inbound row is not in ``history`` yet, so use the turn metadata we
# already carried into the agent run instead of guessing from an older row.
@@ -3605,6 +3610,8 @@ class GatewayTurnMixin:
# distinct from the reply anchor above (None in forum topics). Carry it or two chained
# topic turns with the same text would collide on one obligation id (queued-final-ledger).
next_inbound_id = None
# Same rule as the top-level turn (see _prepare_turn): only self-injected events are machinery.
next_display_kind = "internal_notification" if getattr(pending_event, "internal", False) else None
# See #60671.
if pending_event is not None:
next_source = getattr(pending_event, "source", None) or source
@@ -3677,6 +3684,7 @@ class GatewayTurnMixin:
run_generation=run_generation, _interrupt_depth=_interrupt_depth + 1,
event_message_id=next_message_id, inbound_message_id=next_inbound_id,
channel_prompt=next_channel_prompt, message_type=next_message_type,
persist_user_display_kind=next_display_kind,
)
except asyncio.CancelledError:
await _run_followup_processing_hook(
@@ -3696,7 +3704,11 @@ class GatewayTurnMixin:
# terminal reply, and is never redelivered. A deeper recursion has already set its own id,
# so only fill the key while it is still absent: the innermost turn wins.
if isinstance(merged, dict) and "queued_terminal_inbound_id" not in merged:
merged = {**merged, "queued_terminal_inbound_id": next_inbound_id}
merged = {
**merged,
"queued_terminal_inbound_id": next_inbound_id,
"queued_terminal_display_kind": next_display_kind,
}
return merged
async def _run_agent_cleanup_turn_tasks(

View File

@@ -177,6 +177,52 @@ async def test_queued_human_turn_also_gets_the_visible_fallback():
assert "silence marker" in runner._deliver_queued_first_response.await_args.args[0]
@pytest.mark.asyncio
async def test_queued_terminal_turn_owns_the_silence_verdict(monkeypatch, tmp_path):
"""The chain's LAST turn decides whether a bare marker may vanish, not the opener."""
runner = _runner(monkeypatch, tmp_path)
runner._MAX_INTERRUPT_DEPTH = 8
runner._run_agent = AsyncMock(return_value={"final_response": "NO_REPLY", "messages": []})
runner._is_goal_continuation_event = MagicMock(return_value=False)
runner._session_key_for_source = MagicMock(return_value="agent:main:telegram:group:-1001:12345")
runner._prepare_profile_scoped_inbound_message_text = AsyncMock(return_value="follow-up")
runner._adapter_for_source = MagicMock(return_value=None)
runner._refresh_agent_cache_message_count = AsyncMock()
turn_ctx = SimpleNamespace(
source=_source(), session_id="sid", session_key="agent:main:telegram:group:-1001:12345",
run_generation=1, _interrupt_depth=0, history=[], _status_thread_metadata=None,
context_prompt=None, result_holder=[None])
pending_event = SimpleNamespace(source=_source(), message_id="43", channel_prompt=None,
message_type=None, internal=True)
merged = await gateway_run.GatewayRunner._run_agent_queued_followup(
runner, turn_ctx, adapter=None, pending="hi again", pending_event=pending_event,
response="resp", result={"interrupted": True, "messages": []}, stream_task=None)
assert runner._run_agent.await_args.kwargs["persist_user_display_kind"] == "internal_notification"
assert merged["queued_terminal_display_kind"] == "internal_notification"
def _result(terminal_kind):
return {
"final_response": "[SILENT]", "tools": [], "history_offset": 0, "last_prompt_tokens": 0,
"api_calls": 1, "failed": False, "queued_terminal_inbound_id": "43",
"queued_terminal_display_kind": terminal_kind,
"messages": [{"role": "user", "content": "x"}, {"role": "assistant", "content": "[SILENT]"}],
}
# Human opener, internal terminal turn: silent.
runner = _runner(monkeypatch, tmp_path)
runner._run_agent = AsyncMock(return_value=_result("internal_notification"))
assert await runner._handle_message_with_agent(
_event(), _source(), "agent:main:telegram:group:-1001:12345", 1) == ""
# Internal opener, human terminal turn: visible fallback.
runner = _runner(monkeypatch, tmp_path)
runner._run_agent = AsyncMock(return_value=_result(None))
response = await runner._handle_message_with_agent(
_event(internal=True), _source(), "agent:main:telegram:group:-1001:12345", 1)
assert "silence marker" in response
@pytest.mark.asyncio
async def test_empty_success_still_gets_empty_response_warning(monkeypatch, tmp_path):
runner = _runner(monkeypatch, tmp_path)