From 9fc7f17906eab1dd81ddfdf8a1edeecac1e79940 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Sat, 26 Sep 2026 06:37:36 +0530 Subject: [PATCH] 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. --- gateway/platforms/event.py | 7 +++--- gateway/run_busy.py | 25 ++++++++++++++----- gateway/run_inbound.py | 1 + gateway/run_turn.py | 3 ++- tests/gateway/test_active_turn_recovery.py | 6 +++-- tests/gateway/test_busy_redirect_anchor.py | 28 ++++++++++++++++++++++ 6 files changed, 58 insertions(+), 12 deletions(-) diff --git a/gateway/platforms/event.py b/gateway/platforms/event.py index 9293e21cb9..4e9faa4097 100644 --- a/gateway/platforms/event.py +++ b/gateway/platforms/event.py @@ -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. diff --git a/gateway/run_busy.py b/gateway/run_busy.py index 523aa43474..b7d5f8b1c7 100644 --- a/gateway/run_busy.py +++ b/gateway/run_busy.py @@ -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}'" diff --git a/gateway/run_inbound.py b/gateway/run_inbound.py index 48d793d7a7..f84b9e8796 100644 --- a/gateway/run_inbound.py +++ b/gateway/run_inbound.py @@ -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) diff --git a/gateway/run_turn.py b/gateway/run_turn.py index 81e28dba54..95aa81223f 100644 --- a/gateway/run_turn.py +++ b/gateway/run_turn.py @@ -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): diff --git a/tests/gateway/test_active_turn_recovery.py b/tests/gateway/test_active_turn_recovery.py index 8c781ec2f9..a1b2ba762f 100644 --- a/tests/gateway/test_active_turn_recovery.py +++ b/tests/gateway/test_active_turn_recovery.py @@ -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") diff --git a/tests/gateway/test_busy_redirect_anchor.py b/tests/gateway/test_busy_redirect_anchor.py index d4b9df18cb..3e0b4ae1c4 100644 --- a/tests/gateway/test_busy_redirect_anchor.py +++ b/tests/gateway/test_busy_redirect_anchor.py @@ -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)