fix(gateway): tighten busy re-queue back-off identity check and test matrix
Review cleanups on the drain back-off: - `_drain_after` now requires `guard`: a None default silently meant "release whatever guard is current", which is the guard-swap bug the parameter fixes. - The identity rule is stated once (docstring, wrapped), and the redundant `pending_event is dispatched_event` disjunct is dropped: the same object always has an equal message_id, and when that id is empty its timestamp equals itself, so the remaining comparison already covers it. - Tests drop the `_Adapter` alias and cut the hot-loop matrix from 6 to 4 explicit cases (plain, rewrite with id, rewrite without id, steer). The demotion route does not interact with the identity axis, and each case spends a fixed 1s measuring window.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user