[Bug]: DefaultRequestHandlerV2 keeps the ActiveTask (producer, consumer, 2 dispatchers) alive forever after a direct Message or input-required response
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:
- 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 directMessageresponse 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. - 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
Messageresponse: once theMessagehas 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, asget_or_createalready 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
- Nessun Dockerfile né file Docker Compose
- Ha un modello di pull request
- Leggi la guida per i contributori
Come iniziare
- Leggi tutta la issue e poi la guida ai contributi del progetto.
- Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
- Fai un fork del repository e lavora su un branch.
- Apri una pull request che faccia riferimento al numero della issue.
Altre issue di a2aproject/a2a-python
-
maintainers-only
Difficoltà 2/5 1-3 ore Idoneità per principianti 65/100
a2aproject/a2a-python#805 · 1 commento ·
I maintainer di solito rispondono entro 1 giorno
-
v0.3 gRPC and REST SendMessage return no history when history_length is not setForse già presa @rohityan l’ha presa oggi. Aperta
a2aproject/a2a-python#1305 · 2 assegnatari ·
I maintainer di solito rispondono entro 1 giorno
-
[Feat]: Change Httpx to Httpx2Forse già presa @rohityan l’ha presa 2 giorni fa. Apertacomponent: client status: needs review
a2aproject/a2a-python#1288 · 2 commenti · 1 assegnatario ·
I maintainer di solito rispondono entro 1 giorno
-
[Bug]: Streaming follow-up on an existing task does not begin with a Task; enqueuing the current task drops the follow-up message from historyForse già presa @rohityan l’ha presa 2 giorni fa. Apertacomponent: server status:awaiting response
a2aproject/a2a-python#1285 · 2 commenti · 1 assegnatario ·
I maintainer di solito rispondono entro 1 giorno
-
component: core
a2aproject/a2a-python#1278 · 9 commenti · 1 assegnatario ·
I maintainer di solito rispondono entro 1 giorno
Tutte le issue di a2aproject/a2a-python
Issue simili
-
upstream update
Difficoltà 2/5 1-3 ore Idoneità per principianti 75/100
conan-io/conan-center-index#31098 ·
I maintainer di solito rispondono entro 2 giorni
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 78/100
john-kurkowski/tldextract#382 ·
-
comp/tools duplicate P2 sweeper:risk-compatibility tool/mcp type/bug
Difficoltà 1/5 Meno di un'ora Idoneità per principianti 88/100
NousResearch/hermes-agent#132042 ·
I maintainer di solito rispondono entro 1 giorno
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 78/100
deepset-ai/haystack#13092 ·
I maintainer di solito rispondono entro 1 giorno
-
Difficoltà 1/5 1-3 ore Idoneità per principianti 85/100
feder-cr/invisible_playwright_mcp#1408 ·
I maintainer di solito rispondono entro 1 giorno