Skip to content
Merged
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
35 changes: 23 additions & 12 deletions balancer.go
Original file line number Diff line number Diff line change
Expand Up @@ -185,7 +185,8 @@ type ringBalancer struct {
scStates map[balancer.SubConn]connectivity.State
// scKeys holds the hashring key of every SubConn the resolver gave us.
scKeys map[balancer.SubConn]string
// ringMembers is the set of SubConns on the hashring: only READY ones.
// ringMembers is the set of SubConns on the hashring: every SubConn the
// resolver gave us except those in TRANSIENT_FAILURE.
ringMembers map[balancer.SubConn]struct{}

config *BalancerConfig
Expand Down Expand Up @@ -234,10 +235,11 @@ func (b *ringBalancer) UpdateClientConnState(s balancer.ClientConnState) error {
svcConfig := s.BalancerConfig.(*BalancerConfig)
if b.config == nil || svcConfig.ReplicationFactor != b.config.ReplicationFactor {
b.hashring = hashring.MustNew(b.hasher, svcConfig.ReplicationFactor)
// The new ring starts empty: put every READY SubConn back on it.
// The new ring starts empty: put every SubConn that has not
// failed back on it.
b.ringMembers = make(map[balancer.SubConn]struct{})
for sc, st := range b.scStates {
if st != connectivity.Ready {
if st == connectivity.TransientFailure {
continue
}
if err := b.addToRing(sc); err != nil {
Expand Down Expand Up @@ -282,6 +284,11 @@ func (b *ringBalancer) UpdateClientConnState(s balancer.ClientConnState) error {
b.scKeys[sc] = addr.ServerName + addr.Addr
b.csEvltr.RecordTransition(connectivity.Shutdown, connectivity.Idle)
sc.Connect()
// The backend owns its keys from the start: its connection is
// already on the way, and RPCs picked before READY wait for it.
if err := b.addToRing(sc); err != nil {
return err
}
}
}

Expand Down Expand Up @@ -399,13 +406,16 @@ func (b *ringBalancer) UpdateSubConnState(sc balancer.SubConn, state balancer.Su

b.scStates[sc] = s

// A backend only receives traffic while its connection is READY. Keys
// hashed to a backend that is connecting or failed move to the closest
// ready member instead of waiting on a connection that may never come.
// Ring membership changes on failure, not on every state change. IDLE
// and CONNECTING backends keep their keys: those states normally resolve
// quickly, and a remap on each one would move keys across the fleet on
// every connection recycle. A backend whose connection attempt failed
// leaves the ring, and its keys move to the closest remaining member.
var ringErr error
if s == connectivity.Ready {
switch s {
case connectivity.Ready:
ringErr = b.addToRing(sc)
} else {
case connectivity.TransientFailure, connectivity.Shutdown:
ringErr = b.removeFromRing(sc)
}
if ringErr != nil {
Expand Down Expand Up @@ -453,10 +463,11 @@ var _ balancer.Picker = (*picker)(nil)
// The value stored in CtxKey is hashed into the hashring, and the resulting
// subconnection is used.
//
// The hashring only contains READY subconnections, so a key whose closest
// backend is down is served by the next closest ready backend. When no
// backend is ready the RPC is queued until one becomes ready (or the
// balancer reports TRANSIENT_FAILURE and installs an error picker).
// The hashring contains every resolver-provided backend except those in
// TRANSIENT_FAILURE. The next closest member serves the keys of a failed
// backend. A key whose backend is IDLE or CONNECTING waits for that
// connection. The dial timeout (MinConnectTimeout) bounds the wait: when
// the dial fails, the backend leaves the ring and its keys move.
//
// Spread can be increased to be robust against single node availability
// problems. If spread is greater than 1, a random selection is made from the
Expand Down
43 changes: 28 additions & 15 deletions balancer_readiness_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
"github.com/cespare/xxhash/v2"
"github.com/stretchr/testify/require"
"google.golang.org/grpc"
"google.golang.org/grpc/backoff"
"google.golang.org/grpc/balancer"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/connectivity"
Expand Down Expand Up @@ -39,10 +40,10 @@ func keyHashingTo(t *testing.T, ring *hashring.Ring, memberKey string) []byte {
return nil
}

// TestPickerSkipsSubConnsThatAreNotReady drives the balancer directly: two
// backends come up, then one fails. Keys that hashed to the failed backend
// must be picked on the remaining ready one instead of a dead SubConn.
func TestPickerSkipsSubConnsThatAreNotReady(t *testing.T) {
// TestPickerSkipsFailedSubConns drives the balancer directly: two backends
// become READY, then one fails. Keys that hashed to the failed backend must
// move to the remaining ready one instead of a dead SubConn.
func TestPickerSkipsFailedSubConns(t *testing.T) {
cc := newFakeClientConn()
// Drain state updates so the balancer never blocks on the fake conn.
go func() {
Expand Down Expand Up @@ -96,11 +97,14 @@ func TestPickerSkipsSubConnsThatAreNotReady(t *testing.T) {
require.Same(t, sc2, pick(keyFor2), "key returns to backend 2 once it is ready")
}

// TestRPCsAreNotStuckBehindAConnectingBackend reproduces a peer whose address
// is still resolvable but never completes a connection (a killed pod whose IP
// is still in the endpoint list). RPCs hashed to it must be served by the
// healthy backend instead of waiting for the dial to time out.
func TestRPCsAreNotStuckBehindAConnectingBackend(t *testing.T) {
// TestRPCsFailOverAfterConnectTimeout reproduces a peer whose address is
// still resolvable but never completes a connection (a killed pod whose IP
// is still in the endpoint list). The backend keeps its keys while it is
// CONNECTING. When the dial times out (MinConnectTimeout), the backend
// moves to TRANSIENT_FAILURE and leaves the ring, and the healthy backend
// serves its keys. The failover window is therefore bounded by
// MinConnectTimeout, which the dialer can set.
func TestRPCsFailOverAfterConnectTimeout(t *testing.T) {
// Healthy backend: a real gRPC server with the health service.
healthyLis, err := net.Listen("tcp", "127.0.0.1:0")
require.NoError(t, err)
Expand All @@ -110,7 +114,7 @@ func TestRPCsAreNotStuckBehindAConnectingBackend(t *testing.T) {
t.Cleanup(srv.Stop)

// Black hole: accepts TCP connections but never speaks HTTP/2, so the
// SubConn stays CONNECTING until gRPC's 20s connect timeout.
// SubConn stays CONNECTING until the connect timeout set below.
blackholeLis, err := net.Listen("tcp", "127.0.0.1:0")
require.NoError(t, err)
t.Cleanup(func() { _ = blackholeLis.Close() })
Expand All @@ -135,6 +139,10 @@ func TestRPCsAreNotStuckBehindAConnectingBackend(t *testing.T) {
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithResolvers(rb),
grpc.WithDefaultServiceConfig(svcConfig),
grpc.WithConnectParams(grpc.ConnectParams{
MinConnectTimeout: 300 * time.Millisecond,
Backoff: backoff.DefaultConfig,
}),
)
require.NoError(t, err)
t.Cleanup(func() { _ = conn.Close() })
Expand All @@ -158,10 +166,12 @@ func TestRPCsAreNotStuckBehindAConnectingBackend(t *testing.T) {
// Warm up: the healthy backend is reachable and connected.
require.NoError(t, check(keyForHealthy, 5*time.Second))

// A key hashed to the black-holed backend must still be answered promptly.
err = check(keyForBlackhole, 3*time.Second)
// A key hashed to the black-holed backend waits for the dial, fails over
// when the dial times out at ~300ms, and succeeds well inside the
// deadline.
err = check(keyForBlackhole, 5*time.Second)
if st, ok := status.FromError(err); ok && st.Code() == codes.DeadlineExceeded {
t.Fatalf("RPC hung behind a CONNECTING backend instead of being served by the ready one: %v", err)
t.Fatalf("RPC hung behind a black-holed backend instead of failing over after the connect timeout: %v", err)
}
require.NoError(t, err)
}
Expand Down Expand Up @@ -189,15 +199,18 @@ func ringKeys(b *ringBalancer) []string {
return keys(b.hashring.Members())
}

func TestConnectingSubConnLeavesTheRing(t *testing.T) {
// A reconnecting backend keeps its ring membership. CONNECTING normally
// resolves in milliseconds, and a remap on every reconnect would move the
// member's keys across the fleet each time.
func TestConnectingSubConnKeepsRingMembership(t *testing.T) {
addr := resolver.Address{ServerName: "t", Addr: "1"}
b, _ := readyBalancer(t, addr)
sci, _ := b.subConns.Get(addr)
sc := sci.(balancer.SubConn)
require.Equal(t, []string{"t1"}, ringKeys(b))

b.UpdateSubConnState(sc, balancer.SubConnState{ConnectivityState: connectivity.Connecting})
require.Empty(t, ringKeys(b), "a reconnecting backend must not receive traffic")
require.Equal(t, []string{"t1"}, ringKeys(b), "a reconnecting backend must keep its keys")
require.Equal(t, connectivity.Connecting, b.state)
}

Expand Down
78 changes: 78 additions & 0 deletions membership_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
package consistent

import (
"fmt"
"testing"

"github.com/cespare/xxhash/v2"
"github.com/stretchr/testify/require"
"google.golang.org/grpc/balancer"
"google.golang.org/grpc/connectivity"
"google.golang.org/grpc/resolver"
)

// A SubConn joins the ring when the resolver provides it, before its
// connection is READY. IDLE and CONNECTING backends keep their keyspace.
func TestNewSubConnJoinsRingBeforeReady(t *testing.T) {
cc := newFakeClientConn()
go func() {
for range cc.stateCh {
}
}()

b := NewBuilder(xxhash.Sum64).Build(cc, balancer.BuildOptions{}).(*ringBalancer)
require.NoError(t, b.UpdateClientConnState(balancer.ClientConnState{
ResolverState: resolver.State{Addresses: []resolver.Address{{ServerName: "t", Addr: "1"}}},
BalancerConfig: &BalancerConfig{ReplicationFactor: 100, Spread: 1},
}))

require.Equal(t, []string{"t1"}, ringKeys(b),
"a resolver-provided backend must own its keys before it is READY")
}

// A READY SubConn that moves to IDLE keeps its ring membership and
// reconnects. Servers that limit connection age send GOAWAY on an interval,
// and that produces exactly this transition on a healthy backend.
func TestIdleSubConnKeepsRingMembership(t *testing.T) {
addr := resolver.Address{ServerName: "t", Addr: "1"}
b, _ := readyBalancer(t, addr)
sci, _ := b.subConns.Get(addr)
sc := sci.(*fakeSubConn)
connectsSoFar := sc.connectCalls

b.UpdateSubConnState(sc, balancer.SubConnState{ConnectivityState: connectivity.Idle})

require.Equal(t, []string{"t1"}, ringKeys(b),
"a healthy backend must not lose its keys on a connection recycle")
require.Equal(t, connectsSoFar+1, sc.connectCalls, "IDLE must trigger a reconnect")
}

// A ReplicationFactor change rebuilds the ring. The new ring must contain
// every member except those in TRANSIENT_FAILURE.
func TestReplicationFactorChangeKeepsNonFailedMembers(t *testing.T) {
addrs := []resolver.Address{
{ServerName: "t", Addr: "1"},
{ServerName: "t", Addr: "2"},
{ServerName: "t", Addr: "3"},
}
b, _ := readyBalancer(t, addrs...)

// t2 fails, t3 reconnects: only the failure may cost ring membership.
sci2, _ := b.subConns.Get(addrs[1])
b.UpdateSubConnState(sci2.(balancer.SubConn), balancer.SubConnState{
ConnectivityState: connectivity.TransientFailure,
ConnectionError: fmt.Errorf("refused"),
})
sci3, _ := b.subConns.Get(addrs[2])
b.UpdateSubConnState(sci3.(balancer.SubConn), balancer.SubConnState{
ConnectivityState: connectivity.Connecting,
})
require.ElementsMatch(t, []string{"t1", "t3"}, ringKeys(b))

require.NoError(t, b.UpdateClientConnState(balancer.ClientConnState{
ResolverState: resolver.State{Addresses: addrs},
BalancerConfig: &BalancerConfig{ReplicationFactor: 7, Spread: 1},
}))
require.ElementsMatch(t, []string{"t1", "t3"}, ringKeys(b),
"the rebuilt ring must keep the reconnecting member and exclude the failed one")
}
Loading