package sse import ( "net/http/httptest" "strconv" "strings" "testing" "time" ) func TestEventRingPushReplay(t *testing.T) { ring := newEventRing(5) // 推 3 条 for i := 1; i <= 3; i++ { ring.push(StoredEvent{ ID: string(rune('0' + i)), EventType: "test", Data: []byte(`{"n":` + string(rune('0'+i)) + `}`), Timestamp: time.Now(), }) } // 空 afterID → 首次连接,不回放(缓冲区未满) rec := httptest.NewRecorder() ring.replay("", rec, rec, nil) if rec.Body.Len() > 0 { t.Error("首次连接不应回放事件,实际:", rec.Body.String()) } // 有 afterID → 从下一条开始回放 rec2 := httptest.NewRecorder() ring.replay("1", rec2, rec2, nil) body := rec2.Body.String() if !strings.Contains(body, "id: 2") { t.Error("afterID=1 应该回放 id:2,实际:", body) } if !strings.Contains(body, "id: 3") { t.Error("afterID=1 应该回放 id:3,实际:", body) } if strings.Contains(body, "id: 1") { t.Error("afterID=1 不应回放 id:1,实际:", body) } // 不存在的 afterID → 从头回放全部 rec3 := httptest.NewRecorder() ring.replay("999", rec3, rec3, nil) body3 := rec3.Body.String() if !strings.Contains(body3, "id: 1") { t.Error("不存在的 afterID 应从头回放,实际:", body3) } } func TestEventRingOverflow(t *testing.T) { ring := newEventRing(3) // 推 5 条(超过容量 3,最旧的 2 条被覆盖) for i := 1; i <= 5; i++ { ring.push(StoredEvent{ ID: string(rune('0' + i)), EventType: "test", Data: []byte(`{}`), Timestamp: time.Now(), }) } if !ring.full { t.Fatal("推了 5 条进容量 3 的缓冲区,应该已满") } // afterID=2 已被覆盖 → 找不到位置,从头回放全部 rec := httptest.NewRecorder() ring.replay("2", rec, rec, nil) body := rec.Body.String() if !strings.Contains(body, "id: 3") || !strings.Contains(body, "id: 5") { t.Error("缓冲区溢出后应能回放可用范围,实际:", body) } } func TestEventRingConcurrent(t *testing.T) { ring := newEventRing(100) done := make(chan bool, 10) for i := 0; i < 10; i++ { go func() { for j := 0; j < 200; j++ { ring.push(StoredEvent{ ID: "evt", EventType: "test", Data: []byte(`{}`), Timestamp: time.Now(), }) } done <- true }() } for i := 0; i < 10; i++ { <-done } // 只验证不 panic,不验证内容(并发下顺序无意义) } // TestHeartbeatIntervalIsWellUnderProxyIdleTimeout —— 心跳必须**明显小于**常见代理读超时。 // // 2026-09-15 实测(用户报「每次点击按钮 1-2s 延迟」):他的 SSE 连接每次只活 // 34.6s / 39.4s / 56.9s 就被关闭,而普通 API 只要 30-58ms —— 掐连接的不是我们, // 是中间那层反代的空闲超时,而当时的心跳是 30s,正好与它擦边。 // 这条判据钉的不是"某个数字",而是**量级关系**:心跳要留出余量,不能与超时相当。 func TestHeartbeatIntervalIsWellUnderProxyIdleTimeout(t *testing.T) { if heartbeatInterval >= 20*time.Second { t.Fatalf("心跳间隔 %s 与常见的 30s 代理读超时擦边:实测连接只活 34-57s 就断,"+ "心跳必须明显小于该超时(当前上限取 20s)", heartbeatInterval) } } /* ★ 2026-10-02:Frame 钩子必须在**两条**写路径上都生效。 回放是第三条写 Res 的路径(另两条是 Send / SendWithID)。它原先也把格式写死成 AgentMail 的 SSE 形状 —— MCP 端点要靠它把同一批事件说成 JSON-RPC 方言。 这一格存在的理由:只改 SendWithID 而漏改 replay,**平时完全看不出来** —— 只有在「MCP 客户端带着 Last-Event-ID 重连」时才会收到一堆自己解不开的帧。 那种错要等真实断线才现形,所以只能靠判据钉。 */ func TestCustomFrameAppliesToReplayToo(t *testing.T) { ring := newEventRing(8) for i := 1; i <= 3; i++ { ring.push(StoredEvent{ ID: strconv.Itoa(i), EventType: "new_mail", Data: []byte(`{"n":1}`), Timestamp: time.Now(), }) } // 自定义方言:一眼可辨,且**不带** AgentMail 的 `event:` 行 jsonrpcFrame := func(id, eventType string, data []byte) string { return `{"jsonrpc":"2.0","method":"notifications/message","params":{` + `"id":"` + id + `","type":"` + eventType + `","payload":` + string(data) + `}}` + "\n\n" } rec := httptest.NewRecorder() ring.replay("1", rec, rec, jsonrpcFrame) body := rec.Body.String() if !strings.Contains(body, `"jsonrpc":"2.0"`) { t.Errorf("★ 回放必须用自定义 frame(否则重连的 MCP 客户端收到异种方言):%q", body) } if strings.Contains(body, "event: new_mail") { t.Errorf("★ 回放仍写成 AgentMail 的 SSE 形状 —— Frame 钩子漏接在 replay 路径上:%q", body) } if !strings.Contains(body, `"id":"2"`) { t.Errorf("回放内容不对(应从 afterID=1 之后开始):%q", body) } // 反向对照:nil ⇒ 默认格式一字未变(各桥与 WebUI 走这条) rec2 := httptest.NewRecorder() ring.replay("1", rec2, rec2, nil) if !strings.Contains(rec2.Body.String(), "event: new_mail") { t.Errorf("nil frame 必须回落默认 AgentMail 格式:%q", rec2.Body.String()) } }