A2A / 远程 AgentNote 95
第二章:A2A 流式状态、超时和取消
代码在 deepagentsrc/a2ateach/02streamingandcancel.py。
代码在 deepagent_src/a2a_teach/02_streaming_and_cancel.py。
第 1 章的任务只有最终结果。本章用 SlowExecutor 制造一个可控的长任务,验证 A2A 的三件事:
- Executor 先发布初始
Task,再用TaskUpdater.start_work()发布working状态。 - Executor 发布一个进度 artifact;它不是最终结果,
last_chunk=False表示后面还有内容。 - Client 只等待 250ms。超时后显式发送
CancelTaskRequest,并确认服务端执行器的cancel()被调用、Task 进入TASK_STATE_CANCELED,且不会产生result最终 artifact。
uv run python -m deepagent_src.a2a_teach.02_streaming_and_cancel
预期输出类似:
stream_event: task
stream_event: status_update
status: TASK_STATE_WORKING
stream_event: artifact_update
client_wait: timed out after 250ms
cancel_state: TASK_STATE_CANCELED
a2a streaming, timeout, and cancel ok
超时不等于取消
asyncio.timeout(0.25) 只限制本地 Client 等待流的时间;它不会可靠地停止远端工作。超时后必须把取得的 task_id 传给 CancelTaskRequest,远端 AgentExecutor.cancel() 必须发布 TASK_STATE_CANCELED。
SDK 在收到取消请求时会取消正在执行的 execute() 协程,然后调用 cancel()。所以 execute() 中要让长耗时操作能够响应取消,例如 await 可取消的网络请求、子进程管理或支持取消的模型流。已发布的 progress artifact 会保留在 Task 历史中,取消只阻止后续工作;不能只在 Client 超时后丢弃 HTTP 连接,否则服务端任务可能继续消耗资源。
本章的慢任务使用 asyncio.sleep(5),这是为了稳定验证 A2A 协议事件和取消生命周期,不代表已经验证了某个 LLM 提供商的 generation cancel。对真实 Deep Agent,要把模型流、工具调用和子进程的取消策略单独接入并压测。
何时使用流式 A2A
流式调用适合远端任务持续数秒以上、调用方需要进度或需要随时取消的场景。短小且只关心最终结果的远端调用,使用第 1 章的非流式模式更简单。
相关资源
查看示例代码:deepagent_src/a2a_teach/02_streaming_and_cancel.py
from __future__ import annotations import asyncio import httpx import uvicorn from a2a.client import ClientConfig, create_client from a2a.helpers.proto_helpers import new_text_message, new_text_part from a2a.server.agent_execution import AgentExecutor, RequestContext from a2a.server.request_handlers import DefaultRequestHandler from a2a.server.routes.agent_card_routes import create_agent_card_routes from a2a.server.routes.fastapi_routes import add_a2a_routes_to_fastapi from a2a.server.routes.jsonrpc_routes import create_jsonrpc_routes from a2a.server.tasks import InMemoryTaskStore, TaskUpdater from a2a.types import ( AgentCapabilities, AgentCard, AgentInterface, CancelTaskRequest, Role, SendMessageConfiguration, SendMessageRequest, Task, TaskState, TaskStatus, ) from a2a.utils.constants import TransportProtocol from fastapi import FastAPI AGENT_URL = "http://127.0.0.1:8765" class SlowExecutor(AgentExecutor): """用可控慢任务验证 A2A 生命周期,不依赖模型响应时间。""" def __init__(self) -> None: self.cancel_called = False async def execute(self, context: RequestContext, event_queue) -> None: 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) await updater.start_work() await updater.add_artifact( [new_text_part("progress: remote task is working")], name="progress", last_chunk=False, ) await asyncio.sleep(5) await updater.add_artifact( [new_text_part("final: task completed")], name="result", ) await updater.complete() async def cancel(self, context: RequestContext, event_queue) -> None: self.cancel_called = True updater = TaskUpdater(event_queue, context.task_id, context.context_id) await updater.cancel() def build_a2a_app(executor: AgentExecutor) -> FastAPI: card = AgentCard( name="streaming-cancel-agent", version="0.1.0", description="A streaming A2A teaching agent.", supported_interfaces=[ AgentInterface( protocol_binding=TransportProtocol.JSONRPC, url=f"{AGENT_URL}/", ) ], capabilities=AgentCapabilities(streaming=True), default_input_modes=["text/plain"], default_output_modes=["text/plain"], ) handler = DefaultRequestHandler(executor, InMemoryTaskStore(), card) app = FastAPI() add_a2a_routes_to_fastapi( app, agent_card_routes=create_agent_card_routes(card), jsonrpc_routes=create_jsonrpc_routes(handler, rpc_url="/"), ) return app async def wait_for_server() -> None: async with httpx.AsyncClient() as http_client: for _ in range(50): try: response = await http_client.get( f"{AGENT_URL}/.well-known/agent-card.json" ) except httpx.ConnectError: await asyncio.sleep(0.05) continue if response.is_success: return await asyncio.sleep(0.05) raise RuntimeError("A2A server did not start") async def main() -> None: executor = SlowExecutor() server = uvicorn.Server( uvicorn.Config( build_a2a_app(executor), host="127.0.0.1", port=8765, log_level="error", access_log=False, ) ) server_task = asyncio.create_task(server.serve()) try: await wait_for_server() client = await create_client(AGENT_URL, ClientConfig(streaming=True)) task_id = None event_kinds: list[str] = [] saw_progress = False try: async with asyncio.timeout(0.25): request = SendMessageRequest( message=new_text_message( "请执行一个需要较长时间的远端任务。", role=Role.ROLE_USER, ), configuration=SendMessageConfiguration(), ) async for event in client.send_message(request): kind = event.WhichOneof("payload") event_kinds.append(kind) print("stream_event:", kind) if kind == "task": task_id = event.task.id elif kind == "status_update": task_id = event.status_update.task_id print( "status:", TaskState.Name(event.status_update.status.state), ) elif kind == "artifact_update": saw_progress = True except TimeoutError: print("client_wait: timed out after 250ms") else: raise AssertionError("slow task unexpectedly completed before timeout") assert task_id, "stream did not return a task ID" canceled_task = await client.cancel_task(CancelTaskRequest(id=task_id)) print("cancel_state:", TaskState.Name(canceled_task.status.state)) assert "status_update" in event_kinds, event_kinds assert saw_progress, event_kinds assert executor.cancel_called assert canceled_task.status.state == TaskState.TASK_STATE_CANCELED assert [artifact.name for artifact in canceled_task.artifacts] == ["progress"] print("a2a streaming, timeout, and cancel ok") await client.close() finally: server.should_exit = True await server_task if __name__ == "__main__": asyncio.run(main())