A2A 协议已成功重构为采用 MCP (Model Context Protocol) 架构,将 EventMesh 定位为现代化的 智能体协作总线 (Agent Collaboration Bus)。
EnhancedA2AProtocolAdaptor)jsonrpc 字段自动分发处理逻辑。*.req 事件,属性 mcptype=request。*.resp 事件,属性 mcptype=response。id 映射到 CloudEvent collaborationid 来处理。params._agentId -> CloudEvent 扩展属性 targetagent (P2P)。params._topic -> CloudEvent Subject (Pub/Sub)。_topic 映射到 CloudEvent Subject,支持 O(1) 广播复杂度。message/sendStream 操作,映射为 .stream 事件类型,并通过 _seq -> seq 扩展属性保证顺序。JsonRpcRequest、JsonRpcResponse、JsonRpcError POJO 对象。McpMethods 常量,支持标准操作如 tools/call、resources/read。AgentCard、AgentSkill、AgentInterface、AgentCapabilities 等完整的 Agent 能力描述模型。eventmesh-runtime)完整的独立 HTTP Gateway 服务,桥接外部客户端到 A2A 事件总线。
| 组件 | 职责 |
|---|---|
A2AGatewayServer | Netty HTTP 服务器入口,预注册 mock agent,组装所有组件 |
A2AGatewayHttpHandler | HTTP 请求路由,支持 SSE 流式响应 |
A2AGatewayService | 核心编排:任务提交、响应处理、状态订阅、SSE 推送 |
TaskRegistry | 内存任务状态机 + TTL 自动清理 |
A2APublishSubscribeService | AgentCard 注册、发现、心跳管理 |
InMemoryA2AMessageTransport | 内存 pub/sub 实现(可替换为 EventMesh broker) |
A2ACardHttpHandler | AgentCard CRUD REST 端点 |
A2AClient | Java SDK,提供类型化 API |
| 方法 | 路径 | 说明 |
|---|---|---|
POST | /a2a/tasks?mode=sync | 同步提交任务 |
POST | /a2a/tasks?mode=async | 异步提交任务 |
GET | /a2a/tasks/{taskId} | 查询任务状态 |
DELETE | /a2a/tasks/{taskId} | 取消任务 |
GET | /a2a/tasks/{taskId}/wait | 长轮询等待结果 |
GET | /a2a/tasks/{taskId}/stream | SSE 流式推送状态更新 |
GET | /a2a/agents | 列出已注册 agents |
POST | /a2a/heartbeat | Agent 心跳 |
GET | /a2a/cards/list | 列出所有 AgentCard |
POST | /a2a/cards/card/{org}/{unit}/{agent} | 注册 AgentCard |
ScheduledExecutorService 每 60 秒扫描一次,清理超过 TTL(默认 5 分钟)的终态任务。TaskRegistry(taskTtlMs, cleanupIntervalMs) 构造函数支持自定义调优。InMemoryTransport 同步投递消息,若 transport.publish() 在 pendingTasks.put() 之前执行,handleResponse() 会先于 put() 运行,导致 future 永不完成。pendingTasks.put(taskId, future) 在 transport.publish() 之前执行,并添加注释说明顺序重要性。getTaskStatus() 返回 TaskResult 对象(而非原始 JSON 字符串),listAgents() 返回 List<String>(而非原始 JSON)。TaskResult.data 字段使用 @JsonAlias("result") 注解,兼容服务端 result 字段名。GET /a2a/tasks/{taskId}/streamDefaultHttpContent chunks),返回 null 跳过标准 FullHttpResponse 路径。通过 StatusSubscriber 回调实时推送状态变更。eventmesh-examples/.../demo/README.md,包含架构图、API 表、curl 示例、SDK 用法、运行方式。EnhancedA2AProtocolAdaptorTest 覆盖请求/响应循环、错误处理、通知和批处理。A2ATopicFactoryTest 覆盖 topic 生成与解析。TaskRegistryTest — 任务状态机 + TTL 清理验证InMemoryA2AMessageTransportTest — 内存传输投递A2AGatewayServiceTest — Gateway 服务层A2AGatewayEndToEndTest — 进程内全链路A2AClientServerIntegrationTest — 真实 HTTP 客户端-服务端集成测试McpIntegrationDemoTest、McpPatternsIntegrationTest、McpComprehensiveDemoTest、CloudEventsComprehensiveDemoTestInMemoryA2AMessageTransport,实现生产级部署。targetagent 和 a2amethod 扩展属性实现高级路由规则。methods/list)。TaskRegistry 状态持久化到 Redis/DB,支持崩溃恢复。