fix(gateway): keep swapped command guards and id-less rewrites bounded in drain back-off

The back-off drain's slot-empty exit released whatever guard was current, so a
/stop, /new or /reset guard swapped in during the sleep was deleted, defeating
the #48300 guard-swap protection. Capture the guard owned by the drain at spawn
and release only that; also flush the text debounce buffer before popping the
slot, like every other task exit, so a debounced text isn't orphaned.

A pre_gateway_dispatch rewrite copy of an event with no message_id was never
matched to the dispatched event, so it still hot-looped (#123229). Fall back to
the copied timestamp when message_id is empty (a genuine new message gets a
fresh one); a plain message_id==/timestamp== form breaks the runner's re-queue
of a new object with the same id.

Cap the back-off at 1s: nothing wakes the sleep, so a genuine message merged
into the slot meanwhile waited up to 5s; 1 dispatch/s is still ~250x below the
unbounded loop and needs no new wake plumbing.
This commit is contained in:
kshitijk4poor
2026-09-27 14:40:30 +05:30
committed by kshitij
parent 1288744c79
commit c0f8ac6645

View File

@@ -4587,20 +4587,26 @@ class BasePlatformAdapter(ABC):
self._finish_session_task(session_key, interrupt_event)
_REQUEUE_BACKOFF_INITIAL_SECONDS = 0.25
_REQUEUE_BACKOFF_MAX_SECONDS = 5.0
# Kept at 1s: nothing wakes the back-off sleep, so a genuine message merged into the slot
# meanwhile waits out the remainder; 1 dispatch/s is still ~250x below the unbounded loop.
_REQUEUE_BACKOFF_MAX_SECONDS = 1.0
def _requeue_backoff_delay(self, session_key: str, pending_event: MessageEvent,
dispatched_event: MessageEvent) -> float:
"""Delay before re-dispatching the queued follow-up.
Only the event this task just dispatched coming straight back (same object, or the same
``message_id`` for rewrite-hook copies) backs off: the handler put it back because the
Only the event this task just dispatched coming straight back (same object, or for
rewrite-hook copies the same ``message_id`` / id-less same ``timestamp``) backs off: the handler put it back because the
session is busy elsewhere, and re-dispatching it at once hot-loops for the whole busy
window (#123229). Any other follow-up resets the counter and runs immediately. The first
bounce stays immediate (restart auto-resume relies on one self-bounce), then back off
exponentially to a cap. Defers, never drops."""
same = pending_event is dispatched_event or bool(
pending_event.message_id and pending_event.message_id == dispatched_event.message_id)
# Rewrite-hook copies (dataclasses.replace) keep message_id and timestamp; an id-less
# event falls back to the copied timestamp (a genuine new message gets a fresh one).
same = pending_event is dispatched_event or (
pending_event.message_id == dispatched_event.message_id
and (bool(pending_event.message_id)
or pending_event.timestamp == dispatched_event.timestamp))
if not same:
self._requeue_counts.pop(session_key, None)
return 0.0
@@ -4624,15 +4630,21 @@ class BasePlatformAdapter(ABC):
sleeping, so a cancel/discard during the back-off needs no put-back and can't drop a
newer message."""
self._clear_session_guard(session_key)
# Capture the guard this drain owns now: a /stop//new guard swapped in during the
# back-off must survive the slot-empty exit (#48300).
guard = self._active_sessions.get(session_key)
self._track_session_task(
session_key, asyncio.create_task(self._drain_after(pending_event, session_key, delay)))
session_key,
asyncio.create_task(self._drain_after(pending_event, session_key, delay, guard)))
async def _drain_after(self, pending_event: MessageEvent, session_key: str, delay: float) -> None:
async def _drain_after(self, pending_event: MessageEvent, session_key: str, delay: float,
guard: Optional[asyncio.Event] = None) -> None:
if delay > 0:
await asyncio.sleep(delay)
await self._flush_text_debounce_now(session_key) # as every other task exit does
pending_event = self._pending_messages.pop(session_key, None)
if pending_event is None: # consumed elsewhere during the back-off
self._cleanup_finished_session_task(session_key, self._active_sessions.get(session_key))
self._cleanup_finished_session_task(session_key, guard)
return
await self._process_message_background(pending_event, session_key)