Skip to content

Commit c38ce1d

Browse files
committed
fix: tests
1 parent b50ca84 commit c38ce1d

12 files changed

Lines changed: 93 additions & 65 deletions

core/src/main/java/io/grpc/internal/DelayedClientTransport.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -367,7 +367,8 @@ private class PendingStream extends DelayedStream {
367367
private volatile Status lastPickStatus;
368368
@Nullable private String delayReasonToken;
369369

370-
private PendingStream(PickSubchannelArgs args, ClientStreamTracer[] tracers, @Nullable String initialToken) {
370+
private PendingStream(PickSubchannelArgs args, ClientStreamTracer[] tracers,
371+
@Nullable String initialToken) {
371372
super("connecting_and_lb");
372373
this.args = args;
373374
this.tracers = tracers;

core/src/main/java/io/grpc/internal/ForwardingClientStreamTracer.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,16 @@ public void createPendingStream() {
3939
delegate().createPendingStream();
4040
}
4141

42+
@Override
43+
public void delayStarted(String reasonToken) {
44+
delegate().delayStarted(reasonToken);
45+
}
46+
47+
@Override
48+
public void delayEnded() {
49+
delegate().delayEnded();
50+
}
51+
4252
@Override
4353
public void outboundHeaders() {
4454
delegate().outboundHeaders();

core/src/main/java/io/grpc/internal/PickFirstLoadBalancer.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -38,7 +38,8 @@
3838
* list and sticking to the first that works.
3939
*/
4040
final class PickFirstLoadBalancer extends LoadBalancer {
41-
private static final PickResult CONNECTING_RESULT = PickResult.withNoResult("pick_first:connecting");
41+
private static final PickResult CONNECTING_RESULT =
42+
PickResult.withNoResult("pick_first:connecting");
4243
private final Helper helper;
4344
private Subchannel subchannel;
4445
private ConnectivityState currentState = IDLE;

core/src/test/java/io/grpc/internal/DelayedClientTransportTest.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -782,7 +782,7 @@ public void streamDelayMetrics() {
782782
.thenReturn(PickResult.withNoResult("pick_first:connecting"));
783783

784784
delayedTransport.reprocess(connectingPicker);
785-
ClientStream stream = delayedTransport.newStream(method, headers, callOptions, customTracers);
785+
delayedTransport.newStream(method, headers, callOptions, customTracers);
786786

787787
InOrder inOrder = inOrder(mockTracer);
788788
inOrder.verify(mockTracer).delayStarted("pick_first:connecting");
@@ -812,6 +812,7 @@ public void streamDelayMetrics_cancelled() {
812812

813813
delayedTransport.reprocess(connectingPicker);
814814
ClientStream stream = delayedTransport.newStream(method, headers, callOptions, customTracers);
815+
stream.start(streamListener);
815816

816817
verify(mockTracer).delayStarted("pick_first:connecting");
817818

util/src/main/java/io/grpc/util/ForwardingClientStreamTracer.java

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,16 @@ public void createPendingStream() {
3838
delegate().createPendingStream();
3939
}
4040

41+
@Override
42+
public void delayStarted(String reasonToken) {
43+
delegate().delayStarted(reasonToken);
44+
}
45+
46+
@Override
47+
public void delayEnded() {
48+
delegate().delayEnded();
49+
}
50+
4151
@Override
4252
public void outboundHeaders() {
4353
delegate().outboundHeaders();

util/src/main/java/io/grpc/util/RoundRobinLoadBalancer.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -41,7 +41,8 @@
4141
* EquivalentAddressGroup}s from the {@link NameResolver}.
4242
*/
4343
final class RoundRobinLoadBalancer extends MultiChildLoadBalancer {
44-
private static final PickResult CONNECTING_RESULT = PickResult.withNoResult("round_robin:connecting");
44+
private static final PickResult CONNECTING_RESULT =
45+
PickResult.withNoResult("round_robin:connecting");
4546
private final AtomicInteger sequence = new AtomicInteger(new Random().nextInt());
4647
private SubchannelPicker currentPicker = new FixedResultPicker(CONNECTING_RESULT);
4748

xds/src/main/java/io/grpc/xds/CdsLoadBalancer2.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@
1717
package io.grpc.xds;
1818

1919
import static com.google.common.base.Preconditions.checkNotNull;
20+
import static io.grpc.ConnectivityState.CONNECTING;
2021
import static io.grpc.ConnectivityState.TRANSIENT_FAILURE;
2122
import static io.grpc.xds.XdsLbPolicies.CDS_POLICY_NAME;
2223
import static io.grpc.xds.XdsLbPolicies.PRIORITY_POLICY_NAME;
@@ -119,7 +120,8 @@ public Status acceptResolvedAddresses(ResolvedAddresses resolvedAddresses) {
119120
errorPrefix() + "Unable to find non-dynamic cluster"));
120121
}
121122
// The dynamic cluster must not have loaded yet
122-
helper.updateBalancingState(CONNECTING, new FixedResultPicker(PickResult.withNoResult("cds:discovery_pending")));
123+
helper.updateBalancingState(
124+
CONNECTING, new FixedResultPicker(PickResult.withNoResult("cds:discovery_pending")));
123125
return Status.OK;
124126
}
125127
if (!clusterConfigOr.hasValue()) {

xds/src/main/java/io/grpc/xds/PriorityLoadBalancer.java

Lines changed: 42 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -322,17 +322,11 @@ public void updateBalancingState(final ConnectivityState newState,
322322
}
323323
ConnectivityState oldState = connectivityState;
324324
connectivityState = newState;
325-
picker = new SubchannelPicker() {
326-
@Override
327-
public PickResult pickSubchannel(PickSubchannelArgs args) {
328-
PickResult childResult = newPicker.pickSubchannel(args);
329-
if (!childResult.hasResult() && childResult.getDelayReasonToken() != null) {
330-
return PickResult.withNoResult(
331-
"priority_" + priority + ":" + childResult.getDelayReasonToken());
332-
}
333-
return childResult;
334-
}
335-
};
325+
if (newState == CONNECTING || newState == IDLE) {
326+
picker = new PriorityPicker(newPicker, priority);
327+
} else {
328+
picker = newPicker;
329+
}
336330

337331
if (deletionTimer != null && deletionTimer.isPending()) {
338332
return;
@@ -367,4 +361,41 @@ protected Helper delegate() {
367361
}
368362
}
369363
}
364+
365+
private static final class PriorityPicker extends SubchannelPicker {
366+
private final SubchannelPicker delegate;
367+
private final String priority;
368+
369+
PriorityPicker(SubchannelPicker delegate, String priority) {
370+
this.delegate = com.google.common.base.Preconditions.checkNotNull(delegate, "delegate");
371+
this.priority = com.google.common.base.Preconditions.checkNotNull(priority, "priority");
372+
}
373+
374+
@Override
375+
public PickResult pickSubchannel(PickSubchannelArgs args) {
376+
PickResult childResult = delegate.pickSubchannel(args);
377+
if (!childResult.hasResult() && childResult.getDelayReasonToken() != null) {
378+
return PickResult.withNoResult(
379+
"priority_" + priority + ":" + childResult.getDelayReasonToken());
380+
}
381+
return childResult;
382+
}
383+
384+
@Override
385+
public boolean equals(Object o) {
386+
if (this == o) {
387+
return true;
388+
}
389+
if (o == null || getClass() != o.getClass()) {
390+
return false;
391+
}
392+
PriorityPicker that = (PriorityPicker) o;
393+
return delegate.equals(that.delegate) && priority.equals(that.priority);
394+
}
395+
396+
@Override
397+
public int hashCode() {
398+
return java.util.Objects.hash(delegate, priority);
399+
}
400+
}
370401
}

xds/src/main/java/io/grpc/xds/RingHashLoadBalancer.java

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -465,7 +465,8 @@ public PickResult pickSubchannel(PickSubchannelArgs args) {
465465
}
466466
});
467467

468-
return RING_HASH_CONNECTING_RESULT; // Indicates that this should be retried after backoff
468+
// Indicates that this should be retried after backoff
469+
return RING_HASH_CONNECTING_RESULT;
469470
}
470471
}
471472
} else {

xds/src/test/java/io/grpc/xds/CdsLoadBalancer2Test.java

Lines changed: 11 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -348,19 +348,25 @@ public void discoverDynamicCluster_pending_emitsToken() {
348348
String clusterName = "cluster2";
349349
CdsConfig cdsConfig = new CdsConfig(clusterName, /*dynamic=*/ true);
350350

351-
XdsConfig mockXdsConfig = mock(XdsConfig.class);
352-
when(mockXdsConfig.getClusters()).thenReturn(ImmutableMap.of());
351+
XdsConfig xdsConfig = new XdsConfig(null, null, null, ImmutableMap.of());
353352

354353
loadBalancer.acceptResolvedAddresses(ResolvedAddresses.newBuilder()
355354
.setAddresses(Collections.emptyList())
356355
.setAttributes(Attributes.newBuilder()
357-
.set(XdsAttributes.XDS_CONFIG, mockXdsConfig)
358-
.set(XdsAttributes.XDS_CLUSTER_SUBSCRIPT_REGISTRY, xdsDepManager)
356+
.set(XdsAttributes.XDS_CONFIG, xdsConfig)
357+
.set(
358+
XdsAttributes.XDS_CLUSTER_SUBSCRIPT_REGISTRY,
359+
new XdsConfig.XdsClusterSubscriptionRegistry() {
360+
@Override
361+
public XdsConfig.Subscription subscribeToCluster(String clusterName) {
362+
return mock(XdsConfig.Subscription.class);
363+
}
364+
})
359365
.build())
360366
.setLoadBalancingPolicyConfig(cdsConfig)
361367
.build());
362368

363-
verify(helper).updateBalancingState(eq(CONNECTING), pickerCaptor.capture());
369+
verify(helper).updateBalancingState(eq(ConnectivityState.CONNECTING), pickerCaptor.capture());
364370
PickResult result = pickerCaptor.getValue().pickSubchannel(mock(PickSubchannelArgs.class));
365371
assertThat(result.getDelayReasonToken()).isEqualTo("cds:discovery_pending");
366372
}

0 commit comments

Comments
 (0)