[B-02][critical] broker 大帧名额收到 PUBACK 时从不归还,累计 64 条大帧后 PublishDown 永久阻塞、messageLoops 冻结 #9

Closed
opened 2026-09-30 13:56:51 +08:00 by nixevol · 1 comment
Owner

编号:B-02 严重级:critical 工作线:broker(internal/broker) 来源:消息核心审查 M-01,总审查人已对照 mochi 源码复核
依赖:B-01 (#8)(同文件,先合) 被依赖:C-01 (#32)(消息线删除自己那份名额,即审查 M-02)、B-03 (#10)

现象与影响

  • 对超过 64 KiB 的 QoS 1 下行,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 包本身(mochi server.go:1161-1172),没有载荷,判断永远不成立。名额只在 Publish 出错、QoS 0 或连接断开(releaseAllLarge)时才还。
  • 结果:累计向在线端发出 64 条大帧后,之后所有大帧的 PublishDown 都阻塞。messageLoops 传入的 ctx 没有超时(cmd/nixmsg/serve.go:334-336),全站的推送、到点分发、清理、确认超时处理一起冻结;握手路径在端的上行 worker 里用 context.Background() 推送,这个编号之后的 ack、send 乃至重连后的 hello 全部排在后面。
  • 一条 100 KiB 的消息发给 64 人以上的在线群,就能一次占满名额。PRD F07 允许 256 KiB 正文。

文档依据

DEVELOPMENT 7.5「大帧同时在发不超过 64 个,客户端回 PUBACK(mochi 的 OnQosComplete)、确认超时或连接断开时释放」;DEVIATIONS N1/N2 第 4 条"QoS1 在 OnQosComplete 时释放"——意图明确,实现拿不到载荷。

解决方案

  1. 在 OnQosPublish(cl, out, …) 钩子里(此时出站包带 PacketID 和载荷)把超过 64 KiB 的包 ID 记入 connState.largePIDs。OnQosComplete、OnQosDropped 按 PacketID 查表,命中才归还。在 Provides 里补上这两个钩子。断线时仍全部归还。
  2. 一次 Publish 占了名额却没有产生 inflight(没有订阅者、inflight 已满被静默丢、连接已关),要立即归还:比较 Publish 前后该连接大帧 OnQosPublish 的次数,没有增加就归还。
  3. 若 B-03(每连接异步下发)先落地,名额的获取移到每连接的发送 goroutine 里,归还规则不变。
  4. 更新 DEVIATIONS N1/N2 第 4 条。

方案取舍(三份审查的建议不一致,统一如下)

  • 采用:按 PacketID 在 PUBACK / 丢弃 / 断线时归还(审查 M-01、I-17 的方案),与 DEVELOPMENT 7.5 原意一致,只在 broker 保留一份名额,消息包那份由消息线在 C-01 删除。
  • 不采用:审查 P-01 的"发布完成即归还、在途上限交给消息包的名额"。消息包那份名额本身在断线、撤回、作废时都会泄漏(审查 M-02,附在 C-01),而且在途时间会被拉长到应用层确认(最长 5 分钟),会把卡死点从 broker 挪到消息包。
  • 吸收 P-1 的两点:获取名额必须有界等待,超时按发布失败处理(消息线据此清标记、1 秒后重推);releaseAllLarge 与 largeHeld++ 之间的竞态随 PacketID 表一起消除。

改动文件

internal/broker/hooks.go、internal/broker/broker.go,broker 测试。

交互 / 冲突说明

  • 消息线 C-01 会删除 message 包里重复的一份大帧名额,依赖本 issue 先合入。
  • 与 B-01、B-03 同文件,按 broker 线顺序合入。

验收与测试

  • broker 单测:同一在线客户端连续收 65 个 70 KiB 帧并自动回 PUBACK,第 65 次 PublishDown 1 秒内返回。
  • 不回 PUBACK 的客户端断开后名额全部归还;没有订阅者时发布后名额立即归还。
  • 集成:给 100 人在线群发 100 KiB 消息后,再发一条定时消息,仍按时分发。

问题明细(各区审查原文,证据含文件与行号)

以下是本次复审各区审查报告的原文段落。A、M、I、P、S 开头的是原始发现编号(A 管理后台与网页、M 消息核心、I 身份认证群在线、P 传输平台部署、S SDK)。解决方案以本 issue 上方的"结论与统一方案"为准;原文里的方案与之不一致时,按上方执行。

[M-01] 服务端大帧名额收到 PUBACK 时从不释放,64 条后推送阻塞、每秒循环卡死(跨模块:连接线 broker)

  • 严重级:critical
  • 分类:并发 / 逻辑
  • 现象与影响:
    • 对超过 64 KiB 的 QoS1 下行,PublishDown 会阻塞等待全局 64 个名额之一,本应在客户端回 PUBACK 时释放。
    • 但 mochi 传给 OnQosComplete 的是客户端发来的 PUBACK 包本身,没有 payload,钩子里的大小判断永远直接返回。名额只在该连接断开时才还。
    • 累计 64 条大帧发给仍在线的端之后,后续大帧的 PublishDown 全部阻塞:
      • messageLoops 传的 ctx 没有超时,一旦阻塞就不再到点分发、不再做过期清理、不再处理确认超时。
      • 握手路径在该端的上行 worker 里用 context.Background() 推送,卡住后这个编号之后的 ack、send,乃至重连后的 hello 都排在它后面,设备无法再握手。
      • WakePush 协程 30 秒超时后,clearPushed 复用同一个已过期的 ctx,清不掉「已推送」标记,形成幽灵在途(见 M-15)。
    • 一条 100 KiB 的消息发给 64 人以上的在线群,就能一次占满名额。
  • 证据:
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)
}
  • mochi 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。
  • 文档依据:DEVELOPMENT 7.5「大帧同时在发不超过 64 个,客户端回 PUBACK(mochi 的 OnQosComplete)、确认超时或连接断开时释放」。
  • 为何不是故意设计:DEVIATIONS 连接线 N1/N2 第 4 条写的就是「QoS1 在 OnQosComplete 时释放」。意图明确,只是实现拿不到 payload,结果从不释放。
  • 解决方案:
    1. broker 在 OnQosPublish(cl, out, …) 里(此时出站包带 PacketID 和 payload)把超过 64 KiB 的包 ID 记到 connState.largePIDs。OnQosComplete、OnQosDropped 按 PUBACK 的 PacketID 查表,命中就释放。Provides 里补上这两个钩子。断线时仍然全部释放。
    2. 如果一次 Publish 占了名额却没有产生 inflight(没有订阅者、inflight 已满被静默丢、连接已关),要立即归还。可以在 connState 上用计数器比较 Publish 前后大帧 OnQosPublish 的次数,没增加就释放。
    3. PushPending 调 PublishDown 时包一层 context.WithTimeout(ctx, 5s);失败后的 clearPushed 另起一个短超时的新 ctx。
    4. 更新 DEVIATIONS N 第 4 条。
  • 改动文件:internal/broker/hooks.go、broker.go(连接线目录);internal/app/message/push.go;cmd/nixmsg/serve.go。
  • 交互/冲突风险:和 M-02 是同一类问题,建议一起修。broker 属于连接线目录,要按 TASKS 4.2 走共享改动流程。需要确认会话继承时 mochi 重发 inflight 不会让同一个包 ID 被记两次。
  • 需补测试:
    • broker 单测:同一在线客户端连续收 65 个 70 KiB 的帧并自动回 PUBACK,第 65 个 PublishDown 1 秒内返回。
    • broker 单测:不回 PUBACK 的客户端断开后名额全部归还;没有订阅者时发布后名额立即归还。
    • 集成:给 100 人在线群发 100 KiB 消息后,再发一条定时消息,仍能按时分发。
  • 置信度:代码阅读确定(已对照 mochi 源码)。

[I-17] broker 的大帧名额在收到 PUBACK 时不释放,累计 64 条大帧后全局循环永久阻塞

  • 严重级:critical
  • 分类:并发 / 质量
  • 现象与影响
    • mochi 的 processPuback 传给 OnQosComplete 的是客户端发来的 PUBACK,不带 Payload。hook 里「只处理超过 64 KiB 的包」这个判断永远成立,直接返回,所以名额只在连接断开时才归还。
    • 对保持在线的客户端累计推送 64 条大帧后,下一条大帧的 PublishDown 会一直阻塞,因为 loopCtx 只在停机时才取消。
    • 被阻塞的是 messageLoops 这个全局循环:到点分发、推送、清理、指标采样全部卡住。
  • 证据:broker/hooks.go:256-267;broker/broker.go:223-243;mochi server.go:1161-1170;cmd/nixmsg/serve.go:322-347。
  • 文档依据:DEVELOPMENT 7.5(672)「客户端回 PUBACK(mochi 的 OnQosComplete)、确认超时或连接断开时释放」;DEVIATIONS N1/N2 第 4 条。
  • 为何不是故意设计:N1/N2 第 4 条说的正是在 OnQosComplete 释放,只是拿到的包里没有 Payload,判断不成立。
  • 解决方案:在 OnQosPublish(包进入 inflight 时回调,带 PacketID 和 Payload)里记下大帧的 PacketID,OnQosComplete 和 OnQosDropped 按 PacketID 释放;获取名额改为非阻塞,拿不到就返回错误,让消息线稍后重推。
  • 改动文件:broker/hooks.go、broker/broker.go。
  • 与其他模块的交互/冲突风险:消息线自己也有一份大帧信号量(M2 第 2 条),确认后会释放;改为非阻塞后,消息线按现有「发布失败 1 秒后重推」处理,不需要改。
  • 需补测试:对同一个在线客户端推 65 条 100 KiB 的消息且全部确认,第 65 条能在 1 秒内推出。
  • 置信度:代码阅读确定,没有写测试复现。

[P-01] 大帧名额收到 PUBACK 也不释放:累计 64 条大于 64 KiB 的下行后,推送和调度主循环永久卡死

  • 严重级:critical
  • 分类:并发 / 逻辑
  • 现象与影响:
    • PublishDown 对大于 64 KiB 的帧先占一个 largeSem 名额(全局 64 个)。QoS 1 的帧只在 OnQosComplete 或断线时归还。
    • mochi 调 OnQosComplete 时传的是客户端发来的 PUBACK 包,里面没有 Payload。钩子里 len(pk.Payload) <= largeFrameBytes 永远成立,直接返回,名额永不归还。
    • 只要收件端一直在线,累计 64 条大消息后(F07 允许到 256 KiB),下一条大帧就卡在 b.largeSem <- struct{}{}。卡住时各路径的表现:
      • messageLoops:用的是永不取消的 loopCtx,整条循环永久停住。DispatchDue、所有连接的 PushPending、CleanupOnce、指标采样全部停。
      • WakePush 的 goroutine:卡 30 秒后失败。随后 clearPushed 用的是同一个已过期的 ctx,写库直接失败,投递标记清不掉;不保留的消息 5 分钟后按 not_acked 被丢弃。
      • 上行 worker:发大 resp 时用 context.Background(),该端的上行队列从此不再消费。
    • 另有几条泄漏路径:
      • server.Publish 遇到"没有订阅者"、"发送队列满被丢"、"inflight 已满"都返回 nil,名额挂在连接状态上,直到断线才还。
      • releaseAllLarge 执行后,PublishDown 才做 largeHeld++,这个竞态会永久丢一个名额。
  • 证据:
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。

  • 文档依据:DEVELOPMENT 7.5「拿到名额才发布,客户端回 PUBACK(mochi 的 OnQosComplete)、确认超时或连接断开时释放」;PRD F07、F11。
  • 为何不是故意设计:DEVIATIONS N2 第 4 条写的是 QoS 1「在 OnQosComplete 且 payload>64KiB 时释放」,只接受"客户端永远不回 PUBACK"时名额被占满;实际是回了 PUBACK 也不释放。M2 第 2 条说双重限流"更严,不破坏语义",前提同样不成立。
  • 解决方案:
    1. 推荐方案(改动小):在 broker.PublishDown 里,server.Publish 返回后不论 QoS 都立即 releaseOneLarge,broker 的名额只保护"发布这一步"。同时删掉 largeHeld、releaseAllLarge 和 OnQosComplete 里的名额逻辑。在途数量的上限交给 message 包已有的 largeSem。
    2. 获取名额改成有界等待,例如 select 同时等 ctx.Done() 和 time.After(5*time.Second),超时返回错误。调用方按发布失败处理:清标记,1 秒后重推。
    3. 备选(精确方案):实现 OnQosPublish、OnQosComplete、OnQosDropped,按"客户端 + PacketID"记录大帧并据此释放。OnPublishDropped 和"没有订阅者"的情况在 PublishDown 里同步判定并释放;为了能对上号,同一连接的大帧发布要串行。做完之后 message 包可以去掉自己的名额。
  • 改动文件:internal/broker/broker.go、hooks.go;需要 M 线配合改 internal/app/message/push.go、session.go。
  • 与其他模块的交互/冲突风险:推荐方案依赖 message 包的名额能正确释放,但 M 的 releaseLarge 在断线清标记(message/session.go 的 OnDisconnect)以及撤回、作废、清理结束时都没有调用。断线前没确认的大消息每次都漏一个名额,64 次之后 acquireLarge 永远失败,大消息只清标记、不再重推。所以必须和 M 线一起修,否则卡死点只是从 broker 挪到了 message。
  • 需补测试:
    • broker 单测:用 net.Pipe 客户端自动回 PUBACK,连续发 100 条 70 KiB 的 QoS 1(ctx 2 秒超时),全部成功,结束后 len(b.largeSem)==0。
    • broker 单测:没有订阅者时发布,名额随即归零。
    • 集成测:两端在线互发 80 条 100 KiB 消息,之后再发一条定时消息,仍能按时送达。
  • 置信度:代码阅读确定(已对照 mochi 源码)

复审基线:main 4059a15(2026-09-30)。编号说明、各工作线的合并顺序、共享文件归属见总览 #7。

**编号**:B-02 **严重级**:critical **工作线**:broker(internal/broker) **来源**:消息核心审查 M-01,总审查人已对照 mochi 源码复核 **依赖**:B-01 (#8)(同文件,先合) **被依赖**:C-01 (#32)(消息线删除自己那份名额,即审查 M-02)、B-03 (#10) ### 现象与影响 - 对超过 64 KiB 的 QoS 1 下行,`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 包本身(mochi `server.go:1161-1172`),没有载荷,判断永远不成立。名额只在 Publish 出错、QoS 0 或连接断开(`releaseAllLarge`)时才还。 - 结果:累计向在线端发出 64 条大帧后,之后所有大帧的 `PublishDown` 都阻塞。`messageLoops` 传入的 ctx 没有超时(`cmd/nixmsg/serve.go:334-336`),全站的推送、到点分发、清理、确认超时处理一起冻结;握手路径在端的上行 worker 里用 `context.Background()` 推送,这个编号之后的 ack、send 乃至重连后的 hello 全部排在后面。 - 一条 100 KiB 的消息发给 64 人以上的在线群,就能一次占满名额。PRD F07 允许 256 KiB 正文。 ### 文档依据 DEVELOPMENT 7.5「大帧同时在发不超过 64 个,客户端回 PUBACK(mochi 的 OnQosComplete)、确认超时或连接断开时释放」;DEVIATIONS N1/N2 第 4 条"QoS1 在 OnQosComplete 时释放"——意图明确,实现拿不到载荷。 ### 解决方案 1. 在 `OnQosPublish(cl, out, …)` 钩子里(此时出站包带 PacketID 和载荷)把超过 64 KiB 的包 ID 记入 `connState.largePIDs`。`OnQosComplete`、`OnQosDropped` 按 PacketID 查表,命中才归还。在 `Provides` 里补上这两个钩子。断线时仍全部归还。 2. 一次 Publish 占了名额却没有产生 inflight(没有订阅者、inflight 已满被静默丢、连接已关),要立即归还:比较 Publish 前后该连接大帧 `OnQosPublish` 的次数,没有增加就归还。 3. 若 B-03(每连接异步下发)先落地,名额的获取移到每连接的发送 goroutine 里,归还规则不变。 4. 更新 DEVIATIONS N1/N2 第 4 条。 ### 方案取舍(三份审查的建议不一致,统一如下) - **采用**:按 PacketID 在 PUBACK / 丢弃 / 断线时归还(审查 M-01、I-17 的方案),与 DEVELOPMENT 7.5 原意一致,只在 broker 保留一份名额,消息包那份由消息线在 C-01 删除。 - **不采用**:审查 P-01 的"发布完成即归还、在途上限交给消息包的名额"。消息包那份名额本身在断线、撤回、作废时都会泄漏(审查 M-02,附在 C-01),而且在途时间会被拉长到应用层确认(最长 5 分钟),会把卡死点从 broker 挪到消息包。 - **吸收 P-1 的两点**:获取名额必须有界等待,超时按发布失败处理(消息线据此清标记、1 秒后重推);`releaseAllLarge` 与 `largeHeld++` 之间的竞态随 PacketID 表一起消除。 ### 改动文件 `internal/broker/hooks.go`、`internal/broker/broker.go`,broker 测试。 ### 交互 / 冲突说明 - 消息线 C-01 会删除 message 包里重复的一份大帧名额,依赖本 issue 先合入。 - 与 B-01、B-03 同文件,按 broker 线顺序合入。 ### 验收与测试 - broker 单测:同一在线客户端连续收 65 个 70 KiB 帧并自动回 PUBACK,第 65 次 `PublishDown` 1 秒内返回。 - 不回 PUBACK 的客户端断开后名额全部归还;没有订阅者时发布后名额立即归还。 - 集成:给 100 人在线群发 100 KiB 消息后,再发一条定时消息,仍按时分发。 --- ### 问题明细(各区审查原文,证据含文件与行号) > 以下是本次复审各区审查报告的原文段落。A、M、I、P、S 开头的是原始发现编号(A 管理后台与网页、M 消息核心、I 身份认证群在线、P 传输平台部署、S SDK)。**解决方案以本 issue 上方的"结论与统一方案"为准**;原文里的方案与之不一致时,按上方执行。 #### [M-01] 服务端大帧名额收到 PUBACK 时从不释放,64 条后推送阻塞、每秒循环卡死(跨模块:连接线 broker) - 严重级:critical - 分类:并发 / 逻辑 - 现象与影响: - 对超过 64 KiB 的 QoS1 下行,`PublishDown` 会阻塞等待全局 64 个名额之一,本应在客户端回 PUBACK 时释放。 - 但 mochi 传给 `OnQosComplete` 的是客户端发来的 PUBACK 包本身,没有 payload,钩子里的大小判断永远直接返回。名额只在该连接断开时才还。 - 累计 64 条大帧发给仍在线的端之后,后续大帧的 `PublishDown` 全部阻塞: - `messageLoops` 传的 ctx 没有超时,一旦阻塞就不再到点分发、不再做过期清理、不再处理确认超时。 - 握手路径在该端的上行 worker 里用 `context.Background()` 推送,卡住后这个编号之后的 ack、send,乃至重连后的 hello 都排在它后面,设备无法再握手。 - `WakePush` 协程 30 秒超时后,`clearPushed` 复用同一个已过期的 ctx,清不掉「已推送」标记,形成幽灵在途(见 M-15)。 - 一条 100 KiB 的消息发给 64 人以上的在线群,就能一次占满名额。 - 证据: ```256:267:e:\code\NixMsg\internal\broker\hooks.go 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) } ``` - mochi `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。 - 文档依据:DEVELOPMENT 7.5「大帧同时在发不超过 64 个,客户端回 PUBACK(mochi 的 OnQosComplete)、确认超时或连接断开时释放」。 - 为何不是故意设计:DEVIATIONS 连接线 N1/N2 第 4 条写的就是「QoS1 在 OnQosComplete 时释放」。意图明确,只是实现拿不到 payload,结果从不释放。 - 解决方案: 1. broker 在 `OnQosPublish(cl, out, …)` 里(此时出站包带 PacketID 和 payload)把超过 64 KiB 的包 ID 记到 `connState.largePIDs`。`OnQosComplete`、`OnQosDropped` 按 PUBACK 的 PacketID 查表,命中就释放。`Provides` 里补上这两个钩子。断线时仍然全部释放。 2. 如果一次 Publish 占了名额却没有产生 inflight(没有订阅者、inflight 已满被静默丢、连接已关),要立即归还。可以在 `connState` 上用计数器比较 Publish 前后大帧 `OnQosPublish` 的次数,没增加就释放。 3. `PushPending` 调 `PublishDown` 时包一层 `context.WithTimeout(ctx, 5s)`;失败后的 `clearPushed` 另起一个短超时的新 ctx。 4. 更新 DEVIATIONS N 第 4 条。 - 改动文件:`internal/broker/hooks.go`、`broker.go`(连接线目录);`internal/app/message/push.go`;`cmd/nixmsg/serve.go`。 - 交互/冲突风险:和 M-02 是同一类问题,建议一起修。broker 属于连接线目录,要按 TASKS 4.2 走共享改动流程。需要确认会话继承时 mochi 重发 inflight 不会让同一个包 ID 被记两次。 - 需补测试: - broker 单测:同一在线客户端连续收 65 个 70 KiB 的帧并自动回 PUBACK,第 65 个 `PublishDown` 1 秒内返回。 - broker 单测:不回 PUBACK 的客户端断开后名额全部归还;没有订阅者时发布后名额立即归还。 - 集成:给 100 人在线群发 100 KiB 消息后,再发一条定时消息,仍能按时分发。 - 置信度:代码阅读确定(已对照 mochi 源码)。 #### [I-17] broker 的大帧名额在收到 PUBACK 时不释放,累计 64 条大帧后全局循环永久阻塞 - **严重级**:critical - **分类**:并发 / 质量 - **现象与影响** - mochi 的 `processPuback` 传给 OnQosComplete 的是客户端发来的 PUBACK,不带 Payload。hook 里「只处理超过 64 KiB 的包」这个判断永远成立,直接返回,所以名额只在连接断开时才归还。 - 对保持在线的客户端累计推送 64 条大帧后,下一条大帧的 PublishDown 会一直阻塞,因为 loopCtx 只在停机时才取消。 - 被阻塞的是 messageLoops 这个全局循环:到点分发、推送、清理、指标采样全部卡住。 - **证据**:`broker/hooks.go:256-267`;`broker/broker.go:223-243`;mochi `server.go:1161-1170`;`cmd/nixmsg/serve.go:322-347`。 - **文档依据**:DEVELOPMENT 7.5(672)「客户端回 PUBACK(mochi 的 OnQosComplete)、确认超时或连接断开时释放」;DEVIATIONS N1/N2 第 4 条。 - **为何不是故意设计**:N1/N2 第 4 条说的正是在 OnQosComplete 释放,只是拿到的包里没有 Payload,判断不成立。 - **解决方案**:在 `OnQosPublish`(包进入 inflight 时回调,带 PacketID 和 Payload)里记下大帧的 PacketID,`OnQosComplete` 和 `OnQosDropped` 按 PacketID 释放;获取名额改为非阻塞,拿不到就返回错误,让消息线稍后重推。 - **改动文件**:`broker/hooks.go`、`broker/broker.go`。 - **与其他模块的交互/冲突风险**:消息线自己也有一份大帧信号量(M2 第 2 条),确认后会释放;改为非阻塞后,消息线按现有「发布失败 1 秒后重推」处理,不需要改。 - **需补测试**:对同一个在线客户端推 65 条 100 KiB 的消息且全部确认,第 65 条能在 1 秒内推出。 - **置信度**:代码阅读确定,没有写测试复现。 #### [P-01] 大帧名额收到 PUBACK 也不释放:累计 64 条大于 64 KiB 的下行后,推送和调度主循环永久卡死 - **严重级**:critical - **分类**:并发 / 逻辑 - **现象与影响**: - `PublishDown` 对大于 64 KiB 的帧先占一个 `largeSem` 名额(全局 64 个)。QoS 1 的帧只在 `OnQosComplete` 或断线时归还。 - mochi 调 `OnQosComplete` 时传的是客户端发来的 PUBACK 包,里面没有 Payload。钩子里 `len(pk.Payload) <= largeFrameBytes` 永远成立,直接返回,名额永不归还。 - 只要收件端一直在线,累计 64 条大消息后(F07 允许到 256 KiB),下一条大帧就卡在 `b.largeSem <- struct{}{}`。卡住时各路径的表现: - **messageLoops**:用的是永不取消的 `loopCtx`,整条循环永久停住。DispatchDue、所有连接的 PushPending、CleanupOnce、指标采样全部停。 - **WakePush 的 goroutine**:卡 30 秒后失败。随后 `clearPushed` 用的是同一个已过期的 ctx,写库直接失败,投递标记清不掉;不保留的消息 5 分钟后按 not_acked 被丢弃。 - **上行 worker**:发大 resp 时用 `context.Background()`,该端的上行队列从此不再消费。 - 另有几条泄漏路径: - `server.Publish` 遇到"没有订阅者"、"发送队列满被丢"、"inflight 已满"都返回 nil,名额挂在连接状态上,直到断线才还。 - `releaseAllLarge` 执行后,`PublishDown` 才做 `largeHeld++`,这个竞态会永久丢一个名额。 - **证据**: ```256:267:e:\code\NixMsg\internal\broker\hooks.go func (h *nixHook) OnQosComplete(cl *mqtt.Client, pk packets.Packet) { if len(pk.Payload) <= largeFrameBytes { return } // ... h.b.releaseOneLarge(st) } ``` ```1161:1170:E:\TempData\go\pkg\mod\github.com\mochi-mqtt\server\v2@v2.7.9\server.go 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`。 - **文档依据**:DEVELOPMENT 7.5「拿到名额才发布,客户端回 PUBACK(mochi 的 OnQosComplete)、确认超时或连接断开时释放」;PRD F07、F11。 - **为何不是故意设计**:DEVIATIONS N2 第 4 条写的是 QoS 1「在 OnQosComplete 且 payload>64KiB 时释放」,只接受"客户端永远不回 PUBACK"时名额被占满;实际是回了 PUBACK 也不释放。M2 第 2 条说双重限流"更严,不破坏语义",前提同样不成立。 - **解决方案**: 1. 推荐方案(改动小):在 `broker.PublishDown` 里,`server.Publish` 返回后不论 QoS 都立即 `releaseOneLarge`,broker 的名额只保护"发布这一步"。同时删掉 `largeHeld`、`releaseAllLarge` 和 `OnQosComplete` 里的名额逻辑。在途数量的上限交给 message 包已有的 `largeSem`。 2. 获取名额改成有界等待,例如 `select` 同时等 `ctx.Done()` 和 `time.After(5*time.Second)`,超时返回错误。调用方按发布失败处理:清标记,1 秒后重推。 3. 备选(精确方案):实现 `OnQosPublish`、`OnQosComplete`、`OnQosDropped`,按"客户端 + PacketID"记录大帧并据此释放。`OnPublishDropped` 和"没有订阅者"的情况在 `PublishDown` 里同步判定并释放;为了能对上号,同一连接的大帧发布要串行。做完之后 message 包可以去掉自己的名额。 - **改动文件**:`internal/broker/broker.go`、`hooks.go`;需要 M 线配合改 `internal/app/message/push.go`、`session.go`。 - **与其他模块的交互/冲突风险**:推荐方案依赖 message 包的名额能正确释放,但 M 的 `releaseLarge` 在断线清标记(`message/session.go` 的 OnDisconnect)以及撤回、作废、清理结束时都没有调用。断线前没确认的大消息每次都漏一个名额,64 次之后 `acquireLarge` 永远失败,大消息只清标记、不再重推。所以必须和 M 线一起修,否则卡死点只是从 broker 挪到了 message。 - **需补测试**: - broker 单测:用 net.Pipe 客户端自动回 PUBACK,连续发 100 条 70 KiB 的 QoS 1(ctx 2 秒超时),全部成功,结束后 `len(b.largeSem)==0`。 - broker 单测:没有订阅者时发布,名额随即归零。 - 集成测:两端在线互发 80 条 100 KiB 消息,之后再发一条定时消息,仍能按时送达。 - **置信度**:代码阅读确定(已对照 mochi 源码) --- <sub>复审基线:main `4059a15`(2026-09-30)。编号说明、各工作线的合并顺序、共享文件归属见总览 #7。</sub>
nixevol added the P0-criticallane/brokerreview-2026-09-30 labels 2026-09-30 13:56:51 +08:00
Author
Owner

已合入 origin/main 0c9b459。落地提交 deb2398 fix: 大帧名额按 PacketID 在 PUBACK 与断线时归还 (#9)。

已合入 origin/main `0c9b459`。落地提交 `deb2398` fix: 大帧名额按 PacketID 在 PUBACK 与断线时归还 (#9)。
Sign in to join this conversation.