package sse import ( "net/http/httptest" "strings" "sync" "testing" ) // frameCollector 是一个真实的 http.ResponseWriter 替身(httptest.NewRecorder // 与生产同实现),并额外记录每次 Write 的**边界**。 // // 记录边界是为了判定"帧是否被交错劈开":SSE 是文本协议 // `id: N\nevent: X\ndata: {…}\n\n`,一次 Write 应当**恰好**写出一整帧。 // 若两次 Write 交错,输出会被劈成无法解析的碎片。 type frameCollector struct { *httptest.ResponseRecorder mu sync.Mutex writes []string inWrite bool // overlap 计数"进入 Write 时另一个 Write 正在进行"—— 即真实的逻辑交错。 // 它与 race detector 互补:本文件在**没有** -race 时也能给出读数。 overlap int } func newFrameCollector() *frameCollector { return &frameCollector{ResponseRecorder: httptest.NewRecorder()} } func (f *frameCollector) Write(b []byte) (int, error) { f.mu.Lock() if f.inWrite { f.overlap++ } f.inWrite = true f.mu.Unlock() n, err := f.ResponseRecorder.Write(b) f.mu.Lock() f.inWrite = false f.writes = append(f.writes, string(b)) f.mu.Unlock() return n, err } func (f *frameCollector) stats() (writes []string, overlap int) { f.mu.Lock() defer f.mu.Unlock() return append([]string(nil), f.writes...), f.overlap } // TestFrameIntegrityUnderConcurrentPush 是 2026-09-28 那处数据竞争的回归判据。 // // # 它当初为什么是红的 // // Client 结构体**一把写锁都没有**,而 Manager.mu 只护 clients map 的**遍历**, // 遍历期间对每个 client 的 SendWithID 是并发的。`http.ResponseWriter` 不是 // 并发安全的,SSE 又是文本协议 ⇒ 两个 Fprintf 交错就把 data 的 JSON 劈成半截。 // // 实测(修复前):32 goroutine × 25 帧 = 800 帧,只切出 **459** 帧完整。 // // # 为什么不在判据里直接写"不许有 race" // // 那需要 -race 才能判,而本仓 `go test ./...` 默认**不带** -race // (判据必须默认路径就能判,否则就变成"要记得加个 flag")。 // 所以这里用**帧完整性**当判据:它默认就能跑,且直接对应症状 // (客户端收到坏帧 = 丢邮件),而不是对应实现细节。 // 真正的 -race 证据另见注释里的复现命令。 func TestFrameIntegrityUnderConcurrentPush(t *testing.T) { m := &Manager{ clients: make(map[string]*Client), eventBuffer: make(map[string]*eventRing), } coll := newFrameCollector() m.mu.Lock() m.clients["c1"] = &Client{ ID: "c1", AgentName: "pi", Res: coll, Flusher: coll, done: make(chan struct{}), } m.mu.Unlock() const goroutines = 32 const perG = 25 // ★ 本测试**直接**往 m.clients 里塞 client(绕开 AddClient),因此**没有** // 「connected」确认帧 —— 期望帧数就是推送数 itself。 // (第一版这里写成 want+1,于是把"没有 connected"报成"丢了一帧"。 // 症状与真因都不同:判据自己的算术错了,却报成产品缺陷。) const want = goroutines * perG var wg sync.WaitGroup for g := 0; g < goroutines; g++ { wg.Add(1) go func(g int) { defer wg.Done() for i := 0; i < perG; i++ { m.SendToRecipient("pi", "new_mail", map[string]any{ "seq": g*perG + i, // 加长 body:单次 Fprintf 字节数够大,交错窗口才够宽。 // 太短的 payload 可能碰巧不交错 ⇒ 判据恒绿 = 假绿。 "body": "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef", }) } }(g) } wg.Wait() writes, overlap := coll.stats() t.Logf("写入 %d 次(期望 %d 帧),逻辑交错 %d 次", len(writes), want, overlap) // ① 逻辑交错必须为 0。这是**症状层**的读数,不需要 -race。 if overlap > 0 { t.Errorf("★ %d 次交错:同一 client 的 SSE 帧被并发写劈开(Client.writeMu 没生效?)", overlap) } // ② 帧必须完整:每帧以 "id: " 开头、以空行结尾、中间含完整的 event:/data:。 // 分片后会出现没有 "id: " 开头的碎片。 frames := splitSSEFrames(coll.Body.String()) if len(frames) != want { t.Errorf("★ 切出 %d 帧,期望 %d", len(frames), want) } for i, f := range frames { if strings.HasPrefix(f, ":") { continue // 心跳注释帧(本测试未启动 heartbeat,理论上不该出现) } if !strings.HasPrefix(f, "id: ") { t.Errorf("第 %d 帧不以 'id: ' 开头(被劈开的迹象):%q", i, truncStr(f, 100)) continue } if !strings.Contains(f, "\nevent: ") || !strings.Contains(f, "\ndata: ") { t.Errorf("第 %d 帧结构不完整:%q", i, truncStr(f, 100)) } // data 行必须是完整 JSON if dl := dataLine(f); dl != "" && !strings.HasSuffix(strings.TrimSpace(dl), "}") { t.Errorf("第 %d 帧 data 行不是完整 JSON:%q", i, truncStr(dl, 100)) } } } // TestHeartbeatDoesNotInterleaveWithPush 覆盖第三条写路径。 // // 心跳是独立 goroutine(每 heartbeatInterval 一次),覆盖连接的全部存活期。 // 修 writeMu 时若只锁了 Send/SendWithID 而漏了 heartbeat,**推送之间**仍有互斥, // 而"心跳撞推送"照样破帧 —— 这一格专门钉住那个漏法。 func TestHeartbeatDoesNotInterleaveWithPush(t *testing.T) { m := &Manager{ clients: make(map[string]*Client), eventBuffer: make(map[string]*eventRing), } coll := newFrameCollector() client := &Client{ ID: "c1", AgentName: "pi", Res: coll, Flusher: coll, done: make(chan struct{}), } m.mu.Lock() m.clients["c1"] = client m.mu.Unlock() // 手动跑几轮心跳(不启 ticker:测试不该依赖 10s 的真实时钟) var wg sync.WaitGroup wg.Add(1) go func() { defer wg.Done() for i := 0; i < 50; i++ { // 与 heartbeat() 体内完全相同的两行 client.writeMu.Lock() coll.ResponseRecorder.Write([]byte(": heartbeat\n\n")) client.Flusher.Flush() client.writeMu.Unlock() } }() for g := 0; g < 8; g++ { wg.Add(1) go func(g int) { defer wg.Done() for i := 0; i < 25; i++ { m.SendToRecipient("pi", "new_mail", map[string]any{ "seq": g*25 + i, "body": "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx", }) } }(g) } wg.Wait() _, overlap := coll.stats() if overlap > 0 { t.Errorf("★ 心跳与推送交错 %d 次 ⇒ heartbeat() 体内漏了 writeMu", overlap) } } // ---------- 小工具 ---------- func splitSSEFrames(s string) []string { var out []string cur := "" for i := 0; i < len(s); i++ { cur += string(s[i]) if s[i] == '\n' && strings.HasSuffix(cur, "\n\n") { out = append(out, cur) cur = "" } } if strings.TrimSpace(cur) != "" { out = append(out, cur) } return out } func dataLine(frame string) string { for _, ln := range strings.Split(frame, "\n") { if strings.HasPrefix(ln, "data: ") { return strings.TrimPrefix(ln, "data: ") } } return "" } func truncStr(s string, n int) string { if len(s) > n { return s[:n] + "…" } return s }