Files
MailUI4Agents/gateway/internal/repo/sessionrate.go
JianFeeeee 9d4718a412 feat: SSE Last-Event-ID 补投 + 连接状态指示 + 限速器 DB 化
## 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 两版)。
2026-09-02 14:33:41 +08:00

49 lines
1.6 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 repo
import (
"context"
"time"
"github.com/agentmail/gateway/internal/db"
)
// ---------- 新建会话速率限制 ----------
//
// Agent 可以用 name@path.new 开一串新会话,每条都是全新预算 ——
// 速率限制只压住「短时间内暴开」这个滥用形态,过一个窗口自动恢复。
// ErrSessionRateLimited 表示该 Agent 短时间内新建会话过多。
// (目前未使用,直接返回 retryAfter 由 handler 构造 429 响应)
const (
sessionRateWindow = time.Hour
sessionRateLimit = 20
)
// AllowNewSession 供 handler 调用:Agent 新建会话前先过速率限制。
// 人类用户不走这条路径(手工点「新建邮件」的频率天然受限)。
// 返回 (allowed, retryAfter)。DB 不可用时放行。
func AllowNewSession(ctx context.Context, agentName string) (bool, int) {
if agentName == "" {
return true, 0
}
return RateLimitCheckAndRecord(ctx, "session:"+agentName, sessionRateWindow, sessionRateLimit)
}
// ReleaseNewSession 建会话失败后归还名额。
// DB-backed 方式下记账在 AllowNewSession 里已完成,失败时需手动删除最近一条。
func ReleaseNewSession(ctx context.Context, agentName string) {
if agentName == "" {
return
}
bucket := "session:" + agentName
// 删掉最近一条(建会话失败,那次不该占名额)
_, _ = db.DB.ExecContext(ctx,
`DELETE FROM rate_limits WHERE bucket = $1 AND ts = (
SELECT MAX(ts) FROM rate_limits WHERE bucket = $1
)`, bucket)
}
// SessionRateLimit 暴露窗口内的新建上限,供错误文案使用。
func SessionRateLimit() int { return sessionRateLimit }