鸿蒙智能体开发实战:6.A2A 模式开发接口

0 评论 255 浏览 0 收藏 15 分钟

前言

在上一篇文章中,我们完成了项目的初始化和基础架构搭建。本文将深入介绍 A2A 模式的核心接口开发,包括 JSON-RPC 2.0 协议的路由设计和实现。

JSON-RPC 2.0 协议回顾

JSON-RPC 是一种轻量级的远程过程调用协议,其请求格式如下:

{
  "jsonrpc": "2.0",
  "id": "request-123",
  "method": "message/stream",
  "params": {
    "key": "value"
  }
}

响应格式:

{
  "jsonrpc": "2.0",
  "id": "request-123",
  "result": {
    "data": "response data"
  }
}

统一端点设计

A2A 协议采用统一端点设计,所有 RPC 方法都通过 /agent/message 端点处理:

POST /agent/message
Headers:
  Content-Type: application/json
  agent-session-id: <session-id>

完整的路由实现

1. 消息分发器

# routes/agent_routes.py
from fastapi import APIRouter, Request, Header
from fastapi.responses import JSONResponse, StreamingResponse
import uuid
import json
from typing import OptionalDictAny

router = APIRouter(prefix="/agent", tags=["agent"])

# 共享状态存储
agent_sessions: Dict[strDict[strAny]] = {}
conversation_contexts: Dict[strlist] = {}
task_states: Dict[strDict[strAny]] = {}

@router.post("/message")
async def handle_agent_message(
    request: Request,
    agent_session_id: Optional[str] = Header(None, alias="agent-session-id"),
    x_request_id: Optional[str] = Header(None, alias="X-Request-ID"),
):
    """统一的 Agent 消息处理接口"""
    request_id = x_request_id or str(uuid.uuid4())

    # 解析请求体
    try:
        body = await request.json()
    except json.JSONDecodeError:
        return JSONResponse(
            status_code=400,
            content={"jsonrpc""2.0""id""unknown""error": {"code"400"message""Invalid JSON"}}
        )

    jsonrpc_request = JsonRpcRequest(**body)
    method = jsonrpc_request.method

    # 路由分发
    if method == "initialize":
        return await handle_initialize(jsonrpc_request, request_id)
    elif method == "notifications/initialized":
        return await handle_initialized(jsonrpc_request, agent_session_id, request_id)
    elif method == "message/stream":
        return await handle_message_stream(jsonrpc_request, agent_session_id, request_id)
    elif method == "tasks/cancel":
        return await handle_tasks_cancel(jsonrpc_request, agent_session_id, request_id)
    elif method == "clearContext":
        return await handle_clear_context(jsonrpc_request, agent_session_id, request_id)
    else:
        return JSONResponse(
            status_code=400,
            content=JsonRpcResponse(
                jsonrpc="2.0",
                id=jsonrpc_request.id,
                error={"code": -32601"message"f"Method not found: {method}"}
            ).model_dump()
        )

2. Initialize 方法

初始化会话,获取 agentSessionId

async def handle_initialize(
    request: JsonRpcRequest,
    request_id: str,
) -> JSONResponse:
    """处理 initialize 方法"""
    from datetime import datetime, timedelta

    session_id = uuid.uuid4().hex
    session_ttl = 7 * 24 * 60 * 60  # 7 天

    # 存储会话信息
    agent_sessions[session_id] = {
        "created_at": datetime.utcnow(),
        "expires_at": datetime.utcnow() + timedelta(seconds=session_ttl),
        "status""initialized",
    }

    response = JsonRpcResponse(
        jsonrpc="2.0",
        id=request.id,
        result={
            "version""1.0",
            "agentSessionId": session_id,
            "agentSessionTtl": session_ttl,
        }
    )

    return JSONResponse(content=response.model_dump())

3. Notifications/Initialized 方法

通知服务器初始化完成:

async def handle_initialized(
    request: JsonRpcRequest,
    agent_session_id: Optional[str],
    request_id: str,
) -> JSONResponse:
    """处理 notifications/initialized 方法"""

    if not agent_session_id or agent_session_id not in agent_sessions:
        return JSONResponse(
            status_code=401,
            content={"error""Invalid session ID"}
        )

    # 更新会话状态为 active
    agent_sessions[agent_session_id]["status"] = "active"

    return JSONResponse(content={})

4. Message/Stream 方法

核心的流式消息处理方法,支持 SSE 输出:

from fastapi.responses import StreamingResponse
import asyncio

async def handle_message_stream(
    request: JsonRpcRequest,
    agent_session_id: Optional[str],
    request_id: str,
) -> StreamingResponse:
    """处理 message/stream 方法,支持 SSE 流式输出"""

    params = request.params or {}
    task_id = params.get("id"str(uuid.uuid4()))
    session_id = params.get("sessionId"str(uuid.uuid4()))
    message = params.get("message", {})

    # 提取用户消息文本
    user_text = get_text_from_parts(message.get("parts", []))

    # 存储任务状态
    task_states[task_id] = {
        "sessionId": session_id,
        "status""working",
    }

    async def generate_sse():
        """SSE 流式生成器"""
        try:
            # 步骤 1: 发送 submitted 状态
            yield create_sse_event({
                "taskId": task_id,
                "kind""status-update",
                "status": {"state""submitted"}
            })

            # 步骤 2: 处理用户请求(调用大模型)
            async for chunk in process_user_request(user_text, session_id):
                yield create_sse_event(chunk)

            # 步骤 3: 发送完成状态
            yield create_sse_event({
                "taskId": task_id,
                "kind""status-update",
                "final"True,
                "status": {"state""completed"}
            })

        except Exception as e:
            yield create_sse_event({
                "taskId": task_id,
                "kind""status-update",
                "final"True,
                "status": {"state""failed"},
                "error"str(e)
            })

    return StreamingResponse(
        generate_sse(),
        media_type="text/event-stream",
        headers={
            "Cache-Control""no-cache",
            "Connection""keep-alive",
            "X-Accel-Buffering""no",
        }
    )

def get_text_from_parts(parts: list) -> str:
    """从 Message Parts 中提取文本内容"""
    texts = []
    for part in parts:
        if part.get("kind") == "text" and part.get("text"):
            texts.append(part.get("text"))
    return "\n".join(texts)

def create_sse_event(data: dict) -> str:
    """创建 SSE 事件字符串"""
    event = {
        "jsonrpc""2.0",
        "id": task_id,
        "result": data,
        "error": {"code"0"message""success"}
    }
    return f"data: {json.dumps(event, ensure_ascii=False)}\n\n"

5. Tasks/Cancel 方法

取消正在执行的任务:

async def handle_tasks_cancel(
    request: JsonRpcRequest,
    agent_session_id: Optional[str],
    request_id: str,
) -> JSONResponse:
    """处理 tasks/cancel 方法"""

    params = request.params or {}
    task_id = params.get("id", request.id)

    if task_id in task_states:
        task_states[task_id]["status"] = "canceled"

    return JSONResponse(
        content=JsonRpcResponse(
            jsonrpc="2.0",
            id=request.id,
            result={
                "id": task_id,
                "status": {"state""canceled"}
            }
        ).model_dump()
    )

6. ClearContext 方法

清理对话上下文:

async def handle_clear_context(
    request: JsonRpcRequest,
    agent_session_id: Optional[str],
    request_id: str,
) -> JSONResponse:
    """处理 clearContext 方法"""

    params = request.params or {}
    session_id = params.get("sessionId", request.sessionId)

    if session_id and session_id in conversation_contexts:
        conversation_contexts[session_id] = []

    return JSONResponse(
        content=JsonRpcResponse(
            jsonrpc="2.0",
            id=request.id,
            result={"status": {"state""cleared"}}
        ).model_dump()
    )

接口调用示例

初始化会话

curl -X POST http://localhost:8080/agent/message \
  -H "Content-Type: application/json" \
  -d '{
    "jsonrpc": "2.0",
    "id": 1,
    "method": "initialize"
  }'

响应:

{
  "jsonrpc": "2.0",
  "id": 1,
  "result": {
    "agentSessionId": "8f01f3d172cd4396a0e535ae8aec6687",
    "agentSessionTtl": 604800
  }
}

发送流式消息

curl -X POST http://localhost:8080/agent/message \
  -H "Content-Type: application/json" \
  -H "agent-session-id: 8f01f3d172cd4396a0e535ae8aec6687" \
  -d '{
    "jsonrpc": "2.0",
    "id": 2,
    "method": "message/stream",
    "params": {
      "id": "task-001",
      "sessionId": "session-001",
      "message": {
        "role": "user",
        "parts": [{
          "kind": "text",
          "text": "你好,请帮我介绍一下你自己"
        }]
      }
    }
  }'

小结

本文介绍了 A2A 模式的核心接口开发:

  1. 统一端点设计:所有 RPC 方法通过 /agent/message 路由
  2. JSON-RPC 2.0 协议:标准化的请求响应格式
  3. SSE 流式输出:实时推送任务状态和结果
  4. 会话管理:基于 agentSessionId 的状态跟踪

下一篇文章将详细介绍 A2A 模式的消息规范。

本文由 @少湖说 授权发布于人人都是产品经理。未经作者许可,禁止转载

题图来自Unsplash,基于CC0协议

该文观点仅代表作者本人,人人都是产品经理平台仅提供信息存储空间服务

更多精彩内容,请关注人人都是产品经理微信公众号或下载App
评论
评论请登录
  1. 目前还没评论,等你发挥!