From a75628bc27b18d5570ca9481b27236475d2cd3e1 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Tue, 15 Sep 2026 11:59:36 +0530 Subject: [PATCH] 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> --- gateway/run_turn.py | 16 ++++++- tests/gateway/test_gateway_silence_tokens.py | 46 ++++++++++++++++++++ 2 files changed, 60 insertions(+), 2 deletions(-) diff --git a/gateway/run_turn.py b/gateway/run_turn.py index 338584a936..a4de367a12 100644 --- a/gateway/run_turn.py +++ b/gateway/run_turn.py @@ -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( diff --git a/tests/gateway/test_gateway_silence_tokens.py b/tests/gateway/test_gateway_silence_tokens.py index 50e4f6d976..aabe5da402 100644 --- a/tests/gateway/test_gateway_silence_tokens.py +++ b/tests/gateway/test_gateway_silence_tokens.py @@ -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)