Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
57 changes: 54 additions & 3 deletions sentinel.go
Original file line number Diff line number Diff line change
Expand Up @@ -530,10 +530,61 @@ func (c *sentinelClient) _switchTarget(addr string, isMaster bool) (err error) {
return nil
}

const (
// refreshMaxRetryShift caps the exponent, not the number of retries: the
// delay stops doubling once attempts reach it, but refreshRetry keeps going
// until it succeeds or the client is closed. It only bounds how fast the
// client retries, never whether it does — the client always heals itself.
refreshMaxRetryShift = 9
refreshMaxRetryDelay = time.Second
)

// refreshRetryDelay is the backoff applied between failed refresh attempts.
// It mirrors defaultRetryDelayFn's "equal jitter" scheme, but on a millisecond
// rather than a microsecond base: a refresh is a multi-round-trip operation
// against every known sentinel, so retrying it hundreds of thousands of times
// per second is never useful. It settles at ~512ms-1s.
func refreshRetryDelay(attempts int) time.Duration {
base := 1 << min(refreshMaxRetryShift, attempts)
jitter := util.FastRand(base)
return min(refreshMaxRetryDelay, time.Duration(base+jitter)*time.Millisecond)
}

// refreshRetry keeps refreshing the topology until it succeeds or the client is
// closed.
//
// This used to be an unbounded `goto retry` loop with no delay at all. When
// every sentinel is unreachable — precisely what happens during a failover
// where the primary and a co-located sentinel go down together — refresh()
// fails immediately and the loop spins as fast as the CPU allows. Each
// iteration also dials every sentinel in the list, so one client can saturate a
// core and flood the surviving sentinels with connection attempts exactly while
// they are running the election.
//
// refreshRetry is additionally re-entered from listWatch's error handler, so
// failures compound.
func (c *sentinelClient) refreshRetry() {
retry:
if err := c.refresh(); err != nil {
goto retry
// Stop once Close() has been called, so a client shut down while its
// sentinels are unreachable does not leak this goroutine. The check is at
// the loop head so a Close() during the backoff wait exits without a further
// refresh.
for attempts := 0; atomic.LoadUint32(&c.stop) == 0; attempts++ {
if err := c.refresh(); err == nil {
return
}
c.waitBeforeRetry(refreshRetryDelay(attempts))
}
}

// waitBeforeRetry sleeps for d, but returns early once Close() has been called
// so shutdown does not have to wait out a full backoff interval. There is no
// close channel to select on here, so it polls c.stop in short steps.
func (c *sentinelClient) waitBeforeRetry(d time.Duration) {
const step = 20 * time.Millisecond
for d > 0 && atomic.LoadUint32(&c.stop) == 0 {
s := min(step, d)
time.Sleep(s)
d -= s
}
}

Expand Down
122 changes: 122 additions & 0 deletions sentinel_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3570,3 +3570,125 @@ func TestSentinelSendToReplicasClientPubSub(t *testing.T) {
time.Sleep(time.Millisecond * 100)
}
}

func TestRefreshRetryDelay(t *testing.T) {
// Must never be zero: a zero delay reintroduces the hot spin this backoff
// exists to prevent.
for attempts := 0; attempts < 64; attempts++ {
d := refreshRetryDelay(attempts)
if d <= 0 {
t.Fatalf("attempts %d: delay must be positive, got %v", attempts, d)
}
if d > refreshMaxRetryDelay {
t.Fatalf("attempts %d: delay %v exceeds cap %v", attempts, d, refreshMaxRetryDelay)
}
}
// And it must actually grow, otherwise a long outage still hammers the
// sentinels at the initial rate.
if refreshRetryDelay(0) >= refreshRetryDelay(refreshMaxRetryShift) {
t.Fatalf("delay should increase with attempts: first=%v capped=%v",
refreshRetryDelay(0), refreshRetryDelay(refreshMaxRetryShift))
}
// The jitter has to survive once the base stops growing, which is where a
// long outage spends all of its time. Capping the jittered sum at the same
// magnitude as the base clamps every sample onto the cap, and a fleet of
// clients then retries in lockstep — exactly the thundering herd the jitter
// exists to prevent.
for _, attempts := range []int{refreshMaxRetryShift, refreshMaxRetryShift + 1, 64} {
seen := make(map[time.Duration]struct{})
for range 2000 {
seen[refreshRetryDelay(attempts)] = struct{}{}
}
if len(seen) < 2 {
t.Fatalf("attempts %d: every retry lands on the same delay %v, the jitter is being clamped away",
attempts, refreshRetryDelay(attempts))
}
}
}

// TestRefreshRetryWaitCancelsOnStop pins that the backoff wait is interruptible:
// a Close() mid-wait must return promptly instead of sleeping out the interval.
// It waits an hour, so a plain time.Sleep would hang the test.
func TestRefreshRetryWaitCancelsOnStop(t *testing.T) {
c := &sentinelClient{}
done := make(chan struct{})
go func() {
c.waitBeforeRetry(time.Hour)
close(done)
}()
time.Sleep(30 * time.Millisecond) // let the goroutine enter the wait
atomic.StoreUint32(&c.stop, 1)
select {
case <-done:
case <-time.After(time.Second):
t.Fatal("waitBeforeRetry did not return promptly after stop was set")
}
}

// TestSentinelRefreshRetryBackoff pins the two properties that matter when every
// sentinel is unreachable: refreshRetry must not spin, and it must exit once the
// client is closed rather than leaking a goroutine that retries forever.
func TestSentinelRefreshRetryBackoff(t *testing.T) {
defer ShouldNotLeak(SetupLeakDetection())

var failing atomic.Bool
var attempts atomic.Int64

s0 := &mockConn{
DoFn: func(cmd Completed) RedisResult { return RedisResult{} },
DoMultiFn: func(multi ...Completed) *redisresults {
if failing.Load() {
attempts.Add(1)
return &redisresults{s: []RedisResult{
NewErrorResult(ErrClosing), NewErrorResult(ErrClosing),
}}
}
return &redisresults{s: []RedisResult{
{val: slicemsg('*', []RedisMessage{})},
{val: slicemsg('*', []RedisMessage{strmsg('+', ""), strmsg('+', "1")})},
}}
},
}
m := &mockConn{
DoFn: func(cmd Completed) RedisResult {
return RedisResult{val: slicemsg('*', []RedisMessage{strmsg('+', "master")})}
},
}
client, err := newSentinelClient(
&ClientOption{InitAddress: []string{":0"}},
func(dst string, opt *ClientOption) conn {
switch dst {
case ":0":
return s0
case ":1":
return m
}
return nil
},
newRetryer(defaultRetryDelayFn),
)
if err != nil {
t.Fatalf("unexpected err %v", err)
}

failing.Store(true)
done := make(chan struct{})
go func() {
client.refreshRetry()
close(done)
}()

time.Sleep(200 * time.Millisecond)
// Without backoff this loop managed hundreds of thousands of iterations in
// this window; with it, a few dozen at most.
if n := attempts.Load(); n > 500 {
t.Fatalf("refreshRetry spun %d times in 200ms — backoff not applied", n)
}

client.Close()
select {
case <-done:
case <-time.After(5 * time.Second):
t.Fatal("refreshRetry did not return after Close — goroutine leaked")
}
}
Loading