diff --git a/gateway/platforms/base.py b/gateway/platforms/base.py index c4a3099b19..e59789add4 100644 --- a/gateway/platforms/base.py +++ b/gateway/platforms/base.py @@ -4595,18 +4595,18 @@ class BasePlatformAdapter(ABC): 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 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.""" - # 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)) + Only the event this task just dispatched coming straight back backs off: the same + ``message_id``, or for an id-less event the same ``timestamp`` (rewrite-hook + ``dataclasses.replace`` copies keep both; a genuine new message gets a fresh timestamp). + 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.""" + # The identical object always matches too: its id equals itself, and when empty the + # timestamp comparison does. + same = (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 @@ -4638,7 +4638,7 @@ class BasePlatformAdapter(ABC): 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, - guard: Optional[asyncio.Event] = None) -> None: + guard: Optional[asyncio.Event]) -> None: if delay > 0: await asyncio.sleep(delay) await self._flush_text_debounce_now(session_key) # as every other task exit does diff --git a/tests/gateway/test_busy_requeue_hot_loop.py b/tests/gateway/test_busy_requeue_hot_loop.py index f50e21cd05..99c909dd98 100644 --- a/tests/gateway/test_busy_requeue_hot_loop.py +++ b/tests/gateway/test_busy_requeue_hot_loop.py @@ -31,9 +31,6 @@ from gateway.session import SessionEntry, SessionSource, build_session_key from tests.gateway.restart_test_helpers import RestartTestAdapter -_Adapter = RestartTestAdapter - - def _source() -> SessionSource: return SessionSource(platform=Platform.TELEGRAM, user_id="u1", chat_id="c1", user_name="tester", chat_type="dm") @@ -88,13 +85,15 @@ def _runner_with_running_agent(adapter, *, compression_in_flight): @pytest.mark.asyncio -@pytest.mark.parametrize("busy_mode", ["interrupt", "steer"]) -@pytest.mark.parametrize("rewrite_hook,message_id", [ - (False, "m1"), (True, "m1"), - (True, None), # id-less rewrite copy: only the copied timestamp ties it to the dispatch +@pytest.mark.parametrize("rewrite_hook,message_id,busy_mode", [ + (False, "m1", "interrupt"), + (True, "m1", "interrupt"), + (True, None, "interrupt"), # id-less rewrite copy: only the copied timestamp ties it back + # The demotion route is orthogonal to the identity axis; steer covers the other route. + (False, "m1", "steer"), ]) async def test_requeued_busy_event_does_not_hot_loop(rewrite_hook, message_id, busy_mode): - adapter = _Adapter() + adapter = RestartTestAdapter() runner, agent, sk = _runner_with_running_agent(adapter, compression_in_flight=True) # interrupt: demoted to queue (compression in flight); steer: agent refuses -> queue fallback. runner._busy_input_mode = busy_mode @@ -135,7 +134,7 @@ async def test_requeued_event_runs_once_the_agent_finishes(): """The back-off must defer, not drop: once the running agent is gone the event is processed. And it must key on the runner's demotion, not on ``None``: every streamed turn returns None, so chained genuine follow-ups after it must each dispatch immediately.""" - adapter = _Adapter() + adapter = RestartTestAdapter() runner, agent, sk = _runner_with_running_agent(adapter, compression_in_flight=True) handled, starts, ends = [], [], [] real_handle = runner._handle_message