## SSE Last-Event-ID 补投 EventSource 断线重连时自带 Last-Event-ID 头,但服务端直接忽略了—— 所有断线期间的邮件通知都丢失。用户刷新页面也会错过已推的事件。 改为 per-user 事件环形缓冲区(500 条,~100KB/用户,20 在线 ≈ 2MB): 每次 Broadcast/SendToUser/SendToAgent 同时写入对应用户的缓冲区; AddClient 时取 Last-Event-ID 头,找到该 ID 的位置后从下一条回放。 找不到 ID 说明事件已被覆盖(缓冲区溢出),从头回放全部。 事件 ID 用全局递增序列号(非 UUID),EventSource 的 Last-Event-ID 就是靠这个 ID 记住断点的。 新增测试:缓冲区回放、溢出行为、并发安全(10 goroutine × 200 次 push)、 端到端重连验证(SendToUser → 带 Last-Event-ID 的 AddClient → 补投)。 ## 连接状态指示器 Sidebar 用户头像右下角的小圆点:绿=已连接,黄=连接中,橙=重连中,红=断开。 NarrowNav 底栏也有(移动端)。 SSE 模块新增 onSSEStatus/getSSEStatus 接口,onerror/onopen 驱动状态变化。 状态点用 absolute 定位在头像边缘,不遮挡文字。 ## 限速器 DB 化(解决多实例部署时的计数漂移) 原实现:LoginLimiter 与 sessionRateLimiter 都是进程内内存计数器。 多实例部署时各自独立计数,等效上限变成 N 倍。 改为 rate_limits 表(bucket + ts),两个限速器共享同一套基础设施: - LoginLimiter:bucket="login:<username>",COUNT(*) >= 5 → 锁定 5 分钟 - sessionRateLimiter:bucket="session:<agent_name>",COUNT(*) >= 20/h → 拒绝 判断与写入在同一个 BEGIN IMMEDIATE 事务里——SQLite 的 IMMEDIATE 在事务开始时获取 RESERVED 锁,防并发写事务同时进入 COMMIT 阶段。 实测 80 并发下恰好放行 20 次(旧内存版同样通过,但 DB 版才能多实例共享)。 DB 不可用时放行(宁可放开限速也不能让用户完全无法使用)。 新建 rate_limits 表迁移(SQLite + PG 两版)。
92 lines
2.7 KiB
Go
92 lines
2.7 KiB
Go
package handler
|
||
|
||
import (
|
||
"net/http"
|
||
|
||
"github.com/agentmail/gateway/internal/middleware"
|
||
"github.com/agentmail/gateway/internal/repo"
|
||
"github.com/agentmail/gateway/internal/sse"
|
||
)
|
||
|
||
// GET /api/v1/events/stream
|
||
//
|
||
// 四种凭证,都必须真正验证过身份才能订阅:
|
||
// Authorization: Bearer <agent_key_token> → Agent 通道(密钥认证)
|
||
// X-Agent-Name + X-Agent-Secret → Agent 通道(旧方式,兼容)
|
||
// 登录 Cookie 或 Bearer <user_key_token> → 人类用户通道
|
||
// ?access_token=<token> → 浏览器 EventSource 专用回退
|
||
//
|
||
// 注意不能只凭 X-Agent-Name 就分流:那等于任何人报个名字就能读走别人的新邮件通知。
|
||
// query 令牌仅本端点接受(EventSource 无法带自定义头),其余接口一律要求请求头,
|
||
// 因为 URL 里的令牌会进访问日志与 Referer。
|
||
func SSEStream(w http.ResponseWriter, r *http.Request) {
|
||
agentName, ok := resolveStreamAgent(r)
|
||
if !ok {
|
||
Error(w, http.StatusUnauthorized, "凭证无效")
|
||
return
|
||
}
|
||
|
||
userName := ""
|
||
if agentName == "" {
|
||
u := middleware.OptionalUserWithQuery(r)
|
||
if u == nil {
|
||
Error(w, http.StatusUnauthorized, "not authenticated")
|
||
return
|
||
}
|
||
userName = u.Username
|
||
}
|
||
|
||
client := sse.Default.AddClient(w, r, agentName, userName)
|
||
if client == nil {
|
||
Error(w, http.StatusInternalServerError, "SSE not supported")
|
||
return
|
||
}
|
||
|
||
<-r.Context().Done()
|
||
sse.Default.RemoveClient(client.ID)
|
||
}
|
||
|
||
// resolveStreamAgent 校验 Agent 侧凭证。
|
||
// 返回 ("", true) 表示这不是 Agent 请求,交给人类用户分支;
|
||
// 返回 ("", false) 表示带了 Agent 凭证但验证失败。
|
||
func resolveStreamAgent(r *http.Request) (string, bool) {
|
||
// 密钥认证:Bearer 令牌可能是 Agent 密钥,也可能是用户密钥。
|
||
// 先按 Agent 密钥试,失败就落到人类分支(那里会再按用户密钥试)。
|
||
token := middleware.BearerToken(r)
|
||
if token == "" {
|
||
token = middleware.QueryToken(r) // EventSource 回退
|
||
}
|
||
if token != "" {
|
||
name, err := repo.VerifyAgentKey(r.Context(), token)
|
||
if err == nil && name != "" {
|
||
return name, true
|
||
}
|
||
return "", true
|
||
}
|
||
|
||
name := r.Header.Get("X-Agent-Name")
|
||
if name == "" {
|
||
name = r.URL.Query().Get("agent_name")
|
||
}
|
||
if name == "" {
|
||
return "", true // 非 Agent 请求
|
||
}
|
||
|
||
secret := r.Header.Get("X-Agent-Secret")
|
||
if secret == "" {
|
||
return "", false // 报了名字却没给凭证
|
||
}
|
||
agent, err := repo.VerifyAgent(r.Context(), name, secret)
|
||
if err != nil {
|
||
return "", false
|
||
}
|
||
return agent.Name, true
|
||
}
|
||
|
||
// GET /api/v1/events/status
|
||
func SSEStatus(w http.ResponseWriter, r *http.Request) {
|
||
JSON(w, http.StatusOK, map[string]interface{}{
|
||
"connected_clients": sse.Default.ClientCount(),
|
||
})
|
||
}
|