Chapter 03 · 上手实操

把它跑起来:一个 A2A server,一个 client,端到端

02 讲清了 wire 上发生什么与三种交互——本章把它跑起来。一个 A2A server 暴露一项 skill,一个 client 发现它的 Agent Card、发 SendMessage、收回 Task 与 Artifact。三阶推进:worked(读完整可运行代码)→ partial(补关键 schema 决策)→ open(改成流式)。读完,你不仅能跑通这条链,还能指出 SDK 在哪几处把 wire 协议藏了起来。

本章你将建立的 schema

  • 用 a2a-sdk 1.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 的实际版本与官方文档为准。

A2A client ClientFactory .create(card) send_message() A2A server AgentCard 路由 /.well-known/… AgentExecutor .execute() EventQueue enqueue_event ① GET agent-card.json ② SendMessage(Message: role + parts) ③ Task + Artifact ④ 回流:Task → Artifact 黑线 = client 发起的控制流 · 朱红线 = server 产出的数据回流
图 3.1server + client 的控制流(黑)与数据回流(朱红)。注意:第①步「发现」与第②步「发消息」是两次独立的 HTTP 往返——client 先 GET 那张静态 Card 确认能力与端点,之后才发 SendMessage;产出不是裸字符串,而是 execute() 经 EventQueue 包出的 Artifact,沿朱红线回流。
3.2 Worked server 全给 client 全给 逐行映射回概念 学读代码 3.3 Partial 骨架给定 3 决策留空 streaming/Artifact/input 学做选择 3.4 Open 只给规格 改成流式 SSE 脚手架全撤 学组装系统 撤决策 撤骨架
图 3.2三阶脚手架递减:每往前一阶撤掉一层给定。注意:3.3 抠掉的三个不是任意填空——分别对应一个真正影响 wire 行为的 schema 决策:Card 的 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 传输都带进来了。

setup.sh Bash
# 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 公开发布。先看业务逻辑这半。

upper_executor.py Python
# 未在本机验证 —— 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 起服务。

server.py Python
# 未在本机验证 —— 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 在底层替你聚合。

client.py Python
# 未在本机验证 —— 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。

run.sh Bash
# 终端 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
预期输出(client 终端,示意)

发现 skill: ['uppercase']
收到: HELLO A2A

那行 curl 会先吐出一段 JSON——里面有 name、skills、capabilities、supportedInterfaces 等字段。这就是 client 在第①步读到的同一张 Card;client 只是用代码把这次 GET 做了一遍。

洞察 · SDK 在哪藏了 wire

整段代码里你写的是 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 封装、生命周期中途态上——都不是无关填空。

upper_executor_partial.py Python
# 未在本机验证 —— 在 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:

  • 纯文本结果:一个 text Part(new_text_part),就是 3.2 的做法。
  • 文本 + 结构化对照:再加一个 data Part 装 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)。

forward → 04 失败模式 · 声明与实现不符

决策 (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 逐字生成、长报告这类「产出本身就是渐进的」任务。能力该不该开,取决于产出形态,不是越多越好。
forward → 04 失败模式 · 版本 wire 不兼容

流式响应的 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

先合上代码,把答案写下来,再展开对照。直接展开等于把这一节当又读了一遍。

  1. 你写 client.send_message(...)(snake_case)。这次调用在 wire 上发出的 JSON-RPC 请求里,"method" 字段的字面值是什么?这说明 SDK 与 wire 是什么关系?
  2. 3.2 的 server 端 execute 里,产出为什么不是 return "HELLO A2A",而要走 add_artifact(parts=[...])?对应 01 章哪个概念?
  3. 本例 server 声明了 capabilities.streaming=False,client 却仍用 async for 迭代收 send_message。为什么 client 不能简单写成 result = await client.send_message(...)?
  4. 客户端发来空文本时,把任务置成 TASK_STATE_INPUT_REQUIRED 然后 return,与直接置 TASK_STATE_FAILED,行为上有什么本质区别?
答案(先做完再展开)
  1. 字面值是 SendMessage(PascalCase,v1.0 的 wire 方法名)。Python 方法名 send_message 只是 SDK 的人体工学,SDK 命名 ≠ wire 协议——方法名、TASK_STATE_* 状态值、Card 字段名都由 SDK 在 Python 风格与 wire 风格之间来回翻译。锚 02 章§2.3。
  2. 因为 A2A 的产出是 Artifact(带 artifactId、内容是 Part[]),不是裸字符串。execute 不直接返回值,而是把事件 enqueue 到 EventQueue;结果包成 Artifact 的一个 text Part 发出去。对应 01 章§1.3 的 Artifact = 任务产出。
  3. 因为 1.1.0 的 send_message 签名固定返回 AsyncIterator,把请求/响应和 SSE 流统一成同一个 API:非流式时 SDK 把单个响应也包成只迭代一两次的迭代器。这样调用方代码不必随交互模式(02 章§2.4)改写——3.4 把 server 改流式时,client 一个字都不用动。
  4. INPUT_REQUIRED 是可中断的中途态(02 章§2.5):任务暂停、保住已有 Task 上下文,等客户端补输入后框架带着同一个 Task 再次调 execute、原地续。FAILED 是终态:任务就此结束,客户端得从头发一个新 Task。前者可恢复,后者不可。
进阶挑战 · 刚好够不着

让 client 主动 GetTask 回查,而不是同步等结果

3.2 的 client 是「发了 SendMessage 就地等结果」。改成异步回查风格:client 发完消息、拿到 Task 的 id 后不就地等,而是稍后用 GetTask 主动查这个 Task 现在到哪一步了。要求:

  1. 从 send_message 收回的对象里取出 Task 的 id 与当前 status.state。
  2. 构造一个 GetTaskRequest(name=...)(task 资源标识),调 client.get_task(...) 回查。
  3. 轮询到 status.state 进入终态(TASK_STATE_COMPLETED / FAILED / …)才停,读出 Artifact。
  4. 想清楚:这种「发完不等、稍后回查」对应 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 的方法映射表。