diff --git a/gateway/run_agent_cache.py b/gateway/run_agent_cache.py index 952f6ffa58..462a50ee12 100644 --- a/gateway/run_agent_cache.py +++ b/gateway/run_agent_cache.py @@ -507,7 +507,23 @@ class GatewayAgentCacheMixin: else: await adapter.interrupt_session_activity(session_key, source.chat_id) if adapter and hasattr(adapter, "get_pending_message"): - adapter.get_pending_message(session_key) # consume and discard + # Discard a stale human follow-up (the slot held only user text when /stop started doing + # this, 59575d6a917) — but an internal wake (async-delegation completion, notify+wake) + # shares the slot now and was claim-settled on admission, so dropping it loses it for + # good and the session idles until the next user message (#114456). Leave it parked for + # the adapter's post-command drain; a wake queued behind a discarded human head is + # promoted out of the overflow FIFO for the same reason. Whether a wake may still run + # against a session /new just closed is decided where it is processed + # (_resolve_async_delegation_session fails closed), not here. + parked = adapter.get_pending_message(session_key) + wake = parked if getattr(parked, "internal", False) else None + if wake is None: + overflow = self._overflow_queue(session_key) or [] + wake = next((e for e in overflow if getattr(e, "internal", False)), None) + if wake is not None: + overflow.remove(wake) + if wake is not None: + adapter._pending_messages[session_key] = wake if state is not None: state.persistent.pending_command_text = None if release_running_state: diff --git a/tests/gateway/test_interrupt_keeps_parked_internal_wake.py b/tests/gateway/test_interrupt_keeps_parked_internal_wake.py new file mode 100644 index 0000000000..deec07cd1f --- /dev/null +++ b/tests/gateway/test_interrupt_keeps_parked_internal_wake.py @@ -0,0 +1,104 @@ +"""Invariant: interrupting a session (/stop, /new, /reset) discards a parked human follow-up but +never a parked internal wake (async-delegation completion, notify+wake). See #114456. + +Real ``BasePlatformAdapter`` pending slot + post-command drain, real +``GatewayRunner._interrupt_and_clear_session`` with the reason pairs its three callers pass. +""" +import asyncio + +import pytest + +from gateway.config import GatewayConfig, Platform, PlatformConfig +from gateway.platforms.base import BasePlatformAdapter, SendResult +from gateway.platforms.event import MessageEvent, MessageType +from gateway.run import _INTERRUPT_REASON_RESET, _INTERRUPT_REASON_STOP, GatewayRunner +from gateway.session import SessionSource + +# (interrupt_reason, invalidation_reason) as passed by _busy_stop_command, _handle_stop_command +# (pending sentinel) and _busy_new_command. +_COMMAND_REASONS = [ + (_INTERRUPT_REASON_STOP, "stop_command"), + (_INTERRUPT_REASON_STOP, "stop_command_pending"), + (_INTERRUPT_REASON_RESET, "new_command"), +] + + +class _Adapter(BasePlatformAdapter): + def __init__(self): + super().__init__(PlatformConfig(enabled=True), Platform.TELEGRAM) + self.restarted = [] + + @property + def name(self): + return "telegram" + + async def connect(self, *, is_reconnect=False): + return True + + async def disconnect(self): + pass + + async def send(self, chat_id, content, reply_to=None, metadata=None): + return SendResult(success=True) + + async def get_chat_info(self, chat_id): + return {"id": chat_id, "type": "private"} + + def _start_session_processing(self, event, session_key, *, interrupt_event=None): + self.restarted.append(event) + return True + + +def _gateway(): + adapter = _Adapter() + runner = object.__new__(GatewayRunner) + runner.config = GatewayConfig() + runner.adapters = {Platform.TELEGRAM: adapter} + source = SessionSource(platform=Platform.TELEGRAM, chat_id="c1", chat_type="dm", user_id="u1") + key = adapter._event_session_key(MessageEvent(text="", message_type=MessageType.TEXT, source=source)) + return adapter, runner, source, key + + +def _event(source, text, *, internal): + event = MessageEvent(text=text, message_type=MessageType.TEXT, source=source, internal=internal) + event._gateway_accepted = True + return event + + +async def _run_command(runner, adapter, source, key, reasons): + interrupt_reason, invalidation_reason = reasons + command_guard = asyncio.Event() + adapter._active_sessions[key] = command_guard + await runner._interrupt_and_clear_session( + key, source, interrupt_reason=interrupt_reason, invalidation_reason=invalidation_reason, + ) + await adapter._drain_pending_after_session_command(key, command_guard) + + +@pytest.mark.asyncio +@pytest.mark.parametrize("reasons", _COMMAND_REASONS, ids=[r[1] for r in _COMMAND_REASONS]) +async def test_interrupt_keeps_parked_internal_wake_and_discards_human_followup(reasons): + adapter, runner, source, key = _gateway() + wake = _event(source, "[ASYNC DELEGATION BATCH COMPLETE] 1 task done", internal=True) + adapter._pending_messages[key] = wake + await _run_command(runner, adapter, source, key, reasons) + assert adapter.restarted == [wake] + assert key not in adapter._pending_messages and key not in adapter._active_sessions + + adapter.restarted.clear() + adapter._pending_messages[key] = _event(source, "stale human follow-up", internal=False) + await _run_command(runner, adapter, source, key, reasons) + assert adapter.restarted == [] + assert key not in adapter._pending_messages + + +@pytest.mark.asyncio +async def test_interrupt_promotes_internal_wake_queued_behind_discarded_human_head(): + adapter, runner, source, key = _gateway() + adapter._pending_messages[key] = _event(source, "human head", internal=False) + later_human = _event(source, "second human", internal=False) + wake = _event(source, "[ASYNC DELEGATION BATCH COMPLETE] 1 task done", internal=True) + runner._session_state(key).conversation.queued_events.extend([later_human, wake]) + await _run_command(runner, adapter, source, key, _COMMAND_REASONS[0]) + assert adapter.restarted == [wake] + assert runner._overflow_queue(key) == [later_human] diff --git a/website/docs/developer-guide/gateway-session-lifecycle.md b/website/docs/developer-guide/gateway-session-lifecycle.md index e6c46286fc..0597d9e538 100644 --- a/website/docs/developer-guide/gateway-session-lifecycle.md +++ b/website/docs/developer-guide/gateway-session-lifecycle.md @@ -455,6 +455,13 @@ Called at the drain site after the slot was consumed. If there's an overflow ite ### Clearing Queued events for a session are cleared on `/new` and `/reset` (via `_handle_reset_command`). +`/stop` drops the single-slot follow-up the user sent during the interrupted turn. An +**internal** wake parked in either store (an async-delegation completion notice, a kanban/cron +`notify+wake`) survives all three commands: `_interrupt_and_clear_session` leaves it in the slot +(promoting it out of the overflow when a discarded human follow-up held the slot) so the +post-command drain starts it right away instead of the session idling until the next user +message. Whether a wake pinned to a session that `/new` just closed may still run is decided at +processing time (`_resolve_async_delegation_session`, fail-closed). ### FIFO Invariant