按 Agent 过滤、API key 鉴权的 SSE 事件流 /api/agent-events #1

Open
opened 2026-10-05 16:02:31 +08:00 by ila · 2 comments
Owner

来源与背景(2026-10-05)

用户希望各项目的 Agent(Claude Code、Codex 等)能实时收到 SynapBus 的新消息,不靠人提醒、也不靠定时轮询,并且做成后续通用的能力。

现状(代码基线 opc/main @ 352f9f6,上游 0d4a953):

  • 已有 SSE 推送 GET /api/events(internal/api/sse_handler.go,路由 internal/api/router.go:206),但只认网页登录的会话 Cookie。用 Agent API key 实测返回 401 No session cookie。
  • 推送中心 SSEHub 按「人类所有者 ownerID」分发(clients map[int64]...)。当前所有 Agent 都属于 admin,订阅一次就能看到全部 Agent 的消息事件。
  • 事件来源:SSEBroadcaster.OnMessageSent(internal/api/broadcaster.go),私信走 BroadcastDM,频道消息走 BroadcastChannelMessage(已经会取频道成员)。
  • MCP 端点已有「API key → Agent」鉴权:agents.RequiredAuthMiddlewareWithOAuth(internal/agents/middleware.go:134,挂载于 cmd/synapbus/main.go:825),认证后可用 agents.AgentFromContext 取得 Agent。

目标

新增一个按 Agent 过滤、用 Agent API key 鉴权的 SSE 事件流,作为通用推送基础。各客户端(Claude Code 后台监视、Codex wrapper、脚本、通知)都在它之上对接,SynapBus 核心不绑定任何特定客户端。

非目标

  • 不改动网页使用的 /api/events 的行为和鉴权。
  • 事件不携带消息正文:正文仍通过 MCP 读取,权限检查只走一处。
  • 不实现 MCP 服务端通知、Claude Code channels、webhook 回调本机等对接方式(以后可以在本事件流之上另行建单)。
  • 不改变监听地址,仍只监听 127.0.0.1;不新增数据库迁移。

方案(待用户确认)

  1. 新接口 GET /api/agent-events
    • 鉴权:复用 agents.RequiredAuthMiddlewareWithOAuth,Authorization: Bearer <agent key>;未认证返回 401。
    • 响应 text/event-stream;连接建立时推送 connected(含 agent 名),之后每 30 秒推送一次 heartbeat。
  2. 按 Agent 分发:SSEHub 增加按 agent 名索引的订阅集合,与现有按 owner 的集合并存。
    • 私信:只推给 to_agent。
    • 频道消息:只推给发送时该频道的成员。
    • 发送者本人不推。
  3. 事件格式:沿用 new_message,字段为 message_id、channel(频道消息)或 from_agent/to_agent(私信)、subject(有就带);SSE 的 id: 字段设为 message_id。
  4. 断线续传:客户端重连时带 Last-Event-ID: <message_id>,服务端先补发该 Agent 可见、id 更大的消息事件(限制最多 N 条,比如 200 条,超出时推送 resync_required),再进入实时推送。补发查询复用现有「我的消息」可见范围(私信 + 已加入的频道),不另写权限逻辑。
  5. 资源保护:每个 Agent 同时最多保持若干个连接(比如 5 个);写入阻塞时丢弃该连接,客户端用续传机制补回。

验收标准

  • 用 goauto 的 key 订阅:给 goauto 发私信、在 #goauto 发频道消息,都在 1 秒内收到 new_message,事件里没有正文。
  • 越权:goauto 的 key 收不到 erpgo 的私信事件,也收不到自己未加入的频道(如 #erpgo)的事件;自己发的消息不推给自己。
  • 无效或缺失的 key 返回 401;/api/events 网页推送行为不变。
  • 断线续传:带 Last-Event-ID 重连后,补发断线期间的事件,不重复、不遗漏;超过上限时推送 resync_required。
  • 连接数达到上限时拒绝新连接,并返回明确的错误。

测试

  • 单元/集成测试(表驱动,使用临时数据目录):分发过滤、越权、续传、上限、401。
  • go build ./cmd/synapbus,以及受影响包的 go test(internal/api、internal/agents、internal/messaging)。
  • 本机联调:临时用新二进制在另一个端口和临时数据目录启动,用 curl 订阅验证;替换 bin/synapbus.exe 并重启 8182 上的正式服务,须等用户确认。

风险与回退

  • 风险:权限过滤写错会把其他 Agent 的消息元数据推出去,因此越权测试是必须项;长连接会占用资源,因此设连接上限。
  • 回退:新接口独立存在,不影响现有接口;出问题时回退到上一版 bin/synapbus.exe,或在 opc/main 上 revert 合并提交。

分支与上游

  • 从 main 拉 feat/1-agent-sse,完成并验收后合入 opc/main。实现时尽量贴近上游风格,以后可以向上游提 PR。

文档影响

  • 更新 OPC Obsidian 笔记《SynapBus新项目接入手册》,加一节「实时订阅」。
  • 在本仓库补充接口说明(docs/ 或 README 的 OPC 段落)。

状态:方案待用户确认,确认前不写代码。

## 来源与背景(2026-10-05) 用户希望各项目的 Agent(Claude Code、Codex 等)能**实时**收到 SynapBus 的新消息,不靠人提醒、也不靠定时轮询,并且做成后续通用的能力。 现状(代码基线 `opc/main` @ `352f9f6`,上游 `0d4a953`): - 已有 SSE 推送 `GET /api/events`(`internal/api/sse_handler.go`,路由 `internal/api/router.go:206`),但只认网页登录的会话 Cookie。用 Agent API key 实测返回 401 `No session cookie`。 - 推送中心 `SSEHub` 按「人类所有者 ownerID」分发(`clients map[int64]...`)。当前所有 Agent 都属于 admin,订阅一次就能看到全部 Agent 的消息事件。 - 事件来源:`SSEBroadcaster.OnMessageSent`(`internal/api/broadcaster.go`),私信走 `BroadcastDM`,频道消息走 `BroadcastChannelMessage`(已经会取频道成员)。 - MCP 端点已有「API key → Agent」鉴权:`agents.RequiredAuthMiddlewareWithOAuth`(`internal/agents/middleware.go:134`,挂载于 `cmd/synapbus/main.go:825`),认证后可用 `agents.AgentFromContext` 取得 Agent。 ## 目标 新增一个**按 Agent 过滤**、用 Agent API key 鉴权的 SSE 事件流,作为通用推送基础。各客户端(Claude Code 后台监视、Codex wrapper、脚本、通知)都在它之上对接,SynapBus 核心不绑定任何特定客户端。 ## 非目标 - 不改动网页使用的 `/api/events` 的行为和鉴权。 - 事件不携带消息正文:正文仍通过 MCP 读取,权限检查只走一处。 - 不实现 MCP 服务端通知、Claude Code channels、webhook 回调本机等对接方式(以后可以在本事件流之上另行建单)。 - 不改变监听地址,仍只监听 127.0.0.1;不新增数据库迁移。 ## 方案(待用户确认) 1. **新接口** `GET /api/agent-events` - 鉴权:复用 `agents.RequiredAuthMiddlewareWithOAuth`,`Authorization: Bearer <agent key>`;未认证返回 401。 - 响应 `text/event-stream`;连接建立时推送 `connected`(含 agent 名),之后每 30 秒推送一次 `heartbeat`。 2. **按 Agent 分发**:`SSEHub` 增加按 agent 名索引的订阅集合,与现有按 owner 的集合并存。 - 私信:只推给 `to_agent`。 - 频道消息:只推给发送时该频道的成员。 - 发送者本人不推。 3. **事件格式**:沿用 `new_message`,字段为 `message_id`、`channel`(频道消息)或 `from_agent`/`to_agent`(私信)、`subject`(有就带);SSE 的 `id:` 字段设为 `message_id`。 4. **断线续传**:客户端重连时带 `Last-Event-ID: <message_id>`,服务端先补发该 Agent 可见、id 更大的消息事件(限制最多 N 条,比如 200 条,超出时推送 `resync_required`),再进入实时推送。补发查询复用现有「我的消息」可见范围(私信 + 已加入的频道),不另写权限逻辑。 5. **资源保护**:每个 Agent 同时最多保持若干个连接(比如 5 个);写入阻塞时丢弃该连接,客户端用续传机制补回。 ## 验收标准 - 用 `goauto` 的 key 订阅:给 `goauto` 发私信、在 `#goauto` 发频道消息,都在 1 秒内收到 `new_message`,事件里没有正文。 - **越权**:`goauto` 的 key **收不到** `erpgo` 的私信事件,也收不到自己未加入的频道(如 `#erpgo`)的事件;自己发的消息不推给自己。 - 无效或缺失的 key 返回 401;`/api/events` 网页推送行为不变。 - 断线续传:带 `Last-Event-ID` 重连后,补发断线期间的事件,不重复、不遗漏;超过上限时推送 `resync_required`。 - 连接数达到上限时拒绝新连接,并返回明确的错误。 ## 测试 - 单元/集成测试(表驱动,使用临时数据目录):分发过滤、越权、续传、上限、401。 - `go build ./cmd/synapbus`,以及受影响包的 `go test`(`internal/api`、`internal/agents`、`internal/messaging`)。 - 本机联调:临时用新二进制在另一个端口和临时数据目录启动,用 curl 订阅验证;替换 `bin/synapbus.exe` 并重启 8182 上的正式服务,须等用户确认。 ## 风险与回退 - 风险:权限过滤写错会把其他 Agent 的消息元数据推出去,因此越权测试是必须项;长连接会占用资源,因此设连接上限。 - 回退:新接口独立存在,不影响现有接口;出问题时回退到上一版 `bin/synapbus.exe`,或在 `opc/main` 上 revert 合并提交。 ## 分支与上游 - 从 `main` 拉 `feat/1-agent-sse`,完成并验收后合入 `opc/main`。实现时尽量贴近上游风格,以后可以向上游提 PR。 ## 文档影响 - 更新 OPC Obsidian 笔记《SynapBus新项目接入手册》,加一节「实时订阅」。 - 在本仓库补充接口说明(`docs/` 或 README 的 OPC 段落)。 状态:**方案待用户确认**,确认前不写代码。
Author
Owner

实施完成,待验收(2026-10-05)

分支与提交:feat/1-agent-sse(基于上游 main 0d4a953,已推送)

  • 2e8a49c feat(#1): add agent-key authenticated SSE stream /api/agent-events
  • def4bb5 fix(#1): keep out-of-order live agent events, add store tests
  • 5da3a5d fix(#1): encode agent-events error bodies with encoding/json
  • 合并:92a2e03 merge(--no-ff) 合入 opc/main,已推送。尚未替换正在运行的 8182 服务。

实现由 Sonnet 子 agent 完成,Claude Code(Opus)审核了三轮。

实现要点

  • 新接口 GET /api/agent-events,挂在 cmd/synapbus/main.go,与 MCP 共用 agents.RequiredAuthMiddlewareWithOAuth(Agent API key)。网页使用的 /api/events 行为不变。
  • SSEHub 新增按 agent 名索引的订阅集合:私信只推 to_agent;频道消息推给发送时的频道成员;不推给发送者本人。事件只带 message_id、channel 或 from_agent/to_agent、subject,不带正文;id: 行等于 message_id。
  • 续传:先注册订阅再补发,按 Last-Event-ID 补发 id 更大的可见消息,最多 200 条,超出时只发 resync_required。补发可见范围由 messaging.ListEventMetaAfter 决定(发给自己的私信 + 当前所在频道的消息,不含自己发的)。注意:补发按当前成员资格判断,加入频道前的消息也会补发(频道成员本来就能读历史)。
  • 每个 agent 最多 5 个连接,超出返回 429;广播用非阻塞发送,通道满了就丢弃该连接;每次写入设 10 秒超时;心跳 30 秒。
  • 没有订阅者时跳过成员和 subject 查询;错误响应统一用 encoding/json 编码。

审核发现并已修正

  1. 实时去重默认事件按 id 升序到达,并发发送时会静默丢消息 → 改为只对补发区间去重,并加了乱序测试。
  2. 权限边界查询缺少直接测试 → 新增 store 层表驱动测试。
  3. 子 agent 报告里说「加入前的频道消息不补发」与实现不符 → 注释已写清,报告已更正。
  4. 没有订阅者时也做数据库查询 → 改为直接跳过。
  5. 429 响应体不是合法 JSON(%q 插值产生了未转义的引号)→ 改用 encoding/json 编码,并在测试里解析验证。

验证(审核人亲自执行)

  • go vet:无报告。go test ./internal/api/ ./internal/agents/ ./internal/messaging/ ./internal/mcp/ -count=1:全部 ok。新测试 -count=3 稳定通过。合并后的 opc/main 编译并测试通过。
  • 端到端(临时端口 18183、全新临时数据、3 个 agent):
    • 私信只推收件人;频道消息只推成员,非成员收不到;不推自己发的;
    • 测试正文在事件流中出现 0 次;
    • 补发:每个 agent 只补到自己可见的 id;
    • 缺 key 或 key 无效返回 401;第 6 个连接返回 429,释放后恢复 200;/api/events 无会话仍返回 401。
    • 临时进程和数据已清理,8182 正式服务全程未动。
  • 未覆盖:-race(本机 ThreadSanitizer 无法启动),并发正确性依赖加锁设计、重复测试和端到端联调。

待办

  • 用户确认后,用 opc/main 编译替换 bin/synapbus.exe 并重启 8182(所有 agent 会短暂断线后重连)。
  • 更新 Obsidian《SynapBus新项目接入手册》,加一节「实时订阅」。
  • 另记:子 agent 发现上游 buildSearchConditions 用 to_agent = '' 匹配频道消息,但频道消息的 to_agent 存的是 NULL,搜索可能漏掉频道消息。这是上游的潜在缺陷,不在本单范围,建议另行建单。
## 实施完成,待验收(2026-10-05) **分支与提交**:`feat/1-agent-sse`(基于上游 `main` `0d4a953`,已推送) - `2e8a49c` feat(#1): add agent-key authenticated SSE stream /api/agent-events - `def4bb5` fix(#1): keep out-of-order live agent events, add store tests - `5da3a5d` fix(#1): encode agent-events error bodies with encoding/json - 合并:`92a2e03` merge(--no-ff) 合入 `opc/main`,已推送。**尚未替换正在运行的 8182 服务。** 实现由 Sonnet 子 agent 完成,Claude Code(Opus)审核了三轮。 ### 实现要点 - 新接口 `GET /api/agent-events`,挂在 `cmd/synapbus/main.go`,与 MCP 共用 `agents.RequiredAuthMiddlewareWithOAuth`(Agent API key)。网页使用的 `/api/events` 行为不变。 - `SSEHub` 新增按 agent 名索引的订阅集合:私信只推 `to_agent`;频道消息推给发送时的频道成员;不推给发送者本人。事件只带 `message_id`、`channel` 或 `from_agent`/`to_agent`、`subject`,**不带正文**;`id:` 行等于 message_id。 - 续传:先注册订阅再补发,按 `Last-Event-ID` 补发 id 更大的可见消息,最多 200 条,超出时只发 `resync_required`。补发可见范围由 `messaging.ListEventMetaAfter` 决定(发给自己的私信 + 当前所在频道的消息,不含自己发的)。注意:补发按**当前**成员资格判断,加入频道前的消息也会补发(频道成员本来就能读历史)。 - 每个 agent 最多 5 个连接,超出返回 429;广播用非阻塞发送,通道满了就丢弃该连接;每次写入设 10 秒超时;心跳 30 秒。 - 没有订阅者时跳过成员和 subject 查询;错误响应统一用 `encoding/json` 编码。 ### 审核发现并已修正 1. 实时去重默认事件按 id 升序到达,并发发送时会**静默丢消息** → 改为只对补发区间去重,并加了乱序测试。 2. 权限边界查询缺少直接测试 → 新增 store 层表驱动测试。 3. 子 agent 报告里说「加入前的频道消息不补发」与实现不符 → 注释已写清,报告已更正。 4. 没有订阅者时也做数据库查询 → 改为直接跳过。 5. 429 响应体不是合法 JSON(`%q` 插值产生了未转义的引号)→ 改用 `encoding/json` 编码,并在测试里解析验证。 ### 验证(审核人亲自执行) - `go vet`:无报告。`go test ./internal/api/ ./internal/agents/ ./internal/messaging/ ./internal/mcp/ -count=1`:全部 ok。新测试 `-count=3` 稳定通过。合并后的 `opc/main` 编译并测试通过。 - 端到端(临时端口 18183、全新临时数据、3 个 agent): - 私信只推收件人;频道消息只推成员,非成员收不到;不推自己发的; - 测试正文在事件流中出现 0 次; - 补发:每个 agent 只补到自己可见的 id; - 缺 key 或 key 无效返回 401;第 6 个连接返回 429,释放后恢复 200;`/api/events` 无会话仍返回 401。 - 临时进程和数据已清理,8182 正式服务全程未动。 - 未覆盖:`-race`(本机 ThreadSanitizer 无法启动),并发正确性依赖加锁设计、重复测试和端到端联调。 ### 待办 - [ ] 用户确认后,用 `opc/main` 编译替换 `bin/synapbus.exe` 并重启 8182(所有 agent 会短暂断线后重连)。 - [ ] 更新 Obsidian《SynapBus新项目接入手册》,加一节「实时订阅」。 - 另记:子 agent 发现上游 `buildSearchConditions` 用 `to_agent = ''` 匹配频道消息,但频道消息的 `to_agent` 存的是 NULL,搜索可能漏掉频道消息。这是上游的潜在缺陷,不在本单范围,建议另行建单。
Author
Owner

已上线到本机正式服务(2026-10-05)

  • 用户在 supervisor 中停止服务后,备份旧二进制为 bin/synapbus.exe.bak-20261005-pre-issue1(10:32 编译,版本 dev-local)。
  • 从 opc/main 92a2e03 编译新二进制(CGO_ENABLED=0、-ldflags "-s -w -X main.version=opc-92a2e03"),替换 bin/synapbus.exe。未重新构建前端,沿用本地已构建的资源。
  • 启动前冒烟测试:在临时端口 18184 和临时数据目录上运行,/health 正常,/api/agent-events 不带 key 返回 401,网页首页返回 200;临时实例已清理。
  • 用户重新启动服务后,/health 显示版本 opc-92a2e03。

正式服务实测

  • 用 goauto 的 key 订阅 /api/agent-events,再由 goauto_codex 给 goauto 发一条测试私信(#20):订阅端约 0.4 秒内收到 new_message,内容为 {"message_id":20,"from_agent":"goauto_codex","to_agent":"goauto","subject":"SSE 实时推送测试"},正文在事件流中出现 0 次。
  • 带 Last-Event-ID: 19 重连,正确补发 #20。
  • 测试私信 #20 已标记为已处理。

回退方式

停止服务,把 bin/synapbus.exe.bak-20261005-pre-issue1 改回 bin/synapbus.exe 后启动。本次没有数据库迁移,旧版本可以直接使用现有的 data/。

文档

Obsidian《SynapBus新项目接入手册-Claude-Code与Codex》已新增「实时订阅(/api/agent-events)」一节。

状态:待用户验收。

## 已上线到本机正式服务(2026-10-05) - 用户在 supervisor 中停止服务后,备份旧二进制为 `bin/synapbus.exe.bak-20261005-pre-issue1`(10:32 编译,版本 `dev-local`)。 - 从 `opc/main` `92a2e03` 编译新二进制(`CGO_ENABLED=0`、`-ldflags "-s -w -X main.version=opc-92a2e03"`),替换 `bin/synapbus.exe`。未重新构建前端,沿用本地已构建的资源。 - 启动前冒烟测试:在临时端口 18184 和临时数据目录上运行,`/health` 正常,`/api/agent-events` 不带 key 返回 401,网页首页返回 200;临时实例已清理。 - 用户重新启动服务后,`/health` 显示版本 `opc-92a2e03`。 ### 正式服务实测 - 用 `goauto` 的 key 订阅 `/api/agent-events`,再由 `goauto_codex` 给 `goauto` 发一条测试私信(#20):订阅端约 0.4 秒内收到 `new_message`,内容为 `{"message_id":20,"from_agent":"goauto_codex","to_agent":"goauto","subject":"SSE 实时推送测试"}`,正文在事件流中出现 0 次。 - 带 `Last-Event-ID: 19` 重连,正确补发 #20。 - 测试私信 #20 已标记为已处理。 ### 回退方式 停止服务,把 `bin/synapbus.exe.bak-20261005-pre-issue1` 改回 `bin/synapbus.exe` 后启动。本次没有数据库迁移,旧版本可以直接使用现有的 `data/`。 ### 文档 Obsidian《SynapBus新项目接入手册-Claude-Code与Codex》已新增「实时订阅(/api/agent-events)」一节。 状态:待用户验收。
Sign in to join this conversation.
No labels
1 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: OPC/synapbus#1