| A1 |
Hub /client/ws |
S→C |
auth.ok / message.new / agent.stream / device.kicked …(31 常量) |
Hub PushToConn 内 c.seq.Add(1)(sendMu 临界区) |
纯内存 per-connection atomic.Int64;不落库、不跨连接 |
无 last_seq/resume/cursor;仅重做 Sec-WebSocket-Protocol bearer auth |
服务端无重放;PushToConn 不缓存历史帧。客户端对 agent.stream 单独走 REST gap-fill(GET /web/agent-tasks/:id/events?after_seq=) |
业务键 UPSERT / idempotent-on-apply / watermark(见 api/events.md 表);seq_id 不参与去重 |
at-least-once(fanout + outbox redispatch) |
frame.go:5-17, fanout.go:65-91, conn.go:41-45, handler/ws.go:103-108, api/events.md:27-29 |
| A2 |
Hub /client/ws |
C→S |
typing(唯一合法 C→S) |
n/a |
n/a |
n/a |
n/a |
processIncoming switch 只接受 TypeTyping,其余丢弃;typing 自身 ephemeral 无需去重 |
ephemeral |
handler/ws.go:142-179, frame.go:21 |
| B1 |
Edge /v1/events |
S→C |
run.agent.* / system.gap / error / artifact.* …(EventEnvelope) |
Bus.Publish atomic.AddInt64(&b.seq, 1) |
内存 + 磁盘 eventLog;重启从 eventLog.orderedSeq 末尾续接;history ring 10000 + eventLog 50 MiB(保留 75%) |
query ?cursor=<lastSeq>(或 pageCursor) |
Subscribe(cursor) merge history + eventLog;cursor predates log → 注入 system.gap;handler 收到 gap 发 WS close code 4001 强制客户端 resync |
客户端水位 drop seq<=lastSeq(system.gap 不污染 cursor)+ 业务键 UPSERT |
at-least-once(replay + live fanout + gap→resync) |
bus.go:149-194, bus.go:138-144, bus.go:333-388, eventlog.go:48, handlers_events.go:24-90, handlers_events.go:152-161, handlers.go:110-113, types.go:4 |
| B2 |
Edge /v1/events |
C→S |
ping(clientControl chan) |
n/a |
n/a |
n/a |
n/a |
websocketClientControlResponse 只识别 "type":"ping",回 {type:"pong", ts};其余忽略 |
ephemeral(heartbeat) |
handlers_events.go:180-195 |
| C |
Hub↔Edge |
— |
无 WS;HTTP outbox + callback + relay-via-PushToUser |
n/a(delivery_id 是 outbox journal key,非 wire seq) |
outbox PG journal(pending/sent/retrying/delivered/dead);callback fire-and-forget 3 attempts / 30s timeout / 10s budget |
n/a |
outbox ScanRetryableDeliveries 每 15s;PendingTimeout 30s / SentTimeout 60s / DefaultMaxAttempts 3 / backoff 2s base cap 30s ±25% jitter;relay CreateCommand PushToUser 到 client WS(不是独立 edge WS) |
delivery_id + task_id 幂等;MarkDeliverySent CAS;AckDelivery 幂等 |
at-least-once(journal + redispatch + dead-letter) |
deliveryoutbox/outbox.go:50-140, retry.go:11-49, dispatchsvc/agent_dispatch.go:395-445, relay/relay.go:39-67, hub/callback.go:38-70 |
基线
4b3f236b66cd417b600beb032e0c278891849ce9(2026-08-29)/client/ws+ Edge/v1/events两条 WS 链路;Hub↔Edge 无 WS(HTTP outbox/callback)契约矩阵
/client/wsPushToConn内c.seq.Add(1)(sendMu 临界区)atomic.Int64;不落库、不跨连接agent.stream单独走 REST gap-fill(GET /web/agent-tasks/:id/events?after_seq=)/client/ws/v1/eventsatomic.AddInt64(&b.seq, 1)eventLog.orderedSeq末尾续接;history ring 10000 + eventLog 50 MiB(保留 75%)?cursor=<lastSeq>(或 pageCursor)system.gap;handler 收到 gap 发 WS close code 4001 强制客户端 resyncseq<=lastSeq(system.gap 不污染 cursor)+ 业务键 UPSERT/v1/events"type":"ping",回{type:"pong", ts};其余忽略缺口清单(按严重度排序)
G1 [High] Hub
/client/ws客户端未实现 seq_id gap 检测 → IM/session/device 事件静默丢失fanout.go:64-65注释明示 "dropped frames consume a seq too: the resulting gap is the client-side loss signal",但app/shared/src/hub/hubWS.ts全文不解析seq_id、不维护期望 seq、不检测 gap。buffer-full drop(fanout.go:105-119)和 marshal error 产生的 seq hole 在客户端完全不可见。hub:gap让 UI 显示同步指示器。G2 [High] Hub outbox redispatch 与 WS PushToConn 双通道可对同一 task 产生并发重复投递,Edge 端缺少 delivery_id 级去重契约文档
G3 [Medium] Edge eventLog 50 MiB 截断 + history 10000 条在高频 stream 下窗口过窄,gap 触发全量 resync 代价高
/web/agent-tasks/:id/events也有上限,可能进一步丢中段事件。firstDroppedSeq/lastDroppedSeq(已有字段但 handler 未利用),客户端可精准 REST backfill 而非全量 resync;③ 增加 eventLog 压缩(zstd/jsonlines)延长有效窗口。G4 [Medium] Hub
/client/wsreconnect 无 resume 语义,多端 fanout 下离线期间消息仅靠 REST syncMessages 补齐,syncMessages 自身无 after_seq 游标契约after_seq参数,但 api/events.md 未声明其语义是 "message seq" 还是 "ws seq_id"(实际应是 message 表的内部 seq,与 ws seq_id 无关)。G5 [Low] 两条链路 envelope 字段命名不对称(seq_id vs seq),客户端 parser 易混用
seq_id(frame.go:14),Edge EventEnvelope 用seq(bus.go:163)。前端 hubWS.ts 完全不读 seq_id;eventClient.ts 只读 seq。但若未来 hubWS.ts 加 gap 检测,开发者可能误用seq字段名导致静默失败。seq(breaking change,需版本化)。G6 [Low] relay CreateCommand 经 PushToUser(targetEdgeID) 推送,targetEdgeID 被当作 userID 路由
s.mgr.PushToUser(targetEdgeID, frame)。ws.Manager.byUser 以 userID 为键(fanout.go:128),但此处传入的是 edge device ID。若 edge device ID ≠ userID(通常如此),帧将被 fanout 到零连接,静默丢失。修复切片建议(≤3 条优先级最高)
关联 Issue/PR
审计局限性