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