Files
selfrelease da9c8334d8
CI / lint (push) Has been cancelled
CI / test (push) Has been cancelled
CI / build (push) Has been cancelled
CI / security-scan (push) Has been cancelled
feat: 十轮网关优化 - 安全加固/可观测性/性能/可靠性
- 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 全链路传播
2026-08-03 15:43:11 +08:00

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()
}
}