Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 10 additions & 1 deletion src/a2a/server/request_handlers/default_request_handler_v2.py
Original file line number Diff line number Diff line change
Expand Up @@ -82,16 +82,25 @@
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,
extended_agent_card: AgentCard | None = None,
extended_card_modifier: Callable[
[AgentCard, ServerCallContext], Awaitable[AgentCard]
]
| None = None,
) -> None:
if queue_manager is not None:

Check notice on line 95 in src/a2a/server/request_handlers/default_request_handler_v2.py

View workflow job for this annotation

GitHub Actions / Lint Code Base

Copy/pasted code

see src/a2a/server/request_handlers/default_request_handler.py (97-119)
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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down Expand Up @@ -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."""
Expand Down
Loading