Skip to content

[Bug]: DefaultRequestHandlerV2 keeps the ActiveTask (producer, consumer, 2 dispatchers) alive forever after a direct Message or input-required response #1296

Description

@christophstach

Version

a2a-sdk 1.1.2, 1.1.5, 1.2.0 and 1.2.1 (all show the same result), Python 3.14, standard DefaultRequestHandler (which resolves to DefaultRequestHandlerV2), InMemoryTaskStore.

Summary

For every request whose execution does not end in a terminal task state, the handler leaves the request's ActiveTask running indefinitely: ActiveTask._run_producer, ActiveTask._run_consumer and two EventQueueSource._dispatch_loop tasks stay pending, and the entry stays in the ActiveTaskRegistry. This happens in two common cases:

  1. The executor answers with a direct Message. Per the spec this is the response for simple interactions that do not need task tracking (§3.1.1: the agent "MAY return a direct Message response for simple interactions"; §3.1.2: for a message-only stream "No task tracking or updates are provided"). No follow-up can ever continue such a request, yet its ActiveTask is never released.
  2. The task ends in TASK_STATE_INPUT_REQUIRED. The ActiveTask waits for a follow-up. If the client never sends one, which is common when the multi-turn state lives elsewhere, it waits forever. There is no timeout.

Unlike #1101 / #1121 / #1123, which are about teardown at shutdown or turn boundaries, this happens during normal operation. The pending task count grows linearly with traffic until the process restarts.

Reproduction

Standalone script, SDK only (attached below). It sends 50 requests through DefaultRequestHandler.on_message_send with three executors and counts the asyncio tasks left afterwards:

$ uv run --no-project --with "a2a-sdk[all]==1.2.1" python a2a_activetask_leak_repro.py
a2a-sdk 1.2.1
completed       50 requests ->    0 pending asyncio tasks
message         50 requests ->  200 pending asyncio tasks {'EventQueueSource._dispatch_loop': 100, 'ActiveTask._run_producer': 50, 'ActiveTask._run_consumer': 50}
input-required  50 requests ->  200 pending asyncio tasks {'EventQueueSource._dispatch_loop': 100, 'ActiveTask._run_consumer': 50, 'ActiveTask._run_producer': 50}

1.1.2 and 1.1.5 print the same numbers. The producer is parked on await self._request_queue.get(), and the dispatchers wait on their queues.

a2a_activetask_leak_repro.py
import asyncio
import importlib.metadata
from collections import Counter

from a2a.helpers import new_message, new_text_part
from a2a.server.agent_execution import AgentExecutor, RequestContext
from a2a.server.context import ServerCallContext
from a2a.server.events import EventQueue
from a2a.server.request_handlers import DefaultRequestHandler
from a2a.server.tasks import InMemoryTaskStore, TaskUpdater
from a2a.types import (
    AgentCapabilities, AgentCard, AgentInterface, Message, Role,
    SendMessageRequest, Task, TaskState, TaskStatus,
)

REQUESTS = 50


class Executor(AgentExecutor):
    def __init__(self, mode: str) -> None:
        self.mode = mode

    async def execute(self, context: RequestContext, event_queue: EventQueue) -> None:
        if self.mode == "message":
            await event_queue.enqueue_event(new_message([new_text_part("hello")]))
            return
        await event_queue.enqueue_event(
            Task(id=context.task_id, context_id=context.context_id,
                 status=TaskStatus(state=TaskState.TASK_STATE_SUBMITTED))
        )
        updater = TaskUpdater(event_queue, context.task_id, context.context_id)
        reply = updater.new_agent_message([new_text_part("hello")])
        if self.mode == "completed":
            await updater.complete(reply)
        else:
            await updater.requires_input(reply)

    async def cancel(self, context: RequestContext, event_queue: EventQueue) -> None:
        pass


def card() -> AgentCard:
    return AgentCard(
        name="repro", description="repro", version="0.0.1",
        default_input_modes=["text"], default_output_modes=["text"],
        capabilities=AgentCapabilities(streaming=False),
        supported_interfaces=[AgentInterface(url="http://localhost/", protocol_binding="JSONRPC", protocol_version="1.0")],
    )


async def run(mode: str) -> None:
    handler = DefaultRequestHandler(agent_executor=Executor(mode), task_store=InMemoryTaskStore(), agent_card=card())
    baseline = asyncio.all_tasks()
    for i in range(REQUESTS):
        request = SendMessageRequest(message=Message(message_id=f"m-{i}", role=Role.ROLE_USER, parts=[new_text_part("hi")]))
        await asyncio.wait_for(handler.on_message_send(request, ServerCallContext()), timeout=5)
    await asyncio.sleep(1)
    leftover = asyncio.all_tasks() - baseline - {asyncio.current_task()}
    names = Counter(t.get_coro().__qualname__ for t in leftover)
    print(f"{mode:<15} {REQUESTS} requests -> {len(leftover):>4} pending asyncio tasks {dict(names) if names else ''}")


async def main() -> None:
    print(f"a2a-sdk {importlib.metadata.version('a2a-sdk')}")
    for mode in ("completed", "message", "input-required"):
        await run(mode)


asyncio.run(main())

Expected behavior

  • Direct Message response: once the Message has been delivered (non-streaming result or the single event of a message-only stream), the handler releases the ActiveTask: it stops the producer and consumer, closes the queues, and removes the registry entry. That's the same as what happens today when a task reaches a terminal state.
  • Input-required: the ActiveTask should not have to outlive the request. Since the interrupted task is persisted in the TaskStore, the ActiveTask could be released and recreated from the store when a follow-up arrives, as get_or_create already supports. Alternatively, a configurable idle timeout would bound it.

Impact

The leak is invisible in short tests and only shows up in long-running servers. In our service (a few A2A requests per second per pod, most answered with a direct Message), pending asyncio tasks grew by about 10 per minute per pod, i.e. tens of thousands per day. Memory grew with them. Anything that walks all tasks pays for it too: with a sampling profiler that inspects asyncio tasks, the profiler thread reached a full core after several hours.

Workaround we use (for reference)

We subclass DefaultRequestHandler. We override _setup_active_task to capture the request's ActiveTask (there is no public way to get it), and call ActiveTask.aclose() (from #1105) when on_message_send returns a Message, or when a stream yields only a Message. This brings pending tasks back to the baseline, and the registry entry is removed as well. For input-required tasks we call the public on_cancel_task after an idle timeout. Since #1170 that writes CANCELED and releases the producer.

Relying on the private _setup_active_task hook is fragile, so a fix in the handler, or a public hook to release a request's ActiveTask, would be much appreciated.

Related

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Labels

No labels
No labels

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions