From c0f8ac6645fb2e6455910ced85fffadcd869fe89 Mon Sep 17 00:00:00 2001 From: kshitijk4poor <82637225+kshitijk4poor@users.noreply.github.com> Date: Sun, 27 Sep 2026 14:40:30 +0530 Subject: [PATCH] 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. --- gateway/platforms/base.py | 28 ++++++++++++++++++++-------- 1 file changed, 20 insertions(+), 8 deletions(-) diff --git a/gateway/platforms/base.py b/gateway/platforms/base.py index db1e533c80..c4a3099b19 100644 --- a/gateway/platforms/base.py +++ b/gateway/platforms/base.py @@ -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)