diff --git a/gateway/response_filters.py b/gateway/response_filters.py index b465053248..522ede2893 100644 --- a/gateway/response_filters.py +++ b/gateway/response_filters.py @@ -14,6 +14,10 @@ from typing import Any # error/empty-response path, not silence. LIVE_GATEWAY_SILENT_MARKERS = frozenset({"[SILENT]", "SILENT", "NO_REPLY", "NO REPLY"}) +# only these persisted turn kinds are allowed to disappear when they emit a bare marker. +# ordinary user turns must still get a visible fallback if a model emits one by mistake. +MACHINERY_DISPLAY_KINDS = frozenset({"internal_notification", "model_switch", "auto_continue"}) + # Longer than any marker could plausibly be, even with stray punctuation. _MARKER_LENGTH_CAP = 64 @@ -80,6 +84,25 @@ def is_intentional_silence_agent_result(agent_result: dict | None, response: Any return isinstance(agent_result, dict) and not agent_result.get("failed") and is_intentional_silence_response(response) +def should_swallow_silence( + agent_result: dict | None, + response: Any, + *, + display_kind: Any = None, +) -> bool: + """allow bare silence only for the current synthetic gateway turn. + + the caller passes the current turn's persisted display kind instead of asking + us to infer it from the old transcript. the inbound user row is not in that + transcript yet, and a previous internal row must never authorize a human turn. + """ + return ( + is_intentional_silence_agent_result(agent_result, response) + and isinstance(display_kind, str) + and display_kind in MACHINERY_DISPLAY_KINDS + ) + + def is_partial_silence_marker(text: Any) -> bool: """True while streamed ``text`` could still resolve to a silence marker. diff --git a/gateway/run_turn.py b/gateway/run_turn.py index 38b586bde2..338584a936 100644 --- a/gateway/run_turn.py +++ b/gateway/run_turn.py @@ -51,6 +51,11 @@ _CONTEXT_OVERFLOW_ERROR_PHRASES = ( "payload too large", "input is too long", ) +_UNEXPECTED_SILENCE_REPLY = ( + "⚠️ the model returned only a silence marker for a message that needed a reply. " + "try again or rephrase." +) + def is_context_overflow_failure_result(agent_result: dict, history_len: int) -> bool: """One verdict for "this failed turn is a context overflow", shared by transcript persistence @@ -276,6 +281,14 @@ class GatewayTurnMixin: except Exception: return False + @staticmethod + def _should_swallow_silence(agent_result, response, *, display_kind=None) -> bool: + try: + from gateway.response_filters import should_swallow_silence + return should_swallow_silence(agent_result, response, display_kind=display_kind) + except Exception: + return False + async def _hmwa_resolve_session(self, event, source): """Resolve ``source`` to its session entry (topic recovery, internal-route guards, Telegram topic-binding heal). Returns ``(source, session_entry, session_key)`` or ``None`` to drop @@ -1362,6 +1375,7 @@ class GatewayTurnMixin: async def _hmwa_shape_agent_response( self, agent_result, source, history, session_entry, session_key, _quick_key, run_generation, _run_start_session_id, _platform_name, _msg_start_time, + persist_user_display_kind: Optional[str] = None, ): """Turn the raw agent result into the outbound text: sentinel/silence handling, response logging, resume-pending clear, empty-response normalization, and identity-guarded @@ -1377,6 +1391,21 @@ class GatewayTurnMixin: if _is_gateway_hidden_reasoning_incomplete_turn(agent_result): response = "" _intentional_silence = self._is_intentional_silence(agent_result, response) + if _intentional_silence and not self._should_swallow_silence( + agent_result, response, display_kind=persist_user_display_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. + logger.warning( + "silence marker rejected on a user turn: platform=%s chat=%s", + _platform_name, source.chat_id or "unknown", + ) + _intentional_silence = False + response = _UNEXPECTED_SILENCE_REPLY + # a stream consumer may have seen the marker before the final filter. make sure the + # visible fallback still goes through the normal final-send path. + if isinstance(agent_result, dict): + agent_result["already_sent"] = False # "(empty)" = the model produced no visible content after exhausting all retries. if response == "(empty)" and not _intentional_silence: @@ -2046,6 +2075,7 @@ class GatewayTurnMixin: response, _intentional_silence, agent_messages = await self._hmwa_shape_agent_response( agent_result, source, history, session_entry, session_key, _quick_key, run_generation, _run_start_session_id, _platform_name, _msg_start_time, + persist_user_display_kind=prepared.persist_user_display_kind, ) response = self._hmwa_prepend_reasoning(agent_result, response, source, _intentional_silence) _footer_line = self._hmwa_runtime_footer_line(agent_result, source, _turn_seconds) @@ -3487,11 +3517,22 @@ class GatewayTurnMixin: ) # Same silence predicate as the normal path, else this branch leaks the literal marker. if self._is_intentional_silence(_delivery_result, first_response): - logger.info( - "Queued follow-up for session %s: suppressing intentional silence marker before continuing.", - session_key or "?", - ) - elif first_response: + if self._should_swallow_silence( + _delivery_result, first_response, display_kind=turn_ctx.persist_user_display_kind, + ): + logger.info( + "Queued follow-up for session %s: suppressing intentional silence marker before continuing.", + session_key or "?", + ) + first_response = "" + else: + logger.warning( + "Queued follow-up for session %s: replacing a human-turn silence marker.", + session_key or "?", + ) + first_response = _UNEXPECTED_SILENCE_REPLY + _already_streamed = False + if first_response: logger.info( "Queued follow-up for session %s: final text delivery confirmed; delivering explicit media before continuing." if _already_streamed else diff --git a/tests/gateway/test_gateway_silence_tokens.py b/tests/gateway/test_gateway_silence_tokens.py index a9f7ce445d..50e4f6d976 100644 --- a/tests/gateway/test_gateway_silence_tokens.py +++ b/tests/gateway/test_gateway_silence_tokens.py @@ -1,6 +1,7 @@ """Gateway intentional-silence token behavior.""" from datetime import datetime +from types import SimpleNamespace from unittest.mock import AsyncMock, MagicMock import pytest @@ -12,6 +13,7 @@ from gateway.session import SessionEntry, SessionSource from gateway.response_filters import ( is_intentional_silence_agent_result, is_intentional_silence_response, + should_swallow_silence, ) @@ -24,11 +26,12 @@ def _source(): ) -def _event(): +def _event(*, internal: bool = False): return MessageEvent( text="side chatter", source=_source(), message_id="msg-42", + internal=internal, ) @@ -92,8 +95,17 @@ def test_failed_agent_result_never_counts_as_intentional_silence(): assert not is_intentional_silence_agent_result({"failed": True}, "NO_REPLY") +def test_only_synthetic_turns_can_swallow_silence(): + result = {"failed": False} + for display_kind in ("internal_notification", "model_switch", "auto_continue"): + assert should_swallow_silence(result, "NO_REPLY", display_kind=display_kind) + + for display_kind in (None, "steer", ""): + assert not should_swallow_silence(result, "NO_REPLY", display_kind=display_kind) + + @pytest.mark.asyncio -async def test_silence_token_suppresses_delivery_but_preserves_transcript(monkeypatch, tmp_path): +async def test_human_turn_gets_a_visible_fallback_for_a_silence_marker(monkeypatch, tmp_path): runner = _runner(monkeypatch, tmp_path) runner._run_agent = AsyncMock(return_value={ "final_response": "[SILENT]", @@ -112,12 +124,59 @@ async def test_silence_token_suppresses_delivery_but_preserves_transcript(monkey _event(), _source(), "agent:main:telegram:group:-1001:12345", 1 ) + assert "silence marker" in response + assert "try again or rephrase" in response + + +@pytest.mark.asyncio +async def test_internal_silence_token_suppresses_delivery_but_preserves_transcript(monkeypatch, tmp_path): + runner = _runner(monkeypatch, tmp_path) + runner._run_agent = AsyncMock(return_value={ + "final_response": "[SILENT]", + "messages": [ + {"role": "user", "content": "side chatter"}, + {"role": "assistant", "content": "[SILENT]"}, + ], + "tools": [], + "history_offset": 0, + "last_prompt_tokens": 0, + "api_calls": 1, + "failed": False, + }) + + response = await runner._handle_message_with_agent( + _event(internal=True), _source(), "agent:main:telegram:group:-1001:12345", 1 + ) + assert response == "" appended = [call.args[1] for call in runner.session_store.append_to_transcript.call_args_list] assert {"role": "assistant", "content": "[SILENT]"}.items() <= appended[-1].items() assert [msg["role"] for msg in appended if msg.get("role") in {"user", "assistant"}] == ["user", "assistant"] +@pytest.mark.asyncio +async def test_queued_human_turn_also_gets_the_visible_fallback(): + runner = gateway_run.GatewayRunner(GatewayConfig()) + runner._deliver_queued_first_response = AsyncMock() + turn_ctx = SimpleNamespace( + session_key="agent:main:telegram:group:-1001:12345", + stream_consumer_holder=[None], + persist_user_display_kind=None, + source=_source(), + _status_thread_metadata=None, + event_message_id=None, + inbound_message_id="msg-42", + run_generation=1, + ) + result = {"final_response": "NO_REPLY", "failed": False} + + await runner._run_agent_deliver_first_response( + turn_ctx, None, result, result, None, + ) + + assert "silence marker" in runner._deliver_queued_first_response.await_args.args[0] + + @pytest.mark.asyncio async def test_empty_success_still_gets_empty_response_warning(monkeypatch, tmp_path): runner = _runner(monkeypatch, tmp_path)