Chapter 03 · 上手实操
把它跑起来:一个 A2A server,一个 client,端到端
02 讲清了 wire 上发生什么与三种交互——本章把它跑起来。一个 A2A server 暴露一项 skill,一个 client 发现它的 Agent Card、发 SendMessage、收回 Task 与 Artifact。三阶推进:worked(读完整可运行代码)→ partial(补关键 schema 决策)→ open(改成流式)。读完,你不仅能跑通这条链,还能指出 SDK 在哪几处把 wire 协议藏了起来。
本章你将建立的 schema
- 用
a2a-sdk1.1.0 把 01 章的 Agent Card / Task / Message / Part / Artifact 落成可运行代码——每个对象都映射回它的概念定义。 - server 端 = 子类化
AgentExecutor实现execute、发布一张AgentCard、用 Starlette 路由挂起来;client 端 = 解析 Card → 发SendMessage→ 异步迭代收回结果。 - SDK 在哪藏了 wire:Python 方法是 snake_case
send_message,它在线缆上发出的"method"却是SendMessage——SDK 的人体工学不等于 wire 协议。 - 三个 schema 决策点亲手补:Card 里要不要声明
capabilities.streaming、结果怎么包成 Artifact、input-required状态怎么处理。
本章除 shell 安装命令外的全部代码均标注「未在本机验证」。它们按 a2a-sdk 1.1.0(2026-05-29,实现 spec v1.0)的公开源码逐字核对过导入路径与方法签名,语法完整、可复制粘贴、无 ... 省略;但本教程未在本机实跑这条链。这个领域迭代很快——读到这里时请以 pip show a2a-sdk 的实际版本与官方文档为准。
SendMessage;产出不是裸字符串,而是 execute() 经 EventQueue 包出的 Artifact,沿朱红线回流。capabilities.streaming 声明(§1.2)、结果包成 Artifact(§1.3)、input-required 中途态(§2.5)。3.1环境准备
整章用官方 Python SDK a2a-sdk。版本锁到 1.1.0(2026-05-29 发布,实现 spec v1.0、带 v0.3 兼容模式)。需要 Python 3.10+。最小依赖只有这一个包——它已把服务端的 Starlette 路由、客户端的 httpx 传输都带进来了。
# Python 3.10+
python -m venv .venv && source .venv/bin/activate
# 锁定版本,避免装到 v0.3.x 的旧 API(见下方预防说明)
pip install "a2a-sdk==1.1.0"
# 自检:导入不报错、且版本是 1.1.0 = 装对了
python -c "import a2a; from importlib.metadata import version; print(version('a2a-sdk'))"
装到 v0.3.x 旧版是头号失败源。a2a-sdk 在 1.0 做过破坏性重构:v0.3.x 的客户端入口是 A2AClient.get_client_from_agent_card_url(...),而 1.1.0 改成了 ClientFactory / create_client;服务端的 A2AStarletteApplication 也在 1.0 被路由函数 create_jsonrpc_routes 取代。裸跑 pip install a2a-sdk 可能解析到与本章不符的版本,照着本章代码写就会报 ImportError。永远带上 ==1.1.0(或你确认过的版本),装完先跑那行 python -c 把版本号打出来核对。
3.2Worked · 一个大写转换 agent,端到端
需求:做一个最小的 A2A server,暴露一项 skill「把文本转成大写」。再写一个 client,去发现它的 Agent Card、把一句话发过去、把结果收回来。这条链把 01 章的五个抽象(Agent Card §1.2、Task / Message / Part / Artifact §1.3)和 02 章的 SendMessage(§2.3)一次性落地。
server 端:暴露能力的两件事
一个 A2A server 在代码上拆成两半:一·业务逻辑——子类化 AgentExecutor,在 execute() 里干活、把产出 enqueue 到 EventQueue;二·能力声明——构造一张 AgentCard 公开发布。先看业务逻辑这半。
# 未在本机验证 —— a2a-sdk 1.1.0(spec v1.0)
from a2a.helpers import get_message_text, new_task_from_user_message, new_text_part
from a2a.server.agent_execution import AgentExecutor, RequestContext
from a2a.server.events import EventQueue
from a2a.server.tasks import TaskUpdater
from a2a.types.a2a_pb2 import TaskState
class UpperCaseExecutor(AgentExecutor):
"""把客户端发来的文本转成大写,作为一个 Artifact 返回。"""
async def execute(
self,
context: RequestContext,
event_queue: EventQueue,
) -> None:
# 1. 取出 / 新建一个有状态 Task(对应 01 章 §1.3 的 Task)
if context.current_task:
task = context.current_task
else:
task = new_task_from_user_message(context.message)
await event_queue.enqueue_event(task) # 把新建的 Task 投到事件队列
# 2. 用 TaskUpdater 推进状态:先置 WORKING
updater = TaskUpdater(
event_queue=event_queue,
task_id=task.id,
context_id=task.context_id,
)
await updater.update_status(state=TaskState.TASK_STATE_WORKING)
# 3. 从入站 Message 抽出纯文本,干活(这里就是 .upper())
query = get_message_text(context.message)
result = query.upper() if query else "(空输入)"
# 4. 把结果包成 Artifact(对应 01 章 §1.3 的 Artifact = 任务产出)
await updater.add_artifact(
parts=[new_text_part(text=result, media_type="text/plain")],
name="uppercased",
)
# 5. 终态:COMPLETED(对应 02 章 §2.5 生命周期的终态)
await updater.update_status(state=TaskState.TASK_STATE_COMPLETED)
async def cancel(
self,
context: RequestContext,
event_queue: EventQueue,
) -> None:
# 本例不支持取消;真实长任务应在此置 TASK_STATE_CANCELED
raise NotImplementedError("Cancel is not supported.")
execute(context, event_queue)这是 server 的核心契约:SDK 框架保证同一请求不会并发调它;agent 只管从 context 读入站 Message、把事件 enqueue 到 event_queue。返回即代表本次执行结束。
new_task_from_user_message把入站 Message 包成一个有状态 Task 并 enqueue_event 投出去——这就是 01 章§1.3 说的「服务端为一次托付建出一个 Task」在代码里的样子。
TaskState.TASK_STATE_WORKING / COMPLETED正是 02 章§2.5 状态机里的 wire 值(SCREAMING_SNAKE,不是小写 working)。这里把任务从 WORKING 推到终态 COMPLETED。
add_artifact(parts=[…])产出不是直接 return 一个字符串,而是包成 Artifact 的 Part[](这里一个 text Part)。这对应 01 章§1.3「Artifact 是任务的产出、内容是 Part[]」。
上面 cancel() 里直接 raise NotImplementedError。对一个 50ms 就算完的 .upper(),不实现取消没问题;但换成一个跑几小时的调研 agent,cancel() 该往 event_queue 里发什么,才算「真的取消了」?
展开答案(先停 10 秒再点)
它该把 Task 推进到终态 TASK_STATE_CANCELED——即在 cancel() 里用 TaskUpdater(...).update_status(state=TaskState.TASK_STATE_CANCELED)(或等价的 updater.cancel())发一个 TaskStatusUpdateEvent。原因在 02 章§2.5 的状态机:CANCELED 是四个终态之一,客户端正是靠这个状态变化知道「这个 Task 不会再有产出了」,从而停止等待 / 回查。只是停掉后台计算却不发终态事件,客户端会一直以为任务还在 WORKING、傻等下去。注意 wire 上是美式拼写 CANCELED(v0.x 的 cancelled 已移除)——这也是 SDK 替你对齐 wire 的一处细节。
第二半是能力声明加把它挂起来。构造一张 AgentCard(含 skills、capabilities、服务端点),交给 DefaultRequestHandler,再用两个路由函数生成 Starlette 路由,uvicorn 起服务。
# 未在本机验证 —— a2a-sdk 1.1.0(spec v1.0)
import uvicorn
from starlette.applications import Starlette
from a2a.server.request_handlers import DefaultRequestHandler
from a2a.server.routes import create_agent_card_routes, create_jsonrpc_routes
from a2a.server.tasks import InMemoryTaskStore
from a2a.types import AgentCapabilities, AgentCard, AgentInterface, AgentSkill
from upper_executor import UpperCaseExecutor
PORT = 9999
BASE_URL = f"http://127.0.0.1:{PORT}"
# ── 一项 skill:向别的 agent 声明本 agent 能干这件事(对应 01 章 §1.2 的 skills)──
skill = AgentSkill(
id="uppercase",
name="Uppercase Text",
description="把输入的纯文本整体转成大写。",
input_modes=["text/plain"],
output_modes=["text/plain"],
tags=["text", "demo"],
examples=["hello a2a", "make this loud"],
)
# ── Agent Card:对外的一纸能力声明(对应 01 章 §1.2、02 章 §2.2)──
agent_card = AgentCard(
name="Uppercase Agent",
description="一个把文本转大写的最小 A2A agent。",
version="1.0.0",
default_input_modes=["text/plain"],
default_output_modes=["text/plain"],
# capabilities:声明支持哪些可选能力。这里先不开流式(3.3 让你决定要不要开)
capabilities=AgentCapabilities(streaming=False),
# supported_interfaces:声明走哪种传输、服务端点在哪(这里 JSON-RPC)
supported_interfaces=[
AgentInterface(protocol_binding="JSONRPC", url=BASE_URL),
],
skills=[skill],
)
# ── 把 executor、task store、card 接到请求处理器上 ──
handler = DefaultRequestHandler(
agent_executor=UpperCaseExecutor(),
task_store=InMemoryTaskStore(),
agent_card=agent_card,
)
# ── 生成路由:一组发布 Agent Card,一组处理 JSON-RPC 的 SendMessage 等 ──
routes = []
routes.extend(create_agent_card_routes(agent_card)) # GET /.well-known/agent-card.json
routes.extend(create_jsonrpc_routes(handler, "/")) # POST / —— JSON-RPC 入口
app = Starlette(routes=routes)
if __name__ == "__main__":
uvicorn.run(app, host="127.0.0.1", port=PORT)
AgentSkill / skills=[skill]对应 01 章§1.2 的 skills 字段:声明这个 agent 会哪几手。client 读 Card 时靠它判断对端能不能干这件事。
capabilities=AgentCapabilities(streaming=False)对应 01 章§1.2 的 capabilities 可选能力开关。此处声明不支持流式——3.3 的决策 (a) 就是让你决定要不要把它改成 True,以及那意味着什么。
create_agent_card_routes这一组路由把 Card 发布到 /.well-known/agent-card.json——正是 02 章§2.2 那个固定发现路径。client 第一次 GET 的就是它。
create_jsonrpc_routes(handler, "/")把 JSON-RPC 端点挂在 /。client 发的 SendMessage 落到这里,由 handler 路由进 execute()。
client 端:发现、发消息、收回流
client 这一侧三步走:解析 Agent Card(发现)、构造一个 Message 调 send_message(发)、异步迭代拿回 Task 与 Artifact(收)。注意 1.1.0 的 send_message 永远返回一个异步迭代器——无论 server 流不流式,client 都用 async for 收,由 SDK 在底层替你聚合。
# 未在本机验证 —— a2a-sdk 1.1.0(spec v1.0)
import asyncio
import httpx
from a2a.client import A2ACardResolver, ClientConfig, ClientFactory
from a2a.helpers import get_stream_response_text, new_text_message
from a2a.types.a2a_pb2 import Role, SendMessageRequest
BASE_URL = "http://127.0.0.1:9999"
async def main() -> None:
async with httpx.AsyncClient() as httpx_client:
# 1. 发现:去固定路径 GET 那张 Agent Card(对应 01 §1.2 / 02 §2.2)
resolver = A2ACardResolver(httpx_client=httpx_client, base_url=BASE_URL)
agent_card = await resolver.get_agent_card() # GET /.well-known/agent-card.json
print("发现 skill:", [s.id for s in agent_card.skills])
# 2. 用 Card 造一个 client。streaming=False 与本例 server 的声明一致
config = ClientConfig(streaming=False, httpx_client=httpx_client)
client = ClientFactory(config).create(agent_card)
# 3. 发:构造一个 user 角色的 Message,调 send_message
# 注意:Python 方法名是 snake_case send_message,
# 但它在 wire 上发出的 JSON-RPC "method" 是 PascalCase "SendMessage"
message = new_text_message("hello a2a", role=Role.ROLE_USER)
request = SendMessageRequest(message=message)
# 4. 收:send_message 返回异步迭代器,逐块收回 Task / Artifact
async for chunk in client.send_message(request):
print("收到:", get_stream_response_text(chunk))
await client.close()
if __name__ == "__main__":
asyncio.run(main())
A2ACardResolver.get_agent_card()这是发现那一步——一次普通 HTTP GET 打到 /.well-known/agent-card.json,把 02 章§2.1 端到端时序的第①步落成一行代码。此刻还没发任何任务。
ClientFactory(config).create(agent_card)工厂读 Card 的 supported_interfaces,挑一个双方都支持的传输绑定,造出对应的传输层 client。(便捷写法 create_client(agent=agent_card, client_config=config) 是它的一行封装。)
Role.ROLE_USER + SendMessageRequest对应 01 章§1.3 的 Message:一个 role 加 parts;new_text_message 替你包了 text Part。
async for chunk in client.send_message(request)1.1.0 统一用异步迭代收回——非流式时 SDK 聚合成一两块吐出来,流式时逐帧吐。返回里带着 Task 与 Artifact。
跑起来 + 预期输出
开两个终端。第一个起 server,第二个跑 client。
# 终端 1:起 server(监听 127.0.0.1:9999)
python server.py
# 终端 2:先手动看一眼那张 Card(发现这一步本质就是个 HTTP GET)
curl -s http://127.0.0.1:9999/.well-known/agent-card.json | head -c 400
echo
# 终端 2:再跑 client
python client.py
发现 skill: ['uppercase']
收到: HELLO A2A
那行 curl 会先吐出一段 JSON——里面有 name、skills、capabilities、supportedInterfaces 等字段。这就是 client 在第①步读到的同一张 Card;client 只是用代码把这次 GET 做了一遍。
整段代码里你写的是 Python 方法 client.send_message(...)(snake_case)。但它在线缆上发出的 JSON-RPC 请求,"method" 字段的字面值是 SendMessage(PascalCase)——正是 02 章§2.3 教的那个 v1.0 wire 方法名。SDK 的命名人体工学 ≠ wire 协议:方法名、TASK_STATE_* 状态值、Card 的 JSON 字段名,都是 SDK 在帮你做 Python 风格 ↔ wire 风格的来回翻译。想确认 wire 真长什么样,别只读 SDK 源码——抓一次包,或读 02 章那段 spec 示例。
client 里 send_message 是 async for chunk in ... 迭代着收的。可本例 server 明明声明了 capabilities.streaming=False、也没逐帧推任何东西——为什么 client 这边还要用异步迭代,而不是一句 result = await client.send_message(...) 直接拿一个返回值?
展开答案(先停 10 秒再点)
因为 1.1.0 把两种交互模式统一成了同一个 client API。send_message 的签名固定返回 AsyncIterator,由 SDK 在底层根据「server 支不支持流式 + ClientConfig.streaming 怎么设」自动选走 SSE 流还是普通请求/响应(02 章§2.4 的两种交互)。非流式时,SDK 把单个响应也包成「只迭代一两次」的迭代器吐给你——于是调用方代码不必随交互模式改写。这正是 SDK 抹平 wire 差异的又一处:同一行 async for,底下可能是一次 HTTP 往返,也可能是一条 text/event-stream 长连接。3.4 把 server 改成流式时,client 这段一个字都不用动。
3.3Partial · 补三个关键 schema 决策点
把 3.2 的 server 改成需要你拍板的版本——三处真正影响 wire 行为的 schema 决策抠成 TODO。每个 TODO 先合上页面想 30 秒该怎么选、代价是什么,再展开对照。三个决策分别落在 Card 能力声明、Artifact 封装、生命周期中途态上——都不是无关填空。
# 未在本机验证 —— 在 3.2 基础上抠掉 3 个 schema 决策
from a2a.helpers import get_message_text, new_task_from_user_message, new_text_part
from a2a.server.agent_execution import AgentExecutor, RequestContext
from a2a.server.events import EventQueue
from a2a.server.tasks import TaskUpdater
from a2a.types.a2a_pb2 import TaskState
# ─────────────────────────────────────────────
# 决策 (a) · Card 里要不要声明 capabilities.streaming = True?
# 3.2 声明的是 False。如果这个 agent 要逐字吐结果(像 LLM 流式输出),
# 该把它改成 True 吗?改了之后,server 端 execute 里要相应多做什么?
# (这个开关在 server.py 的 AgentCapabilities(...) 里,见下方答案)
# ─────────────────────────────────────────────
class UpperCaseExecutor(AgentExecutor):
async def execute(
self,
context: RequestContext,
event_queue: EventQueue,
) -> None:
if context.current_task:
task = context.current_task
else:
task = new_task_from_user_message(context.message)
await event_queue.enqueue_event(task)
updater = TaskUpdater(
event_queue=event_queue,
task_id=task.id,
context_id=task.context_id,
)
await updater.update_status(state=TaskState.TASK_STATE_WORKING)
query = get_message_text(context.message)
# ─────────────────────────────────────────────
# 决策 (c) · input-required 中途态怎么处理?
# 如果客户端发来空文本(没东西可转大写),该怎么办?
# 直接 FAILED?还是把任务暂停成 INPUT_REQUIRED 等它补一句?
# ─────────────────────────────────────────────
if not query:
# await updater.update_status(state=???) # 该置成哪个状态?然后 return 让出控制权?
return
result = query.upper()
# ─────────────────────────────────────────────
# 决策 (b) · 结果怎么包成 Artifact?
# 下面这行只发了一个 text Part。如果还想附一份「原文 + 大写」的
# 结构化对照(JSON),该往 parts 里加什么 kind 的 Part?
# ─────────────────────────────────────────────
await updater.add_artifact(
parts=[new_text_part(text=result, media_type="text/plain")], # 还要再加一个 Part 吗?
name="uppercased",
)
await updater.update_status(state=TaskState.TASK_STATE_COMPLETED)
async def cancel(self, context: RequestContext, event_queue: EventQueue) -> None:
raise NotImplementedError("Cancel is not supported.")
决策 (a) 答案 · Card 里要不要声明 capabilities.streaming
看产出形态,不是默认开:
- 短任务、一次性产出(如本例
.upper()):声明streaming=False就够。结果一步算完,没有「逐字吐」的中间过程,开流式只是徒增复杂度。 - 逐字 / 分块产出(如 LLM 流式生成、长报告边写边发):声明
capabilities=AgentCapabilities(streaming=True)。但声明只是一半——server 端execute还得真的多次add_artifact(..., append=True)或多发TaskStatusUpdateEvent,否则声明了流式却一次性吐,等于没流。 - 代价:声明
True是对 client 的承诺。client 读到后可能选走 SSE 长连接(02 章§2.4);server 却没真流式地发,行为就和声明不符。能力声明是契约,不是装饰。
填空:本例保持 AgentCapabilities(streaming=False);要改流式得连 execute 一起改(这正是 3.4 的开放练习)。锚 01 章 Card 能力声明(§1.2)+ 02 章交互模式(§2.4)。
决策 (b) 答案 · 结果怎么包成 Artifact
一个 Artifact 可以装多个 Part,按内容形态选 kind:
- 纯文本结果:一个
textPart(new_text_part),就是 3.2 的做法。 - 文本 + 结构化对照:再加一个
dataPart 装 JSON。01 章§1.3 说过 Part 靠kind分text/file/data,一条消息 / 一个 Artifact 能混装多种——这正是模态无关设计的用处。用a2a.helpers的new_data_part(data={...})。 - 多个独立产出(如同时给报告和图):可以发多次
add_artifact,每次一个独立artifactId;也可以一个 Artifact 塞多个 Part。前者语义是「几份成果」,后者是「一份成果的几个部分」。
填空(文本 + 结构化对照):
# 未在本机验证
from a2a.helpers import new_data_part, new_text_part
await updater.add_artifact(
parts=[
new_text_part(text=result, media_type="text/plain"),
new_data_part(data={"original": query, "upper": result}),
],
name="uppercased",
)
锚 01 章 Part / Artifact(§1.3)。
决策 (c) 答案 · input-required 中途态怎么处理
缺必要输入时,用 INPUT_REQUIRED 暂停、等补,而不是直接 FAILED:
- 置
TASK_STATE_INPUT_REQUIRED并 return:02 章§2.5 讲过,INPUT_REQUIRED是可中断的中途态,不是终态。server 把任务停在这里、捎一条「请补一句要转大写的文本」的 message,然后execute返回让出控制权;client 补完输入后框架会带着同一个 Task 再次调execute,任务原地续,不重开。 - 为什么不直接
FAILED:FAILED是终态,任务就此结束,客户端得从头再发一个新 Task。对「只是少了点输入」这种可恢复情况,INPUT_REQUIRED保住了已有的 Task 上下文,体验和效率都更好。 - 代价 / 注意:中途态把控制权交还给客户端等待补输入——这个能力很强,但也正是钓鱼攻击的入口(恶意 agent 伪造
INPUT_REQUIRED骗用户重交敏感信息,04 章§4.6)。
填空:
# 未在本机验证
from a2a.helpers import new_text_message
if not query:
await updater.update_status(
state=TaskState.TASK_STATE_INPUT_REQUIRED,
message=new_text_message("请补一句要转成大写的文本。"),
)
return # 让出控制权,等客户端补输入后框架再次调 execute
锚 02 章生命周期中途态(§2.5)。
决策 (a) 里若 Card 声明了 streaming=True 却没真流式地发,或决策 (c) 用 INPUT_REQUIRED 暂停后再没人来续——任务就永远卡在中途态。能力声明与实际行为不一致、中途态无兜底,都是 04 章会展开的协议-实现耦合与失败模式(§4.8)。
3.4Open · 把它改成流式
3.2 是一次性产出(请求/响应)。开放练习把同一个 agent 改成流式:server 声明 capabilities.streaming=True,在 execute 里逐块把结果发出去——一个字符一个字符(或一词一词)地 emit,让 client 能边收边显示,像看 LLM 打字。
这条练习要求你用到 01 章的 ≥2 个概念——Artifact / Part §1.3(分块产出仍是 Artifact 的 Part)和 Card 的 capabilities §1.2(声明 streaming)——以及 02 章的某个权衡:§2.4 推送/流式 vs 请求/响应,或 §2.6 表里「有状态长任务 + 推送」那一行(流式拿到实时进度,代价是要维持长连接、处理重连)。先自己写完,再对照参考实现要点。
SSE 流式的 wire 形态(02 章§2.4):server 先发 Task,再连发若干 TaskStatusUpdateEvent(进度)和 TaskArtifactUpdateEvent(产出分块)。SDK 里这两类事件由 TaskUpdater.update_status(...) 和带 append=True / last_chunk=True 的 TaskUpdater.add_artifact(...) 发出。client 端不用改——3.2 的 async for 会自动逐帧收(这正是上面那道 predict 的结论)。
参考实现要点 · 流式版(写完自己的版本再展开)
第一步,Card 声明流式(server.py 里):
capabilities=AgentCapabilities(streaming=True) # 把 3.2 的 False 改成 True
第二步,execute 里分块 emit:把整段 .upper() 拆成块,用同一个 artifact_id 连发,append=True 表示「拼到上一块后面」,最后一块 last_chunk=True:
# 未在本机验证 —— 流式版 execute 的核心改动
import asyncio
from a2a.helpers import get_message_text, new_task_from_user_message, new_text_part
from a2a.server.tasks import TaskUpdater
from a2a.types.a2a_pb2 import TaskState
async def execute(self, context, event_queue) -> None:
if context.current_task:
task = context.current_task
else:
task = new_task_from_user_message(context.message)
await event_queue.enqueue_event(task)
updater = TaskUpdater(
event_queue=event_queue,
task_id=task.id,
context_id=task.context_id,
)
await updater.update_status(state=TaskState.TASK_STATE_WORKING)
text = (get_message_text(context.message) or "").upper()
artifact_id = "out-1" # 同一 Artifact,分多块
for i, ch in enumerate(text):
await updater.add_artifact(
parts=[new_text_part(text=ch, media_type="text/plain")],
artifact_id=artifact_id,
name="uppercased",
append=(i > 0), # 第一块不 append,其后都拼接
last_chunk=(i == len(text) - 1), # 标出最后一块
)
await asyncio.sleep(0.05) # 模拟逐字产出的节奏
await updater.update_status(state=TaskState.TASK_STATE_COMPLETED)
关键决策说明:
- 为什么用同一个
artifact_id+append:流式产出的语义是「一份成果分多块陆续到」,不是「多份成果」。同一artifact_id让 client 知道这些块属于同一个 Artifact,append=True告诉它拼接而非替换,last_chunk=True标出收尾。锚 01 章 Artifact / Part(§1.3)。 - client 为什么不用改:3.2 的 client 已是
async for chunk in client.send_message(...)。server 一旦流式,同一个迭代会逐帧吐回每个TaskArtifactUpdateEvent——SDK 抹平了请求/响应与 SSE 的差异(上面 predict 的结论)。把ClientConfig(streaming=True)打开,让 client 优先走流式即可。 - 权衡(锚 02 章§2.4 / §2.6):流式换来「实时进度、边收边显示」,代价是 server 与 client 要维持一条
text/event-stream长连接,断了得靠 resubscribe 续;对一个 50ms 就算完的.upper()其实不值——流式真正的用武之地是 LLM 逐字生成、长报告这类「产出本身就是渐进的」任务。能力该不该开,取决于产出形态,不是越多越好。
流式响应的 Content-Type 在版本间有差异(v1.0.1 倾向 application/a2a+json,某些 v0.x 实现对 application/json 与 text/event-stream 处理不一致)。一个 v0.x 的 client 连 v1.0 的流式 server,哪怕都用 JSON-RPC,也可能因 Content-Type 协商失败而收不到流——这是 04 章版本 wire 不兼容失败模式(§4.7)的一种。
§本章 self-check
先合上代码,把答案写下来,再展开对照。直接展开等于把这一节当又读了一遍。
- 你写
client.send_message(...)(snake_case)。这次调用在 wire 上发出的 JSON-RPC 请求里,"method"字段的字面值是什么?这说明 SDK 与 wire 是什么关系? - 3.2 的 server 端
execute里,产出为什么不是return "HELLO A2A",而要走add_artifact(parts=[...])?对应 01 章哪个概念? - 本例 server 声明了
capabilities.streaming=False,client 却仍用async for迭代收send_message。为什么 client 不能简单写成result = await client.send_message(...)? - 客户端发来空文本时,把任务置成
TASK_STATE_INPUT_REQUIRED然后return,与直接置TASK_STATE_FAILED,行为上有什么本质区别?
答案(先做完再展开)
- 字面值是
SendMessage(PascalCase,v1.0 的 wire 方法名)。Python 方法名send_message只是 SDK 的人体工学,SDK 命名 ≠ wire 协议——方法名、TASK_STATE_*状态值、Card 字段名都由 SDK 在 Python 风格与 wire 风格之间来回翻译。锚 02 章§2.3。 - 因为 A2A 的产出是 Artifact(带
artifactId、内容是Part[]),不是裸字符串。execute不直接返回值,而是把事件 enqueue 到EventQueue;结果包成 Artifact 的一个textPart 发出去。对应 01 章§1.3 的 Artifact = 任务产出。 - 因为 1.1.0 的
send_message签名固定返回AsyncIterator,把请求/响应和 SSE 流统一成同一个 API:非流式时 SDK 把单个响应也包成只迭代一两次的迭代器。这样调用方代码不必随交互模式(02 章§2.4)改写——3.4 把 server 改流式时,client 一个字都不用动。 INPUT_REQUIRED是可中断的中途态(02 章§2.5):任务暂停、保住已有 Task 上下文,等客户端补输入后框架带着同一个 Task 再次调execute、原地续。FAILED是终态:任务就此结束,客户端得从头发一个新 Task。前者可恢复,后者不可。
让 client 主动 GetTask 回查,而不是同步等结果
3.2 的 client 是「发了 SendMessage 就地等结果」。改成异步回查风格:client 发完消息、拿到 Task 的 id 后不就地等,而是稍后用 GetTask 主动查这个 Task 现在到哪一步了。要求:
- 从
send_message收回的对象里取出 Task 的id与当前status.state。 - 构造一个
GetTaskRequest(name=...)(task 资源标识),调client.get_task(...)回查。 - 轮询到
status.state进入终态(TASK_STATE_COMPLETED/FAILED/ …)才停,读出 Artifact。 - 想清楚:这种「发完不等、稍后回查」对应 02 章三种交互里的哪一种思路?它和 webhook 推送的差别在哪?
提示(卡住再展开)
从 a2a.types.a2a_pb2 导入 GetTaskRequest;client.get_task(request) 返回一个 Task,读它的 status.state 判终态、读 artifacts 取产出。轮询用 while + asyncio.sleep,并务必设一个最大轮询次数/超时兜底,否则任务卡在中途态时会死循环。这对应 02 章§2.4 的「客户端轮询」思路——与 webhook 推送的关键差别:轮询是客户端反复主动问(多数次拿到「没变化」,浪费且有延迟),推送是服务端在有变化时反向 POST(省往返,但引入反向信任,客户端要验来源真伪,04 章§4.4)。GetTask 的 wire 方法名同样是 PascalCase(GetTask),见 02 章§2.3 的方法映射表。