diff --git a/src/a2a/server/request_handlers/default_request_handler_v2.py b/src/a2a/server/request_handlers/default_request_handler_v2.py index 872a3bfa2..372f11a10 100644 --- a/src/a2a/server/request_handlers/default_request_handler_v2.py +++ b/src/a2a/server/request_handlers/default_request_handler_v2.py @@ -82,7 +82,7 @@ def __init__( # noqa: PLR0913 task_store: TaskStore, agent_card: AgentCard, queue_manager: Any - | None = None, # Kept for backward compat in signature + | None = None, # Accepted for signature compat; ignored in v2 (warns) push_config_store: PushNotificationConfigStore | None = None, push_sender: PushNotificationSender | None = None, request_context_builder: RequestContextBuilder | None = None, @@ -92,6 +92,15 @@ def __init__( # noqa: PLR0913 ] | None = None, ) -> None: + if queue_manager is not None: + logger.warning( + 'A queue_manager was passed to DefaultRequestHandlerV2, but it ' + 'is not used: v2 delegates event streaming to an in-memory ' + 'ActiveTaskRegistry, so custom or distributed QueueManager ' + 'implementations are ignored. For multi-replica event ' + 'streaming, either use LegacyRequestHandler or route ' + 'subscription requests to the replica holding the task.' + ) self.agent_executor = agent_executor self.task_store = task_store self._agent_card = agent_card diff --git a/tests/server/request_handlers/test_default_request_handler_v2.py b/tests/server/request_handlers/test_default_request_handler_v2.py index b276fb77a..9c2e18bb2 100644 --- a/tests/server/request_handlers/test_default_request_handler_v2.py +++ b/tests/server/request_handlers/test_default_request_handler_v2.py @@ -20,7 +20,7 @@ ) from a2a.server.agent_execution.active_task_registry import ActiveTaskRegistry from a2a.server.context import ServerCallContext -from a2a.server.events import EventQueue +from a2a.server.events import EventQueue, InMemoryQueueManager from a2a.server.events.event_queue_v2 import EventQueueSource from a2a.server.request_handlers import DefaultRequestHandlerV2 from a2a.server.tasks import ( @@ -140,6 +140,36 @@ def test_init_default_dependencies(): assert handler._request_context_builder._task_store == task_store +def test_init_warns_when_queue_manager_passed(caplog): + """A caller-supplied queue_manager is not honored in v2, so passing one + must emit a warning instead of being silently ignored (issue #1135).""" + queue_manager = InMemoryQueueManager() + with caplog.at_level(logging.WARNING): + DefaultRequestHandlerV2( + agent_executor=MockAgentExecutor(), + task_store=InMemoryTaskStore(), + agent_card=create_default_agent_card(), + queue_manager=queue_manager, + ) + assert any( + 'queue_manager' in record.message and record.levelno == logging.WARNING + for record in caplog.records + ) + + +def test_init_no_warning_without_queue_manager(caplog): + """No warning is emitted when queue_manager is omitted.""" + with caplog.at_level(logging.WARNING): + DefaultRequestHandlerV2( + agent_executor=MockAgentExecutor(), + task_store=InMemoryTaskStore(), + agent_card=create_default_agent_card(), + ) + assert not any( + 'queue_manager' in record.message for record in caplog.records + ) + + @pytest.mark.asyncio async def test_on_get_task_not_found(): """Test on_get_task when task_store.get returns None."""