## 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 两版)。
49 lines
1.6 KiB
Go
49 lines
1.6 KiB
Go
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 }
|