What happened?
When a client cancels a task while it is streaming that task, the cancel itself works. The client gets TASK_STATE_CANCELED, the stream receives the CANCELED update, and the stream closes. But the task's background producer (the asyncio task named producer:<task id>) never finishes. It stays running until the whole request handler is shut down with aclose(). Because the consumer task waits for the producer, it stays running too, and the task is never removed from the ActiveTaskRegistry.
On a long-running server, every cancelled task that had a stream open leaves two tasks and one registry entry behind until the process stops.
This happens with any agent. It does not matter whether the agent's cancel() publishes CANCELED itself or leaves that to the SDK.
Version: a2a-sdk 1.2.1. The same code is on main today.
Cause
When a subscriber reads from its queue, it is expected to mark each item as done with task_done(). In ActiveTask.subscribe, two internal events, _RequestStarted and _RequestCompleted, are skipped with continue before the try/finally that calls tapped_queue.task_done(). Those two items are never marked done.
if isinstance(event, _RequestCompleted):
if request_id is not None and event.request_id == request_id:
return
continue # never marked done
elif isinstance(event, _RequestStarted):
continue # never marked done
try:
yield event
finally:
tapped_queue.task_done()
When a task finishes normally, this does no harm: the _RequestCompleted event makes the subscriber leave before the producer cleans up.
A cancel stops the agent before _RequestCompleted is sent. The producer then cleans up with await self._event_queue_subscribers.close(immediate=False), which waits until every subscriber queue has marked all of its items done. The stream's queue still holds the unmarked _RequestStarted, so the wait never ends. When the hang happens, the subscriber queue is empty but still has one item not marked done.
How to reproduce
pip install "a2a-sdk==1.2.1", then run:
import asyncio
import uuid
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.a2a_pb2 import (
AgentCapabilities,
AgentCard,
CancelTaskRequest,
Message,
Part,
Role,
SendMessageRequest,
Task,
TaskState,
TaskStatus,
)
class SlowAgent(AgentExecutor):
"""Starts a task, then works "forever" until it is cancelled."""
async def execute(self, context: RequestContext, event_queue: EventQueue) -> None:
await event_queue.enqueue_event(
Task(
id=context.task_id,
context_id=context.context_id,
status=TaskStatus(state=TaskState.TASK_STATE_SUBMITTED),
)
)
await TaskUpdater(event_queue, context.task_id, context.context_id).start_work()
await asyncio.Event().wait() # long-running work
async def cancel(self, context: RequestContext, event_queue: EventQueue) -> None:
await TaskUpdater(event_queue, context.task_id, context.context_id).cancel()
def producers_still_running() -> list[str]:
return [
t.get_name()
for t in asyncio.all_tasks()
if t.get_name().startswith("producer:") and not t.done()
]
async def main() -> None:
handler = DefaultRequestHandler(
SlowAgent(),
InMemoryTaskStore(),
AgentCard(capabilities=AgentCapabilities(streaming=True)),
)
call = ServerCallContext()
request = SendMessageRequest(
message=Message(
message_id=str(uuid.uuid4()), role=Role.ROLE_USER, parts=[Part(text="hi")]
)
)
# 1. A client streams the task.
stream = handler.on_message_send_stream(request, call)
first = await anext(stream)
task_id = first.id
print("streaming task", task_id)
# 2. While it streams, the client cancels it.
async def cancel_soon() -> None:
await asyncio.sleep(0.2)
result = await handler.on_cancel_task(CancelTaskRequest(id=task_id), call)
print("cancel answered:", TaskState.Name(result.status.state))
canceller = asyncio.create_task(cancel_soon())
# 3. The stream receives the cancel and ends normally.
async for event in stream:
print("stream got:", type(event).__name__, TaskState.Name(event.status.state))
print("stream closed")
await canceller
# 4. Everything is done, yet the task's producer is still running.
await asyncio.sleep(2)
print("producers still running after 2 s:", producers_still_running())
# 5. Only shutting down the whole handler frees it.
await handler.aclose()
print("after handler.aclose():", producers_still_running())
asyncio.run(main())
Output:
streaming task 1247c447-a6ed-4fdd-b700-88b8e66248c0
stream got: TaskStatusUpdateEvent TASK_STATE_WORKING
cancel answered: TASK_STATE_CANCELED
stream got: TaskStatusUpdateEvent TASK_STATE_CANCELED
stream closed
producers still running after 2 s: ['producer:1247c447-a6ed-4fdd-b700-88b8e66248c0']
after handler.aclose(): []
Expected: the list on the "after 2 s" line is empty.
The result is the same if SlowAgent.cancel() does nothing (pass), so that the SDK writes CANCELED itself.
Suggested fix
Mark the two internal events as done before skipping them, in ActiveTask.subscribe:
if isinstance(event, _RequestStarted | _RequestCompleted):
tapped_queue.task_done()
if isinstance(event, _RequestCompleted):
...
With this change applied to 1.2.1, the script above prints producers still running after 2 s: [].
Related
What happened?
When a client cancels a task while it is streaming that task, the cancel itself works. The client gets
TASK_STATE_CANCELED, the stream receives the CANCELED update, and the stream closes. But the task's background producer (the asyncio task namedproducer:<task id>) never finishes. It stays running until the whole request handler is shut down withaclose(). Because the consumer task waits for the producer, it stays running too, and the task is never removed from theActiveTaskRegistry.On a long-running server, every cancelled task that had a stream open leaves two tasks and one registry entry behind until the process stops.
This happens with any agent. It does not matter whether the agent's
cancel()publishes CANCELED itself or leaves that to the SDK.Version:
a2a-sdk1.2.1. The same code is onmaintoday.Cause
When a subscriber reads from its queue, it is expected to mark each item as done with
task_done(). InActiveTask.subscribe, two internal events,_RequestStartedand_RequestCompleted, are skipped withcontinuebefore thetry/finallythat callstapped_queue.task_done(). Those two items are never marked done.When a task finishes normally, this does no harm: the
_RequestCompletedevent makes the subscriber leave before the producer cleans up.A cancel stops the agent before
_RequestCompletedis sent. The producer then cleans up withawait self._event_queue_subscribers.close(immediate=False), which waits until every subscriber queue has marked all of its items done. The stream's queue still holds the unmarked_RequestStarted, so the wait never ends. When the hang happens, the subscriber queue is empty but still has one item not marked done.How to reproduce
pip install "a2a-sdk==1.2.1", then run:Output:
Expected: the list on the "after 2 s" line is empty.
The result is the same if
SlowAgent.cancel()does nothing (pass), so that the SDK writes CANCELED itself.Suggested fix
Mark the two internal events as done before skipping them, in
ActiveTask.subscribe:With this change applied to 1.2.1, the script above prints
producers still running after 2 s: [].Related
await self._event_queue_subscribers.close(immediate=False). feat(server): add aclose() to drain ActiveTask background tasks (#1101) #1105 closed it by addingaclose(), which frees the stuck producer when the server shuts down. This issue is the cause of that stuck state after a cancel, and the fix above prevents it while the server is running.