消息¶
A2A 消息使用 JSON-RPC 2.0 envelope。MOI 在 envelope 顶层增加 Agent 选择器,params.message 则承载用户消息。普通调用和流式调用的消息结构相同,差别在于方法名与响应传输方式。
构造用户消息¶
下面的请求选择 explore Agent,并发送一个文本 Part:
{
"agent_code": "explore",
"jsonrpc": "2.0",
"id": "req_01J0A",
"method": "message/send",
"params": {
"message": {
"kind": "message",
"role": "user",
"messageId": "msg_01J0A",
"parts": [
{
"kind": "text",
"text": "汇总本季度各区域销售额"
}
]
}
}
}
ID 各有不同用途:
JSON-RPC
id关联请求和响应;每次调用都应唯一。messageId标识用户消息;重试同一逻辑消息时不要随意换 ID。服务端返回的 Task
id标识一次执行。contextId标识可供后续消息复用的对话上下文。
如需让 Agent 继续之前的上下文,把服务端返回的 contextId 放进下一条 message:
{
"kind": "message",
"role": "user",
"messageId": "msg_01J0B",
"contextId": "ctx_123",
"parts": [
{
"kind": "text",
"text": "只保留华东和华南,并按销售额降序排列"
}
]
}
不要使用 JSON-RPC id 或 Task id 代替 contextId。是否还需带上一轮 Task id 取决于具体 Agent;没有明确契约时只传服务端返回的上下文字段。
非流式调用¶
message/send 返回一个 JSON-RPC 文档。成功响应的 result 通常是 A2A Task,但客户端仍应先判断是否存在 error,再根据 result.kind 解析。
import os
import uuid
import requests
base_url = os.environ["MOI_BASE_URL"].rstrip("/")
payload = {
"agent_code": "explore",
"jsonrpc": "2.0",
"id": f"req_{uuid.uuid4().hex}",
"method": "message/send",
"params": {
"message": {
"kind": "message",
"role": "user",
"messageId": f"msg_{uuid.uuid4().hex}",
"parts": [{"kind": "text", "text": "列出本月销售额最高的五个产品"}],
}
},
}
response = requests.post(
f"{base_url}/newmoi/agents/a2a",
json=payload,
headers={"moi-key": os.environ["MOI_API_KEY"]},
timeout=120,
)
response.raise_for_status()
document = response.json()
if "error" in document:
raise RuntimeError(document["error"])
result = document["result"]
task_id = result.get("id") if result.get("kind") == "task" else None
context_id = result.get("contextId")
不要仅以 HTTP 200 作为完成信号。响应可能只创建了仍处于 working 的 Task,最终结果需要通过 Task 状态或流式事件确认。
流式调用¶
将方法改为 message/stream 后,请求体保持不变,响应类型变为 text/event-stream:
event: message
id: 12
data: {"jsonrpc":"2.0","id":"req_01J0A","result":{"kind":"status-update",...}}
一个事件可以包含以下几类结果:
首个成功事件通常包含 Task,使客户端尽早取得稳定的 Task
id。status-update更新 Task 状态。artifact-update传递文本、结构化数据或其他 Agent 产物。某些终止路径可能直接返回 A2A Message;解析器不应假定每个事件都是 Task 更新。
流中还可能出现 : heartbeat 注释。SSE 客户端必须忽略注释和空行,不要把心跳当作 JSON 解析。
下面的最小 Python 读取器保存事件序号,并支持多行 data:
import json
import requests
last_event_id = None
event_type = "message"
data_lines = []
with requests.post(
f"{base_url}/newmoi/agents/a2a",
json={**payload, "method": "message/stream"},
headers={
"moi-key": os.environ["MOI_API_KEY"],
"Accept": "text/event-stream",
},
stream=True,
timeout=(10, None),
) as response:
response.raise_for_status()
for line in response.iter_lines(decode_unicode=True):
if line == "":
if data_lines:
event = json.loads("\n".join(data_lines))
handle_a2a_event(event_type, last_event_id, event)
event_type, data_lines = "message", []
continue
if line.startswith(":"):
continue
if line.startswith("event:"):
event_type = line[6:].strip()
elif line.startswith("id:"):
last_event_id = line[3:].strip()
elif line.startswith("data:"):
data_lines.append(line[5:].lstrip())
handle_a2a_event 是应用自己的投影函数:按 result.kind 更新 Task 状态、追加 Artifact,并在持久化业务结果的同时保存 last_event_id。不要直接把未知 Data Part 当作文本展示。
超时、断线与重试¶
连接超时和读取超时应分开设置;长任务不适合使用普通 JSON 请求的短读取超时。
流在终态之前断开时,不要合成
completed或failed。使用最后保存的事件序号调用tasks/resubscribe。只有能够确认请求未到达服务端,或应用使用稳定的消息 ID/幂等策略时,才自动重发
message/stream。否则可能创建重复 Task。JSON-RPC
error、非 2xx HTTP 响应和 SSE 读取错误是三种不同故障,应分别记录。