fix(gateway): /stop, /new and /reset keep a parked internal wake instead of discarding it
_interrupt_and_clear_session popped the adapter's single pending slot and dropped whatever it
held ("consume and discard", 59575d6a91). That was right when the slot only ever carried the
user's stale follow-up text, but internal wakes (async-delegation completion notices,
kanban/cron notify+wake) now park in the same slot and are claim-settled the moment the adapter
admits them, so the pop lost them for good: the drain that runs after the command found an
empty slot and the session idled until the next user message (#114456, ~6 min stall in the
reported session).
Now the human follow-up is still discarded, but an internal wake stays parked (promoted out of
the overflow FIFO when a discarded human head occupied the slot) and the post-command drain
starts it immediately. Applies to every caller of the helper — /stop (busy fast path, handler,
pending sentinel, thread sibling), /new and /reset — since a wake that arrives a second after
/new runs against the fresh session anyway; whether a completion pinned to the closed session
may run stays with _resolve_async_delegation_session (fail-closed).
Salvaged from #114538 (@whyyagswhy, earliest filer): the pop-and-re-park mechanism is theirs;
widened here to the /new and /reset callers and the overflow promotion, and the invariant tests
rewritten against the real adapter drain. #114540 (@JoaoMarcos44) reached the same fix
independently; its stop-reason taxonomy and overflow analysis informed the class coverage.
Co-authored-by: joaomarcos <joaomarcosdias444@gmail.com>
This commit is contained in:
@@ -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:
|
||||
|
||||
104
tests/gateway/test_interrupt_keeps_parked_internal_wake.py
Normal file
104
tests/gateway/test_interrupt_keeps_parked_internal_wake.py
Normal file
@@ -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]
|
||||
@@ -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
|
||||
|
||||
|
||||
Reference in New Issue
Block a user