fix(gateway): preserve replies on human silence markers

This commit is contained in:
MKanso
2026-09-14 16:16:09 +01:00
committed by kshitij
parent 5871d750bf
commit 5ea8fb2b78
3 changed files with 130 additions and 7 deletions

View File

@@ -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.

View File

@@ -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

View File

@@ -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)