fix(gateway): fire processing hooks for runner-drained queued follow-ups

A message that arrives while a turn is already running never gets the
processing-start acknowledgement — the 👀 read receipt on Slack, and the
equivalent in-progress reaction on Discord, Telegram, Feishu, Matrix,
Signal and Photon. It is not added-then-removed; the hook is never called.

`_run_processing_hook("on_processing_start", …)` has exactly one call site,
inside `BasePlatformAdapter._process_message_background` (base.py:5403).
`handle_message` takes the busy branch at base.py:5165, parks the event and
returns at base.py:5316 — above `_start_session_processing`, which is the
only thing that spawns `_process_message_background`. The parked event is
then drained in-band by the runner (`_dequeue_pending_event`, run.py:23277)
and replayed through a recursive `_run_agent` that touches no adapter hooks.
Because `get_pending_message` pops, the adapter's own drain (base.py:5823)
— the one path that would fire the hook — finds an empty slot. The gap is
structural, not a race, and it is shared by every mode (`queue`,
`interrupt`, steer-demoted-to-queue), by `/queue`, by photo-burst and
text-debounce flushes, and by voice drains.

Fire the existing hook pair around the recursive call. The follow-up now
gets the same lifecycle an idle-session message already gets, on every
platform, through one shared site.

Details that shaped the placement:

- Fired after the depth-cap requeue (run.py:23372) and after every discard
  and early return in the block, so no path can strand an in-progress
  marker: from that point on, control either reaches the recursion or
  raises, and both close the hook.
- Fired before `_refresh_agent_cache_message_count` so that re-baseline
  stays adjacent to the recursive call it exists to protect — inserting a
  reaction round-trip between them would widen the window in which the
  cross-process coherence guard (#45966) can trip on our own writes and
  rebuild the agent, destroying the prompt-cache prefix #46237 preserves.
- The hook adapter is resolved from the follow-up's own source, not the
  completing turn's: a multiplexed gateway can route it to a different
  profile's adapter, and only that instance holds the per-message reaction
  state.
- Gated on a truthy `message_id`, which is what every adapter's own hook
  already checks. Synthetic drains (`/goal` continuations, wake-ups, CLI
  hand-offs) carry no id and stay silent; `interrupt_message` and leftover
  `/steer` carry no event at all.
- Cancellation maps to CANCELLED rather than FAILURE, matching
  _process_message_background — Telegram clears the marker on CANCELLED and
  Signal deliberately leaves it, so the distinction is load-bearing.

Outcome is SUCCESS unless the recursion raises. That mirrors the existing
non-queued contract, where a run returning `failed: True` still delivers a
diagnostic message and reports SUCCESS; making the outcome track agent
failure is a separate change that should apply to both paths at once.

No new config, no new env var, no new hook, no change to message
construction or role alternation.

Fixes #72502

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
Mira Solari
2026-07-26 23:06:43 -07:00
committed by Teknium
parent f7c79efbac
commit 8828a3db37
2 changed files with 328 additions and 12 deletions

View File

@@ -2750,6 +2750,7 @@ from gateway.platforms.base import (
EphemeralReply,
MessageEvent,
MessageType,
ProcessingOutcome,
_prefix_within_utf16_limit,
_reply_anchor_for_event,
build_auto_tts_output_path,
@@ -3321,6 +3322,36 @@ def _dequeue_pending_event(adapter, session_key: str) -> MessageEvent | None:
return adapter.get_pending_message(session_key)
async def _run_followup_processing_hook(
adapter,
event: MessageEvent | None,
hook_name: str,
*args,
) -> None:
"""Fire a platform processing-lifecycle hook for a runner-drained follow-up.
A message that arrives mid-turn is parked in the adapter's pending slot and
drained in-band by ``_run_agent``, never by
``BasePlatformAdapter._process_message_background`` — which owns the only
other call site for these hooks. Without firing them here, the read-receipt
reaction every adapter renders from ``on_processing_start`` is silently
skipped for queued, interrupting, and steer-demoted messages.
No-ops unless there is a real inbound platform message to acknowledge:
interrupt text and leftover ``/steer`` carry no event at all, and synthetic
drains (``/goal`` continuations, wake-ups, CLI hand-offs) carry no
``message_id`` — the same field every adapter's own hook already gates on.
"""
if event is None or adapter is None:
return
if not getattr(event, "message_id", None):
return
run_hook = getattr(adapter, "_run_processing_hook", None)
if not callable(run_hook):
return
await run_hook(hook_name, event, *args)
_INTERRUPT_REASON_STOP = "Stop requested"
_INTERRUPT_REASON_RESET = "Session reset requested"
_INTERRUPT_REASON_TIMEOUT = "Execution timed out (inactivity)"
@@ -30558,6 +30589,28 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
except Exception:
pass
# Acknowledge the follow-up the way an idle-session message is
# already acknowledged. This in-band drain is the only place a
# queued/interrupting message ever runs — base.py's
# _process_message_background, which owns the sole other call
# site for the processing hooks, is never entered for it — so
# without this every platform that renders a read receipt from
# on_processing_start silently skips mid-turn messages.
# Resolve the adapter from the follow-up's OWN source: a
# multiplexed gateway can route it to a different profile's
# adapter than the turn we are completing, and only that
# instance holds the per-message reaction state. Fired here
# rather than below so the cache re-baseline stays adjacent to
# the recursive call it exists to protect.
_hook_adapter = (
self._adapter_for_source(next_source)
if pending_event is not None
else None
)
await _run_followup_processing_hook(
_hook_adapter, pending_event, "on_processing_start",
)
# Re-baseline the cached agent's message_count snapshot before
# recursing into the in-band queued (/queue) follow-up turn.
# The first turn has completed and flushed its own user +
@@ -30574,18 +30627,40 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
# what the follow-up's guard will consult. Fail-safe in helper.
await self._refresh_agent_cache_message_count(session_key, session_id)
followup_result = await self._run_agent(
message=next_message,
context_prompt=context_prompt,
history=updated_history,
source=next_source,
session_id=session_id,
session_key=next_session_key,
run_generation=run_generation,
_interrupt_depth=_interrupt_depth + 1,
event_message_id=next_message_id,
channel_prompt=next_channel_prompt,
message_type=next_message_type,
try:
followup_result = await self._run_agent(
message=next_message,
context_prompt=context_prompt,
history=updated_history,
source=next_source,
session_id=session_id,
session_key=next_session_key,
run_generation=run_generation,
_interrupt_depth=_interrupt_depth + 1,
event_message_id=next_message_id,
channel_prompt=next_channel_prompt,
message_type=next_message_type,
)
except asyncio.CancelledError:
# Matches _process_message_background: a cancelled turn is
# not a failure, and adapters that special-case CANCELLED
# (Telegram clears the marker, Signal leaves it) rely on the
# distinction. Best-effort — a re-delivered cancellation
# can pre-empt the await, exactly as it can in base.py.
await _run_followup_processing_hook(
_hook_adapter, pending_event, "on_processing_complete",
ProcessingOutcome.CANCELLED,
)
raise
except BaseException:
await _run_followup_processing_hook(
_hook_adapter, pending_event, "on_processing_complete",
ProcessingOutcome.FAILURE,
)
raise
await _run_followup_processing_hook(
_hook_adapter, pending_event, "on_processing_complete",
ProcessingOutcome.SUCCESS,
)
return _preserve_queued_followup_history_offset(result, followup_result)
finally:

View File

@@ -0,0 +1,241 @@
"""Processing-hook parity for queued follow-up turns.
A message that arrives while a turn is already running is parked in the
adapter's ``_pending_messages`` slot and drained *in-band* by
``GatewayRunner._run_agent`` rather than by
``BasePlatformAdapter._process_message_background``. The runner-side drain
must still fire the ``on_processing_start`` / ``on_processing_complete``
lifecycle hooks, otherwise every platform that renders a read-receipt
reaction from those hooks (Slack 👀, Discord, Telegram, Feishu, Matrix,
Signal, ...) silently skips the acknowledgement for mid-turn messages.
"""
import importlib
import sys
import types
from types import SimpleNamespace
import pytest
from gateway.config import Platform, PlatformConfig
from gateway.platforms.base import (
BasePlatformAdapter,
MessageEvent,
MessageType,
ProcessingOutcome,
SendResult,
)
from gateway.session import SessionSource
class HookRecordingAdapter(BasePlatformAdapter):
"""Adapter that records the processing-hook lifecycle it is driven through."""
def __init__(self):
super().__init__(PlatformConfig(enabled=True, token="***"), Platform.TELEGRAM)
self.started: list = []
self.completed: list = []
async def connect(self) -> bool:
return True
async def disconnect(self) -> None:
return None
async def send(self, chat_id, content, reply_to=None, metadata=None) -> SendResult:
return SendResult(success=True, message_id="sent-1")
async def send_typing(self, chat_id, metadata=None) -> None:
return None
async def stop_typing(self, chat_id) -> None:
return None
async def get_chat_info(self, chat_id: str):
return {"id": chat_id}
async def on_processing_start(self, event: MessageEvent) -> None:
self.started.append(getattr(event, "message_id", None))
async def on_processing_complete(self, event, outcome) -> None:
self.completed.append((getattr(event, "message_id", None), outcome))
class _TwoTurnAgent:
calls: list = []
def __init__(self, **kwargs):
self.tools = []
def run_conversation(self, message, conversation_history=None, task_id=None):
type(self).calls.append(message)
return {
"final_response": f"done-{len(type(self).calls)}",
"messages": [],
"api_calls": 1,
}
class _RaisingSecondTurnAgent:
calls: list = []
def __init__(self, **kwargs):
self.tools = []
def run_conversation(self, message, conversation_history=None, task_id=None):
type(self).calls.append(message)
if len(type(self).calls) >= 2:
raise RuntimeError("boom in the queued follow-up turn")
return {
"final_response": "done-1",
"messages": [],
"api_calls": 1,
}
def _make_runner(adapter):
gateway_run = importlib.import_module("gateway.run")
runner = object.__new__(gateway_run.GatewayRunner)
runner.adapters = {adapter.platform: adapter}
runner._voice_mode = {}
runner._prefill_messages = []
runner._ephemeral_system_prompt = ""
runner._reasoning_config = None
runner._provider_routing = {}
runner._fallback_model = None
runner._session_db = None
runner._running_agents = {}
runner._session_run_generation = {}
runner.hooks = SimpleNamespace(loaded_hooks=False)
runner.config = SimpleNamespace(
thread_sessions_per_user=False,
group_sessions_per_user=False,
stt_enabled=False,
)
runner._model = "openai/gpt-4.1-mini"
runner._base_url = None
return runner
def _install_fake_agent(monkeypatch, tmp_path, agent_cls):
fake_dotenv = types.ModuleType("dotenv")
fake_dotenv.load_dotenv = lambda *args, **kwargs: None
monkeypatch.setitem(sys.modules, "dotenv", fake_dotenv)
fake_run_agent = types.ModuleType("run_agent")
fake_run_agent.AIAgent = agent_cls
monkeypatch.setitem(sys.modules, "run_agent", fake_run_agent)
gateway_run = importlib.import_module("gateway.run")
monkeypatch.setattr(gateway_run, "_hermes_home", tmp_path)
monkeypatch.setattr(
gateway_run, "_resolve_runtime_agent_kwargs", lambda: {"api_key": "***"}
)
SESSION_KEY = "agent:main:telegram:dm:4242"
def _source():
return SessionSource(platform=Platform.TELEGRAM, chat_id="4242", chat_type="dm")
@pytest.mark.asyncio
async def test_queued_followup_fires_processing_hooks(monkeypatch, tmp_path):
"""The runner-drained follow-up gets the same start/complete hooks as a
message that arrives while the session is idle."""
_TwoTurnAgent.calls = []
_install_fake_agent(monkeypatch, tmp_path, _TwoTurnAgent)
adapter = HookRecordingAdapter()
runner = _make_runner(adapter)
adapter._pending_messages[SESSION_KEY] = MessageEvent(
text="the follow-up",
message_type=MessageType.TEXT,
source=_source(),
message_id="queued-1",
)
result = await runner._run_agent(
message="the first turn",
context_prompt="",
history=[],
source=_source(),
session_id="sess-hooks",
session_key=SESSION_KEY,
)
# The follow-up really did run in-band.
assert result["final_response"] == "done-2"
assert _TwoTurnAgent.calls == ["the first turn", "the follow-up"]
# ...and it was acknowledged through the lifecycle hooks.
assert adapter.started == ["queued-1"]
assert adapter.completed == [("queued-1", ProcessingOutcome.SUCCESS)]
@pytest.mark.asyncio
async def test_queued_followup_failure_completes_the_hook(monkeypatch, tmp_path):
"""A follow-up turn that blows up still closes its hook, so a platform
never strands a 'still working' marker on the user's message."""
_RaisingSecondTurnAgent.calls = []
_install_fake_agent(monkeypatch, tmp_path, _RaisingSecondTurnAgent)
adapter = HookRecordingAdapter()
runner = _make_runner(adapter)
adapter._pending_messages[SESSION_KEY] = MessageEvent(
text="the doomed follow-up",
message_type=MessageType.TEXT,
source=_source(),
message_id="queued-2",
)
with pytest.raises(RuntimeError):
await runner._run_agent(
message="the first turn",
context_prompt="",
history=[],
source=_source(),
session_id="sess-hooks-failure",
session_key=SESSION_KEY,
)
assert adapter.started == ["queued-2"]
assert adapter.completed == [("queued-2", ProcessingOutcome.FAILURE)]
@pytest.mark.asyncio
async def test_synthetic_followup_is_not_acknowledged(monkeypatch, tmp_path):
"""Drains with no inbound platform message — /goal continuations, wake-ups,
CLI hand-offs — carry no message_id and must stay silent: there is nothing
on the platform to react to."""
_TwoTurnAgent.calls = []
_install_fake_agent(monkeypatch, tmp_path, _TwoTurnAgent)
adapter = HookRecordingAdapter()
runner = _make_runner(adapter)
adapter._pending_messages[SESSION_KEY] = MessageEvent(
text="synthetic continuation",
message_type=MessageType.TEXT,
source=_source(),
message_id=None,
)
result = await runner._run_agent(
message="the first turn",
context_prompt="",
history=[],
source=_source(),
session_id="sess-hooks-synthetic",
session_key=SESSION_KEY,
)
# It still ran — we only suppressed the acknowledgement, not the turn.
assert result["final_response"] == "done-2"
assert _TwoTurnAgent.calls == ["the first turn", "synthetic continuation"]
assert adapter.started == []
assert adapter.completed == []