2026-08-15 14:38:14 +08:00
|
|
|
package middleware
|
2026-05-23 20:32:12 +08:00
|
|
|
|
|
|
|
|
import (
|
2026-05-24 15:33:27 +08:00
|
|
|
"log"
|
2026-05-23 20:32:12 +08:00
|
|
|
"net"
|
|
|
|
|
"net/http"
|
2026-08-26 00:36:22 +08:00
|
|
|
"os"
|
2026-09-21 11:14:48 +08:00
|
|
|
"sort"
|
2026-08-26 00:36:22 +08:00
|
|
|
"strconv"
|
2026-05-23 20:32:12 +08:00
|
|
|
"strings"
|
|
|
|
|
"sync"
|
|
|
|
|
"time"
|
|
|
|
|
|
|
|
|
|
"github.com/volcano-tts/tts-api/common"
|
2026-08-15 13:30:35 +08:00
|
|
|
"github.com/volcano-tts/tts-api/metrics"
|
2026-08-26 00:36:22 +08:00
|
|
|
"github.com/volcano-tts/tts-api/setting"
|
2026-05-23 20:32:12 +08:00
|
|
|
)
|
|
|
|
|
|
|
|
|
|
type RateLimiter struct {
|
|
|
|
|
requests map[string][]time.Time
|
|
|
|
|
mutex sync.Mutex
|
|
|
|
|
limit int
|
|
|
|
|
window time.Duration
|
|
|
|
|
lastCleanup time.Time
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
var (
|
|
|
|
|
GlobalRateLimiter *RateLimiter
|
|
|
|
|
ConcurrencySem chan struct{}
|
2026-08-26 00:36:22 +08:00
|
|
|
// trustedProxyHops controls how X-Forwarded-For (XFF) is parsed when the
|
|
|
|
|
// direct connection comes from a private IP (i.e., we're behind a reverse
|
|
|
|
|
// proxy). Two modes are supported, switched by this single value:
|
|
|
|
|
//
|
|
|
|
|
// HEURISTIC MODE (trustedProxyHops == 0, the default):
|
|
|
|
|
// Walk XFF from the end, return the first PUBLIC IP. Skips private
|
|
|
|
|
// and loopback hops automatically. Works for ~90% of deployments
|
|
|
|
|
// without the operator needing to know the exact number of proxy
|
|
|
|
|
// hops. Trade-off in multi-hop: rate limiting is per-CDN-edge rather
|
|
|
|
|
// than per-real-client, which is "good enough" for abuse protection
|
|
|
|
|
// but not for fine-grained per-user quotas.
|
|
|
|
|
//
|
|
|
|
|
// PRECISE MODE (trustedProxyHops > 0):
|
|
|
|
|
// Count back N hops from the end of XFF and return that value. Gives
|
|
|
|
|
// precise per-real-client rate limiting even in multi-hop setups
|
|
|
|
|
// (e.g., Cloudflare + nginx). Operator MUST set this to the number
|
|
|
|
|
// of trusted reverse proxies between this service and the client.
|
|
|
|
|
//
|
|
|
|
|
// Both modes walk from the END of the XFF chain. The first value is
|
|
|
|
|
// client-controllable; trusting it would let attackers bypass IP rate
|
|
|
|
|
// limiting by sending a forged X-Forwarded-For header.
|
|
|
|
|
trustedProxyHops = 0
|
2026-05-23 20:32:12 +08:00
|
|
|
)
|
|
|
|
|
|
|
|
|
|
func InitRateLimiter() {
|
2026-08-26 00:36:22 +08:00
|
|
|
switch v := os.Getenv("TRUSTED_PROXY_HOPS"); {
|
|
|
|
|
case v == "":
|
|
|
|
|
log.Printf("TRUSTED_PROXY_HOPS 未设置,使用默认启发式模式(XFF 链尾第一个公网 IP)")
|
|
|
|
|
default:
|
|
|
|
|
n, err := strconv.Atoi(v)
|
|
|
|
|
switch {
|
|
|
|
|
case err != nil || n < 0 || n > 10:
|
|
|
|
|
log.Printf("警告: TRUSTED_PROXY_HOPS=%q 无效(需 0-10 的整数),回退到默认启发式模式", v)
|
|
|
|
|
case n == 0:
|
|
|
|
|
// "0" 或 "00" 等被 Atoi 解析为 0 的形式都归到启发式模式,
|
|
|
|
|
// 避免日志出现"精确模式, 信任 0 跳"这种自相矛盾的输出。
|
|
|
|
|
log.Printf("已配置 TRUSTED_PROXY_HOPS=%d(启发式模式,等同默认)", n)
|
|
|
|
|
default:
|
|
|
|
|
trustedProxyHops = n
|
|
|
|
|
log.Printf("已配置 TRUSTED_PROXY_HOPS=%d(精确模式,信任 %d 跳反代)", n, n)
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-05-23 20:32:12 +08:00
|
|
|
GlobalRateLimiter = &RateLimiter{
|
|
|
|
|
requests: make(map[string][]time.Time),
|
|
|
|
|
limit: common.RateLimitRequests,
|
|
|
|
|
window: common.RateLimitWindow,
|
|
|
|
|
}
|
|
|
|
|
ConcurrencySem = make(chan struct{}, common.MaxConcurrentRequests)
|
2026-08-26 00:36:22 +08:00
|
|
|
|
|
|
|
|
// 同步到 setting 包,供 LogStartupSummary 展示
|
|
|
|
|
setting.TrustedProxyHops = trustedProxyHops
|
2026-05-23 20:32:12 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (rl *RateLimiter) Allow(key string) bool {
|
|
|
|
|
rl.mutex.Lock()
|
|
|
|
|
defer rl.mutex.Unlock()
|
|
|
|
|
|
|
|
|
|
now := time.Now()
|
|
|
|
|
cutoff := now.Add(-rl.window)
|
|
|
|
|
|
|
|
|
|
if now.Sub(rl.lastCleanup) > common.CleanupInterval {
|
|
|
|
|
rl.cleanup()
|
|
|
|
|
rl.lastCleanup = now
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
timestamps := rl.requests[key]
|
|
|
|
|
valid := make([]time.Time, 0, len(timestamps))
|
|
|
|
|
for _, ts := range timestamps {
|
|
|
|
|
if ts.After(cutoff) {
|
|
|
|
|
valid = append(valid, ts)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if len(valid) >= rl.limit {
|
|
|
|
|
rl.requests[key] = valid
|
2026-08-15 13:30:35 +08:00
|
|
|
metrics.RateLimitRejected.Inc(nil)
|
2026-05-23 20:32:12 +08:00
|
|
|
return false
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
valid = append(valid, now)
|
|
|
|
|
rl.requests[key] = valid
|
|
|
|
|
return true
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (rl *RateLimiter) cleanup() {
|
|
|
|
|
cutoff := time.Now().Add(-rl.window)
|
|
|
|
|
for k, v := range rl.requests {
|
|
|
|
|
valid := make([]time.Time, 0, len(v))
|
|
|
|
|
for _, ts := range v {
|
|
|
|
|
if ts.After(cutoff) {
|
|
|
|
|
valid = append(valid, ts)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
if len(valid) == 0 {
|
|
|
|
|
delete(rl.requests, k)
|
|
|
|
|
} else {
|
|
|
|
|
rl.requests[k] = valid
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-05-24 15:33:27 +08:00
|
|
|
|
|
|
|
|
if len(rl.requests) > common.MaxRateLimiterEntries {
|
|
|
|
|
log.Printf("警告: 限流器条目数 %d 超过上限 %d,触发强制清理", len(rl.requests), common.MaxRateLimiterEntries)
|
2026-09-21 11:14:48 +08:00
|
|
|
// 【修复】原版用 `for k := range rl.requests` 删,Go map 遍历顺序随机,
|
|
|
|
|
// 会随机删掉活跃用户(其条目 timestamp 仍在窗口内),导致该用户下次请求
|
|
|
|
|
// 拿到新配额 — 攻击者可用大量伪造 IP 撑爆 map 触发清理,反而"清洗"
|
|
|
|
|
// 自己留的活跃条目,绕过限流。
|
|
|
|
|
// 修复:按"最近一次请求时间(lastTs)"升序排序,删最旧的(最可能已离开/低频),
|
|
|
|
|
// 保留最活跃用户,语义符合"限流器只淘汰冷条目"的预期。
|
|
|
|
|
// 排序复杂度 O(n log n),但只在超 10w 条目时触发,代价可接受。
|
|
|
|
|
type entry struct {
|
|
|
|
|
key string
|
|
|
|
|
lastTs time.Time
|
|
|
|
|
}
|
|
|
|
|
entries := make([]entry, 0, len(rl.requests))
|
|
|
|
|
for k, v := range rl.requests {
|
|
|
|
|
// 走到这里 v 一定非空(cleanup 第一阶段会把空 timestamps 删掉),
|
|
|
|
|
// 取最后一个 timestamp 作为"最近活跃时间"。
|
|
|
|
|
lastTs := v[len(v)-1]
|
|
|
|
|
entries = append(entries, entry{key: k, lastTs: lastTs})
|
|
|
|
|
}
|
|
|
|
|
sort.Slice(entries, func(i, j int) bool {
|
|
|
|
|
return entries[i].lastTs.Before(entries[j].lastTs)
|
|
|
|
|
})
|
|
|
|
|
for _, e := range entries {
|
2026-05-24 15:33:27 +08:00
|
|
|
if len(rl.requests) <= common.MaxRateLimiterEntries/2 {
|
|
|
|
|
break
|
|
|
|
|
}
|
2026-09-21 11:14:48 +08:00
|
|
|
delete(rl.requests, e.key)
|
2026-05-24 15:33:27 +08:00
|
|
|
}
|
|
|
|
|
}
|
2026-05-23 20:32:12 +08:00
|
|
|
}
|
|
|
|
|
|
2026-06-26 19:12:50 +08:00
|
|
|
var privateCIDRs []*net.IPNet
|
|
|
|
|
|
|
|
|
|
func init() {
|
|
|
|
|
for _, cidr := range []string{
|
|
|
|
|
"10.0.0.0/8",
|
|
|
|
|
"172.16.0.0/12",
|
|
|
|
|
"192.168.0.0/16",
|
|
|
|
|
"127.0.0.0/8",
|
|
|
|
|
"169.254.0.0/16",
|
|
|
|
|
"::1/128",
|
|
|
|
|
"fc00::/7",
|
|
|
|
|
"fe80::/10",
|
|
|
|
|
} {
|
|
|
|
|
_, ipNet, _ := net.ParseCIDR(cidr)
|
|
|
|
|
privateCIDRs = append(privateCIDRs, ipNet)
|
2026-05-23 20:32:12 +08:00
|
|
|
}
|
2026-06-26 19:12:50 +08:00
|
|
|
}
|
2026-05-23 20:32:12 +08:00
|
|
|
|
2026-06-26 19:12:50 +08:00
|
|
|
func isPrivateIP(ipStr string) bool {
|
|
|
|
|
ip := net.ParseIP(ipStr)
|
|
|
|
|
if ip == nil {
|
|
|
|
|
return false
|
|
|
|
|
}
|
|
|
|
|
if ip.IsLoopback() || ip.IsLinkLocalUnicast() || ip.IsLinkLocalMulticast() {
|
|
|
|
|
return true
|
2026-05-23 20:32:12 +08:00
|
|
|
}
|
2026-06-26 19:12:50 +08:00
|
|
|
for _, cidr := range privateCIDRs {
|
|
|
|
|
if cidr.Contains(ip) {
|
|
|
|
|
return true
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
return false
|
|
|
|
|
}
|
2026-05-23 20:32:12 +08:00
|
|
|
|
2026-06-26 19:12:50 +08:00
|
|
|
func GetClientIP(r *http.Request) string {
|
|
|
|
|
directIP, _, err := net.SplitHostPort(r.RemoteAddr)
|
2026-05-23 20:32:12 +08:00
|
|
|
if err != nil {
|
2026-06-26 19:12:50 +08:00
|
|
|
directIP = r.RemoteAddr
|
2026-05-23 20:32:12 +08:00
|
|
|
}
|
2026-06-26 19:12:50 +08:00
|
|
|
|
|
|
|
|
if isPrivateIP(directIP) {
|
2026-08-26 00:36:22 +08:00
|
|
|
// Parse X-Forwarded-For when there's a reverse proxy in front (direct
|
|
|
|
|
// connection is from a private IP). Both modes walk from the END of
|
|
|
|
|
// the chain so that the client-controllable first value cannot be
|
|
|
|
|
// used to spoof a different client IP for rate limit bypass.
|
2026-06-26 19:12:50 +08:00
|
|
|
if xff := r.Header.Get("X-Forwarded-For"); xff != "" {
|
2026-08-26 00:36:22 +08:00
|
|
|
parts := strings.Split(xff, ",")
|
|
|
|
|
if trustedProxyHops > 0 {
|
|
|
|
|
// PRECISE MODE: count back N hops from end. Real client IP
|
|
|
|
|
// sits at index (len(parts) - N). Walk backwards to skip
|
|
|
|
|
// any malformed values; if chain is shorter than expected,
|
|
|
|
|
// fall through to the first valid IP in the chain.
|
|
|
|
|
target := len(parts) - trustedProxyHops
|
|
|
|
|
if target < 0 {
|
|
|
|
|
target = 0
|
|
|
|
|
}
|
|
|
|
|
for i := target; i >= 0; i-- {
|
|
|
|
|
ip := strings.TrimSpace(parts[i])
|
|
|
|
|
if net.ParseIP(ip) != nil {
|
|
|
|
|
return ip
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
} else {
|
|
|
|
|
// HEURISTIC MODE (default): walk from end, return first
|
|
|
|
|
// PUBLIC IP. Skips private/loopback hops that come from
|
|
|
|
|
// internal proxies between the public-facing proxy and us.
|
|
|
|
|
for i := len(parts) - 1; i >= 0; i-- {
|
|
|
|
|
ip := strings.TrimSpace(parts[i])
|
|
|
|
|
if parsed := net.ParseIP(ip); parsed != nil && !isPrivateIP(ip) {
|
|
|
|
|
return ip
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-06-26 19:12:50 +08:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
if xri := strings.TrimSpace(r.Header.Get("X-Real-IP")); xri != "" {
|
|
|
|
|
if net.ParseIP(xri) != nil {
|
|
|
|
|
return xri
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return directIP
|
2026-05-23 20:32:12 +08:00
|
|
|
}
|