134 lines
2.3 KiB
Go
134 lines
2.3 KiB
Go
package nixmsg
|
||
|
||
import (
|
||
"math/rand"
|
||
"sync"
|
||
"time"
|
||
)
|
||
|
||
// 重连 / 限速重交共用标称间隔:第 n 次 min(1s×2^(n-1), 30s),n≥1。
|
||
func nominalDelay(n int) time.Duration {
|
||
if n < 1 {
|
||
return 0
|
||
}
|
||
if n > 6 {
|
||
return 30 * time.Second
|
||
}
|
||
d := time.Second
|
||
for i := 1; i < n; i++ {
|
||
d *= 2
|
||
if d >= 30*time.Second {
|
||
return 30 * time.Second
|
||
}
|
||
}
|
||
return d
|
||
}
|
||
|
||
var backoffJitter = withJitter
|
||
|
||
func withJitter(d time.Duration) time.Duration {
|
||
if d <= 0 {
|
||
return 0
|
||
}
|
||
f := 0.7 + rand.Float64()*0.6
|
||
return time.Duration(float64(d) * f)
|
||
}
|
||
|
||
// reconnectBackoff 只维护一个连续失败计数 n。
|
||
// 应用调用 Connect 后的第一次连接不等待;之后第 n 次等待 nominalDelay(n)×抖动。
|
||
type reconnectBackoff struct {
|
||
mu sync.Mutex
|
||
n int
|
||
skipFirst bool
|
||
online bool
|
||
onlineAt time.Time
|
||
counted bool
|
||
}
|
||
|
||
func newReconnectBackoff() *reconnectBackoff {
|
||
return &reconnectBackoff{skipFirst: true}
|
||
}
|
||
|
||
// Func 供 autopaho;忽略其 attempt,只用本对象的 n。
|
||
func (b *reconnectBackoff) Func(attempt int) time.Duration {
|
||
_ = attempt
|
||
return b.NextWait()
|
||
}
|
||
|
||
func (b *reconnectBackoff) NextWait() time.Duration {
|
||
return b.nextWait(true)
|
||
}
|
||
|
||
func (b *reconnectBackoff) nextWait(jitter bool) time.Duration {
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
b.counted = false
|
||
if b.skipFirst {
|
||
b.skipFirst = false
|
||
return 0
|
||
}
|
||
n := b.n
|
||
if n < 1 {
|
||
n = 1
|
||
}
|
||
d := nominalDelay(n)
|
||
if jitter {
|
||
return backoffJitter(d)
|
||
}
|
||
return d
|
||
}
|
||
|
||
func (b *reconnectBackoff) MarkOnline() {
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
b.online = true
|
||
b.onlineAt = time.Now()
|
||
b.counted = false
|
||
}
|
||
|
||
func (b *reconnectBackoff) MarkOffline() {
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
if b.counted {
|
||
return
|
||
}
|
||
b.counted = true
|
||
was := b.online
|
||
onlineAt := b.onlineAt
|
||
b.online = false
|
||
if !was {
|
||
b.n++
|
||
if b.n < 1 {
|
||
b.n = 1
|
||
}
|
||
return
|
||
}
|
||
if time.Since(onlineAt) >= 60*time.Second {
|
||
b.n = 1
|
||
return
|
||
}
|
||
b.n++
|
||
if b.n < 1 {
|
||
b.n = 1
|
||
}
|
||
}
|
||
|
||
func (b *reconnectBackoff) N() int {
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
return b.n
|
||
}
|
||
|
||
func (b *reconnectBackoff) setOnlineAtForTest(t time.Time) {
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
b.online = true
|
||
b.onlineAt = t
|
||
}
|
||
|
||
// computeBackoffDelay 保留给旧单测:忽略 base,按 n 计算标称延迟。
|
||
func computeBackoffDelay(base time.Duration, n int) time.Duration {
|
||
_ = base
|
||
return nominalDelay(n)
|
||
}
|