From 09fb544b7d522ce30c12172addf5399011466bf2 Mon Sep 17 00:00:00 2001 From: Nixevol Date: Wed, 30 Sep 2026 14:58:48 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E6=B5=8B=E8=AF=95=20WS=20=E5=AE=A2?= =?UTF-8?q?=E6=88=B7=E7=AB=AF=E6=8C=89=E5=AD=97=E8=8A=82=E6=B5=81=E6=8B=86?= =?UTF-8?q?=E5=8C=85=E5=B9=B6=E6=81=A2=E5=A4=8D=E7=BE=A4=E4=BA=8B=E4=BB=B6?= =?UTF-8?q?=E5=90=8C=E6=AD=A5=E4=B8=8B=E5=8F=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/DEVIATIONS.md | 21 +++- internal/app/group/app.go | 21 ++-- test/accept/ACCEPTANCE.md | 4 +- test/accept/rest_accept_test.go | 54 +++------- test/harness/mqtt.go | 47 ++++++-- test/harness/mqtt_test.go | 185 ++++++++++++++++++++++++++++++++ 6 files changed, 265 insertions(+), 67 deletions(-) create mode 100644 test/harness/mqtt_test.go diff --git a/docs/DEVIATIONS.md b/docs/DEVIATIONS.md index 5e9e6be..233ea8e 100644 --- a/docs/DEVIATIONS.md +++ b/docs/DEVIATIONS.md @@ -1454,21 +1454,24 @@ - 实际做法:WebSocket 升级成功与 TCP dial 成功后 `SetDeadline(time.Time{})`,避免长会话在 dial timeout 到期后读写全部失败。 - 原因:F10 等短宽限仍需跨数秒保持连接;未清 deadline 时旧 10s dial 会在会话中途使 Recv 失败,表现为 `timeout waiting resp`。 - 备选方案:每次读写刷新 deadline(更繁琐)。 - - 影响:跨线改了 harness;行为仅更正测试客户端,不改产品。 + - 影响:跨线改了 harness;行为仅更正测试客户端,不改产品。 + - **T-01 注**:清 deadline 的做法仍保留。当时把部分 `timeout waiting resp` 归因于服务端死锁的判断已被 issue #3 复审(T-01)取代。 3. **F15 带密建群用独立短生命周期进程** - 原条款:拉进群须当次带对话密码。 - 实际做法:主会话用 `group.create` 无密断言失败;带密成功在干净进程上立刻建群。 - 原因:与第 4 条同一死锁,补测时先用隔离进程覆盖校验路径。 - 备选方案:仅依赖第 4 条修复后在同一长会话上测 `group.add`。 - - 影响:验收覆盖仍成立。 + - 影响:验收覆盖仍成立。 + - **已被 T-01(issue #3 复审)取代**:不是服务端死锁;带密拉人已改回同一长会话 `group.add`。 4. **群事件 `emit` 改为异步 PublishDown** - 原条款:群变更向成员推 `group_event`(QoS 0)。 - 实际做法:`internal/app/group/app.go` 的 `emit` 在独立 goroutine 里延迟约 20ms 再 `PublishDown`,让上行 worker 先把 `resp` 推完。 - 原因:同一连接上 `group.create`/`group.add` 同步向本连接注入下行时,与 mochi InlineClient 互相等待,`resp` 回不去(`TestUplinkDMOfflineGroupRecall` 在清掉测试客户端 dial deadline 后稳定复现)。 - 备选方案:broker 层对 Inline 发布做无锁队列。 - - 影响:`group_event` 可能略晚于 `resp` 到达;业务结果仍以 `resp` 为准。 + - 影响:`group_event` 可能略晚于 `resp` 到达;业务结果仍以 `resp` 为准。 + - **已被 T-01(issue #3 复审)取代**:20ms 是误诊绕过;`emit` 已改回同步 QoS 0。`feat/fix-3-downlink-deadlock` 仍不合入。 ### fix-issue-1 @@ -1517,6 +1520,8 @@ ### 死锁未修(issue #3) +> **已被 T-01 取代(2026-09-30)**:下列「服务端 InlineClient 死锁」结论经实证不成立。真实根因是测试用 WebSocket 客户端 `test/harness/mqtt.go` 不按字节流拆包、并发写不加锁。保留原文供追溯。`feat/fix-3-downlink-deadlock` 与工作树 `E:\code\NixMsg-wt\fix3` 仍不合入。 + issue #3 未关闭,`feat/fix-3-downlink-deadlock` 未合入 `main`。下面是核对过的调用链、三次尝试和仍留在 `main` 上的绕过。不改产品行为。 1. **现象** @@ -1683,3 +1688,13 @@ issue #3 未关闭,`feat/fix-3-downlink-deadlock` 未合入 `main`。下面是 - 原因:未转义时搜 `e_ab` 会命中 `exab…`,搜 `%` 返回全部端。 - 备选方案:在 `internal/protocol.ValidEndpointID` 加保留字,使开通/导入一并拒绝(否决,本波不改 protocol 与 `internal/admin/endpoints.go`)。后台开通与批量导入仍可能使用 `inline`,留给后续波次。 - 影响:目录下划线按字面匹配;自助注册不能占用 `inline`。 + +### 复审修复 T-01 + +1. **测试 WS 客户端按字节流拆包并加写锁;群 emit 改回同步 QoS0** + - 日期:2026-09-30 + - 原条款:MQTT 5.0 §6 [MQTT-6.0.0-3](不得假设控制包与 WebSocket 帧对齐);PRD F16 群事件尽力推送、没有延迟 20ms;issue #3 复审「结论与统一方案」。 + - 实际做法:`test/harness/mqtt.go` 增加接收缓冲与写锁;`Recv` 按 MQTT 剩余长度拆包,opcode `0x2`/`0x0` 载荷都追加,`0x9` 在写锁内回 PONG;`Send` 与 PONG 持锁且整帧一次 `Write`。`internal/app/group/app.go` 的 `emit` 去掉 goroutine 与 `time.Sleep(20ms)`,在调用方上下文同步 `PublishDown`(QoS 0)。F15 去掉独立短进程,在同一长会话用 `group.add` + 当次对话密码。未合入 `feat/fix-3-downlink-deadlock`。 + - 原因:`resp` 回不来是测试客户端丢掉合帧里的后续 MQTT 包,以及并发写把发往服务端的流写乱;不是服务端死锁。 + - 备选方案:合入 fix3 的 broker `OnPublish`/延后队列(否决,基于错误假设,还会改 PUBACK 语义)。 + - 影响:跨线改了总控目录 `test/harness` 与身份线 `emit`;不改 `internal/broker/`。所有经 `DialMQTTWebSocket` 的测试客户端一并受益。 diff --git a/internal/app/group/app.go b/internal/app/group/app.go index 73e60d4..92537a5 100644 --- a/internal/app/group/app.go +++ b/internal/app/group/app.go @@ -892,7 +892,6 @@ func (a *App) emit(ctx context.Context, recipients []string, groupID, event, end if a.down == nil { return } - _ = ctx frame := protocol.GroupEvent{ V: protocol.Version, Type: protocol.TypeGroupEvent, GroupID: groupID, Event: event, EndpointID: endpointID, AtMs: atMs, @@ -901,20 +900,14 @@ func (a *App) emit(ctx context.Context, recipients []string, groupID, event, end if encErr != nil { return } - // 异步且略推迟:必须让处理该端上行的 worker 先 PublishDown resp。 - // 若与 resp 同时向本连接注入 group_event,会与 mochi InlineClient 互相等待。 - ids := append([]string(nil), recipients...) - go func() { - time.Sleep(20 * time.Millisecond) - seen := map[string]struct{}{} - for _, id := range ids { - if _, ok := seen[id]; ok { - continue - } - seen[id] = struct{}{} - _ = a.down.PublishDown(context.Background(), id, "", payload, port.PublishOpts{QoS: 0}) + seen := map[string]struct{}{} + for _, id := range recipients { + if _, ok := seen[id]; ok { + continue } - }() + seen[id] = struct{}{} + _ = a.down.PublishDown(ctx, id, "", payload, port.PublishOpts{QoS: 0}) + } } func encodeFrame(v any) ([]byte, error) { diff --git a/test/accept/ACCEPTANCE.md b/test/accept/ACCEPTANCE.md index 8f6c565..2a1ea25 100644 --- a/test/accept/ACCEPTANCE.md +++ b/test/accept/ACCEPTANCE.md @@ -4,6 +4,8 @@ 汇总:通过 23,失败 0,未测 0 +T-01(issue #3)更正:先前把 `group.create` / `group.add`「收不到 resp」写成服务端死锁,实际是测试 WS 客户端不按字节流拆包、并发写不加锁。F15 带密拉人已在同一长会话用 `group.add` 覆盖。 + | 编号 | 一句话 | 结果 | 备注 | |---|---|---|---| | F01 | 批量开通整批校验、停用、删除群主转让、删除后同编号重开不串数据 | 通过 | 已测:开通一端、错误密码 MQTT 拒绝、正确密码可连;批量校验/停用/删除转让/同号重开未在本用例穷尽 | @@ -20,7 +22,7 @@ | F12 | 延迟窗口内撤回对方收不到 | 通过 | 已测:延迟窗口内撤回对方无 msg/revoked | | F13 | 未推送必撤成功;群部分确认得到部分撤回 | 通过 | 已测:未推送前撤回成功;群部分撤回未覆盖 | | F14 | 回执能补送给当时离线的发送方 | 通过 | 已测:发送方离线期间对方确认,发送方重连后补到 state=accepted 回执 | -| F15 | 输一次记住、改密失效、回复免密、进群仍要密码、防多账号轮流猜 | 通过 | 已测:不带密拒绝、带对后第二条免密、改密失效、对方先发可免密回、拉群须当次密码、5 账号×10 错触发总数锁后正确密也 rate_limited 且已有授权仍可发 | +| F15 | 输一次记住、改密失效、回复免密、进群仍要密码、防多账号轮流猜 | 通过 | 已测:不带密拒绝、带对后第二条免密、改密失效、对方先发可免密回、同一长会话 group.add 带密拉人成功、5 账号×10 错触发总数锁后正确密也 rate_limited 且已有授权仍可发 | | F16 | 群主权限、退出后不再收到、解散后同编号新群不收旧消息 | 通过 | 已测:建群并拉成员后可群发;群主权限/退出/解散同号等未穷尽 | | F17 | 后台管端、管注册、管群、查记录,响应里没有正文;API 令牌可用且不能越权 | 通过 | 已测:管理登录、错误密码锁定、无 CSRF 被拒 / 有 CSRF 可通过;管端开通见 F01;管注册见 F23;令牌越权/查记录无正文等未穷尽 | | F18 | 送达后正文消失;记录天数 0 时连记录消失;防重仍在 | 通过 | 已测:确认后 message_bodies 为空;同号重试不再投递;record_retention_days=0 完成后 status=not_found 且消息行消失 | diff --git a/test/accept/rest_accept_test.go b/test/accept/rest_accept_test.go index 2a78238..7096de6 100644 --- a/test/accept/rest_accept_test.go +++ b/test/accept/rest_accept_test.go @@ -614,9 +614,6 @@ func runF15(t *testing.T, set func(string, report.Status, string)) { } // 拉进群仍要当次带对话密码(已有单聊授权不能代替)。 - // 注:真实进程上 group.add+talk_password,以及长会话后再 group.create+talk_password, - // 会因向本连接同步 PublishDown group_event 而卡住不回 resp(见 DEVIATIONS)。 - // 无密失败在本会话用 create 覆盖;带密成功在独立短生命周期进程上覆盖(同校验路径)。 alice.Request(t, map[string]any{ "v": 1, "type": "send", "rid": "f15s6", "id": "f15-reauth", "to": map[string]any{"kind": "endpoint", "id": "f15bob001"}, @@ -650,9 +647,18 @@ func runF15(t *testing.T, set func(string, report.Status, string)) { t.Errorf("expected member fail: %+v", addNo.Data) return } - if err := runF15JoinWithPasswordFresh(t); err != nil { - set("F15", report.StatusFail, "带密拉人建群: "+err.Error()) - t.Error(err) + addYes := alice.Request(t, map[string]any{ + "v": 1, "type": "group.add", "rid": "f15g2", "group_id": "g_f15a", + "members": []map[string]any{{"id": "f15bob001", "talk_password": "talk-secret-2"}}, + }) + if !addYes.OK { + set("F15", report.StatusFail, fmt.Sprintf("带密拉人失败: %+v", addYes)) + t.Errorf("group.add pw: %+v", addYes) + return + } + if fails := memberFailures(addYes.Data); len(fails) > 0 { + set("F15", report.StatusFail, fmt.Sprintf("带密拉人仍失败: %+v", addYes.Data)) + t.Errorf("group.add failed: %+v", addYes.Data) return } @@ -700,7 +706,7 @@ func runF15(t *testing.T, set func(string, report.Status, string)) { t.Errorf("still: %+v", still) return } - set("F15", report.StatusPass, "已测:不带密拒绝、带对后第二条免密、改密失效、对方先发可免密回、拉群须当次密码、5 账号×10 错触发总数锁后正确密也 rate_limited 且已有授权仍可发") + set("F15", report.StatusPass, "已测:不带密拒绝、带对后第二条免密、改密失效、对方先发可免密回、同一长会话 group.add 带密拉人成功、5 账号×10 错触发总数锁后正确密也 rate_limited 且已有授权仍可发") } func runF18(t *testing.T, set func(string, report.Status, string)) { @@ -826,40 +832,6 @@ func runF18(t *testing.T, set func(string, report.Status, string)) { set("F18", report.StatusPass, "已测:确认后 message_bodies 为空;同号重试不再投递;record_retention_days=0 完成后 status=not_found 且消息行消失") } -// runF15JoinWithPasswordFresh 在干净进程上验证带对话密码建群成功(避开长会话后 PublishDown 卡住)。 -func runF15JoinWithPasswordFresh(t *testing.T) error { - t.Helper() - srv, err := harness.Start(harness.Options{}) - if err != nil { - return fmt.Errorf("harness: %w", err) - } - defer func() { _ = srv.Stop() }() - ac := accept.AdminLogin(t, srv) - accept.CreateEndpoint(t, ac, "f15jalice", epPassword) - accept.CreateEndpoint(t, ac, "f15jbob01", epPassword) - alice := accept.MQTTLogin(t, srv.HTTPBase, "f15jalice", epPassword) - defer alice.Close() - bob := accept.MQTTLogin(t, srv.HTTPBase, "f15jbob01", epPassword) - defer bob.Close() - setTalk := bob.Request(t, map[string]any{ - "v": 1, "type": "self.talk_password", "rid": "jtp1", "talk_password": "join-secret", - }) - if !setTalk.OK { - return fmt.Errorf("设对话密码失败: %+v", setTalk) - } - addYes := alice.Request(t, map[string]any{ - "v": 1, "type": "group.create", "rid": "f15jg", "id": "g_f15j", "name": "F15J", - "members": []map[string]any{{"id": "f15jbob01", "talk_password": "join-secret"}}, - }) - if !addYes.OK { - return fmt.Errorf("带密建群失败: %+v", addYes) - } - if fails := memberFailures(addYes.Data); len(fails) > 0 { - return fmt.Errorf("带密建群仍失败: %+v", addYes.Data) - } - return nil -} - func drainReceipts(s *accept.MQTTSession, d time.Duration) { deadline := time.Now().Add(d) for time.Now().Before(deadline) { diff --git a/test/harness/mqtt.go b/test/harness/mqtt.go index 0efc129..ddbc5c5 100644 --- a/test/harness/mqtt.go +++ b/test/harness/mqtt.go @@ -12,6 +12,7 @@ import ( "net/http" "net/url" "strings" + "sync" "time" ) @@ -164,25 +165,35 @@ func (c *tcpMQTT) RemoteAddr() net.Addr { type wsMQTT struct { conn net.Conn r *bufio.Reader + buf []byte + wmu sync.Mutex } func (c *wsMQTT) Send(packet []byte) error { + c.wmu.Lock() + defer c.wmu.Unlock() return writeWSClientBinary(c.conn, packet) } func (c *wsMQTT) Recv() ([]byte, error) { for { + if pkt, n, ok := splitMQTTPacket(c.buf); ok { + c.buf = c.buf[n:] + return pkt, nil + } payload, opcode, err := readWSFrame(c.r) if err != nil { return nil, err } switch opcode { - case 0x2: - return payload, nil + case 0x0, 0x2: + c.buf = append(c.buf, payload...) case 0x8: return nil, io.EOF case 0x9: + c.wmu.Lock() _ = writeWSClientControl(c.conn, 0xA, payload) + c.wmu.Unlock() case 0xA: continue default: @@ -265,17 +276,37 @@ func writeWSClientFrame(w io.Writer, opcode byte, payload []byte) error { header = append(header, ext[:]...) } header = append(header, mask...) - masked := make([]byte, n) + frame := make([]byte, len(header)+n) + copy(frame, header) for i := 0; i < n; i++ { - masked[i] = payload[i] ^ mask[i%4] + frame[len(header)+i] = payload[i] ^ mask[i%4] } - if _, err := w.Write(header); err != nil { - return err - } - _, err := w.Write(masked) + _, err := w.Write(frame) return err } +// splitMQTTPacket 从缓冲头部切出一个完整 MQTT 控制包。不够一包时返回 ok=false。 +func splitMQTTPacket(b []byte) ([]byte, int, bool) { + if len(b) < 2 { + return nil, 0, false + } + rem, mult := 0, 1 + for i := 1; i < len(b) && i <= 4; i++ { + rem += int(b[i]&127) * mult + if b[i]&128 == 0 { + total := 1 + i + rem + if total < 0 || len(b) < total { + return nil, 0, false + } + pkt := make([]byte, total) + copy(pkt, b[:total]) + return pkt, total, true + } + mult *= 128 + } + return nil, 0, false +} + func readWSFrame(r *bufio.Reader) (payload []byte, opcode byte, err error) { b0, err := r.ReadByte() if err != nil { diff --git a/test/harness/mqtt_test.go b/test/harness/mqtt_test.go new file mode 100644 index 0000000..3b7e051 --- /dev/null +++ b/test/harness/mqtt_test.go @@ -0,0 +1,185 @@ +package harness + +import ( + "bufio" + "bytes" + "encoding/binary" + "io" + "net" + "sync" + "testing" + "time" +) + +var pingReq = []byte{0xC0, 0x00} + +func TestSplitMQTTPacket(t *testing.T) { + t.Parallel() + if _, _, ok := splitMQTTPacket(nil); ok { + t.Fatal("empty") + } + if _, _, ok := splitMQTTPacket([]byte{0xC0}); ok { + t.Fatal("incomplete header") + } + three := concat(pingReq, pingReq, []byte{0xD0, 0x00}) + pkt, n, ok := splitMQTTPacket(three) + if !ok || n != 2 || !bytes.Equal(pkt, pingReq) { + t.Fatalf("first pkt=%x n=%d ok=%v", pkt, n, ok) + } + pkt, n, ok = splitMQTTPacket(three[n:]) + if !ok || n != 2 || !bytes.Equal(pkt, pingReq) { + t.Fatalf("second pkt=%x n=%d ok=%v", pkt, n, ok) + } + pkt, n, ok = splitMQTTPacket(three[4:]) + if !ok || n != 2 || pkt[0] != 0xD0 { + t.Fatalf("third pkt=%x n=%d ok=%v", pkt, n, ok) + } +} + +func TestWSMQTTRecvSplitsCoalescedPackets(t *testing.T) { + t.Parallel() + c, peer := pipeWS(t) + three := concat(pingReq, pingReq, []byte{0xD0, 0x00}) + errCh := make(chan error, 1) + go func() { errCh <- writeWSServerFrame(peer, 0x2, three, true) }() + + got := make([][]byte, 0, 3) + for i := 0; i < 3; i++ { + pkt, err := c.Recv() + if err != nil { + t.Fatalf("recv %d: %v", i, err) + } + got = append(got, pkt) + } + if err := <-errCh; err != nil { + t.Fatalf("write: %v", err) + } + want := [][]byte{pingReq, pingReq, {0xD0, 0x00}} + for i := range want { + if !bytes.Equal(got[i], want[i]) { + t.Fatalf("pkt %d = %x want %x", i, got[i], want[i]) + } + } +} + +func TestWSMQTTRecvContinuation(t *testing.T) { + t.Parallel() + c, peer := pipeWS(t) + errCh := make(chan error, 1) + go func() { + if err := writeWSServerFrame(peer, 0x2, []byte{0xC0}, false); err != nil { + errCh <- err + return + } + errCh <- writeWSServerFrame(peer, 0x0, []byte{0x00}, true) + }() + pkt, err := c.Recv() + if err != nil { + t.Fatalf("recv: %v", err) + } + if err := <-errCh; err != nil { + t.Fatalf("write: %v", err) + } + if !bytes.Equal(pkt, pingReq) { + t.Fatalf("pkt=%x", pkt) + } +} + +func TestWSMQTTSendConcurrentLocked(t *testing.T) { + t.Parallel() + c, peer := pipeWS(t) + const nSenders = 2 + const perSender = 1000 + want := nSenders * perSender + + var wg sync.WaitGroup + errCh := make(chan error, nSenders) + for i := 0; i < nSenders; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for j := 0; j < perSender; j++ { + if err := c.Send(pingReq); err != nil { + errCh <- err + return + } + } + }() + } + + r := bufio.NewReader(peer) + _ = peer.SetDeadline(time.Now().Add(15 * time.Second)) + var buf []byte + got := 0 + for got < want { + payload, opcode, err := readWSFrame(r) + if err != nil { + t.Fatalf("read ws after %d pkts: %v", got, err) + } + if opcode != 0x0 && opcode != 0x2 { + continue + } + buf = append(buf, payload...) + for { + pkt, n, ok := splitMQTTPacket(buf) + if !ok { + break + } + buf = buf[n:] + if !bytes.Equal(pkt, pingReq) { + t.Fatalf("decoded %x after %d", pkt, got) + } + got++ + } + } + wg.Wait() + select { + case err := <-errCh: + t.Fatalf("send: %v", err) + default: + } + if got != want { + t.Fatalf("got=%d want=%d leftover=%d", got, want, len(buf)) + } +} + +func pipeWS(t *testing.T) (*wsMQTT, net.Conn) { + t.Helper() + a, b := net.Pipe() + t.Cleanup(func() { + _ = a.Close() + _ = b.Close() + }) + _ = a.SetDeadline(time.Now().Add(15 * time.Second)) + _ = b.SetDeadline(time.Now().Add(15 * time.Second)) + return &wsMQTT{conn: a, r: bufio.NewReader(a)}, b +} + +func concat(parts ...[]byte) []byte { + return bytes.Join(parts, nil) +} + +func writeWSServerFrame(w io.Writer, opcode byte, payload []byte, fin bool) error { + b0 := opcode & 0x0f + if fin { + b0 |= 0x80 + } + header := []byte{b0} + n := len(payload) + switch { + case n < 126: + header = append(header, byte(n)) + case n <= 65535: + header = append(header, 126, byte(n>>8), byte(n)) + default: + var ext [8]byte + binary.BigEndian.PutUint64(ext[:], uint64(n)) + header = append(header, 127) + header = append(header, ext[:]...) + } + frame := make([]byte, len(header)+n) + copy(frame, header) + copy(frame[len(header):], payload) + _, err := w.Write(frame) + return err +}