Files
MailUI4Agents/server/internal/sse/manager.go
JianFeeeee aeb1f4116b fix(WebUI): SSE 订阅跟着账号凭证走 + 断线重放 + 兜底轮询
用户报:**页面停留不动,新邮件不自动同步**(手动刷新能看到)。

## 根因一(主因):SSE 连接不跟着账号走

`App.tsx` 的 effect 依赖是 `[phase]`,而切号(`accountStore.setActive`)
只换 `api/config` 的 base/token、**不改 phase** ⇒ SSE 连接仍绑旧账号的凭证:
旧账号的新邮件照收,新账号的一封都不推。而 `fetchInbox` 走**新**凭证 ⇒ 数据是新的。
⇒ 表现正是「不自动同步,但手动刷新能看到」。

修法:effect 依赖加上「当前凭证身份」(base + token)。
不在切号处显式重建订阅 —— 那要改所有调用点、漏一处就不刷新;
凭证变化的**唯一发生地**是 api/config,从那里取身份更可靠。

★ 身份**不含 user**:同一账号重新登录 token 变了,那个账号的邮件仍该收
  (服务端按 user/agent 绑通道,见 sse.bufferKey);
  按 base+token 判只会让「同账号换令牌」多触发一次重连(无害)。

## 根因二:断线重连不重放

服务端一直支持按 Last-Event-ID 回放(ring.replay,500 条缓冲),
EventSource 断线后**本来会自己重连并带该头**。但这里的 onerror 主动
`close(false)` 再 `open()` —— **换了 EventSource 对象**,
而 Last-Event-ID 是浏览器为**那个对象**记的 ⇒ 服务端拿不到 ⇒ 不回放。

EventSource 不能设请求头 ⇒ 游标只能进 query,服务端相应要读
`?lastEventId=`(**两侧都要改,缺一半都不生效且没有任何东西会红**)。
服务端写成 query 优先、header 兜底 —— header 保留给 Agent 侧(curl/SDK)。
⚠ query 会进访问日志;游标是自增数字(不是令牌),与「令牌不进日志」的约定不同级。

`onerror` 区分两种重连:
- 断线(凭证没变)⇒ 带游标,服务端回放断线期间的事件
- 切号(凭证变了)⇒ **必须不带** —— 拿旧账号的 id 去问新账号会搅乱事件流

## 根因三:连接静默但不再收数据

SSE 只在**真的断开**时触发 onerror。有一类故障它看不见:
连接还在、TCP 没断、却不再收数据(代理静默丢包 / NAT 超时 /
中间设备挂死长连接)。两端都认为正常 ⇒ 不重连 ⇒ 页面停留就再也不同步。

补 `lib/inboxFallbackPoll.ts` 作为冗余通道:
- 探针 `getInbox('all', 1)` **只要 total**(全量重拉会让接口与渲染无谓抖动)
- 首轮只建基线不触发;探针失败**不重置基线**(否则一次抖动会变成「下一轮假装有变化」)
- inFlight 去重,慢网络下不叠请求
- 页面隐藏时暂停,恢复可见**立刻探一次**(用户往往正是「切回来发现没更新」才报的)
- 切号时 resetPollBaseline:新账号 total 与旧账号无关,不丢会白拉一次
- 间隔 30s:远大于 SSE 的秒级延迟(正常时纯冗余),又短到挂死最多 30s 被发现

## 判据

12 格(sse-credentials 5 + inbox-fallback-poll 7)。九个变异全部经得起:
依赖退回 [phase] / 重连不带游标 / 切号也带旧游标 / 服务端不读 query /
catch 重置基线 / cleanup 漏停轮询 / 凭证依赖丢失 / 探针拉全量 / 恢复可见不立即探。

★ 一处判据自身缺陷被变异抓出来并修掉:第 3 格原先只查
  `url += \`${sep}lastEventId=…\`` 这行**文本存在**,把 `if (lastEventId)`
  改成 `if (false)` 后照样绿 —— 正则匹配文本,缺陷在控制流。
  补了条件本身的断言才红。与「catch 里不得重置基线」是同一类教训。
2026-10-02 10:46:45 +08:00

437 lines
13 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 的位置,从下一条开始回放。
func (r *eventRing) replay(afterID string, flush http.Flusher, res http.ResponseWriter) {
r.mu.Lock()
defer r.mu.Unlock()
if afterID == "" {
return // 首次连接,不回放
}
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.Fprintf(res, "id: %s\nevent: %s\ndata: %s\n\n", 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
}
// 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 二者恰其一)
func (m *Manager) AddClient(res http.ResponseWriter, r *http.Request, agentName, userName 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{}),
}
// 设置 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.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()
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()
fmt.Fprintf(c.Res, "id: %s\nevent: %s\ndata: %s\n\n", 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()
}
}
}