From ca3434421b19f45be1295f87e0bd8d8a6dd5b327 Mon Sep 17 00:00:00 2001 From: Nixevol Date: Wed, 30 Sep 2026 19:29:37 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E5=81=9C=E6=9C=BA=E6=8E=92=E7=A9=BA?= =?UTF-8?q?=E8=AE=A1=E5=85=A5=E6=AD=A3=E5=9C=A8=E5=8F=91=E9=80=81=E7=9A=84?= =?UTF-8?q?=E5=B8=A7=E5=B9=B6=E4=B8=BA=E6=96=AD=E5=BC=80=E9=A2=84=E7=95=99?= =?UTF-8?q?=E6=97=B6=E9=97=B4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/DEVIATIONS.md | 3 ++- internal/broker/b04_b06_b08_test.go | 21 +++++++++++++++++++++ internal/broker/broker.go | 22 ++++++++++++++++------ internal/broker/downlink.go | 2 ++ 4 files changed, 41 insertions(+), 7 deletions(-) diff --git a/docs/DEVIATIONS.md b/docs/DEVIATIONS.md index a6fb37b..f77cb0a 100644 --- a/docs/DEVIATIONS.md +++ b/docs/DEVIATIONS.md @@ -1796,4 +1796,5 @@ issue #3 未关闭,`feat/fix-3-downlink-deadlock` 未合入 `main`。下面是 - 实际做法:带断开的下行帧用本帧 `OnPacketSent` 完成信号(优先 packet id,否则按载荷匹配),不再用连接级 `sentPub` 总数。`Shutdown` 在 ctx 未取消时先等下行队列与 `wirePending` 排空,再 `DisconnectClient` 发 `0x8B` 并在截止前等连接拆掉;ctx 已取消则发完即 `Close`。`serve` 仍给 5 秒预算并记录非超时错误。未合 `feat/fix-3-downlink-deadlock`。 - 原因:前面 PUBLISH 的 `OnPacketSent` 会让总数等待提前返回;`Shutdown` 对 ctx 非阻塞 select 使 5 秒预算用不上,有 outbound 积压时 `0x8B` 只进 outbuf 随 `Stop` 丢掉。 - 备选方案:恢复固定 `Sleep`(否决);改 `PublishDown` 签名(否决)。 - - 影响:队列/outbound 有积压时 fatal、logout 先到客户端再断开;停机在预算内尽量发出 `0x8B`,超时返回 ctx 错误而非空等。 + - 影响:队列/outbound 有积压时 fatal、logout 先到客户端再断开;停机在预算内尽量发出 `0x8B`,超时返回 ctx 错误而非空等。 + - 补强:`sendOne` 从取出帧到返回前记 `inSend`,排空等待把它算上(含等大帧名额、尚未 `wirePending++` 的窗口)。有截止时间时排空最多用到截止前 1 秒,剩下的时间留给 `0x8B` 写出;排空没完成也不会因此跳过这段等待。 diff --git a/internal/broker/b04_b06_b08_test.go b/internal/broker/b04_b06_b08_test.go index 27e1424..1ce71a4 100644 --- a/internal/broker/b04_b06_b08_test.go +++ b/internal/broker/b04_b06_b08_test.go @@ -274,6 +274,27 @@ func TestShutdownCancelledContextReturnsQuickly(t *testing.T) { } } +func TestWaitConnsQuietSeesInSend(t *testing.T) { + b, err := New(Options{}) + if err != nil { + t.Fatal(err) + } + defer func() { _ = b.Close() }() + st := &connState{downCh: make(chan downItem, 1)} + st.inSend.Store(1) + ctx, cancel := context.WithTimeout(context.Background(), 40*time.Millisecond) + defer cancel() + if b.waitConnsQuiet(ctx, []*connState{st}) { + t.Fatal("inSend should keep shutdown from treating the conn as quiet") + } + st.inSend.Store(0) + ctx2, cancel2 := context.WithTimeout(context.Background(), 200*time.Millisecond) + defer cancel2() + if !b.waitConnsQuiet(ctx2, []*connState{st}) { + t.Fatal("quiet when inSend is 0 and queues are empty") + } +} + func TestEffectivePayloadLimitSubtractsOverhead(t *testing.T) { got := EffectivePayloadLimit(200, 0) if got != 200-packetOverheadBudget { diff --git a/internal/broker/broker.go b/internal/broker/broker.go index b3fa1c5..01d2256 100644 --- a/internal/broker/broker.go +++ b/internal/broker/broker.go @@ -142,6 +142,7 @@ type connState struct { downDone chan struct{} downBytes atomic.Int64 wirePending atomic.Int64 // Publish 入 mochi outbound 后、OnPacketSent 前 + inSend atomic.Int32 // downLoop 已取出帧、尚未从 sendOne 返回 mu sync.Mutex // 带断开的下行帧:只等本帧 OnPacketSent,不用连接级计数。 @@ -237,8 +238,8 @@ func (b *Broker) Close() error { } // Shutdown 向所有连接发 MQTT 5 0x8B 后关闭。完整 HTTP 停机顺序见 L-03。 -// ctx 未取消时先等下行队列与 wirePending 排空,再 DisconnectClient(此时 outbound 空, -// 0x8B 直写套接字),并在截止前等连接拆掉;ctx 已取消则发完即 Close,不等待。 +// ctx 未取消时先等下行队列、正在 sendOne 的帧与 wirePending 排空(若有截止时间则预留约 1 秒), +// 再 DisconnectClient,并在截止前等连接拆掉;ctx 已取消则发完即 Close,不等待。 func (b *Broker) Shutdown(ctx context.Context) error { if b.closed.Load() { return nil @@ -275,17 +276,26 @@ func (b *Broker) Shutdown(ctx context.Context) error { var waitErr error if !alreadyCancelled { - if !b.waitConnsQuiet(ctx, states) { + // 留出约 1 秒给 0x8B 写出,避免排空把整个截止时间用完后立刻 Close。 + quietCtx := ctx + cancelQuiet := func() {} + if dl, ok := ctx.Deadline(); ok && time.Until(dl) > time.Second { + quietCtx, cancelQuiet = context.WithDeadline(ctx, dl.Add(-time.Second)) + } + if !b.waitConnsQuiet(quietCtx, states) && ctx.Err() != nil { waitErr = ctx.Err() } + cancelQuiet() } for _, cl := range clients { _ = b.server.DisconnectClient(cl, packets.ErrServerShuttingDown) } - if !alreadyCancelled && waitErr == nil { - waitErr = b.waitConnsGone(ctx) + if !alreadyCancelled && ctx.Err() == nil { + if err := b.waitConnsGone(ctx); err != nil { + waitErr = err + } } closeErr := b.Close() @@ -299,7 +309,7 @@ func (b *Broker) waitConnsQuiet(ctx context.Context, states []*connState) bool { for { quiet := true for _, st := range states { - if st.wirePending.Load() > 0 { + if st.inSend.Load() > 0 || st.wirePending.Load() > 0 { quiet = false break } diff --git a/internal/broker/downlink.go b/internal/broker/downlink.go index edb3ee8..3e24cf6 100644 --- a/internal/broker/downlink.go +++ b/internal/broker/downlink.go @@ -118,6 +118,8 @@ func (st *connState) sendOne(b *Broker, item downItem) { st.signalSent(item) return } + st.inSend.Add(1) + defer st.inSend.Add(-1) large := len(item.payload) > largeFrameBytes if large { if err := b.acquireLarge(context.Background()); err != nil {