Hacktoberfest 2026: le issue che i maintainer hanno segnato per ottobre, aperte e adatte ai principianti. Sfoglia le issue Hacktoberfest

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

Aperta
#1,296 0 commenti 0 reazioni 2 assegnatari Vedi su GitHub

I maintainer di solito rispondono entro 1 giorno

@rohityan ci sta già lavorando.

Dal 2/10/2026.

Valutazione

Questa issue non è ancora stata valutata.

Descrizione

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
  • #1101 / #1105: aclose() for teardown at shutdown, which the workaround reuses
  • #1121, #1123: producer lifecycle at turn boundaries and teardown
  • #1170: cancel now writes a terminal state for parked input-required tasks
  • #1188: ActiveTaskRegistry state after input_required
Lingua principale
Python
Stelle
2.2k
Fork
499
Merge medio
3g 20h
PR unite (30g)
32

Preparare l'ambiente

Come iniziare

  1. Leggi tutta la issue e poi la guida ai contributi del progetto.
  2. Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
  3. Fai un fork del repository e lavora su un branch.
  4. Apri una pull request che faccia riferimento al numero della issue.

Altre issue di a2aproject/a2a-python

Tutte le issue di a2aproject/a2a-python

Issue simili

Altre issue su Python

Ricevi le nuove issue nella tua casella

Un breve riepilogo di issue GitHub adatte ai principianti.