diff --git a/core/src/main/java/io/grpc/internal/PickFirstLoadBalancer.java b/core/src/main/java/io/grpc/internal/PickFirstLoadBalancer.java index 705cfdbc1dc..b92cc2a48a4 100644 --- a/core/src/main/java/io/grpc/internal/PickFirstLoadBalancer.java +++ b/core/src/main/java/io/grpc/internal/PickFirstLoadBalancer.java @@ -107,6 +107,9 @@ public void handleNameResolutionError(Status error) { } private void processSubchannelState(Subchannel subchannel, ConnectivityStateInfo stateInfo) { + if (subchannel != this.subchannel) { + return; + } ConnectivityState newState = stateInfo.getState(); if (newState == SHUTDOWN) { return; @@ -161,6 +164,7 @@ private void updateBalancingState(ConnectivityState state, SubchannelPicker pick public void shutdown() { if (subchannel != null) { subchannel.shutdown(); + subchannel = null; } } diff --git a/core/src/test/java/io/grpc/internal/PickFirstLoadBalancerTest.java b/core/src/test/java/io/grpc/internal/PickFirstLoadBalancerTest.java index 1c3182237d3..ab6af82ade2 100644 --- a/core/src/test/java/io/grpc/internal/PickFirstLoadBalancerTest.java +++ b/core/src/test/java/io/grpc/internal/PickFirstLoadBalancerTest.java @@ -591,6 +591,29 @@ public void requestConnection() { verify(mockSubchannel, times(2)).requestConnection(); } + @Test + public void ignoreStaleSubchannelStateChange() throws Exception { + loadBalancer.acceptResolvedAddresses( + ResolvedAddresses.newBuilder().setAddresses(servers).setAttributes(affinity).build()); + verify(mockSubchannel).start(stateListenerCaptor.capture()); + SubchannelStateListener oldListener = stateListenerCaptor.getValue(); + + // Name resolution error occurs, shutting down old subchannel and resetting subchannel reference + loadBalancer.handleNameResolutionError(Status.UNAVAILABLE); + + // New resolution result arrives, creating a new subchannel + Subchannel newSubchannel = org.mockito.Mockito.mock(Subchannel.class); + when(mockHelper.createSubchannel(any())).thenReturn(newSubchannel); + loadBalancer.acceptResolvedAddresses( + ResolvedAddresses.newBuilder().setAddresses(servers).setAttributes(affinity).build()); + + // Old subchannel (e.g. during delayed shutdown) fires READY state update + oldListener.onSubchannelState(ConnectivityStateInfo.forNonError(READY)); + + // Verify that the LB state was NOT updated to READY with the old subchannel + verify(mockHelper, never()).updateBalancingState(eq(READY), any(SubchannelPicker.class)); + } + private static class FakeSocketAddress extends SocketAddress { final String name;