package core import ( "math" "sync" "time" ) // AdaptiveTimeout 基于 RTT 采样的自适应超时计算器 // 算法:timeout = mean(RTT) + 4 * stddev(RTT),clamp 到 [min, max] // 冷启动阶段(样本不足)返回用户配置的固定超时 type AdaptiveTimeout struct { mu sync.Mutex samples []float64 // 环形缓冲区,单位 ms pos int // 写入位置 count int // 已采集总数 size int // 缓冲区容量 minTO time.Duration maxTO time.Duration warmup int // 冷启动所需最小样本数 cachedTO time.Duration dirty bool } // NewAdaptiveTimeout 创建自适应超时计算器 // maxTimeout: 用户配置的超时上限(即原始固定超时) func NewAdaptiveTimeout(maxTimeout time.Duration) *AdaptiveTimeout { // minTO: 自适应超时下限,取 max(500ms, maxTimeout/5) // 依据:高并发下 TCP 握手存在尾延迟(OS 调度抖动、backlog 溢出、端口竞争), // 过低的下限会导致开放端口被误判为关闭(issue #503) minTO := maxTimeout / 5 if minTO < 500*time.Millisecond { minTO = 500 * time.Millisecond } return &AdaptiveTimeout{ samples: make([]float64, 64), size: 64, minTO: minTO, maxTO: maxTimeout, warmup: 10, } } // Record 记录一次成功连接的 RTT func (a *AdaptiveTimeout) Record(rtt time.Duration) { a.mu.Lock() a.samples[a.pos%a.size] = float64(rtt.Milliseconds()) a.pos++ a.count++ a.dirty = true a.mu.Unlock() } // Timeout 获取当前推荐超时值 // 样本不足时返回 maxTO(冷启动) // 锁外执行均值/标准差计算,减少锁持有时间 func (a *AdaptiveTimeout) Timeout() time.Duration { a.mu.Lock() if a.count < a.warmup { a.mu.Unlock() return a.maxTO } if !a.dirty { cached := a.cachedTO a.mu.Unlock() return cached } n := a.size if a.count < a.size { n = a.count } // 拷贝样本到本地,释放锁后再计算 localSamples := make([]float64, n) start := a.pos % a.size if a.count < a.size { copy(localSamples, a.samples[:n]) } else { copy(localSamples[:a.size-start], a.samples[start:]) copy(localSamples[a.size-start:], a.samples[:start]) } a.mu.Unlock() // 锁外计算 var sum float64 for _, s := range localSamples { sum += s } mean := sum / float64(n) var variance float64 for _, s := range localSamples { d := s - mean variance += d * d } stddev := math.Sqrt(variance / float64(n)) ms := mean + 4*stddev to := time.Duration(ms) * time.Millisecond if to < a.minTO { to = a.minTO } if to > a.maxTO { to = a.maxTO } // 短暂加锁更新缓存 a.mu.Lock() a.cachedTO = to a.dirty = false a.mu.Unlock() return to }