test(gateway): harden busy re-queue hot-loop regression

Measure the 1s window after the first dispatch so a cold first handler call under load
cannot false-green the no-fix mutation; parametrize over interrupt demotion and steer
fallback; cover cancel during back-off keeping the event queued and the /stop tail
replaying it. Reuse RestartTestAdapter instead of a duplicate adapter.
This commit is contained in:
kshitijk4poor
2026-09-27 14:03:02 +05:30
committed by kshitij
parent b6e6f00790
commit 1288744c79

View File

@@ -26,37 +26,19 @@ from unittest.mock import AsyncMock, MagicMock
import pytest
from gateway.config import GatewayConfig, Platform, PlatformConfig
from gateway.platforms.base import BasePlatformAdapter, SendResult
from gateway.platforms.event import MessageEvent
from gateway.session import SessionEntry, SessionSource, build_session_key
from tests.gateway.restart_test_helpers import RestartTestAdapter
class _Adapter(BasePlatformAdapter):
class _Adapter(RestartTestAdapter):
def __init__(self):
super().__init__(PlatformConfig(enabled=True), Platform.TELEGRAM)
self.sent = []
super().__init__()
self.typing_calls = 0
@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):
self.sent.append(content)
return SendResult(success=True)
async def send_typing(self, chat_id, metadata=None):
self.typing_calls += 1
async def get_chat_info(self, chat_id):
return {"id": chat_id, "type": "private"}
def _source() -> SessionSource:
return SessionSource(platform=Platform.TELEGRAM, user_id="u1", chat_id="c1",
@@ -112,12 +94,18 @@ def _runner_with_running_agent(adapter, *, compression_in_flight):
@pytest.mark.asyncio
@pytest.mark.parametrize("busy_mode", ["interrupt", "steer"])
@pytest.mark.parametrize("rewrite_hook", [False, True])
async def test_requeued_busy_event_does_not_hot_loop(rewrite_hook):
async def test_requeued_busy_event_does_not_hot_loop(rewrite_hook, busy_mode):
adapter = _Adapter()
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
agent.steer.return_value = False
busy_calls = []
async def handler(event):
busy_calls.append(time.monotonic())
# A pre_gateway_dispatch "rewrite" hands the runner a dataclasses.replace copy, so the
# event it re-queues is a different object than the one the adapter dispatched.
if rewrite_hook:
@@ -129,10 +117,15 @@ async def test_requeued_busy_event_does_not_hot_loop(rewrite_hook):
# The new adapter holds no guard for the session (the reconnect replaced it mid-turn).
assert sk not in adapter._active_sessions
await adapter.handle_message(MessageEvent(text="still there?", source=_source(), message_id="m1"))
# Measure after the first dispatch so a cold first handler call can't eat the window.
for _ in range(1000):
if busy_calls:
break
await asyncio.sleep(0.01)
base_dispatches, base_typing = len(busy_calls), adapter.typing_calls
await asyncio.sleep(1.0)
dispatches = runner._session_has_compression_in_flight.await_count
typing = adapter.typing_calls
dispatches = len(busy_calls) - base_dispatches
typing = adapter.typing_calls - base_typing
await adapter.cancel_background_tasks()
agent.interrupt.assert_not_called()
@@ -163,7 +156,11 @@ async def test_requeued_event_runs_once_the_agent_finishes():
adapter.set_message_handler(handler)
await adapter.handle_message(MessageEvent(text="queued msg", source=_source(), message_id="m2"))
await asyncio.sleep(0.6)
# A cancel during the back-off (e.g. /stop) must not lose the queued event: it stays pending.
await adapter.cancel_session_processing(sk, discard_pending=False)
assert adapter._pending_messages[sk].text == "queued msg"
runner._running_agents.pop(sk) # the long turn finishes
await adapter._drain_pending_after_session_command(sk, asyncio.Event()) # /stop tail replays it
for _ in range(1000):
if handled:
break