[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
Los mantenedores suelen responder en 2 días
@rohityan ya está trabajando en esto.
Desde el 30/9/2026.
Evaluación
Este issue todavía no se ha evaluado.
Descripción
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)
default_request_handler_v2.py#L466:on_message_send_streamsubscribes withinclude_initial_task=False, so the handler never emits the stored task. The leadingTaskhas to come from the executor.active_task.py#L274-L283: when the executor enqueues aTaskand the task already exists,_handle_initial_tasklogs the ERROR, ignores the event, and still setsself.message_to_save = None. The follow-up message was stored inmessage_to_savewhen the request started, so_handle_task_modification_eventno longer adds it to history.active_task.py#L291-L299: theInvalidAgentResponseErrorcheck added in #979 ("Agent should enqueue Task before …") only fires when no task exists yet, so a continuation without a leadingTaskpasses silently.test_end_to_end.py#L589-L596:test_end_to_end_input_requiredcurrently expects the follow-up stream to bestatus_update, artifact_update, status_update.
Related
- #965 was fixed for new tasks by #979.
- The earlier attempt #964 injected
Task(SUBMITTED)in the handler, but only whennot params.message.task_id, so continuations were out of scope there too. - #1188 concerns resuming from a stale snapshot after
input_required. It's related but different: that report is about multi-replica state, and this one reproduces on a single in-memory replica.
Possible directions (maintainers will know better)
- Handler: for a message that carries
task_idof an existing task, emit the stored task as the first stream event, ason_subscribe_to_taskalready does withinclude_initial_task=True. Ideally that snapshot would already include the follow-up message. - Consumer: when the executor enqueues a
Taskwhose id matches the existing task, treat it as "emit current state": keep or applymessage_to_saveand 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
- Lenguaje dominante
- Python
- Estrellas
- 2.2k
- Forks
- 509
- Merge medio
- 3 d 18 h
- PR fusionados (30 d)
- 44
Preparar el entorno
- Sin Dockerfile ni archivo de Docker Compose
- Tiene una plantilla de pull request
- Leer la guía de contribución
Primeros pasos
- Lee el issue completo y luego la guía de contribución del proyecto.
- Comenta en el issue que vas a ocuparte — evita que dos personas hagan lo mismo.
- Haz un fork del repositorio y trabaja en una rama.
- Abre un pull request que haga referencia al número del issue.
Más de a2aproject/a2a-python
-
[Bug]: REST task/request id sanitizationPosiblemente ocupada @Linux2010 la tomó hace 102 días. Abiertomaintainers-only
Dificultad 2/5 1-3 horas Aptitud para principiantes 65/100
a2aproject/a2a-python#805 · 1 comentario ·
Los mantenedores suelen responder en 2 días
-
[Feat]: Cluster mode: detect and recover tasks abandoned by a crashed replica (heartbeat/lease)Posiblemente ocupada @rohityan la tomó hace 1 día. Abierto
a2aproject/a2a-python#1324 · 2 asignados ·
Los mantenedores suelen responder en 2 días
-
[Bug]: Cluster mode: SubscribeToTask on a replica holding a paused ActiveTask returns a stale INPUT_REQUIRED snapshot and closesPosiblemente ocupada @rohityan la tomó hace 1 día. Abierto
a2aproject/a2a-python#1323 · 2 asignados ·
Los mantenedores suelen responder en 2 días
-
[Bug]: After a streamed task is cancelled, its background producer never finishesPosiblemente ocupada @rohityan la tomó hace 1 día. Abierto
a2aproject/a2a-python#1322 · 2 asignados ·
Los mantenedores suelen responder en 2 días
-
v0.3 gRPC and REST SendMessage without configuration run non-blockingPosiblemente ocupada @rohityan la tomó hace 2 días. Abierto
a2aproject/a2a-python#1321 · 2 asignados ·
Los mantenedores suelen responder en 2 días
Todos los issues de a2aproject/a2a-python
Issues similares
-
first
Dificultad 2/5 1-3 horas Aptitud para principiantes 72/100
AcademySoftwareFoundation/rmtc#54 · 1 comentario ·
-
feature/cohorts feature/feature-flags team/feature-flags
Dificultad 2/5 1-3 horas Aptitud para principiantes 74/100
Los mantenedores suelen responder en 1 día
-
License examples/ as MITPosiblemente ocupada @PGrayCS la tomó hoy. Abiertodocumentation enhancement example good first issue
Dificultad 2/5 1-3 horas Aptitud para principiantes 84/100
speedyk-005/yasbd-lib#383 ·
Los mantenedores suelen responder en 1 día
-
Dificultad 2/5 1-3 horas Aptitud para principiantes 68/100
interactions-py/interactions.py#1827 ·
-
Managed start can fail when OpenVMM reads its control capability before NVX writes itPosiblemente ocupada @ppenna la tomó hoy. Abiertobug
Dificultad 2/5 1-3 horas Aptitud para principiantes 76/100
Los mantenedores suelen responder en 1 día