
1. 为什么要在本地搭一套 A2A MCP 骨架如果你最近在折腾 Agent 系统大概率会撞上两个词A2A 和 MCP。A2A 解决的是「代理和代理之间怎么对话」MCP 解决的是「代理和工具、数据之间怎么对话」。前者是编排骨干后者是工具接入层。听起来分工清晰但真到写代码时很多人会卡在同一个地方协议消息到底长什么样SSE 通道怎么把流式结果透传出去JSON-RPC 的请求和响应字段怎么对齐。这篇就是来解决这个问题的。我会用 Python 从零搭一个最小可运行骨架一个 Router 负责按任务类型分发三个子 AgentASR / Legal / TTS各自暴露/a2a/task入口再加一个 Mock MCP Server 用 JSON-RPC 风格提供工具调用。全链路支持 SSE 流式你启动之后发一条 JSON-RPC 请求就能看到progress、delta、final事件逐步回包。适合谁看已经会写 FastAPI、想快速理解 A2A 与 MCP 协作方式的开发者正在做多 Agent 编排、需要一套可复制骨架的人以及想搞清楚 SSE 和 JSON-RPC 在 Agent 场景里怎么配合的同学。整套代码不依赖任何外部服务本地五个终端就能跑通。另外骨架里所有需要调模型或工具的地方我都会用 TaoToken 的统一 Key 来配置。这样你后面把 Mock 换成真实 LLM、ASR、TTS 时不用改一堆环境变量一个 Key 走通。2. TaoToken 前置统一 Key 与配置项在动手写代码之前先把「钥匙」准备好。TaoToken 的作用是把模型调用和工具调用的鉴权收敛到一个 Key 上这样 Router 和各个 Agent 不用各自维护一套凭证。官网入口在这里https://taotoken.net/?utm_sourcetaotoken_aicg_blog_endutm_mediumcsdnutm_campaignrewriteutm_content 注册后进控制台拿 Key。拿到 Key 之后建议先建一个.env文件把所有地址和凭证集中管理。骨架里我会用python-dotenv读取这样本地调试和后面部署都不用改代码。# .env.sample 复制为 .env 后按需修改 # TaoToken 统一 Key模型与工具调用共用 TAOTOKEN_API_KEYsk-你的key TAOTOKEN_BASE_URLhttps://taotoken.net/api # Router 将任务路由到各个 Agent 的地址 ASR_AGENT_URLhttp://127.0.0.1:8011 LEGAL_AGENT_URLhttp://127.0.0.1:8012 TTS_AGENT_URLhttp://127.0.0.1:8013 # MCP Server 地址示例使用 mock MCP_ENDPOINThttp://127.0.0.1:8099/mcp这里有个细节值得说TAOTOKEN_BASE_URL我写的是https://taotoken.net/api注意这个地址不带任何查询参数是纯 API 入口。而官网首页那个带 UTM 的链接是给人点的代码里请求不要带 UTM否则有些网关会当成异常参数处理。如果你后面要接 Coding Plan 做长期编码或 Agent 任务可以在控制台里单独开一个 planKey 还是同一个只是额度策略不同。接入文档在 https://taotoken.net/doc API Keys 管理页在 https://taotoken.net/api-keys 这两个页面建议先收藏排障时会反复用到。注意.env不要提交到 Git。骨架里我会加.gitignore但你自己也要养成习惯Key 泄露比代码泄露麻烦得多。3. 可复制配置目录结构与核心文件先把目录结构定下来后面所有代码都按这个放。这套结构的好处是common放协议层agents放业务代理tools放工具服务职责清晰新增一个 Agent 或一个工具都不用动其他目录。. ├─ requirements.txt ├─ .env.sample ├─ .gitignore ├─ config.toml ├─ settings.json └─ src/ ├─ common/ │ ├─ a2a.py │ └─ mcp.py ├─ router.py ├─ agents/ │ ├─ asr_agent.py │ ├─ legal_agent.py │ └─ tts_agent.py └─ tools/ └─ mock_mcp_server.py依赖清单很轻四个核心包就够# requirements.txt fastapi0.111.0 uvicorn[standard]0.30.0 httpx0.27.2 pydantic2.8.2 python-dotenv1.0.1接下来是config.toml把路由表和超时策略放这里避免硬编码在 Python 里。这样你改路由前缀不用重启代码逻辑只改配置。# config.toml [router] host 0.0.0.0 port 8000 timeout_seconds 60 [router.routes] asr. http://127.0.0.1:8011 legal. http://127.0.0.1:8012 tts. http://127.0.0.1:8013 [mcp] endpoint http://127.0.0.1:8099/mcp timeout_seconds 30 [taotoken] base_url https://taotoken.net/apisettings.json则用来放 Agent 级别的元信息比如每个 Agent 支持的任务类型和描述。这个文件后面可以扩展成 Agent 注册表配合/.well-known/agents做服务发现。{ agents: [ { name: asr_agent, url: http://127.0.0.1:8011, task_types: [asr.transcribe], description: 语音转写代理 }, { name: legal_agent, url: http://127.0.0.1:8012, task_types: [legal.analyze], description: 法律分析代理 }, { name: tts_agent, url: http://127.0.0.1:8013, task_types: [tts.speak], description: 语音合成代理 } ], mcp_server: { endpoint: http://127.0.0.1:8099/mcp, transport: http-jsonrpc } }配置就绪后先写协议层。src/common/a2a.py定义 A2A 的任务与事件模型以及 SSE 编解码。这里的关键是A2AEvent的event字段它决定了调用方怎么解析流。# src/common/a2a.py from __future__ import annotations import uuid from typing import Any, AsyncIterator, Dict from pydantic import BaseModel, Field A2A_MIME_SSE text/event-stream class A2ATask(BaseModel): id: str Field(default_factorylambda: str(uuid.uuid4())) type: str payload: Dict[str, Any] {} stream: bool True class A2AEvent(BaseModel): task_id: str event: str data: Dict[str, Any] {} def sse_encode(event: A2AEvent) - bytes: line data: event.model_dump_json() \n\n return line.encode(utf-8) async def sse_stream_generator(generator: AsyncIterator[A2AEvent]): async for ev in generator: yield sse_encode(ev) class A2AError(RuntimeError): passsrc/common/mcp.py是最小 MCP 客户端用 HTTP JSON-RPC 实现。真实 MCP 支持 stdio 和 SSE 传输但 HTTP 版最容易本地跑通也方便你后面替换。# src/common/mcp.py from __future__ import annotations import uuid from typing import Any, Dict, List import httpx class MCPError(RuntimeError): pass class MCPClient: def __init__(self, endpoint: str, timeout: float 30.0): self.endpoint endpoint.rstrip(/) self.timeout timeout self._client httpx.AsyncClient(timeouttimeout) async def _rpc(self, method: str, params: Dict[str, Any]) - Any: req { jsonrpc: 2.0, id: str(uuid.uuid4()), method: method, params: params, } r await self._client.post(self.endpoint, jsonreq) r.raise_for_status() data r.json() if error in data and data[error]: raise MCPError(str(data[error])) return data.get(result) async def list_tools(self) - List[Dict[str, Any]]: return await self._rpc(tools/list_tools, {}) async def call(self, tool: str, args: Dict[str, Any]) - Dict[str, Any]: return await self._rpc(tools/call, {tool: tool, args: args}) async def aclose(self): await self._client.aclose()Router 是 A2A 的入口按task.type前缀查路由表然后把下游的 SSE 字节流原样转发。这里用aiter_bytes()而不是aiter_text()是为了避免编码转换破坏 SSE 分块边界。# src/router.py from __future__ import annotations import os from typing import AsyncIterator import httpx from dotenv import load_dotenv from fastapi import FastAPI, HTTPException from fastapi.responses import StreamingResponse, JSONResponse from common.a2a import A2ATask load_dotenv() ROUTE_TABLE { asr.: os.getenv(ASR_AGENT_URL, http://127.0.0.1:8011), legal.: os.getenv(LEGAL_AGENT_URL, http://127.0.0.1:8012), tts.: os.getenv(TTS_AGENT_URL, http://127.0.0.1:8013), } app FastAPI(titleA2A Router) def _pick_agent(task_type: str) - str: for prefix, url in ROUTE_TABLE.items(): if task_type.startswith(prefix): return url raise HTTPException(status_code400, detailfNo agent for task type: {task_type}) app.post(/a2a/task) async def route_task(task: A2ATask): agent_base _pick_agent(task.type) agent_endpoint f{agent_base}/a2a/task async with httpx.AsyncClient(timeoutNone) as client: if task.stream: resp await client.post( agent_endpoint, jsontask.model_dump(), headers{accept: text/event-stream}, ) if resp.status_code ! 200: raise HTTPException(status_coderesp.status_code, detailresp.text) async def forward() - AsyncIterator[bytes]: async for chunk in resp.aiter_bytes(): yield chunk return StreamingResponse(forward(), media_typetext/event-stream) else: resp await client.post(agent_endpoint, jsontask.model_dump()) return JSONResponse(status_coderesp.status_code, contentresp.json())三个 Agent 的结构类似我拿 Legal Agent 举例它内部会调 MCP 的kb.search和llm.complete并把结果分段流出。注意from __future__ import annotations是 Python 3.12 的写法如果你用 3.10 或 3.11改成from __future__ import annotations即可。# src/agents/legal_agent.py from __future__ import annotations import os import asyncio from typing import AsyncIterator from dotenv import load_dotenv from fastapi import FastAPI, HTTPException from fastapi.responses import StreamingResponse, JSONResponse from common.a2a import A2ATask, A2AEvent, sse_stream_generator from common.mcp import MCPClient load_dotenv() MCP_ENDPOINT os.getenv(MCP_ENDPOINT, http://127.0.0.1:8099/mcp) app FastAPI(titleLegal Expert Agent) async def _legal_stream(mcp: MCPClient, query: str, task_id: str) - AsyncIterator[A2AEvent]: yield A2AEvent(task_idtask_id, eventprogress, data{msg: legal_started}) kb await mcp.call(kb.search, {query: query, top_k: 3}) yield A2AEvent(task_idtask_id, eventprogress, data{kb_hits: kb.get(hits, [])}) prompt f根据以下材料做初步法律分析\n{kb.get(context, )}\n---\n问题{query}\n llm await mcp.call(llm.complete, {prompt: prompt, temperature: 0.2}) text llm.get(text, ) for chunk in [text[i:i60] for i in range(0, len(text), 60)]: await asyncio.sleep(0.05) yield A2AEvent(task_idtask_id, eventdelta, data{text: chunk}) yield A2AEvent(task_idtask_id, eventfinal, data{text: text}) app.post(/a2a/task) async def handle(task: A2ATask): if task.type ! legal.analyze: raise HTTPException(status_code400, detailfUnsupported task type: {task.type}) query task.payload.get(query) if not query: raise HTTPException(status_code400, detailpayload.query required) mcp MCPClient(MCP_ENDPOINT) if task.stream: gen _legal_stream(mcp, query, task.id) return StreamingResponse(sse_stream_generator(gen), media_typetext/event-stream) else: kb await mcp.call(kb.search, {query: query, top_k: 3}) llm await mcp.call(llm.complete, {prompt: str(kb), temperature: 0.2}) return JSONResponse(llm)Mock MCP Server 用 FastAPI 实现 JSON-RPC 端点注册四个工具asr.transcribe、kb.search、llm.complete、tts.speak。每个工具返回结构化的result方便 Agent 解析。# src/tools/mock_mcp_server.py from __future__ import annotations import random from typing import Any, Dict from fastapi import FastAPI from pydantic import BaseModel app FastAPI(titleMock MCP Server (JSON-RPC over HTTP)) TOOLS [ {name: asr.transcribe, args: {audio_url: str}}, {name: kb.search, args: {query: str, top_k: int}}, {name: llm.complete, args: {prompt: str, temperature: float}}, {name: tts.speak, args: {text: str}}, ] class JSONRPCRequest(BaseModel): jsonrpc: str id: str method: str params: Dict[str, Any] | None None class JSONRPCResponse(BaseModel): jsonrpc: str 2.0 id: str result: Dict[str, Any] | None None error: Dict[str, Any] | None None app.post(/mcp) async def mcp(req: JSONRPCRequest) - JSONRPCResponse: try: if req.method tools/list_tools: return JSONRPCResponse(idreq.id, result{tools: TOOLS}) if req.method tools/call: params req.params or {} tool params.get(tool) args params.get(args, {}) if tool asr.transcribe: text f(mock asr) transcribed from {args.get(audio_url,?)} return JSONRPCResponse(idreq.id, result{text: text}) if tool kb.search: hits [ {id: fdoc{i}, score: round(random.random(), 3), snippet: fsnippet {i}} for i in range(1, (args.get(top_k, 3) 1)) ] context \n.join(h[snippet] for h in hits) return JSONRPCResponse(idreq.id, result{hits: hits, context: context}) if tool llm.complete: prompt args.get(prompt, ) return JSONRPCResponse(idreq.id, result{text: f(mock llm) summary for: {prompt[:60]}...}) if tool tts.speak: text args.get(text, ) chunks [faudio-bytes-chunk-{i} for i in range(5)] return JSONRPCResponse(idreq.id, result{chunks: chunks, desc: f(mock) tts of {text[:20]}...}) return JSONRPCResponse(idreq.id, error{message: funknown tool: {tool}}) return JSONRPCResponse(idreq.id, error{message: funknown method: {req.method}}) except Exception as e: return JSONRPCResponse(idreq.id, error{message: str(e)})4. 验证请求启动服务并观察 SSE 回包代码写完后启动顺序有讲究先起 Mock MCP再起三个 Agent最后起 Router。因为 Agent 启动时会读MCP_ENDPOINTRouter 启动时会读三个 Agent 地址顺序反了不影响启动但第一次请求可能连不上。python -m venv .venv source .venv/bin/activate pip install -r requirements.txt cp .env.sample .env # 终端 1Mock MCP Server uvicorn src.tools.mock_mcp_server:app --host 0.0.0.0 --port 8099 --reload # 终端 2ASR Agent uvicorn src.agents.asr_agent:app --host 0.0.0.0 --port 8011 --reload # 终端 3Legal Agent uvicorn src.agents.legal_agent:app --host 0.0.0.0 --port 8012 --reload # 终端 4TTS Agent uvicorn src.agents.tts_agent:app --host 0.0.0.0 --port 8013 --reload # 终端 5Router uvicorn src.router:app --host 0.0.0.0 --port 8000 --reload五个终端都起来后先单独验证 MCP Server 的 JSON-RPC 是否正常。这一步很关键因为后面 Agent 调工具全靠它。curl -s -X POST http://127.0.0.1:8099/mcp \ -H Content-Type: application/json \ -d { jsonrpc: 2.0, id: req-1, method: tools/list_tools, params: {} } | python -m json.tool正常会返回四个工具的列表。接着验证 Router 的 SSE 流式转发让 Router 把任务派给 Legal Agentcurl -N -H Accept: text/event-stream \ -H Content-Type: application/json \ -d { type: legal.analyze, payload: {query: 公司违法辞退赔偿怎么计算}, stream: true } \ http://127.0.0.1:8000/a2a/task你会看到类似这样的回包事件按progress→delta→final顺序逐步到达data: {task_id:...,event:progress,data:{msg:legal_started}} data: {task_id:...,event:progress,data:{kb_hits:[{id:doc1,score:0.42,snippet:snippet 1}]}} data: {task_id:...,event:delta,data:{text:(mock llm) summary for: 根据以下材料...}} data: {task_id:...,event:final,data:{text:(mock llm) summary for: ...}}如果你把stream改成falseRouter 会走非流式分支直接返回一个 JSON 响应。两种模式共用同一个/a2a/task入口这是 A2A 设计里比较舒服的一点。再验证一下 ASR Agent 的流式增量它会把转写文本按 token 逐个吐出curl -N -H Accept: text/event-stream \ -H Content-Type: application/json \ -d { type: asr.transcribe, payload: {audio_url: http://example.com/a.wav}, stream: true } \ http://127.0.0.1:8000/a2a/task到这里A2A 的编排链路和 MCP 的工具调用链路就都跑通了。你可以打开模型对话页面 https://taotoken.net/chat 用同一个 Key 对比一下真实模型输出和 Mock 输出的差异方便后面替换。5. 本篇常见错排查跑这套骨架时最容易踩的坑集中在端口、编码和协议字段三块。我按实际遇到的频率排一下。第一个坑SSE 回包被缓冲看不到流式效果。如果你用aiter_text()转发或者中间加了await resp.aread()SSE 会被一次性读完再返回curl -N也看不到逐条事件。正确做法是用aiter_bytes()原样转发并且 Router 的StreamingResponse不要加Content-Length。另外curl必须带-N否则 curl 自己会缓冲。第二个坑JSON-RPC 的id类型不一致。骨架里id用的是字符串如果你在客户端传了数字某些严格校验的服务端会报Invalid Request。建议统一用str(uuid.uuid4())避免类型歧义。jsonrpc字段必须是2.0少一个字符都会解析失败。第三个坑from __future__ import annotations写错。这个语句必须是文件第一行docstring 之后而且 Python 3.12 之前要写成from __future__ import annotations。如果你看到SyntaxError: future feature annotations is not defined先检查 Python 版本再检查这行有没有被注释或放错位置。第四个坑Agent 连不上 MCP Server。表现是httpx.ConnectError或MCPError。先确认 8099 端口在监听再确认.env里的MCP_ENDPOINT没有多余斜杠。MCPClient里做了rstrip(/)但如果你写成http://127.0.0.1:8099/mcp/拼接后可能变成/mcp//部分框架会 404。第五个坑Router 路由不到 Agent。检查task.type的前缀是否在ROUTE_TABLE里。比如legal.analyze匹配legal.但如果你写成legal_analyze就会命中No agent for task type。前缀匹配是大小写敏感的Legal.analyze也不行。第六个坑TaoToken Key 没生效。如果你把 Mock 换成了真实模型调用报 401 或 403先确认.env里TAOTOKEN_API_KEY没有引号包裹再确认TAOTOKEN_BASE_URL是https://taotoken.net/api而不是带 UTM 的首页地址。Key 的额度和管理在 https://taotoken.net/api-keys 接入细节看 https://taotoken.net/doc 。提示排障时建议先单独测 MCP Server再测单个 Agent最后测 Router。逐层验证比一上来就发全链路请求高效得多。6. 下一步把 Mock 换成真实能力骨架跑通之后替换成真实能力其实只改两处。第一处是MCPClient的endpoint把 Mock Server 换成你真实的 MCP 工具服务第二处是 Agent 里的工具名和参数比如把llm.complete换成你实际用的模型调用把asr.transcribe换成真实 ASR 服务。如果你要长期跑编码类或 Agent 类任务建议在 TaoToken 控制台开一个 Coding Plan地址是 https://taotoken.net/coding-plan Key 还是同一个但额度策略更适合长时间运行。Claude Code 相关的接入可以参考 https://taotoken.net/claudecode 里面有针对 Anthropic 协议的配置说明。再往后的升级方向我在骨架里留了几个扩展点Router 层加 Bearer 鉴权做租户隔离Agent 内部对 MCP 调用加重试和熔断把 MCP 的 HTTP 传输换成 SSE 或 stdio实现端到端全链路流式用/.well-known/agents做 Agent 动态注册。这些都不需要推翻现有结构加在对应层就行。最后说个实际经验这套骨架最大的价值不是代码本身而是它把 A2A 和 MCP 的边界划清楚了。A2A 管任务怎么流转MCP 管工具怎么调用两者通过 Agent 内部的MCPClient解耦。你后面换模型、换工具、加 Agent都不会互相牵连。先把这条链路跑顺再往上堆复杂度会稳很多。