编号:B-03 严重级:high 工作线:broker(internal/broker) 来源:总审查人实证 依赖:B-01 (#8)、B-02 (#9)(同文件,先合) 被依赖:B-04 (#11)、B-06 (#13);与消息线 C-01 (#32) 配合(接口不变,可并行)
mochi 的 WritePacket 在写网络期间一直持有客户端锁 cl.Lock()(mochi clients.go:602-631),而 QoS 1 发布要先 NextPacketID 取包号,同样要这把锁(clients.go:276-278,由 server.go:1072 publishToClient 调用)。客户端不读数据、TCP 发送缓冲写满时,写出方卡在网络写上并持锁,所有向这个客户端发布 QoS 1 的调用都跟着卡住。mochi 每读到一个入站包就刷新读写超时(clients.go:262-270、373),客户端只要还在发心跳,写就可以无限期阻塞。
WritePacket
cl.Lock()
clients.go:602-631
NextPacketID
clients.go:276-278
server.go:1072 publishToClient
clients.go:262-270
373
NixMsg 里在共享或其他端线程上同步调用 PublishDown(QoS 1)的地方:
PublishDown
cmd/nixmsg/serve.go:334-337
messageLoops
PushPending
Session.fatalKick
mochi 层复现:客户端订阅后停止读取、每 2 秒发一次 PINGREQ;服务端连续发 256 KiB 的 QoS 1。填满 socket 缓冲(约 17 条、4.3 MB)后,第 18 次 server.Publish 阻塞超过 10 秒,心跳仍在持续推后写超时。
server.Publish
DEVELOPMENT 7.5 要求"每个已握手连接一个推送循环",各端互不影响;PRD §8 要求 1000 在线端下的吞吐和延迟。现有实现让一个慢端影响全站,DEVIATIONS 无相关说明。
在 broker 层做每连接异步下发,让 PublishDown 永不阻塞调用方:
connState
OnSessionEstablished
OnDisconnect
ErrBackpressure
Session.publishJSON
port.Downlink
消息线只需确认把 ErrBackpressure 当作发布失败:清掉"已推送"标记并 1 秒后重推(internal/app/message/push.go:225-231 已有此路径),回执与 revoked 路径同样处理——这一小处改动由消息线在 C-01 里一并完成,不改接口签名。
internal/app/message/push.go:225-231
internal/broker/broker.go、hooks.go、新增 downlink.go、session.go(发布走新通道);internal/app/port/port.go(仅注释与哨兵错误);docs/DEVELOPMENT.md §5/§7.5;broker 测试。
internal/broker/broker.go
hooks.go
downlink.go
session.go
internal/app/port/port.go
docs/DEVELOPMENT.md
以下是本次复审各区审查报告的原文段落。A、M、I、P、S 开头的是原始发现编号(A 管理后台与网页、M 消息核心、I 身份认证群在线、P 传输平台部署、S SDK)。解决方案以本 issue 上方的"结论与统一方案"为准;原文里的方案与之不一致时,按上方执行。
loopCtx
func messageLoops(ctx context.Context, msgApp *message.App, conns *message.MemoryConns, db *store.DB, hashPool auth.HashPool, met *metrics.Registry) { t := time.NewTicker(time.Second) // ... case <-t.C: nowMs := time.Now().UnixMilli() if _, err := msgApp.DispatchDue(ctx, nowMs, 100); err != nil { slog.Error("dispatch due", "err", err) } for ep, live := range conns.Snapshot() { if err := msgApp.PushPending(ctx, ep, live.ConnID); err != nil {
func (cl *Client) NextPacketID() (i uint32, err error) { cl.Lock() defer cl.Unlock()
写包时持有同一把锁的位置:mochi clients.go:602-608;QoS 1 调用 NextPacketID 的位置:mochi server.go:1064-1072。
clients.go:602-608
server.go:1064-1072
cmd/nixmsg/serve.go
loops.go
复审基线:main 4059a15(2026-09-30)。编号说明、各工作线的合并顺序、共享文件归属见总览 #7。
4059a15
已合入 origin/main 0c9b459。落地提交 91e887b fix: 完成 broker 复审 B-03 至 B-12 (#10)。
0c9b459
91e887b
No dependencies set.
The note is not visible to the blocked user.
编号:B-03 严重级:high 工作线:broker(internal/broker) 来源:总审查人实证
依赖:B-01 (#8)、B-02 (#9)(同文件,先合) 被依赖:B-04 (#11)、B-06 (#13);与消息线 C-01 (#32) 配合(接口不变,可并行)
现象与影响
mochi 的
WritePacket在写网络期间一直持有客户端锁cl.Lock()(mochiclients.go:602-631),而 QoS 1 发布要先NextPacketID取包号,同样要这把锁(clients.go:276-278,由server.go:1072 publishToClient调用)。客户端不读数据、TCP 发送缓冲写满时,写出方卡在网络写上并持锁,所有向这个客户端发布 QoS 1 的调用都跟着卡住。mochi 每读到一个入站包就刷新读写超时(clients.go:262-270、373),客户端只要还在发心跳,写就可以无限期阻塞。NixMsg 里在共享或其他端线程上同步调用
PublishDown(QoS 1)的地方:cmd/nixmsg/serve.go:334-337的messageLoops串行对所有在线端PushPending:一个前台卡死但心跳仍在的手机,就能让全站推送、到点分发、清理停止;Session.fatalKick同步发布),都会被目标端拖住;实证
mochi 层复现:客户端订阅后停止读取、每 2 秒发一次 PINGREQ;服务端连续发 256 KiB 的 QoS 1。填满 socket 缓冲(约 17 条、4.3 MB)后,第 18 次
server.Publish阻塞超过 10 秒,心跳仍在持续推后写超时。为何判定为真实缺陷
DEVELOPMENT 7.5 要求"每个已握手连接一个推送循环",各端互不影响;PRD §8 要求 1000 在线端下的吞吐和延迟。现有实现让一个慢端影响全站,DEVIATIONS 无相关说明。
解决方案
在 broker 层做每连接异步下发,让
PublishDown永不阻塞调用方:connState增加有界 FIFO(按帧数与字节数双上限,例如 256 帧 / 16 MiB)和一个发送 goroutine:在OnSessionEstablished启动,在OnDisconnect停止并丢弃剩余帧。PublishDown只做同步检查(连接存在、大小上限、broker 未关闭),然后非阻塞入队;队列满时立即返回新的哨兵错误ErrBackpressure,不阻塞。发送 goroutine 按顺序调用server.Publish,同一连接内保持 FIFO 顺序。Session.publishJSON(resp、fatal、logout 响应)走同一通道。port.Downlink的注释里写明:返回 nil 只表示已入队;ErrBackpressure与其他发布失败同样处理。消息线只需确认把
ErrBackpressure当作发布失败:清掉"已推送"标记并 1 秒后重推(internal/app/message/push.go:225-231已有此路径),回执与 revoked 路径同样处理——这一小处改动由消息线在 C-01 里一并完成,不改接口签名。改动文件
internal/broker/broker.go、hooks.go、新增downlink.go、session.go(发布走新通道);internal/app/port/port.go(仅注释与哨兵错误);docs/DEVELOPMENT.md§5/§7.5;broker 测试。交互 / 冲突说明
PublishDown签名不变为契约,可并行开发。审查 P-07 对慢客户端拖住全局的分析附在下方。验收与测试
PublishDown在 100 ms 内返回并送达。PublishDown返回ErrBackpressure,且调用耗时小于 10 ms。问题明细(各区审查原文,证据含文件与行号)
[P-07] messageLoops 单 goroutine 串行:一个不读数据的客户端能拖住全局调度,DispatchDue 每秒最多 100 条
loopCtx。NextPacketID,它要拿cl.Lock()。写包时持有同一把锁的位置:mochi
clients.go:602-608;QoS 1 调用 NextPacketID 的位置:mochiserver.go:1064-1072。cmd/nixmsg/serve.go(可以拆出loops.go)复审基线:main
4059a15(2026-09-30)。编号说明、各工作线的合并顺序、共享文件归属见总览 #7。已合入 origin/main
0c9b459。落地提交91e887bfix: 完成 broker 复审 B-03 至 B-12 (#10)。