## 这半边解决什么
工具面(POST)只解决「接入方**问**」。这一条解决「服务端**说**」:
邮件投递时把 new_mail / session_update 推给接入方,让它**拉起对话** ——
与各桥靠 /api/v1/events/stream 收信是同一件事,只是方言不同:
桥: id: 7\nevent: new_mail\ndata: {…}\n\n
MCP: {"jsonrpc":"2.0","method":"notifications/message","params":{…}}
## 为什么复用 sse.Manager 而不是另起一套
Manager 里那些东西**都是踩过坑才对的**:writeMu 串行化(2026-09-28 -race
实测 http.ResponseWriter 并发写会把 JSON 劈成半截,800 帧只切出 459 个完整)、
Last-Event-ID 回放(宁可重复也不丢失)、心跳(反代按空闲 30-58s 掐连接)、
环形缓冲上限、断线清理。复制一份等于把那些坑再踩一遍,
而两边的修复从此各走各的。
代价是 `sse.Client` 多了一个可选 `Frame` 钩子:
**nil = AgentMail 原格式,各桥与 WebUI 行为一字未变**(默认值即历史行为)。
## ★ 回放是第三条写路径,漏了就只在断线时现形
`Send` / `SendWithID` / `replay` 是三条写 Res 的路径。原先**三条都把格式写死**,
只改前两条的话:MCP 客户端**平时**一切正常,只有带 `Last-Event-ID` 重连时
才会收到一批自己解不开的帧 —— 同一个连接上两种方言。
判据 `TestCustomFrameAppliesToReplayToo` 专门钉这条,并带反向对照
(nil 帧必须回落 AgentMail 格式)。
`Frame` 必须在**注册时**传入(`AddClientWithFrame`),不能事后设 ——
回放发生在「先写响应、再注册」的前半段,事后设只影响之后推来的事件。
原先 `AddClient` 保留为薄封装,各桥与 WebUI 调用点一字未改。
## 判据(6 格)
Frame 是 JSON-RPC 2.0 通知 + 帧完整性(单事件、\n\n 结尾)
payload 原样嵌入(不是 JSON 字符串)—— 再 marshal 会让客户端解析两次
event_id / event_type 必带(前者是 Last-Event-ID 续传的依据)
Accept 判定(含 q 值、大小写)
匿名 GET → 401(不能变成静默的匿名订阅)
缺 Accept → 406(接错的客户端会静默收不到东西)
## 顺带修:TestAdvanceRecurrenceLunar 的时区缺陷(★ 今天第三次假红)
全量测试红了,查下来是**我今天早些时候改判据时引入的**,与本次改动无关。
农历换算必须按**本地公历日**算(`AdvanceRecurrence` 里那句
`eventTime.In(time.Local)` 就是这条规则)。库里读回的 EventTime 是 **UTC**
(DSN 用 `_timezone=UTC`),UTC 比本地晚 8 小时,跨零点时农历日差一天:
start (Local) = 2026-10-04 农历日 24
after (UTC) = 2026-11-01 16:00 农历日 23 ← 断言没换算时区(错)
after.In(Local) = 2026-11-02 00:00 农历日 24 ← 正确
服务端代码一直是对的,是判据没照做。失败信息里现在打印时区,
免得下次要重新推导一遍。变异验证:去掉 `.In(time.Local)` → 红 1 ✓
(这条判据是农历的第三次假红了:3459605「断言要求不存在的农历日」、
今天早些「起点写死日期 + advanceToFuture 跳过过期月份」、现在「没换算时区」——
三次都是判据自己写错,代码三次都对。它依赖 Local 时区与「今天」,
天生脆弱,值得记着。)
## 验证
go test ./... 14 包全绿
go test ./internal/sse/ 含新判据绿
go test ./internal/mcp/ 6 格新判据 + 原 19 格全绿
164 lines
5.2 KiB
Go
164 lines
5.2 KiB
Go
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())
|
||
}
|
||
}
|