编号:B-02 严重级:critical 工作线:broker(internal/broker) 来源:消息核心审查 M-01,总审查人已对照 mochi 源码复核 依赖:B-01 (#8)(同文件,先合) 被依赖:C-01 (#32)(消息线删除自己那份名额,即审查 M-02)、B-03 (#10)
Broker.PublishDown
internal/broker/broker.go:223-232
OnQosComplete
len(pk.Payload) > 64KiB
internal/broker/hooks.go:256-267
OnQosComplete(cl, pk)
pk
server.go:1161-1172
releaseAllLarge
PublishDown
messageLoops
cmd/nixmsg/serve.go:334-336
context.Background()
DEVELOPMENT 7.5「大帧同时在发不超过 64 个,客户端回 PUBACK(mochi 的 OnQosComplete)、确认超时或连接断开时释放」;DEVIATIONS N1/N2 第 4 条"QoS1 在 OnQosComplete 时释放"——意图明确,实现拿不到载荷。
OnQosPublish(cl, out, …)
connState.largePIDs
OnQosDropped
Provides
OnQosPublish
largeHeld++
internal/broker/hooks.go、internal/broker/broker.go,broker 测试。
internal/broker/hooks.go
internal/broker/broker.go
以下是本次复审各区审查报告的原文段落。A、M、I、P、S 开头的是原始发现编号(A 管理后台与网页、M 消息核心、I 身份认证群在线、P 传输平台部署、S SDK)。解决方案以本 issue 上方的"结论与统一方案"为准;原文里的方案与之不一致时,按上方执行。
WakePush
clearPushed
func (h *nixHook) OnQosComplete(cl *mqtt.Client, pk packets.Packet) { if len(pk.Payload) <= largeFrameBytes { return } h.b.connsMu.RLock() st := h.b.byClient[cl] h.b.connsMu.RUnlock() if st == nil { return } h.b.releaseOneLarge(st) }
processPuback
s.hooks.OnQosComplete(cl, pk)
internal/broker/broker.go:223-242
PushPending
internal/app/message/push.go:225-227
connState
context.WithTimeout(ctx, 5s)
broker.go
internal/app/message/push.go
cmd/nixmsg/serve.go
broker/hooks.go:256-267
broker/broker.go:223-243
server.go:1161-1170
cmd/nixmsg/serve.go:322-347
broker/hooks.go
broker/broker.go
largeSem
len(pk.Payload) <= largeFrameBytes
b.largeSem <- struct{}{}
loopCtx
server.Publish
func (h *nixHook) OnQosComplete(cl *mqtt.Client, pk packets.Packet) { if len(pk.Payload) <= largeFrameBytes { return } // ... h.b.releaseOneLarge(st) }
func (s *Server) processPuback(cl *Client, pk packets.Packet) error { // ... if ok := cl.State.Inflight.Delete(pk.PacketID); ok { // [MQTT-4.3.2-5] cl.State.Inflight.IncreaseSendQuota() atomic.AddInt64(&s.Info.Inflight, -1) s.hooks.OnQosComplete(cl, pk) }
其余位置:broker.go:223-243(阻塞获取名额,只有 QoS 0 立即释放);serve.go:330-338(ctx 为 loopCtx);message/push.go:225-230、549-560;uplink.go:298。
broker.go:223-243
serve.go:330-338
message/push.go:225-230
549-560
uplink.go:298
broker.PublishDown
releaseOneLarge
largeHeld
select
ctx.Done()
time.After(5*time.Second)
OnPublishDropped
hooks.go
session.go
releaseLarge
message/session.go
acquireLarge
len(b.largeSem)==0
复审基线:main 4059a15(2026-09-30)。编号说明、各工作线的合并顺序、共享文件归属见总览 #7。
4059a15
已合入 origin/main 0c9b459。落地提交 deb2398 fix: 大帧名额按 PacketID 在 PUBACK 与断线时归还 (#9)。
0c9b459
deb2398
No dependencies set.
The note is not visible to the blocked user.
编号:B-02 严重级:critical 工作线:broker(internal/broker) 来源:消息核心审查 M-01,总审查人已对照 mochi 源码复核
依赖:B-01 (#8)(同文件,先合) 被依赖:C-01 (#32)(消息线删除自己那份名额,即审查 M-02)、B-03 (#10)
现象与影响
Broker.PublishDown先阻塞获取全局 64 个名额之一(internal/broker/broker.go:223-232),设计上应在客户端回 PUBACK 时归还(DEVELOPMENT 7.5)。OnQosComplete钩子里,并按len(pk.Payload) > 64KiB判断(internal/broker/hooks.go:256-267)。但 mochi 调OnQosComplete(cl, pk)时传入的pk是客户端发来的 PUBACK 包本身(mochiserver.go:1161-1172),没有载荷,判断永远不成立。名额只在 Publish 出错、QoS 0 或连接断开(releaseAllLarge)时才还。PublishDown都阻塞。messageLoops传入的 ctx 没有超时(cmd/nixmsg/serve.go:334-336),全站的推送、到点分发、清理、确认超时处理一起冻结;握手路径在端的上行 worker 里用context.Background()推送,这个编号之后的 ack、send 乃至重连后的 hello 全部排在后面。文档依据
DEVELOPMENT 7.5「大帧同时在发不超过 64 个,客户端回 PUBACK(mochi 的 OnQosComplete)、确认超时或连接断开时释放」;DEVIATIONS N1/N2 第 4 条"QoS1 在 OnQosComplete 时释放"——意图明确,实现拿不到载荷。
解决方案
OnQosPublish(cl, out, …)钩子里(此时出站包带 PacketID 和载荷)把超过 64 KiB 的包 ID 记入connState.largePIDs。OnQosComplete、OnQosDropped按 PacketID 查表,命中才归还。在Provides里补上这两个钩子。断线时仍全部归还。OnQosPublish的次数,没有增加就归还。方案取舍(三份审查的建议不一致,统一如下)
releaseAllLarge与largeHeld++之间的竞态随 PacketID 表一起消除。改动文件
internal/broker/hooks.go、internal/broker/broker.go,broker 测试。交互 / 冲突说明
验收与测试
PublishDown1 秒内返回。问题明细(各区审查原文,证据含文件与行号)
[M-01] 服务端大帧名额收到 PUBACK 时从不释放,64 条后推送阻塞、每秒循环卡死(跨模块:连接线 broker)
PublishDown会阻塞等待全局 64 个名额之一,本应在客户端回 PUBACK 时释放。OnQosComplete的是客户端发来的 PUBACK 包本身,没有 payload,钩子里的大小判断永远直接返回。名额只在该连接断开时才还。PublishDown全部阻塞:messageLoops传的 ctx 没有超时,一旦阻塞就不再到点分发、不再做过期清理、不再处理确认超时。context.Background()推送,卡住后这个编号之后的 ack、send,乃至重连后的 hello 都排在它后面,设备无法再握手。WakePush协程 30 秒超时后,clearPushed复用同一个已过期的 ctx,清不掉「已推送」标记,形成幽灵在途(见 M-15)。server.go:1161-1172:processPuback调用的是s.hooks.OnQosComplete(cl, pk),这里的pk就是 PUBACK。internal/broker/broker.go:223-242:阻塞获取名额,只有 Publish 出错或 QoS0 时才释放。cmd/nixmsg/serve.go:334-336:循环里串行调用PushPending,ctx 无超时。internal/app/message/push.go:225-227:发布失败后clearPushed用的还是同一个 ctx。OnQosPublish(cl, out, …)里(此时出站包带 PacketID 和 payload)把超过 64 KiB 的包 ID 记到connState.largePIDs。OnQosComplete、OnQosDropped按 PUBACK 的 PacketID 查表,命中就释放。Provides里补上这两个钩子。断线时仍然全部释放。connState上用计数器比较 Publish 前后大帧OnQosPublish的次数,没增加就释放。PushPending调PublishDown时包一层context.WithTimeout(ctx, 5s);失败后的clearPushed另起一个短超时的新 ctx。internal/broker/hooks.go、broker.go(连接线目录);internal/app/message/push.go;cmd/nixmsg/serve.go。PublishDown1 秒内返回。[I-17] broker 的大帧名额在收到 PUBACK 时不释放,累计 64 条大帧后全局循环永久阻塞
processPuback传给 OnQosComplete 的是客户端发来的 PUBACK,不带 Payload。hook 里「只处理超过 64 KiB 的包」这个判断永远成立,直接返回,所以名额只在连接断开时才归还。broker/hooks.go:256-267;broker/broker.go:223-243;mochiserver.go:1161-1170;cmd/nixmsg/serve.go:322-347。OnQosPublish(包进入 inflight 时回调,带 PacketID 和 Payload)里记下大帧的 PacketID,OnQosComplete和OnQosDropped按 PacketID 释放;获取名额改为非阻塞,拿不到就返回错误,让消息线稍后重推。broker/hooks.go、broker/broker.go。[P-01] 大帧名额收到 PUBACK 也不释放:累计 64 条大于 64 KiB 的下行后,推送和调度主循环永久卡死
PublishDown对大于 64 KiB 的帧先占一个largeSem名额(全局 64 个)。QoS 1 的帧只在OnQosComplete或断线时归还。OnQosComplete时传的是客户端发来的 PUBACK 包,里面没有 Payload。钩子里len(pk.Payload) <= largeFrameBytes永远成立,直接返回,名额永不归还。b.largeSem <- struct{}{}。卡住时各路径的表现:loopCtx,整条循环永久停住。DispatchDue、所有连接的 PushPending、CleanupOnce、指标采样全部停。clearPushed用的是同一个已过期的 ctx,写库直接失败,投递标记清不掉;不保留的消息 5 分钟后按 not_acked 被丢弃。context.Background(),该端的上行队列从此不再消费。server.Publish遇到"没有订阅者"、"发送队列满被丢"、"inflight 已满"都返回 nil,名额挂在连接状态上,直到断线才还。releaseAllLarge执行后,PublishDown才做largeHeld++,这个竞态会永久丢一个名额。其余位置:
broker.go:223-243(阻塞获取名额,只有 QoS 0 立即释放);serve.go:330-338(ctx 为 loopCtx);message/push.go:225-230、549-560;uplink.go:298。broker.PublishDown里,server.Publish返回后不论 QoS 都立即releaseOneLarge,broker 的名额只保护"发布这一步"。同时删掉largeHeld、releaseAllLarge和OnQosComplete里的名额逻辑。在途数量的上限交给 message 包已有的largeSem。select同时等ctx.Done()和time.After(5*time.Second),超时返回错误。调用方按发布失败处理:清标记,1 秒后重推。OnQosPublish、OnQosComplete、OnQosDropped,按"客户端 + PacketID"记录大帧并据此释放。OnPublishDropped和"没有订阅者"的情况在PublishDown里同步判定并释放;为了能对上号,同一连接的大帧发布要串行。做完之后 message 包可以去掉自己的名额。internal/broker/broker.go、hooks.go;需要 M 线配合改internal/app/message/push.go、session.go。releaseLarge在断线清标记(message/session.go的 OnDisconnect)以及撤回、作废、清理结束时都没有调用。断线前没确认的大消息每次都漏一个名额,64 次之后acquireLarge永远失败,大消息只清标记、不再重推。所以必须和 M 线一起修,否则卡死点只是从 broker 挪到了 message。len(b.largeSem)==0。复审基线:main
4059a15(2026-09-30)。编号说明、各工作线的合并顺序、共享文件归属见总览 #7。已合入 origin/main
0c9b459。落地提交deb2398fix: 大帧名额按 PacketID 在 PUBACK 与断线时归还 (#9)。