# a2a-broker **Repository Path**: fa0/a2a-broker ## Basic Information - **Project Name**: a2a-broker - **Description**: No description available - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: main - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-06-28 - **Last Updated**: 2026-07-07 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # xz-a2a-broker WebSocket **A2A(agent2agent)服务端**,参考 `../broker`(MCP broker)实现,用 TypeScript 编写。 居中做三件事:① 让远程 A2A Agent 连进来并发现其能力;② 接收 xz-chat 的任务派发;③ 把远程 Agent 的流式进度发成 `inject_system_message` 经 Redis 推回 xz-chat。 ``` xz-chat ──HTTP 派发──► a2a-broker ──A2A/WS message/stream──► 远程 Agent ▲ │ ▲ │ └──Redis inject◄──────────┘ └────────── 流式进度 ◄──────────────┘ ``` > 设计来源:`xz-chat/docs/a2a-server.md`。三段传输:**xz-chat↔broker = HTTP(去程)+ Redis(回程)**,**broker↔Agent = WebSocket**。 ## 运行 ```bash npm install cp .env.example .env # 可选:按部署环境改端口 / 监听地址 / Redis npm run dev # tsx 直接跑 src(开发) # 或 npm run build && npm start ``` 需要一个可用的 Redis(`config/server.json` 或 `.env` 里配 host/port);该 Redis 必须与 xz-chat 同一实例。 ## 环境变量(`.env`) 启动时自动加载项目根目录的 `.env`(`dotenv`),对 `npm run dev` 和 `npm start` 都生效。**优先级:环境变量 / `.env` > `config/server.json` > 默认值**;未设置的变量回退到 JSON 配置。命令行直接传(如 `PORT=9000 npm start`)优先于 `.env`。 | 变量 | 覆盖 | 默认 | 说明 | | --- | --- | --- | --- | | `PORT` | `port` | `3100` | HTTP + WS 监听端口 | | `HOST` | —(仅监听地址) | `0.0.0.0` | 监听网卡:`0.0.0.0`=所有网卡(允许外网);`127.0.0.1`=仅本机 | | `REDIS_HOST` | `redis.host` | — | Redis 主机 | | `REDIS_PORT` | `redis.port` | — | Redis 端口 | | `REDIS_PASSWORD` | `redis.password` | — | Redis 密码(不设则用 JSON 值,可为空) | | `REDIS_DB` | `redis.db` | `0` | Redis 库号 | | `DEBUG` | — | — | debug 日志命名空间(设 `a2a-broker` 开启,留空关闭) | `.env` 已被 `.gitignore` 忽略,不会提交;模板见 `.env.example`。 ## 配置 `config/server.json`(热重载) > `port` 与 `redis.*` 可被环境变量覆盖(见上「环境变量」);其余字段仅走 JSON。 | 字段 | 说明 | | --- | --- | | `port` | HTTP + WS 监听端口(默认 3100) | | `request_timeout` | 对 Agent 的单次请求超时(ms,取卡片 / tasks-cancel 用;不约束流式任务) | | `task_idle_timeout` | 任务空闲超时(ms,默认 120000):连续这么久没有进度则判定 Agent 卡死,停止转发并注入「任务超时」 | | `inject_min_interval` | 中间进度注入节流(ms,默认 1000):同一任务每秒最多注入一条 `delay` 进度,间隔内连续来的直接忽略丢弃;终态(`auto`)必达,不受此限 | | `auth.mode` | `none`(dev,从 query 取身份)/ `jwt`(ES256 校验 `?token=`) | | `redis` | Redis 连接(发布 `inject_system_message`) | | `blacklist_ip` | WS 接入 IP 黑名单 | ## HTTP 接口(面向 xz-chat) > 所有 HTTP 路由与 WS 升级都挂在 **`/a2a`** 前缀下(便于反代 / 共用网关)。下文路径已含前缀。 - `GET /a2a/agents?agent_id=` —— 能力发现。**按 `agent_id` 过滤**:只返回绑定到该 xz-chat 智能体的远程 Agent(外加未绑定的全局 Agent);不带 `agent_id` 则返回全部(调试用)。返回在线 Agent + 由 Agent Card 聚合的 **LLM 工具定义**: ```json { "agents": [{ "a2a_id": "travel", "name": "出行助手", "skills": [{ "id": "train", "name": "火车票查询" }], "tool": { "name": "a2a_travel", "description": "…", "inputSchema": { "type":"object","properties":{"task":{"type":"string"}},"required":["task"] } } }] } ``` - `POST /a2a/tasks/dispatch` —— 派发任务,**即时返回 ack**: ```jsonc // 入参 { "task": "帮我查明天到上海的高铁并比价", "targetAgent": "travel", "agentId": 123, "chatId": "abc", "macAddress": "..." } // 返回 { "taskId": "uuid", "status": "accepted" } // 或 { "status": "no_agent" } ``` - `POST /a2a/tasks/cancel` —— 按 `chatId` 取消该会话在跑的任务:`{ "chatId": "abc" }` → `{ "canceled": 2 }` - `GET /a2a/health` / `GET /a2a/stats` ## 回程:Redis 事件(与 xz-chat activity-contract 对齐) 每段进度 `PUBLISH` 到频道 `agent:events:`(`agentId` 为 xz-chat 侧的 agentId): ```jsonc // 中间进度(working)→ delay:只插入历史、不触发回复,等用户下次说话再带出 { "type": "inject_system_message", "agentId": 123, "chatId": "abc", "text": "[出行助手] 已查到 6 个车次…", "mode": "delay" } // 最终结果(completed/failed/timeout,终态)→ auto:在安全时机主动播报 { "type": "inject_system_message", "agentId": 123, "chatId": "abc", "text": "[出行助手] 已订好备选,最优 G1234。", "mode": "auto" } ``` xz-chat 现成订阅 → `DeviceConnection.onDeviceEvent` → `ChatSession.injectSystemEvent`。 > **按任务状态选 mode**(见 `onAgentEvent`):非终态进度(`TASK_STATE_WORKING`)用 **`delay`**——只插入对话历史、不打断、不触发回复,等用户下次说话再带出;终态(`TASK_STATE_COMPLETED` 及 `failed`/`canceled`/超时)用 **`auto`**——最终结果在安全时机主动讲给用户。这样「边跑只攒、跑完才说」。 > **中间进度节流**(见 `allowInject`):`delay` 进度按 `inject_min_interval`(默认 1000ms)**每任务每秒最多注入一条**,窗口内连续来的进度直接**忽略丢弃**(leading-edge,丢弃不补发),避免 Agent 刷屏。**终态 `auto` 必达**,不走节流。 ## 远程 Agent 接入约定(WS pipe,参考 mcp_pipe.py) 拓扑同 MCP:Agent 端用一个 **ws pipe**(参考 `mcp-calculator/mcp_pipe.py`)把本地 A2A 服务(stdio JSON-RPC)**拨号连到 broker**——broker 是 WS 服务端,Agent 在 pipe 后面(无公网 HTTP)。pipe 只透传:WS→子进程 stdin、stdout→WS。broker↔Agent 全程走行分隔 JSON-RPC 2.0(一条消息一行 JSON,与 mcp_pipe 一致)。 连接 URL(pipe 的 endpoint):`ws://:3100/a2a?a2a_id=<名字>&agent_id=`(**WS 也在 `/a2a` 前缀下**)。 - `a2a_id`:远程 Agent 的名字(= 派发 targetAgent)。 - `agent_id`(可选):把本 Agent **绑定到某个 xz-chat 智能体 `agent.id`**——只有该智能体的设备能在 `/a2a/agents` 发现它、也只有它能派发;**留空则全局可见**(所有设备都能用)。 - `auth=jwt` 时改带 `?token=`,`a2a_id`/`agent_id` 从 token 的 claim 取;旧的 `?agentId=` 仍兼容。 - **能力发现(broker→Agent)**:broker 连上后**主动发请求**取 Agent Card(不是 HTTP 拉 well-known —— pipe 后的 Agent 拉不到): `{ "jsonrpc":"2.0","id":1,"method":"agent/authenticatedExtendedCard","params":{} }` Agent 回 `{ "jsonrpc":"2.0","id":1,"result":{ …AgentCard(含 skills)… } }`,broker 据此建工具定义。(Agent 若直连 HTTP 可达,可在 URL 带 `&cardUrl=` 让 broker 改走 HTTP 拉。) - **任务下发(broker→Agent)**:标准 A2A JSON-RPC ```json { "jsonrpc":"2.0","id":1,"method":"message/stream", "params":{ "message":{ "role":"user","parts":[{"kind":"text","text":""}],"messageId":"" } } } ``` - **进度回流(Agent→broker)**:同一 `id` 回多条事件(`Task` / `status-update` / `artifact-update`),最后一条带 `final:true` 或终态(`completed/failed/canceled`): ```json { "jsonrpc":"2.0","id":1,"result":{ "kind":"status-update","taskId":"t1", "status":{ "state":"working","message":{ "role":"agent","parts":[{"kind":"text","text":"已查到 6 个车次…"}] } },"final":false } } ``` broker 抽取其中文本,逐条经 Redis 注入对话。 - **取消(broker→Agent)**:`{ "jsonrpc":"2.0","id":N,"method":"tasks/cancel","params":{"id":""} }` - **心跳**:WS 层 ping/pong,60s 一轮。 ## 文件结构 ``` src/ index.ts 入口 A2aBroker.ts 主服务(Express HTTP + WS,任务编排,Redis 回推) AgentConnection.ts 单个远程 Agent 的 WS 连接(A2A message/stream 流式) a2a.ts A2A 事件/卡片工具(拉 well-known、抽文本、判终态、建工具定义) redis.ts inject_system_message 发布器 config.ts 配置热重载 types.ts 类型 ```