Files
MailUI4Agents/server/internal/sse/manager.go
JianFeeeee 560c462768 feat(mcp): GET /api/v1/mcp —— 投递侧事件流(让接入方被动收信,不用轮询)
## 这半边解决什么

工具面(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 格全绿
2026-10-02 15:37:34 +08:00

500 lines
16 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package sse
import (
"encoding/json"
"fmt"
"net/http"
"sync"
"time"
"github.com/google/uuid"
)
// eventRing 是单用户事件的有界环形缓冲区。
//
// EventSource 断线重连时自带 Last-Event-ID 头:服务端据此回放断线期间的事件。
// 没有它,重连后永远看不到断线期间收到的邮件 —— 而这正是实时协作的体验核心。
//
// 缓冲区大小 500 条:一条事件约 200B(typical),500 条 ≈ 100KB/用户。
// 20 个在线用户 ≈ 2MB,远低于 OOM 风险。
type eventRing struct {
mu sync.Mutex
events []StoredEvent
cap int
head int // 下一次写入的位置
full bool
}
// StoredEvent 是缓冲区中的单条事件。
type StoredEvent struct {
ID string // 自增序列号,EventSource 的 Last-Event-ID 值
EventType string
Data []byte
Timestamp time.Time
}
func newEventRing(cap int) *eventRing {
return &eventRing{events: make([]StoredEvent, cap), cap: cap}
}
// push 追加一条事件到缓冲区。满了就覆盖最旧的。
func (r *eventRing) push(evt StoredEvent) {
r.mu.Lock()
defer r.mu.Unlock()
r.events[r.head] = evt
r.head = (r.head + 1) % r.cap
if r.head == 0 && !r.full {
r.full = true
}
}
// replay 从 afterID 之后的所有事件回放给 ResponseWriter。
// afterID 为空时:缓冲区未满不回放(首次连接无历史);满了也不回放
// (首次连接的 EventSource 不传 Last-Event-ID)。
// afterID 非空时:找到该 ID 的位置,从下一条开始回放。
//
// ★ frame 必传(2026-10-02):回放是**第三条**写 Res 的路径,它原本也把格式
// 写死成 AgentMail 的 SSE 形状。MCP 连接断线重连时走的就是这条路 ——
// 不传 frame 会让它收到一堆自己的客户端解不开的帧(同一个连接上两种方言)。
// 这种错只在「MCP 客户端带 Last-Event-ID 重连」时现形,平时完全看不见。
func (r *eventRing) replay(afterID string, flush http.Flusher, res http.ResponseWriter,
frame func(id, eventType string, data []byte) string) {
r.mu.Lock()
defer r.mu.Unlock()
if afterID == "" {
return // 首次连接,不回放
}
if frame == nil {
frame = defaultFrame
}
start := -1
total := r.cap
for i := 0; i < r.cap; i++ {
idx := (r.head + i) % r.cap
if r.events[idx].ID == afterID {
start = (idx + 1) % r.cap
break
}
}
if start == -1 {
// afterID 不在缓冲区里(已被覆盖或从未存在),
// 回放缓冲区里所有事件 —— 宁可重复也不丢失
start = 0
if !r.full {
total = r.head
}
} else {
// 从 start 开始到 head 结束
total = r.head - start
if total < 0 {
total += r.cap
}
}
for i := 0; i < total; i++ {
idx := (start + i) % r.cap
evt := &r.events[idx]
if evt.ID == "" {
continue
}
fmt.Fprint(res, frame(evt.ID, evt.EventType, evt.Data))
}
flush.Flush()
}
// Client 是一个 SSE 连接客户端
type Client struct {
ID string
AgentName string // 非空 = Agent 侧连接
UserName string // 非空 = 已登录人类用户的前端连接
Res http.ResponseWriter
Flusher http.Flusher
done chan struct{}
// ★ writeMu 串行化对 Res 的**每一次**写入。
//
// 为什么必需(2026-09-28 实测,-race 证实):
// http.ResponseWriter **不是并发安全**的,而本结构原先**一把写锁都没有**。
// Manager.mu 只护 `clients` map 的**遍历**,遍历期间的 `c.SendWithID` 是并发的 ——
// 任何两条并发请求都会同时向同一个 client 写。
//
// 生产上会打中的三条路径:
// ① handler/permission.go:412-414 —— `SendToAgent(perm.AgentName, …)` 紧接
// `SendToUser(user.Username, …)`,两个不同 HTTP 请求(两个 goroutine)命中同一账号;
// ② 任意两条并发邮件:一封投给 B,B 的插件回信进 C 的 handler,而 A 的
// `notify.Recipients` 还没跑完;
// ③ `heartbeat` 那条 goroutine 每 10s 写一次(见 heartbeatInterval 的注释),
// 与推送撞车的概率随在线时长线性上升。
//
// 症状:SSE 是 `id: N\nevent: X\ndata: {…}\n\n` 的**文本协议**,两个 Fprintf
// 交错 ⇒ data 的 JSON 被劈成半截 ⇒ 客户端 EventSource 收到坏帧、丢事件。
// 实测 32 goroutine × 25 帧 = 800 帧,只切出 **459** 帧完整。
//
// 为什么不能靠上层串行化:推送方有 5 个入口(SendToUser / SendToAgent /
// SendToRecipient / Broadcast / replay),要保证"同一个 client 的所有写互斥",
// 责任只能落在 client 自己身上 —— 那是唯一能覆盖**全部**写者的位置。
writeMu sync.Mutex
// Frame 可选:把一条事件渲染成**要写进流里的字节**。
//
// 为什么需要它(2026-10-02,MCP 端点要把同一批事件说成别的协议):
// 本管理器原本把帧格式**写死**成 AgentMail 的 SSE 形状。而 MCP
// (Streamable HTTP)要求服务端→客户端的方向用 **JSON-RPC 通知**经 SSE 下发
// —— 同一根管子、两种方言。
//
// 把格式抽成一个函数而不是复制一份 Manager:心跳、Last-Event-ID 回放、
// 断线清理、`writeMu` 串行化、容量上限、事件环缓冲 —— 这些是**踩过坑才对的**
// (见 writeMu 的注释)。复制一份等于把那些坑再踩一遍,
// 而两边的修复会各走各的。
//
// nil = 用默认的 AgentMail 格式(各桥与 WebUI 走这条,行为一字未变)。
Frame func(id, eventType string, data []byte) string
}
// defaultFrame 是 AgentMail 自己的 SSE 帧格式(带 id)。
func defaultFrame(id, eventType string, data []byte) string {
return fmt.Sprintf("id: %s\nevent: %s\ndata: %s\n\n", id, eventType, data)
}
// fill 是**唯一**把事件写进 Res 的地方(心跳除外的两条推送路径共用)。
//
// 抽出来的意义:让「默认格式」与「自定义格式」在代码上对称 —— 将来改默认格式时,
// 两边的差异一眼可见;而漏掉其中一条路径就会让某种连接收到半截方言。
func (c *Client) fill(id, eventType string, data []byte) {
if c.Frame != nil {
fmt.Fprint(c.Res, c.Frame(id, eventType, data))
return
}
fmt.Fprint(c.Res, defaultFrame(id, eventType, data))
}
// Manager 管理所有 SSE 客户端连接
type Manager struct {
mu sync.RWMutex
clients map[string]*Client
// eventBuffer:per-user/agent 的事件环形缓冲区,供 Last-Event-ID 回放。
// Key 是 userName(人类)或 agentName(Agent),二者共享一个 map。
// 不是连接级别的 —— 同一用户断线重连后仍能从同一个缓冲区拿到断线期间的事件。
eventBuffer map[string]*eventRing
bufMu sync.RWMutex
seqCounter uint64 // 全局递增序列号,用作事件 ID
seqMu sync.Mutex
}
// Default 是全局 SSE 管理器
var Default = &Manager{
clients: make(map[string]*Client),
eventBuffer: make(map[string]*eventRing),
}
const eventBufferCap = 500 // 每用户最多保留 500 条事件
// nextEventID 生成下一个全局递增的事件 ID
func (m *Manager) nextEventID() string {
m.seqMu.Lock()
defer m.seqMu.Unlock()
m.seqCounter++
return fmt.Sprintf("%d", m.seqCounter)
}
// getOrCreateRing 获取或创建用户的环形缓冲区
func (m *Manager) getOrCreateRing(key string) *eventRing {
if key == "" {
return nil
}
m.bufMu.RLock()
ring, ok := m.eventBuffer[key]
m.bufMu.RUnlock()
if ok {
return ring
}
m.bufMu.Lock()
defer m.bufMu.Unlock()
// double-check
if ring, ok = m.eventBuffer[key]; ok {
return ring
}
ring = newEventRing(eventBufferCap)
m.eventBuffer[key] = ring
return ring
}
// AddClient 注册一个新 SSE 客户端(agentName 与 userName 二者恰其一)
// AddClient 注册一个 SSE 客户端,用**默认** AgentMail 帧格式。
//
// 各桥与 WebUI 走这条(行为与历史完全一致)。需要别的方言的调用方
// (MCP 端点把同一批事件说成 JSON-RPC 通知)用 AddClientWithFrame。
func (m *Manager) AddClient(res http.ResponseWriter, r *http.Request, agentName, userName string) *Client {
return m.AddClientWithFrame(res, r, agentName, userName, nil)
}
// AddClientWithFrame 注册一个 SSE 客户端,用调用方给的帧格式。
//
// ★ frame 必须在**构造时**传入,不能注册后再设。
//
// 原因:Last-Event-ID 回放发生在“先写响应、再注册”的**前半段**(见下面那段
// 关于持 writeMu 的注释)。事后挂 Frame 只能影响之后推来的事件,
// 而重连回放那一批仍会走默认格式 —— 同一个连接上出现两种方言,
// 且**只在真实断线重连时现形**。frame=nil 表示默认格式。
func (m *Manager) AddClientWithFrame(res http.ResponseWriter, r *http.Request,
agentName, userName string, frame func(id, eventType string, data []byte) string) *Client {
flusher, ok := res.(http.Flusher)
if !ok {
return nil
}
id := uuid.New().String()[:8]
client := &Client{
ID: id,
AgentName: agentName,
UserName: userName,
Res: res,
Flusher: flusher,
done: make(chan struct{}),
Frame: frame,
}
// 设置 SSE 响应头
res.Header().Set("Content-Type", "text/event-stream")
res.Header().Set("Cache-Control", "no-cache")
res.Header().Set("Connection", "keep-alive")
res.Header().Set("X-Accel-Buffering", "no")
// Last-Event-ID 回放:EventSource 断线重连时自带这个头,
// 服务端据此把断线期间的事件补上 —— 否则重连后永远看不到那段时间的邮件。
//
// ★ 持 writeMu 写入(尽管此时**按构造就是单写者**):
// 回放发生在下面 `m.clients[id] = client` **之前** —— 此刻还没有任何 goroutine
// 拿得到这个 client 的指针,所以本身上就是安全的。持锁是为了让「对 Res 的写入
// 一律经由 writeMu」成为**结构上**的纪律:将来有人把注册提前、或把回放挪到
// 注册之后(很自然的一个改动),没上锁的版本会**静默**退化成并发写。
// 锁在这里零成本,而它买的正是「后人改顺序也不会破」这件事。
// ★ 2026-10-02:lastID 也接受 query(`?lastEventId=`)。
//
// 为什么需要:EventSource **不能设请求头**,而 WebUI 的重连逻辑是
// 「close 掉旧对象再新建一个」(为了 UI 能显示「重连中」——
// readyState 在网络断开时不一定及时反映)。
// 换对象就换掉了浏览器为**那个对象**记的 Last-Event-ID,
// 于是断线期间的事件永远不会被回放:用户报的现象就是
// 「页面停留不动、新邮件不自动同步」。
//
// 两处取值:query 优先(客户端显式给的),否则用标准头
// (Agent 侧如 curl/SDK 可能仍走标准重连语义)。
lastID := r.URL.Query().Get("lastEventId")
if lastID == "" {
lastID = r.Header.Get("Last-Event-ID")
}
key := m.bufferKey(userName, agentName)
if ring := m.getOrCreateRing(key); ring != nil && lastID != "" {
client.writeMu.Lock()
ring.replay(lastID, flusher, res, client.Frame)
client.writeMu.Unlock()
}
m.mu.Lock()
m.clients[id] = client
m.mu.Unlock()
// 发送连接确认(带 id 让客户端知道自己的 ID)
evtID := m.nextEventID()
client.SendWithID(evtID, "connected", map[string]string{"id": id})
// 启动心跳
go m.heartbeat(client)
fmt.Printf("[SSE] Client connected: %s (agent=%q user=%q) lastID=%q\n", id, agentName, userName, lastID)
return client
}
// bufferKey 返回缓冲区 key:优先 userName(人类),其次 agentName(Agent)
func (m *Manager) bufferKey(userName, agentName string) string {
if userName != "" {
return "u:" + userName
}
if agentName != "" {
return "a:" + agentName
}
return ""
}
// RemoveClient 移除一个客户端
func (m *Manager) RemoveClient(id string) {
m.mu.Lock()
if c, ok := m.clients[id]; ok {
close(c.done)
delete(m.clients, id)
fmt.Printf("[SSE] Client disconnected: %s\n", id)
}
m.mu.Unlock()
}
// SendToAgent 向指定 Agent 名的所有客户端推送事件
func (m *Manager) SendToAgent(agentName, eventType string, data interface{}) {
if agentName == "" {
return
}
// 写入缓冲区
evtID := m.nextEventID()
raw, _ := json.Marshal(data)
if ring := m.getOrCreateRing(m.bufferKey("", agentName)); ring != nil {
ring.push(StoredEvent{ID: evtID, EventType: eventType, Data: raw, Timestamp: time.Now()})
}
m.mu.RLock()
defer m.mu.RUnlock()
for _, c := range m.clients {
if c.AgentName == agentName {
c.SendWithID(evtID, eventType, data)
}
}
}
// SendToUser 向指定人类用户的所有前端连接推送事件
func (m *Manager) SendToUser(userName, eventType string, data interface{}) {
if userName == "" {
return
}
evtID := m.nextEventID()
raw, _ := json.Marshal(data)
if ring := m.getOrCreateRing(m.bufferKey(userName, "")); ring != nil {
ring.push(StoredEvent{ID: evtID, EventType: eventType, Data: raw, Timestamp: time.Now()})
}
m.mu.RLock()
defer m.mu.RUnlock()
for _, c := range m.clients {
if c.UserName == userName {
c.SendWithID(evtID, eventType, data)
}
}
}
// SendToRecipient 根据收件人名同时尝试 Agent 通道与人类用户通道
func (m *Manager) SendToRecipient(name, eventType string, data interface{}) {
if name == "" {
return
}
evtID := m.nextEventID()
raw, _ := json.Marshal(data)
// 同时写两个缓冲区(人类或 Agent,或两者都有)
if ring := m.getOrCreateRing(m.bufferKey(name, "")); ring != nil {
ring.push(StoredEvent{ID: evtID, EventType: eventType, Data: raw, Timestamp: time.Now()})
}
if ring := m.getOrCreateRing(m.bufferKey("", name)); ring != nil {
ring.push(StoredEvent{ID: evtID, EventType: eventType, Data: raw, Timestamp: time.Now()})
}
m.mu.RLock()
defer m.mu.RUnlock()
for _, c := range m.clients {
if c.AgentName == name || c.UserName == name {
c.SendWithID(evtID, eventType, data)
}
}
}
// Broadcast 向所有客户端广播事件(心跳、系统通知等)
func (m *Manager) Broadcast(eventType string, data interface{}) {
evtID := m.nextEventID()
raw, _ := json.Marshal(data)
// 广播写入所有用户的缓冲区(确保任何用户重连都能回放)
m.bufMu.RLock()
for key, ring := range m.eventBuffer {
ring.push(StoredEvent{ID: evtID, EventType: eventType, Data: raw, Timestamp: time.Now()})
_ = key // key 仅用于日志,此处不需
}
m.bufMu.RUnlock()
m.mu.RLock()
defer m.mu.RUnlock()
for _, c := range m.clients {
c.SendWithID(evtID, eventType, data)
}
}
// ClientCount 返回当前连接数
func (m *Manager) ClientCount() int {
m.mu.RLock()
defer m.mu.RUnlock()
return len(m.clients)
}
// Send 向单个客户端发送事件(无 ID)
func (c *Client) Send(eventType string, data interface{}) {
defer func() { recover() }()
jsonData, err := json.Marshal(data)
if err != nil {
return
}
c.writeMu.Lock()
defer c.writeMu.Unlock()
if c.Frame != nil {
fmt.Fprint(c.Res, c.Frame("", eventType, jsonData))
} else {
fmt.Fprintf(c.Res, "event: %s\ndata: %s\n\n", eventType, jsonData)
}
c.Flusher.Flush()
}
// SendWithID 向单个客户端发送带 ID 的事件
func (c *Client) SendWithID(id, eventType string, data interface{}) {
defer func() { recover() }()
jsonData, err := json.Marshal(data)
if err != nil {
return
}
c.writeMu.Lock()
defer c.writeMu.Unlock()
c.fill(id, eventType, jsonData)
c.Flusher.Flush()
}
// heartbeatInterval —— 心跳间隔。
//
// ★ 2026-09-15 实测(用户报「每次点击按钮 1-2s 延迟」):他的 SSE 连接每次只活
// 34.6s / 39.4s / 56.9s 就被关闭(网关日志里 /events/stream 的耗时即连接寿命),
// 而普通 API 只要 30-58ms —— 说明不是服务端慢,是**连接被中间反代按空闲超时掐掉**,
// 而我们的心跳是 30s,正好与那个超时擦边:晚一点就被判空闲。
//
// 心跳必须**明显小于**常见的 30s/60s 代理读超时,而不是与它相当。10s 留了三倍余量,
// 代价只是每 10s 一个 16 字节的注释帧。
const heartbeatInterval = 10 * time.Second
// heartbeat 定期发送心跳保活
func (m *Manager) heartbeat(client *Client) {
ticker := time.NewTicker(heartbeatInterval)
defer ticker.Stop()
for {
select {
case <-client.done:
return
case <-ticker.C:
defer func() { recover() }()
// ★ 同一把 writeMu:心跳是本结构里**第三条**写 Res 的路径。
// 漏了它就等于"推送之间互斥、心跳不参与"—— 而心跳每 10s 一次、
// 覆盖连接的全部存活期,撞上推送是必然事件(见 writeMu 的注释 ③)。
client.writeMu.Lock()
fmt.Fprintf(client.Res, ": heartbeat\n\n")
client.Flusher.Flush()
client.writeMu.Unlock()
}
}
}