消息

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 请求的短读取超时。

  • 流在终态之前断开时,不要合成 completedfailed。使用最后保存的事件序号调用 tasks/resubscribe

  • 只有能够确认请求未到达服务端,或应用使用稳定的消息 ID/幂等策略时,才自动重发 message/stream。否则可能创建重复 Task。

  • JSON-RPC error、非 2xx HTTP 响应和 SSE 读取错误是三种不同故障,应分别记录。