fix(gateway): steering an addressed message into a running turn keeps the silence fallback
Busy steer (the busy-mode and priority paths, and /steer) pushes a new message into the running turn, which then answers it without a turn of its own. Redirect already folded the message's reply_expected into the turn; steer did not, so an unaddressed opener steered by an @mention could still end on a bare silence marker. All of them now go through _fold_into_running_turn. Also: stale comments on MessageEvent and the queued-terminal silence verdict, and a crash-recovery case for an unaddressed row.
This commit is contained in:
@@ -86,10 +86,11 @@ class MessageEvent:
|
||||
metadata: Dict[str, Any] = field(default_factory=dict)
|
||||
timestamp: datetime = field(default_factory=datetime.now)
|
||||
# May this event resolve gateway commands / control prompts? Proactive plugin events set False
|
||||
# so untrusted payload text stays conversational. Kept last for positional compat.
|
||||
# so untrusted payload text stays conversational. New fields append after it (positional compat).
|
||||
allow_gateway_control: bool = True
|
||||
# Whether this inbound turn was addressed to this bot. False means the adapter admitted a
|
||||
# free-response or peer-addressed message, None means the adapter cannot determine it.
|
||||
# Was this message addressed to this bot? False lets a bare silence marker stand (the adapter
|
||||
# knows the message was meant for someone else); None means unknown and keeps the visible
|
||||
# fallback, like True.
|
||||
reply_expected: Optional[bool] = None
|
||||
|
||||
# Process-local admission receipt, never routing metadata or execution acknowledgement.
|
||||
|
||||
@@ -593,7 +593,9 @@ class GatewayBusySessionMixin:
|
||||
steered = self._try_agent_verb(
|
||||
running_agent, "steer", steer_text, session_key, event=event
|
||||
)
|
||||
if not steered:
|
||||
if steered:
|
||||
self._fold_into_running_turn(running_agent, session_key, event)
|
||||
else:
|
||||
effective_mode = "queue"
|
||||
elif (
|
||||
effective_mode == "interrupt" and plain_text and agent_live
|
||||
@@ -636,22 +638,32 @@ class GatewayBusySessionMixin:
|
||||
"""
|
||||
if not self._try_agent_verb(running_agent, "redirect", text, session_key, event=event):
|
||||
return False
|
||||
turn = self._session_state(session_key).turn
|
||||
if turn.agent is not running_agent:
|
||||
turn = self._fold_into_running_turn(running_agent, session_key, event)
|
||||
if turn is None:
|
||||
return True # a newer turn already owns the slot; never re-anchor it
|
||||
anchor = self._reply_anchor_for_event(event)
|
||||
inbound_id = str(event.message_id) if event.message_id else None
|
||||
if turn.event is not None and turn.event is not event:
|
||||
turn.event.reply_anchor_override = anchor
|
||||
turn.event.ledger_message_id = inbound_id
|
||||
turn.event.absorb_reply_expected(event)
|
||||
if turn.ctx is not None:
|
||||
turn.ctx.event_message_id = anchor
|
||||
turn.ctx.inbound_message_id = inbound_id
|
||||
if turn.event is not None:
|
||||
turn.ctx.reply_expected = turn.event.reply_expected
|
||||
return True
|
||||
|
||||
def _fold_into_running_turn(self, running_agent, session_key: str, event: MessageEvent):
|
||||
"""The running turn now answers *event* too (steer, redirect): if *event* was addressed to
|
||||
the bot, a bare silence marker must not end the turn. Returns the turn, or None when a newer
|
||||
turn already owns the slot."""
|
||||
turn = self._session_state(session_key).turn
|
||||
if turn.agent is not running_agent:
|
||||
return None
|
||||
if turn.event is not None and turn.event is not event:
|
||||
turn.event.absorb_reply_expected(event)
|
||||
if turn.ctx is not None:
|
||||
turn.ctx.reply_expected = turn.event.reply_expected
|
||||
return turn
|
||||
|
||||
async def _interrupt_running_agent_for_busy_event(self, event: MessageEvent, adapter, running_agent) -> None:
|
||||
"""Interrupt mode: abort in-flight tool calls; the agent loop exits at its next check point."""
|
||||
from gateway.run import _build_media_placeholder
|
||||
@@ -1061,6 +1073,7 @@ class GatewayBusySessionMixin:
|
||||
return f"⚠️ Steer failed: {exc}"
|
||||
if not accepted:
|
||||
return "Steer rejected (empty payload)."
|
||||
self._fold_into_running_turn(running_agent, quick_key, event)
|
||||
preview = steer_text[:60] + ("..." if len(steer_text) > 60 else "")
|
||||
target = "run and its active subagent(s)" if self._agent_has_active_subagents(running_agent) else "run"
|
||||
return f"⏩ Steer queued into current {target} — arrives after the next tool call: '{preview}'"
|
||||
|
||||
@@ -642,6 +642,7 @@ class GatewayInboundMixin:
|
||||
except Exception as exc:
|
||||
logger.warning("PRIORITY steer failed for session %s: %s", _quick_key, exc)
|
||||
if steered:
|
||||
self._fold_into_running_turn(running_agent, _quick_key, event)
|
||||
logger.debug("PRIORITY steer for session %s", _quick_key)
|
||||
return
|
||||
logger.debug("PRIORITY steer-fallback-to-queue for session %s", _quick_key)
|
||||
|
||||
@@ -1540,7 +1540,8 @@ class GatewayTurnMixin:
|
||||
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.
|
||||
# opened the chain: an internal follow-up, or a message not addressed to the bot, may go
|
||||
# silent; any other human one must not.
|
||||
_silence_kind = agent_result.get("queued_terminal_display_kind", persist_user_display_kind)
|
||||
_silence_reply_expected = agent_result.get("queued_terminal_reply_expected", reply_expected)
|
||||
if _intentional_silence and not silence_allowed(_silence_kind, _silence_reply_expected):
|
||||
|
||||
@@ -523,11 +523,13 @@ _WAKE = {"display_kind": "internal_notification"}
|
||||
("disk is 91% full", {**_WAKE, "display_metadata": {"notification_category": "diagnostic"}}, []),
|
||||
("NO_REPLY", {}, ["⚠️ The model returned only a silence marker for a message that needed a reply. "
|
||||
"Try again or rephrase."]),
|
||||
("NO_REPLY", {"display_metadata": {"reply_expected": False}}, []),
|
||||
])
|
||||
async def test_unclean_restart_never_redelivers_a_reply_live_delivery_suppressed(tmp_path, reply, prompt, owed):
|
||||
"""A crash-left reply is owed exactly what live delivery would have sent: nothing for a silence
|
||||
marker on a machinery turn or a muted diagnostic wake (and the finished turn is not resumed), the
|
||||
unexpected-silence notice for a human turn, never the raw marker."""
|
||||
marker on a machinery turn, a muted diagnostic wake or a message the adapter reported as not
|
||||
addressed to the bot (and the finished turn is not resumed), the unexpected-silence notice for
|
||||
any other human turn, never the raw marker."""
|
||||
from gateway.delivery_ledger import sweep_recoverable
|
||||
|
||||
(Path(os.environ["HERMES_HOME"]) / "config.yaml").write_text("display: {suppress_warning_notifications: true}\n", encoding="utf-8")
|
||||
|
||||
@@ -25,6 +25,9 @@ class Receiver:
|
||||
def redirect(self, text):
|
||||
return self.accept
|
||||
|
||||
def steer(self, text):
|
||||
return self.accept
|
||||
|
||||
|
||||
def _running_turn(runner, key, receiver):
|
||||
source = SessionSource(platform=Platform.TELEGRAM, chat_id="c1", user_id="u1", chat_type="dm")
|
||||
@@ -74,3 +77,28 @@ async def test_refused_or_foreign_redirect_leaves_the_anchor_on_the_opening_mess
|
||||
opening2, ctx2, redirecting2, _ = _running_turn(runner, "key2", Receiver())
|
||||
assert await runner._resolve_busy_steer_or_redirect(redirecting2, "key2", "interrupt", displaced)
|
||||
assert _reply_anchor_for_event(opening2) == "A" and ctx2.event_message_id == "A"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@pytest.mark.parametrize(
|
||||
"route", ["busy_interrupt", "priority_interrupt", "busy_steer", "priority_steer", "slash_steer"])
|
||||
async def test_a_turn_that_takes_in_an_addressed_message_keeps_the_silence_fallback(route):
|
||||
"""Redirect and steer fold a new message into the running turn, which then answers it too: if
|
||||
that message was addressed to the bot, a bare silence marker must not end the turn silently."""
|
||||
runner = GatewayRunner(config=GatewayConfig())
|
||||
receiver = Receiver()
|
||||
opening, ctx, incoming, source = _running_turn(runner, "key", receiver)
|
||||
opening.reply_expected, incoming.reply_expected = False, True
|
||||
|
||||
if route == "priority_interrupt":
|
||||
await runner._hm_busy_interrupt(incoming, source, receiver, "key")
|
||||
elif route == "priority_steer":
|
||||
runner._hm_busy_steer(incoming, receiver, "key")
|
||||
elif route == "slash_steer":
|
||||
incoming.text = "/steer " + incoming.text
|
||||
assert (await runner._busy_steer_command(incoming, "key", source)).startswith("⏩")
|
||||
else:
|
||||
outcome = await runner._resolve_busy_steer_or_redirect(incoming, "key", route[5:], receiver)
|
||||
assert outcome.redirected or outcome.steered
|
||||
|
||||
assert (opening.reply_expected, ctx.reply_expected) == (True, True)
|
||||
|
||||
Reference in New Issue
Block a user