da9c8334d8
- SSE Keepalive Ping (15s心跳防止代理断连) - Timing HTTP 头 (X-Timing-Queue/Inference/Total-Ms) - Adapter Request-ID 传播到后端 - Session 清理日志回调 - Server 安全加固 (ReadHeaderTimeout/MaxHeaderBytes 防 slowloris) - Usage Tracker 数据保留清理 (retentionDays + 定期清理) - Config Reload 后 Adapter Registry 更新 (RegisterIfAbsent + RWMutex) - Rate Limiter 空闲 Bucket 清理 (30分钟过期) - Shutdown Drain 超时可配置 (ShutdownDrainSeconds) - Config 模型字段校验增强 (provider/endpoint/actual_model) - Auth 过期 Key 自动清理 (5分钟扫描) - Admin API Rate Limiting - Adapter Health Check 独立超时 (每个 adapter 3s) - TCP 连接阶段超时 (DialContext 5s + KeepAlive 30s) - 幂等键缓存、审计日志、Gzip 中间件、CORS Expose Headers - Backpressure 响应头、熔断器 Prometheus 指标 - 连接池优化、Trace-ID 全链路传播
186 lines
4.3 KiB
Go
186 lines
4.3 KiB
Go
package ratelimit
|
|
|
|
import (
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// TokenBucket 是一个简单的令牌桶限流器。
|
|
type TokenBucket struct {
|
|
capacity float64 // 桶容量
|
|
refillRate float64 // 每秒补充令牌数
|
|
tokens float64 // 当前令牌数
|
|
lastRefill time.Time // 上次补充时间
|
|
lastAccess time.Time // 最后访问时间,用于空闲清理
|
|
mu sync.Mutex
|
|
}
|
|
|
|
// NewTokenBucket 创建一个令牌桶。
|
|
// capacity 为桶容量(突发上限),refillPerSecond 为每秒补充速率。
|
|
func NewTokenBucket(capacity float64, refillPerSecond float64) *TokenBucket {
|
|
now := time.Now()
|
|
return &TokenBucket{
|
|
capacity: capacity,
|
|
refillRate: refillPerSecond,
|
|
tokens: capacity, // 初始满桶
|
|
lastRefill: now,
|
|
lastAccess: now,
|
|
}
|
|
}
|
|
|
|
// Allow 尝试消耗 1 个令牌,返回是否允许。
|
|
func (tb *TokenBucket) Allow() bool {
|
|
return AllowN(tb, 1)
|
|
}
|
|
|
|
// AllowN 尝试消耗 n 个令牌,返回是否允许。
|
|
func AllowN(tb *TokenBucket, n float64) bool {
|
|
tb.mu.Lock()
|
|
defer tb.mu.Unlock()
|
|
|
|
now := time.Now()
|
|
elapsed := now.Sub(tb.lastRefill).Seconds()
|
|
tb.tokens += elapsed * tb.refillRate
|
|
if tb.tokens > tb.capacity {
|
|
tb.tokens = tb.capacity
|
|
}
|
|
tb.lastRefill = now
|
|
tb.lastAccess = now
|
|
|
|
if tb.tokens >= n {
|
|
tb.tokens -= n
|
|
return true
|
|
}
|
|
return false
|
|
}
|
|
|
|
// Tokens 返回当前可用令牌数(近似值)。
|
|
func (tb *TokenBucket) Tokens() float64 {
|
|
tb.mu.Lock()
|
|
defer tb.mu.Unlock()
|
|
return tb.tokens
|
|
}
|
|
|
|
// Limiter 管理 per-app 令牌桶限流器。
|
|
type Limiter struct {
|
|
mu sync.RWMutex
|
|
buckets map[string]*TokenBucket // app_id -> bucket
|
|
capacity float64
|
|
refill float64
|
|
}
|
|
|
|
// NewLimiter 创建一个 per-app 限流管理器。
|
|
// capacity 为每个 app 的桶容量,refillPerSecond 为每秒补充速率。
|
|
// idleTimeout 为空闲 app bucket 的过期时间,超过此时间未被访问的 bucket 将被清理。
|
|
func NewLimiter(capacity float64, refillPerSecond float64) *Limiter {
|
|
l := &Limiter{
|
|
buckets: make(map[string]*TokenBucket),
|
|
capacity: capacity,
|
|
refill: refillPerSecond,
|
|
}
|
|
go l.cleanupIdle()
|
|
return l
|
|
}
|
|
|
|
// Allow 检查指定 app 是否被限流。
|
|
func (l *Limiter) Allow(appID string) bool {
|
|
l.mu.RLock()
|
|
bucket, ok := l.buckets[appID]
|
|
l.mu.RUnlock()
|
|
|
|
if !ok {
|
|
bucket = NewTokenBucket(l.capacity, l.refill)
|
|
l.mu.Lock()
|
|
// 双检查,防止竞态
|
|
if existing, ok := l.buckets[appID]; ok {
|
|
bucket = existing
|
|
} else {
|
|
l.buckets[appID] = bucket
|
|
}
|
|
l.mu.Unlock()
|
|
}
|
|
|
|
return bucket.Allow()
|
|
}
|
|
|
|
// RateLimitInfo 包含限流信息,用于响应头。
|
|
type RateLimitInfo struct {
|
|
Allowed bool
|
|
Limit int
|
|
Remaining int
|
|
}
|
|
|
|
// AllowWithInfo 检查限流并返回详情(用于 X-RateLimit-* 响应头)。
|
|
func (l *Limiter) AllowWithInfo(appID string) RateLimitInfo {
|
|
l.mu.RLock()
|
|
bucket, ok := l.buckets[appID]
|
|
l.mu.RUnlock()
|
|
|
|
if !ok {
|
|
bucket = NewTokenBucket(l.capacity, l.refill)
|
|
l.mu.Lock()
|
|
if existing, ok := l.buckets[appID]; ok {
|
|
bucket = existing
|
|
} else {
|
|
l.buckets[appID] = bucket
|
|
}
|
|
l.mu.Unlock()
|
|
}
|
|
|
|
allowed := bucket.Allow()
|
|
tokens := int(bucket.Tokens())
|
|
|
|
return RateLimitInfo{
|
|
Allowed: allowed,
|
|
Limit: int(l.capacity),
|
|
Remaining: tokens,
|
|
}
|
|
}
|
|
|
|
// SetLimit 为指定 app 设置自定义限流参数。
|
|
func (l *Limiter) SetLimit(appID string, capacity, refillPerSecond float64) {
|
|
l.mu.Lock()
|
|
defer l.mu.Unlock()
|
|
l.buckets[appID] = NewTokenBucket(capacity, refillPerSecond)
|
|
}
|
|
|
|
// GetLimit 返回指定 app 的限流配置(burst 和 rate_per_minute)。
|
|
func (l *Limiter) GetLimit(appID string) (burst int, ratePerMinute int) {
|
|
l.mu.RLock()
|
|
bucket, ok := l.buckets[appID]
|
|
l.mu.RUnlock()
|
|
|
|
if !ok {
|
|
return int(l.capacity), int(l.refill * 60)
|
|
}
|
|
return int(bucket.capacity), int(bucket.refillRate * 60)
|
|
}
|
|
|
|
// Remove 移除指定 app 的限流器。
|
|
func (l *Limiter) Remove(appID string) {
|
|
l.mu.Lock()
|
|
defer l.mu.Unlock()
|
|
delete(l.buckets, appID)
|
|
}
|
|
|
|
// cleanupIdle 定期清理空闲超过 30 分钟的 app bucket,防止内存泄漏。
|
|
func (l *Limiter) cleanupIdle() {
|
|
ticker := time.NewTicker(5 * time.Minute)
|
|
defer ticker.Stop()
|
|
idleThreshold := 30 * time.Minute
|
|
|
|
for range ticker.C {
|
|
cutoff := time.Now().Add(-idleThreshold)
|
|
l.mu.Lock()
|
|
for appID, bucket := range l.buckets {
|
|
bucket.mu.Lock()
|
|
idle := bucket.lastAccess.Before(cutoff)
|
|
bucket.mu.Unlock()
|
|
if idle {
|
|
delete(l.buckets, appID)
|
|
}
|
|
}
|
|
l.mu.Unlock()
|
|
}
|
|
}
|