diff --git a/balancer.go b/balancer.go index 23c7e06..97cb5a2 100644 --- a/balancer.go +++ b/balancer.go @@ -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 @@ -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 { @@ -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 + } } } @@ -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 { @@ -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 diff --git a/balancer_readiness_test.go b/balancer_readiness_test.go index c11fc66..c2ba824 100644 --- a/balancer_readiness_test.go +++ b/balancer_readiness_test.go @@ -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" @@ -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() { @@ -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) @@ -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() }) @@ -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() }) @@ -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) } @@ -189,7 +199,10 @@ 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) @@ -197,7 +210,7 @@ func TestConnectingSubConnLeavesTheRing(t *testing.T) { 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) } diff --git a/membership_test.go b/membership_test.go new file mode 100644 index 0000000..25132cc --- /dev/null +++ b/membership_test.go @@ -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") +}