[L-03][high] 有 MQTT 连接时停机退不出,还可能向已关闭通道发送而 panic #22

Closed
opened 2026-09-30 13:56:55 +08:00 by nixevol · 2 comments
Owner

编号:L-03 严重级:high 工作线:监听与 HTTP(internal/listener、internal/httpx、serve 的 HTTP 装配) 来源:审查 P-06
依赖:B-08 (#15)(Broker.Shutdown)、D-02 (#28)(写队列关闭安全)、C-01 (#32)(可等待退出的消息循环) 被依赖:无

结论与统一方案

采用审查 P-06 的方案,按目录拆给三条线、由本 issue 负责总装:

  1. broker:Broker.Shutdown(ctx),在 B-08 实现。
  2. store:写队列关闭安全,在 D-02 实现(Do 不向已关闭通道发送,关闭后返回 ErrQueueClosed)。
  3. 本 issue:
    • listener 拆成 StopAccept()(关 TCP 监听、对 HTTP 调 Shutdown)和带超时的 Wait(ctx);记录交给 OnMQTT 的连接,超时后强制关闭。
    • serve 在 <-ctx.Done() 之后立即调用 stop(),让第二次 Ctrl+C / SIGTERM 能直接结束进程。
    • 停机顺序:StopAccept → 取消并等待消息循环(C-01 提供可等待的循环)→ Drain(10 秒)→ brk.Shutdown(5 秒)→ 用剩余时间再 Drain 一次 → db.Close()。
    • compose 加 stop_grace_period: 30s;OPS 写明 systemd TimeoutStopSec 至少 30 秒。

断开时还没回 resp 的请求,SDK 会按原消息号重新提交,防重保证不重复。SDK 需要把 DISCONNECT 0x8B 当作可重试,由 SDK 线核对。

改动文件

cmd/nixmsg/serve.go(runServe 末尾的停机段)、internal/listener/server.go、deploy/docker-compose.yml、docs/OPS.md。

与其他问题的交互 / 冲突说明

  • 依赖 B-08、D-02、C-01,最后合入。
  • serve.go 只改停机段:messageLoops 由 C-01 改,staticFileHandler 由 L-05 改,listener 装配段由 L-07 改,admin.Deps 由 L-04 与 H-02 改。

验收与测试

  • 保持一个裸 TCP 和一个 WS 已登录客户端不关,cancel ctx,runServe 在 15 秒内返回、没有 panic、客户端收到 0x8B。
  • Shutdown 期间并发发上行不 panic;并发 Do 与 Close 1000 次不 panic。

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

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

[P-06] 有 MQTT 连接时停机退不出,broker.Close 还可能 panic

  • 严重级:high
  • 分类:并发 / 逻辑
  • 现象与影响:
    1. listener 卡住:lnSrv.Close() 最后会 wg.Wait(),而裸 TCP/TLS 的 MQTT 连接一直阻塞在 OnMQTT → AttachTCP 里,这些 goroutine 都在 wg 中,没有人去关它们(server.go:245-249、270-275)。只要有一台裸 TCP 设备在线,停机就卡在这一步,连 Drain 都不会执行。
    2. broker 卡住并可能 panic:即使走到了 brk.Close():
      • 它先关闭所有上行队列的通道,再调 mochi.Server.Close()。
      • mochi 只断开注册过的监听器上的客户端,然后 ClientsWg.Wait();我们的连接都是 EstablishConnection("tcp"/"ws") 接进来的,不属于任何注册监听器,同样没人断开,于是继续卡住。
      • 这期间任何 WS 客户端再发一条上行,就会向已关闭的通道发送,进程 panic。
    3. 关库时的竞态:db.Close() 会关闭写队列的通道,而 WakePush、重推计时器和没被等待退出的 messageLoops 可能还在 Queue.Do 里向它发送,也可能 panic(store/queue.go:78-90、296-303)。
    4. 信号被吞:signal.NotifyContext 的 stop 要等 cmdServe 返回才调用,卡住期间再按 Ctrl+C 或再发 SIGTERM 都会被吞掉,只能靠 SIGKILL(Docker 默认 10 秒后,systemd 默认 90 秒后)。SIGKILL 会留下 -wal/-shm 文件,和 P-19 叠加有损坏风险。
  • 证据:
	<-ctx.Done()
	loopCancel()
	_ = lnSrv.Close()
	drainCtx, drainCancel := context.WithTimeout(context.Background(), 10*time.Second)
	defer drainCancel()
	if drainErr := db.Queue.Drain(drainCtx); drainErr != nil && !errors.Is(drainErr, context.DeadlineExceeded) {
		slog.Error("write queue drain", "err", drainErr)
	}
	return nil
func (b *Broker) Close() error {
	if b.closed.Swap(true) {
		return nil
	}
	b.queuesMu.Lock()
	for _, q := range b.queues {
		q.close()
	}
	b.queuesMu.Unlock()
	return b.server.Close()
}
func (s *Server) Close() error {
	close(s.done)
	s.Log.Info("gracefully stopping server")
	s.Listeners.CloseAll(s.closeListenerClients)

其余位置:mochi listeners/listeners.go:121-135(只关注册过的监听器,然后 ClientsWg.Wait());queue.go:33-39(向可能已关闭的通道发送)。

  • 文档依据:DEVELOPMENT 7.8「收到停止信号后先停止接受新连接,等写队列里已排队的操作提交完(最多 10 秒),再断开所有连接并退出」。
  • 为何不是故意设计:DEVIATIONS P1 第 4 条写明"有 MQTT 长连接后需 N/P 联调停机路径",是未完成项。现有集成测试都在 cancel 之前先关掉了 MQTT 客户端,所以没暴露出来。
  • 解决方案:
    1. broker 新增 Shutdown(ctx):
      • 置 closed:OnConnect 直接拒绝新连接,PublishDown 返回已关闭。
      • 遍历 b.server.Clients.GetAll()(跳过内联客户端),调 DisconnectClient(cl, packets.ErrServerShuttingDown)(原因码 0x8B)。
      • 在 goroutine 里调 server.Close(),用 ctx 限时等待。
      • 最后停止上行 worker。不要关闭数据通道,改用 done 通道加 WaitGroup,发送时在 select 里同时看 done,从根上消除向已关闭通道发送的 panic。
    2. listener 拆成 StopAccept()(关闭 TCP 监听、对 HTTP 调 Shutdown)和带超时的 Wait(ctx);记录交给 OnMQTT 的连接,超时后强制关闭。
    3. serve:
      • <-ctx.Done() 之后立刻调 stop(),让第二次信号能直接结束进程。
      • 停机顺序:StopAccept → 取消并等待 messageLoops → Drain(10 秒)→ brk.Shutdown(5 秒)→ 用剩余时间再 Drain 一次 → db.Close()。
    4. store.Queue:用读写锁(或"不关闭 ch、只关闭 stop 通道")保证 Do 不会向已关闭的通道发送;关闭之后 Do 返回 ErrQueueClosed。
    5. compose 加 stop_grace_period: 30s;OPS 写明 systemd 的 TimeoutStopSec 至少 30 秒。
  • 改动文件:cmd/nixmsg/serve.go、internal/broker/broker.go、queue.go、internal/listener/server.go、internal/store/queue.go、deploy/docker-compose.yml、docs/OPS.md
  • 与其他模块的交互/冲突风险:断开时还没回 resp 的请求,SDK 会按原消息号重新提交,防重机制保证不会重复。SDK 需要把 0x8B 当作可重试(见待核实第 7 条)。
  • 需补测试:
    • serve 测试:保持一个裸 TCP 和一个 WS 已登录客户端不关,cancel ctx,确认 runServe 在 15 秒内返回、没有 panic、客户端收到 0x8B。
    • broker 测试:Shutdown 期间并发发上行,不 panic。
    • store 测试:并发调用 Do 和 Close 1000 次。
  • 置信度:代码阅读确定

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

**编号**:L-03 **严重级**:high **工作线**:监听与 HTTP(internal/listener、internal/httpx、serve 的 HTTP 装配) **来源**:审查 P-06 **依赖**:B-08 (#15)(Broker.Shutdown)、D-02 (#28)(写队列关闭安全)、C-01 (#32)(可等待退出的消息循环) **被依赖**:无 ### 结论与统一方案 采用审查 P-06 的方案,按目录拆给三条线、由本 issue 负责总装: 1. broker:`Broker.Shutdown(ctx)`,在 B-08 实现。 2. store:写队列关闭安全,在 D-02 实现(`Do` 不向已关闭通道发送,关闭后返回 `ErrQueueClosed`)。 3. 本 issue: - listener 拆成 `StopAccept()`(关 TCP 监听、对 HTTP 调 Shutdown)和带超时的 `Wait(ctx)`;记录交给 `OnMQTT` 的连接,超时后强制关闭。 - serve 在 `<-ctx.Done()` 之后立即调用 `stop()`,让第二次 Ctrl+C / SIGTERM 能直接结束进程。 - 停机顺序:StopAccept → 取消并等待消息循环(C-01 提供可等待的循环)→ Drain(10 秒)→ `brk.Shutdown`(5 秒)→ 用剩余时间再 Drain 一次 → `db.Close()`。 - compose 加 `stop_grace_period: 30s`;OPS 写明 systemd `TimeoutStopSec` 至少 30 秒。 断开时还没回 resp 的请求,SDK 会按原消息号重新提交,防重保证不重复。SDK 需要把 DISCONNECT 0x8B 当作可重试,由 SDK 线核对。 ### 改动文件 `cmd/nixmsg/serve.go`(`runServe` 末尾的停机段)、`internal/listener/server.go`、`deploy/docker-compose.yml`、`docs/OPS.md`。 ### 与其他问题的交互 / 冲突说明 - 依赖 B-08、D-02、C-01,最后合入。 - serve.go 只改停机段:`messageLoops` 由 C-01 改,`staticFileHandler` 由 L-05 改,listener 装配段由 L-07 改,`admin.Deps` 由 L-04 与 H-02 改。 ### 验收与测试 - 保持一个裸 TCP 和一个 WS 已登录客户端不关,cancel ctx,`runServe` 在 15 秒内返回、没有 panic、客户端收到 0x8B。 - Shutdown 期间并发发上行不 panic;并发 `Do` 与 `Close` 1000 次不 panic。 --- ### 问题明细(各区审查原文,证据含文件与行号) > 以下是本次复审各区审查报告的原文段落。A、M、I、P、S 开头的是原始发现编号(A 管理后台与网页、M 消息核心、I 身份认证群在线、P 传输平台部署、S SDK)。**解决方案以本 issue 上方的"结论与统一方案"为准**;原文里的方案与之不一致时,按上方执行。 #### [P-06] 有 MQTT 连接时停机退不出,broker.Close 还可能 panic - **严重级**:high - **分类**:并发 / 逻辑 - **现象与影响**: 1. **listener 卡住**:`lnSrv.Close()` 最后会 `wg.Wait()`,而裸 TCP/TLS 的 MQTT 连接一直阻塞在 `OnMQTT → AttachTCP` 里,这些 goroutine 都在 wg 中,没有人去关它们(`server.go:245-249`、`270-275`)。只要有一台裸 TCP 设备在线,停机就卡在这一步,连 Drain 都不会执行。 2. **broker 卡住并可能 panic**:即使走到了 `brk.Close()`: - 它先关闭所有上行队列的通道,再调 `mochi.Server.Close()`。 - mochi 只断开注册过的监听器上的客户端,然后 `ClientsWg.Wait()`;我们的连接都是 `EstablishConnection("tcp"/"ws")` 接进来的,不属于任何注册监听器,同样没人断开,于是继续卡住。 - 这期间任何 WS 客户端再发一条上行,就会向已关闭的通道发送,进程 panic。 3. **关库时的竞态**:`db.Close()` 会关闭写队列的通道,而 WakePush、重推计时器和没被等待退出的 messageLoops 可能还在 `Queue.Do` 里向它发送,也可能 panic(`store/queue.go:78-90`、`296-303`)。 4. **信号被吞**:`signal.NotifyContext` 的 stop 要等 cmdServe 返回才调用,卡住期间再按 Ctrl+C 或再发 SIGTERM 都会被吞掉,只能靠 SIGKILL(Docker 默认 10 秒后,systemd 默认 90 秒后)。SIGKILL 会留下 `-wal`/`-shm` 文件,和 P-19 叠加有损坏风险。 - **证据**: ```311:319:e:\code\NixMsg\cmd\nixmsg\serve.go <-ctx.Done() loopCancel() _ = lnSrv.Close() drainCtx, drainCancel := context.WithTimeout(context.Background(), 10*time.Second) defer drainCancel() if drainErr := db.Queue.Drain(drainCtx); drainErr != nil && !errors.Is(drainErr, context.DeadlineExceeded) { slog.Error("write queue drain", "err", drainErr) } return nil ``` ```179:189:e:\code\NixMsg\internal\broker\broker.go func (b *Broker) Close() error { if b.closed.Swap(true) { return nil } b.queuesMu.Lock() for _, q := range b.queues { q.close() } b.queuesMu.Unlock() return b.server.Close() } ``` ```1496:1500:E:\TempData\go\pkg\mod\github.com\mochi-mqtt\server\v2@v2.7.9\server.go func (s *Server) Close() error { close(s.done) s.Log.Info("gracefully stopping server") s.Listeners.CloseAll(s.closeListenerClients) ``` 其余位置:mochi `listeners/listeners.go:121-135`(只关注册过的监听器,然后 `ClientsWg.Wait()`);`queue.go:33-39`(向可能已关闭的通道发送)。 - **文档依据**:DEVELOPMENT 7.8「收到停止信号后先停止接受新连接,等写队列里已排队的操作提交完(最多 10 秒),再断开所有连接并退出」。 - **为何不是故意设计**:DEVIATIONS P1 第 4 条写明"有 MQTT 长连接后需 N/P 联调停机路径",是未完成项。现有集成测试都在 cancel 之前先关掉了 MQTT 客户端,所以没暴露出来。 - **解决方案**: 1. broker 新增 `Shutdown(ctx)`: - 置 closed:`OnConnect` 直接拒绝新连接,`PublishDown` 返回已关闭。 - 遍历 `b.server.Clients.GetAll()`(跳过内联客户端),调 `DisconnectClient(cl, packets.ErrServerShuttingDown)`(原因码 0x8B)。 - 在 goroutine 里调 `server.Close()`,用 ctx 限时等待。 - 最后停止上行 worker。不要关闭数据通道,改用 done 通道加 WaitGroup,发送时在 select 里同时看 done,从根上消除向已关闭通道发送的 panic。 2. listener 拆成 `StopAccept()`(关闭 TCP 监听、对 HTTP 调 Shutdown)和带超时的 `Wait(ctx)`;记录交给 OnMQTT 的连接,超时后强制关闭。 3. serve: - `<-ctx.Done()` 之后立刻调 `stop()`,让第二次信号能直接结束进程。 - 停机顺序:StopAccept → 取消并等待 messageLoops → Drain(10 秒)→ `brk.Shutdown`(5 秒)→ 用剩余时间再 Drain 一次 → `db.Close()`。 4. `store.Queue`:用读写锁(或"不关闭 ch、只关闭 stop 通道")保证 `Do` 不会向已关闭的通道发送;关闭之后 `Do` 返回 `ErrQueueClosed`。 5. compose 加 `stop_grace_period: 30s`;OPS 写明 systemd 的 `TimeoutStopSec` 至少 30 秒。 - **改动文件**:`cmd/nixmsg/serve.go`、`internal/broker/broker.go`、`queue.go`、`internal/listener/server.go`、`internal/store/queue.go`、`deploy/docker-compose.yml`、`docs/OPS.md` - **与其他模块的交互/冲突风险**:断开时还没回 resp 的请求,SDK 会按原消息号重新提交,防重机制保证不会重复。SDK 需要把 0x8B 当作可重试(见待核实第 7 条)。 - **需补测试**: - serve 测试:保持一个裸 TCP 和一个 WS 已登录客户端不关,cancel ctx,确认 runServe 在 15 秒内返回、没有 panic、客户端收到 0x8B。 - broker 测试:Shutdown 期间并发发上行,不 panic。 - store 测试:并发调用 `Do` 和 `Close` 1000 次。 - **置信度**:代码阅读确定 --- <sub>复审基线:main `4059a15`(2026-09-30)。编号说明、各工作线的合并顺序、共享文件归属见总览 #7。</sub>
nixevol added the P1-highlane/listenerreview-2026-09-30 labels 2026-09-30 13:56:55 +08:00
Author
Owner

L-03 已在 feat/fix-l03-shutdown 落地,未合入 main。提交 9a17c6a26b (#22)

合成 origin/feat/fix-message-c01-c03 + origin/feat/fix-listener-l01-l07(broker / store 已含于 C-01–C-03)。未改 PublishDown 签名。

做法:

  • listener:StopAccept() 关 TCP 监听并对 HTTP 调 Shutdown;Wait(ctx) 超时后强关仍阻塞在 OnMQTT 的连接
  • serve:<-ctx.Done() 后立刻 stop();顺序 StopAccept → 取消并等待消息循环 → Drain 10s → brk.Shutdown 5s → 再 Drain → Wait → db.Close()
  • compose stop_grace_period: 30s;OPS 写明 systemd TimeoutStopSec 至少 30 秒

验证(listen :0、临时目录):

  • TestStopAcceptDoesNotWaitForMQTT 通过
  • TestServeShutdownWithLiveMQTT 通过:已登录 WS + TCP 客户端,cancel 后 runServe 15s 内返回,DISCONNECT 0x8B
  • go test ./cmd/nixmsg ./internal/listener ./internal/broker ./internal/store ./internal/app/message -count=1 通过

task check 在合成分支上仍有来自 C-01 / L-01 / store 的既有 gofmt/govet,非本提交引入。TestQ2AcceptAndReport 在 C-01–C-03 基线上即因回执路径超时失败,非 L-03 引入。

L-03 已在 `feat/fix-l03-shutdown` 落地,未合入 main。提交 https://git.asio.asia/nixevol/NixMsg/commit/9a17c6a26bd7b09f875d664a48c5caf59cfcb09f (#22) 合成 `origin/feat/fix-message-c01-c03` + `origin/feat/fix-listener-l01-l07`(broker / store 已含于 C-01–C-03)。未改 `PublishDown` 签名。 做法: - listener:`StopAccept()` 关 TCP 监听并对 HTTP 调 `Shutdown`;`Wait(ctx)` 超时后强关仍阻塞在 `OnMQTT` 的连接 - serve:`<-ctx.Done()` 后立刻 `stop()`;顺序 StopAccept → 取消并等待消息循环 → Drain 10s → `brk.Shutdown` 5s → 再 Drain → `Wait` → `db.Close()` - compose `stop_grace_period: 30s`;OPS 写明 systemd `TimeoutStopSec` 至少 30 秒 验证(listen `:0`、临时目录): - `TestStopAcceptDoesNotWaitForMQTT` 通过 - `TestServeShutdownWithLiveMQTT` 通过:已登录 WS + TCP 客户端,cancel 后 `runServe` 15s 内返回,DISCONNECT `0x8B` - `go test ./cmd/nixmsg ./internal/listener ./internal/broker ./internal/store ./internal/app/message -count=1` 通过 `task check` 在合成分支上仍有来自 C-01 / L-01 / store 的既有 gofmt/govet,非本提交引入。`TestQ2AcceptAndReport` 在 C-01–C-03 基线上即因回执路径超时失败,非 L-03 引入。
Author
Owner

已合入 origin/main 0c9b459。落地提交 b40ef5c fix: 停机先停接受并等待循环再断开 MQTT (#22)。

已合入 origin/main `0c9b459`。落地提交 `b40ef5c` fix: 停机先停接受并等待循环再断开 MQTT (#22)。
Sign in to join this conversation.