fix(gateway): don't hand a completion to complete-only adapters; cover raw-envelope events
Self-review of the previous commit found three defects in it. All three are about the *guard*, not the mechanism. 1. Complete-only adapters were handed an unpaired completion. Google Chat (google_chat/adapter.py:2766) and webhook (webhook.py:901) implement on_processing_complete WITHOUT on_processing_start, and theirs is end-of-cycle teardown, not a reaction: Google Chat reaps the typing card, patching it to "(no reply)"/"(interrupted)", and webhook ends the per-delivery session. The drain fires before the follow-up's reply is delivered — delivery happens after the whole chain unwinds back into _process_message_background — so a Google Chat space would get a permanent "(no reply)" tombstone on every queued follow-up, and the real answer would then land as a separate message. Now gated on the adapter actually overriding on_processing_start: we bracket, so both halves must be ours. 2. The message_id gate made the fix a no-op on Signal, which the previous commit message claimed to fix. SignalAdapter never sets message_id (signal.py:749-766) — its hook keys off raw_message["sender"] and ["timestamp_ms"] via _extract_reaction_target — and Discord's start hook reads raw_message too. Gate is now message_id OR raw_message; synthetic drains still carry neither, so /goal continuations stay silent. 3. Cancellation was classified unconditionally as CANCELLED. base.py:5862-5868 only reports CANCELLED for a task in _expected_cancelled_tasks (/stop, /new, /reset, adapter cleanup) and downgrades anything else to FAILURE. Signal and Matrix deliberately LEAVE the marker in place on CANCELLED, so the previous version would strand exactly the marker this PR exists to clear. Now mirrors base.py via _followup_cancel_outcome(). Also moves _refresh_agent_cache_message_count inside the try. It awaits DB I/O and guards it with `except Exception`, which does not catch cancellation, so a /stop landing there escaped both handlers and stranded the marker. Ordering is unchanged, so the prompt-cache adjacency the previous commit describes still holds. Two new tests, both failing against the previous commit: test_complete_only_adapter_is_left_alone and test_raw_envelope_only_followup_is_acknowledged. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
121
gateway/run.py
121
gateway/run.py
@@ -3322,6 +3322,59 @@ def _dequeue_pending_event(adapter, session_key: str) -> MessageEvent | None:
|
||||
return adapter.get_pending_message(session_key)
|
||||
|
||||
|
||||
def _followup_processing_hooks_apply(adapter, event: MessageEvent | None) -> bool:
|
||||
"""Whether a runner-drained follow-up should be bracketed by the hooks.
|
||||
|
||||
Two conditions, both necessary:
|
||||
|
||||
* **A real inbound platform message to acknowledge.** Adapters key their
|
||||
in-progress marker either off ``message_id`` (Slack, Telegram, Feishu,
|
||||
Matrix, Photon) or off the raw envelope — Signal reads
|
||||
``raw_message["sender"]``/``["timestamp_ms"]`` and never sets
|
||||
``message_id`` at all, Discord reads ``raw_message`` — so either field
|
||||
means there is something to react to. Synthetic drains (``/goal``
|
||||
continuations, wake-ups, CLI hand-offs, startup auto-resume) carry
|
||||
neither and must stay silent.
|
||||
* **The adapter overrides ``on_processing_start``.** We *bracket*, so both
|
||||
halves have to belong to us. An adapter that implements only
|
||||
``on_processing_complete`` — Google Chat reaps its typing card there,
|
||||
webhook ends its per-delivery session — would otherwise be handed a
|
||||
completion for a turn whose reply has not been delivered yet, because
|
||||
delivery happens after the whole drain chain unwinds back into
|
||||
``_process_message_background``.
|
||||
"""
|
||||
if adapter is None or event is None:
|
||||
return False
|
||||
if not (getattr(event, "message_id", None) or getattr(event, "raw_message", None)):
|
||||
return False
|
||||
start_hook = getattr(type(adapter), "on_processing_start", None)
|
||||
return start_hook is not None and start_hook is not BasePlatformAdapter.on_processing_start
|
||||
|
||||
|
||||
def _followup_cancel_outcome(adapter) -> ProcessingOutcome:
|
||||
"""Classify a cancelled follow-up exactly as ``_process_message_background``
|
||||
does: only cancels the adapter itself routed (``/stop``, ``/new``,
|
||||
``/reset``, adapter cleanup) count as CANCELLED; anything else is a failure.
|
||||
|
||||
The distinction is load-bearing rather than cosmetic — Signal and Matrix
|
||||
deliberately leave the in-progress marker in place on CANCELLED, so
|
||||
reporting an *unexpected* cancellation as CANCELLED would strand it.
|
||||
"""
|
||||
expected = getattr(adapter, "_expected_cancelled_tasks", None)
|
||||
if expected is None:
|
||||
return ProcessingOutcome.FAILURE
|
||||
try:
|
||||
current = asyncio.current_task()
|
||||
except RuntimeError:
|
||||
current = None
|
||||
if current is None:
|
||||
return ProcessingOutcome.FAILURE
|
||||
try:
|
||||
return ProcessingOutcome.CANCELLED if current in expected else ProcessingOutcome.FAILURE
|
||||
except TypeError:
|
||||
return ProcessingOutcome.FAILURE
|
||||
|
||||
|
||||
async def _run_followup_processing_hook(
|
||||
adapter,
|
||||
event: MessageEvent | None,
|
||||
@@ -3334,17 +3387,12 @@ async def _run_followup_processing_hook(
|
||||
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.
|
||||
reaction adapters render 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.
|
||||
See ``_followup_processing_hooks_apply`` for when this is a no-op.
|
||||
"""
|
||||
if event is None or adapter is None:
|
||||
return
|
||||
if not getattr(event, "message_id", None):
|
||||
if not _followup_processing_hooks_apply(adapter, event):
|
||||
return
|
||||
run_hook = getattr(adapter, "_run_processing_hook", None)
|
||||
if not callable(run_hook):
|
||||
@@ -30594,7 +30642,7 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
|
||||
# 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
|
||||
# without this every adapter 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
|
||||
@@ -30611,23 +30659,30 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
|
||||
_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 +
|
||||
# assistant rows to the SessionDB, so the cross-process
|
||||
# coherence guard (#45966) — which this recursive _run_agent
|
||||
# call re-enters — would otherwise see the grown on-disk count
|
||||
# against the stale build-time snapshot and rebuild the agent
|
||||
# on THIS process's OWN writes, destroying the prompt-cache
|
||||
# prefix #46237 was merged to preserve. The existing
|
||||
# re-baseline in _handle_message_with_agent only runs after the
|
||||
# whole _run_agent chain unwinds — too late for the in-band
|
||||
# follow-up. Use the same (session_key, session_id) the
|
||||
# recursive call runs under so the snapshot matches exactly
|
||||
# what the follow-up's guard will consult. Fail-safe in helper.
|
||||
await self._refresh_agent_cache_message_count(session_key, session_id)
|
||||
|
||||
# Everything from here to the recursion runs inside the try, so
|
||||
# that once the marker is on the message every exit closes it —
|
||||
# including a /stop landing on the re-baseline's DB await, which
|
||||
# would otherwise strand the marker (that helper guards its own
|
||||
# I/O with `except Exception`, which does not catch cancellation).
|
||||
try:
|
||||
# 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 + assistant rows to the SessionDB, so the
|
||||
# cross-process coherence guard (#45966) — which this
|
||||
# recursive _run_agent call re-enters — would otherwise see
|
||||
# the grown on-disk count against the stale build-time
|
||||
# snapshot and rebuild the agent on THIS process's OWN
|
||||
# writes, destroying the prompt-cache prefix #46237 was
|
||||
# merged to preserve. The existing re-baseline in
|
||||
# _handle_message_with_agent only runs after the whole
|
||||
# _run_agent chain unwinds — too late for the in-band
|
||||
# follow-up. Use the same (session_key, session_id) the
|
||||
# recursive call runs under so the snapshot matches exactly
|
||||
# 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,
|
||||
@@ -30642,14 +30697,16 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
|
||||
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.
|
||||
# Classified the same way _process_message_background does:
|
||||
# an *expected* cancel (/stop, /new, /reset, cleanup) is
|
||||
# CANCELLED, anything else is a failure. Signal and Matrix
|
||||
# deliberately leave the marker in place on CANCELLED, so
|
||||
# misclassifying here would strand it. 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,
|
||||
_followup_cancel_outcome(_hook_adapter),
|
||||
)
|
||||
raise
|
||||
except BaseException:
|
||||
|
||||
@@ -239,3 +239,75 @@ async def test_synthetic_followup_is_not_acknowledged(monkeypatch, tmp_path):
|
||||
|
||||
assert adapter.started == []
|
||||
assert adapter.completed == []
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_raw_envelope_only_followup_is_acknowledged(monkeypatch, tmp_path):
|
||||
"""Signal never sets message_id — its hook keys off the raw envelope
|
||||
(sender + timestamp_ms) — and Discord's reads raw_message. An event
|
||||
carrying only a raw envelope is still a real inbound message."""
|
||||
_TwoTurnAgent.calls = []
|
||||
_install_fake_agent(monkeypatch, tmp_path, _TwoTurnAgent)
|
||||
|
||||
adapter = HookRecordingAdapter()
|
||||
runner = _make_runner(adapter)
|
||||
|
||||
adapter._pending_messages[SESSION_KEY] = MessageEvent(
|
||||
text="signal-shaped follow-up",
|
||||
message_type=MessageType.TEXT,
|
||||
source=_source(),
|
||||
message_id=None,
|
||||
raw_message={"sender": "+15550100", "timestamp_ms": 1700000000000},
|
||||
)
|
||||
|
||||
await runner._run_agent(
|
||||
message="the first turn",
|
||||
context_prompt="",
|
||||
history=[],
|
||||
source=_source(),
|
||||
session_id="sess-hooks-raw",
|
||||
session_key=SESSION_KEY,
|
||||
)
|
||||
|
||||
assert adapter.started == [None]
|
||||
assert adapter.completed == [(None, ProcessingOutcome.SUCCESS)]
|
||||
|
||||
|
||||
class CompleteOnlyAdapter(HookRecordingAdapter):
|
||||
"""Google Chat and webhook implement on_processing_complete WITHOUT
|
||||
on_processing_start; theirs is end-of-cycle teardown (reap the typing
|
||||
card / end the delivery session), not a reaction."""
|
||||
|
||||
on_processing_start = BasePlatformAdapter.on_processing_start
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_complete_only_adapter_is_left_alone(monkeypatch, tmp_path):
|
||||
"""We bracket, so both halves must belong to us. An adapter that only
|
||||
implements the completion half must not be handed a completion here: at
|
||||
this point the follow-up's reply has not been delivered yet, so its
|
||||
teardown would fire against a live turn."""
|
||||
_TwoTurnAgent.calls = []
|
||||
_install_fake_agent(monkeypatch, tmp_path, _TwoTurnAgent)
|
||||
|
||||
adapter = CompleteOnlyAdapter()
|
||||
runner = _make_runner(adapter)
|
||||
|
||||
adapter._pending_messages[SESSION_KEY] = MessageEvent(
|
||||
text="the follow-up",
|
||||
message_type=MessageType.TEXT,
|
||||
source=_source(),
|
||||
message_id="queued-3",
|
||||
)
|
||||
|
||||
result = await runner._run_agent(
|
||||
message="the first turn",
|
||||
context_prompt="",
|
||||
history=[],
|
||||
source=_source(),
|
||||
session_id="sess-hooks-complete-only",
|
||||
session_key=SESSION_KEY,
|
||||
)
|
||||
|
||||
assert result["final_response"] == "done-2"
|
||||
assert adapter.completed == []
|
||||
|
||||
Reference in New Issue
Block a user