[B-03][high] 慢客户端写阻塞时,向它发布 QoS1 会卡住调用线程,messageLoops 与其他端的处理被连带冻结 #10

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

编号: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),客户端只要还在发心跳,写就可以无限期阻塞。

NixMsg 里在共享或其他端线程上同步调用 PublishDown(QoS 1)的地方:

  • cmd/nixmsg/serve.go:334-337 的 messageLoops 串行对所有在线端 PushPending:一个前台卡死但心跳仍在的手机,就能让全站推送、到点分发、清理停止;
  • 接收方确认后给发送方推回执、撤回时给接收方推 revoked、管理接口作废后推 revoked 和 fatal(Session.fatalKick 同步发布),都会被目标端拖住;
  • PRD 允许 256 KiB 正文、推送窗口 32,单连接最多约 8 MiB 待写,移动弱网下很容易写满缓冲。

实证

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 永不阻塞调用方:

  1. connState 增加有界 FIFO(按帧数与字节数双上限,例如 256 帧 / 16 MiB)和一个发送 goroutine:在 OnSessionEstablished 启动,在 OnDisconnect 停止并丢弃剩余帧。
  2. PublishDown 只做同步检查(连接存在、大小上限、broker 未关闭),然后非阻塞入队;队列满时立即返回新的哨兵错误 ErrBackpressure,不阻塞。发送 goroutine 按顺序调用 server.Publish,同一连接内保持 FIFO 顺序。
  3. 大帧名额(B-02)的获取移到发送 goroutine 中。
  4. Session.publishJSON(resp、fatal、logout 响应)走同一通道。
  5. 在 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 测试。

交互 / 冲突说明

  • 与消息线 C-01(每连接推送 worker、拆分 messageLoops)互补:B-03 保证任何调用方都不会被慢端卡住,C-01 保证推送的窗口与顺序。两边以 PublishDown 签名不变为契约,可并行开发。审查 P-07 对慢客户端拖住全局的分析附在下方。
  • B-04 的"写出后再断开"原语依赖本 issue 的每连接发送队列。

验收与测试

  • broker 测试:一个慢客户端(不读、定时 PINGREQ)加一个正常客户端;慢客户端写阻塞期间,对正常客户端的 PublishDown 在 100 ms 内返回并送达。
  • 慢客户端队列满时 PublishDown 返回 ErrBackpressure,且调用耗时小于 10 ms。
  • 同一连接连续 1000 帧到达顺序与发布顺序一致;断开后发送 goroutine 退出(goroutine 计数回落)。
  • 集成:一个端停止读取时,其他端之间的单聊仍能在 1 秒内送达。

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

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

[P-07] messageLoops 单 goroutine 串行:一个不读数据的客户端能拖住全局调度,DispatchDue 每秒最多 100 条

  • 严重级:medium
  • 分类:设计 / 并发 / 与文档不符
  • 现象与影响:
    • serve 用一个 1 秒的 ticker 串行执行 DispatchDue(100) → 逐个连接 PushPending → CleanupOnce → 指标采样,全部用没有超时的 loopCtx。
    • 慢客户端卡住全局:
      • QoS 1 发布要经过 mochi 的 NextPacketID,它要拿 cl.Lock()。
      • 该客户端的写循环在向对端写包时一直持有同一把锁。对端不读数据(移动网络卡住、TCP 零窗口、恶意客户端)时,写入会一直阻塞到写 deadline,也就是 1.5 倍心跳,最长 900 秒。
      • 这段时间里整条循环停住:定时消息不分发,确认超时和到期清理不执行,指标冻结在看起来正常的旧值上。
      • 和 P-1 叠加就是永久卡住。
    • 吞吐:DispatchDue 每轮最多 100 条,而每轮间隔至少是 max(1 秒, 整轮耗时)。持续每秒超过 100 条定时或延迟消息时,积压会无上限增长,到点发送越来越晚。
  • 证据:
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。

  • 文档依据:DEVELOPMENT 7.4「调度循环按最早的 send_at 定时唤醒」、7.5「每个已握手连接一个推送循环」;PRD 第 8 节吞吐 200 条/秒,F09、F11。
  • 为何不是故意设计:M2 第 4 条把起循环交给接线方,L-WIRE 第 1 条只写了"周期 DispatchDue/PushPending/CleanupOnce",没有记录串行结构和它的后果。这里不讨论 #3 的定性,只说同一把锁在真实慢客户端下造成的全局停顿。
  • 解决方案:serve 拆成几个独立的 goroutine,都用 WaitGroup 管理,停机时等待退出(配合 P-6):
    • 调度:每轮循环调 DispatchDue,直到返回数小于上限或本轮用满 500 毫秒,每次调用带 5 秒超时。
    • 清理:每秒一次,单独一个 goroutine。
    • 指标:10 到 15 秒采样一次。
    • 推送:按连接做"同一端同时只跑一个"的 worker,总并发用信号量限制(例如 32),每次调用带 5 秒超时,只推已握手的连接(配合 P-2)。
  • 改动文件:cmd/nixmsg/serve.go(可以拆出 loops.go)
  • 与其他模块的交互/冲突风险:不改 message 包的接口。同一端的 PushPending 要保证不并发执行,避免重复查询。
  • 需补测试:
    • 用假的 Downlink 让端 X 阻塞 60 秒,同时确认其他端照常推送、定时消息照常分发。
    • 500 条同一秒到点的定时消息在 2 秒内全部分发。
  • 置信度:循环结构和吞吐上限为代码阅读确定;慢客户端卡住全局由 mochi 源码推出,未复现

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

**编号**: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`),客户端只要还在发心跳,写就可以无限期阻塞。 NixMsg 里在共享或其他端线程上同步调用 `PublishDown`(QoS 1)的地方: - `cmd/nixmsg/serve.go:334-337` 的 `messageLoops` 串行对所有在线端 `PushPending`:一个前台卡死但心跳仍在的手机,就能让全站推送、到点分发、清理停止; - 接收方确认后给发送方推回执、撤回时给接收方推 revoked、管理接口作废后推 revoked 和 fatal(`Session.fatalKick` 同步发布),都会被目标端拖住; - PRD 允许 256 KiB 正文、推送窗口 32,单连接最多约 8 MiB 待写,移动弱网下很容易写满缓冲。 ### 实证 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` 永不阻塞调用方: 1. `connState` 增加有界 FIFO(按帧数与字节数双上限,例如 256 帧 / 16 MiB)和一个发送 goroutine:在 `OnSessionEstablished` 启动,在 `OnDisconnect` 停止并丢弃剩余帧。 2. `PublishDown` 只做同步检查(连接存在、大小上限、broker 未关闭),然后非阻塞入队;队列满时立即返回新的哨兵错误 `ErrBackpressure`,不阻塞。发送 goroutine 按顺序调用 `server.Publish`,同一连接内保持 FIFO 顺序。 3. 大帧名额(B-02)的获取移到发送 goroutine 中。 4. `Session.publishJSON`(resp、fatal、logout 响应)走同一通道。 5. 在 `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 测试。 ### 交互 / 冲突说明 - 与消息线 C-01(每连接推送 worker、拆分 messageLoops)互补:B-03 保证任何调用方都不会被慢端卡住,C-01 保证推送的窗口与顺序。两边以 `PublishDown` 签名不变为契约,可并行开发。审查 P-07 对慢客户端拖住全局的分析附在下方。 - B-04 的"写出后再断开"原语依赖本 issue 的每连接发送队列。 ### 验收与测试 - broker 测试:一个慢客户端(不读、定时 PINGREQ)加一个正常客户端;慢客户端写阻塞期间,对正常客户端的 `PublishDown` 在 100 ms 内返回并送达。 - 慢客户端队列满时 `PublishDown` 返回 `ErrBackpressure`,且调用耗时小于 10 ms。 - 同一连接连续 1000 帧到达顺序与发布顺序一致;断开后发送 goroutine 退出(goroutine 计数回落)。 - 集成:一个端停止读取时,其他端之间的单聊仍能在 1 秒内送达。 --- ### 问题明细(各区审查原文,证据含文件与行号) > 以下是本次复审各区审查报告的原文段落。A、M、I、P、S 开头的是原始发现编号(A 管理后台与网页、M 消息核心、I 身份认证群在线、P 传输平台部署、S SDK)。**解决方案以本 issue 上方的"结论与统一方案"为准**;原文里的方案与之不一致时,按上方执行。 #### [P-07] messageLoops 单 goroutine 串行:一个不读数据的客户端能拖住全局调度,DispatchDue 每秒最多 100 条 - **严重级**:medium - **分类**:设计 / 并发 / 与文档不符 - **现象与影响**: - serve 用一个 1 秒的 ticker 串行执行 DispatchDue(100) → 逐个连接 PushPending → CleanupOnce → 指标采样,全部用没有超时的 `loopCtx`。 - **慢客户端卡住全局**: - QoS 1 发布要经过 mochi 的 `NextPacketID`,它要拿 `cl.Lock()`。 - 该客户端的写循环在向对端写包时一直持有同一把锁。对端不读数据(移动网络卡住、TCP 零窗口、恶意客户端)时,写入会一直阻塞到写 deadline,也就是 1.5 倍心跳,最长 900 秒。 - 这段时间里整条循环停住:定时消息不分发,确认超时和到期清理不执行,指标冻结在看起来正常的旧值上。 - 和 P-1 叠加就是永久卡住。 - **吞吐**:DispatchDue 每轮最多 100 条,而每轮间隔至少是 max(1 秒, 整轮耗时)。持续每秒超过 100 条定时或延迟消息时,积压会无上限增长,到点发送越来越晚。 - **证据**: ```322:338:e:\code\NixMsg\cmd\nixmsg\serve.go 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 { ``` ```276:278:E:\TempData\go\pkg\mod\github.com\mochi-mqtt\server\v2@v2.7.9\clients.go 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`。 - **文档依据**:DEVELOPMENT 7.4「调度循环按最早的 send_at 定时唤醒」、7.5「每个已握手连接一个推送循环」;PRD 第 8 节吞吐 200 条/秒,F09、F11。 - **为何不是故意设计**:M2 第 4 条把起循环交给接线方,L-WIRE 第 1 条只写了"周期 DispatchDue/PushPending/CleanupOnce",没有记录串行结构和它的后果。这里不讨论 #3 的定性,只说同一把锁在真实慢客户端下造成的全局停顿。 - **解决方案**:serve 拆成几个独立的 goroutine,都用 WaitGroup 管理,停机时等待退出(配合 P-6): - **调度**:每轮循环调 DispatchDue,直到返回数小于上限或本轮用满 500 毫秒,每次调用带 5 秒超时。 - **清理**:每秒一次,单独一个 goroutine。 - **指标**:10 到 15 秒采样一次。 - **推送**:按连接做"同一端同时只跑一个"的 worker,总并发用信号量限制(例如 32),每次调用带 5 秒超时,只推已握手的连接(配合 P-2)。 - **改动文件**:`cmd/nixmsg/serve.go`(可以拆出 `loops.go`) - **与其他模块的交互/冲突风险**:不改 message 包的接口。同一端的 PushPending 要保证不并发执行,避免重复查询。 - **需补测试**: - 用假的 Downlink 让端 X 阻塞 60 秒,同时确认其他端照常推送、定时消息照常分发。 - 500 条同一秒到点的定时消息在 2 秒内全部分发。 - **置信度**:循环结构和吞吐上限为代码阅读确定;慢客户端卡住全局由 mochi 源码推出,未复现 --- <sub>复审基线:main `4059a15`(2026-09-30)。编号说明、各工作线的合并顺序、共享文件归属见总览 #7。</sub>
nixevol added the P1-highlane/brokerreview-2026-09-30 labels 2026-09-30 13:56:51 +08:00
Author
Owner

已合入 origin/main 0c9b459。落地提交 91e887b fix: 完成 broker 复审 B-03 至 B-12 (#10)。

已合入 origin/main `0c9b459`。落地提交 `91e887b` fix: 完成 broker 复审 B-03 至 B-12 (#10)。
Sign in to join this conversation.