Skip to content

[Bug]: After a streamed task is cancelled, its background producer never finishes #1322

Description

@carstenpiepel

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

Activity

  1. added
    component: serverIssues related to frameworks for agent execution, HTTP/event handling, database persistence logic.
    on Oct 8, 2026
  2. added theissue type on Oct 8, 2026
  3. rohityan commented on Oct 9, 2026

    @rohityan

    Hi @carstenpiepel , thanks for the detailed reproduction script. I was able to reproduce this issue. Your root cause analysis is correct.

    As a workaround after exhausting the cancelled task's event stream, find and explicitly call .cancel() on the residual producer:<task_id> asyncio task so the background consumer unblocks and prunes the task from ActiveTaskRegistry without requiring a full server restart.

    Your proposed fix of calling tapped_queue.task_done() for _RequestStarted and _RequestCompleted in ActiveTask.subscribe is the right approach, and we would welcome a pull request from you implementing it along with a regression test. Let me know if you are interested in contributing I could assign this issue to you.

  4. carstenpiepel commented on Oct 9, 2026

    @carstenpiepel
    Author

    Hi @rohityan, thanks for your prompt response. I am happy to contribute a fix. Please assign the issue to @carstenpiepel. Thank you.

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

Metadata

Metadata

Assignees

Labels

component: serverIssues related to frameworks for agent execution, HTTP/event handling, database persistence logic.

Type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions