Skip to content

[Bug]: Streaming follow-up on an existing task does not begin with a Task; enqueuing the current task drops the follow-up message from history #1285

Description

@marcinbelczewski

What happened?

When a client continues an existing task with SendStreamingMessage (for example, after TASK_STATE_INPUT_REQUIRED), the stream does not begin with a Task object. The spec requires one, in §3.1.2 Send Streaming Message, Behavior 2:

If the agent returns a Task, the stream MUST begin with the Task object, followed by zero or more TaskStatusUpdateEvent or TaskArtifactUpdateEvent objects.

The executor cannot fix this from its side. Enqueueing the current task first, as the v1.0 migration guide shows (task = context.current_task or new_task_from_user_message(...) followed by enqueue_event(task)), makes the stream start with a Task. But the follow-up user message is then never written to the task history, and an ERROR is logged.

Reproduced on a2a-sdk 1.1.5 and 1.2.1. I found no change to the affected code on main since 1.2.1.

Expected: a streamed follow-up on an existing task begins with a Task (its current state), and the follow-up user message is added to Task.history.

Actual (reproduction below):

Executor 2nd stream history after 2nd turn
enqueues Task only for a new task (as in tests/integration/test_end_to_end.py) status_update, status_update, no leading Task user 1, agent, user 2
always enqueues context.current_task first (migration guide) task, status_update, status_update user 1, agent. User 2 is missing, and ERROR … already exists. Ignoring task replacement. is logged

Where it comes from (v1.2.1)

Related

Possible directions (maintainers will know better)

  1. Handler: for a message that carries task_id of an existing task, emit the stored task as the first stream event, as on_subscribe_to_task already does with include_initial_task=True. Ideally that snapshot would already include the follow-up message.
  2. Consumer: when the executor enqueues a Task whose id matches the existing task, treat it as "emit current state": keep or apply message_to_save and don't log an ERROR. That would make the migration-guide pattern correct.

I'm happy to test a fix. Which direction do you prefer?

Reproduction

Self-contained. Needs a2a-sdk[http-server] (and httpx). Run python repro.py for row 1 and python repro.py --always for row 2.

import asyncio
import logging
import sys

import httpx
from a2a.client import ClientConfig, ClientFactory
from a2a.helpers.proto_helpers import new_task_from_user_message
from a2a.server.agent_execution import AgentExecutor, RequestContext
from a2a.server.events import EventQueue
from a2a.server.request_handlers import DefaultRequestHandler
from a2a.server.routes import create_agent_card_routes, create_jsonrpc_routes
from a2a.server.tasks import TaskUpdater
from a2a.server.tasks.inmemory_task_store import InMemoryTaskStore
from a2a.types import (
    AgentCapabilities,
    AgentCard,
    AgentInterface,
    GetTaskRequest,
    Message,
    Part,
    Role,
    SendMessageRequest,
    TaskState,
)
from a2a.utils import TransportProtocol
from starlette.applications import Starlette

ALWAYS_ENQUEUE_TASK = '--always' in sys.argv


class Executor(AgentExecutor):
    async def execute(self, context: RequestContext, event_queue: EventQueue) -> None:
        task = context.current_task or new_task_from_user_message(context.message)
        if ALWAYS_ENQUEUE_TASK or context.current_task is None:
            await event_queue.enqueue_event(task)
        updater = TaskUpdater(event_queue, task.id, task.context_id)
        await updater.update_status(TaskState.TASK_STATE_WORKING)
        if context.current_task is None:
            await updater.update_status(
                TaskState.TASK_STATE_INPUT_REQUIRED,
                message=updater.new_agent_message([Part(text='Which one?')]),
            )
        else:
            await updater.complete()

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


async def main() -> None:
    card = AgentCard(
        name='repro',
        description='repro',
        version='1.0.0',
        capabilities=AgentCapabilities(streaming=True),
        default_input_modes=['text/plain'],
        default_output_modes=['text/plain'],
        supported_interfaces=[
            AgentInterface(protocol_binding=TransportProtocol.JSONRPC, url='http://testserver')
        ],
    )
    handler = DefaultRequestHandler(
        agent_executor=Executor(), task_store=InMemoryTaskStore(), agent_card=card
    )
    app = Starlette(
        routes=[
            *create_agent_card_routes(agent_card=card, card_url='/'),
            *create_jsonrpc_routes(request_handler=handler, rpc_url='/'),
        ]
    )
    http = httpx.AsyncClient(transport=httpx.ASGITransport(app=app), base_url='http://testserver')
    client = ClientFactory(
        ClientConfig(httpx_client=http, supported_protocol_bindings=[TransportProtocol.JSONRPC])
    ).create(card)

    async def send(text: str, task_id: str | None = None) -> list:
        message = Message(role=Role.ROLE_USER, message_id=f'msg-{text}', parts=[Part(text=text)])
        if task_id:
            message.task_id = task_id
        return [e async for e in client.send_message(SendMessageRequest(message=message))]

    first = await send('first')
    task_id = first[0].task.id
    second = await send('second', task_id)
    task = await client.get_task(GetTaskRequest(id=task_id))

    print('first stream: ', [e.WhichOneof('payload') for e in first])
    print('second stream:', [e.WhichOneof('payload') for e in second])
    print('final state:  ', TaskState.Name(task.status.state))
    print('history:      ', [(Role.Name(m.role), m.message_id) for m in task.history])


if __name__ == '__main__':
    logging.basicConfig(level=logging.ERROR, format='%(levelname)s %(name)s: %(message)s')
    asyncio.run(main())

Environment: a2a-sdk 1.1.5 (PyPI) and 1.2.1 (tag v1.2.1), Python 3.14, macOS, JSON-RPC transport, InMemoryTaskStore.

Relevant log output

$ python repro.py
first stream:  ['task', 'status_update', 'status_update']
second stream: ['status_update', 'status_update']
final state:   TASK_STATE_COMPLETED
history:       [('ROLE_USER', 'msg-first'), ('ROLE_AGENT', '<uuid>'), ('ROLE_USER', 'msg-second')]

$ python repro.py --always
ERROR a2a.server.agent_execution.active_task: Task <uuid> already exists. Ignoring task replacement.
first stream:  ['task', 'status_update', 'status_update']
second stream: ['task', 'status_update', 'status_update']
final state:   TASK_STATE_COMPLETED
history:       [('ROLE_USER', 'msg-first'), ('ROLE_AGENT', '<uuid>')]

Code of Conduct

  • I agree to follow this project's Code of Conduct

Activity

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.status:awaiting response

Type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions