Skip to content

Commit c52c2aa

Browse files
committed
Address review comments on PR grpc#12807: inline CONNECTING_RESULT in RoundRobinLoadBalancer and refine PendingStream delay telemetry synchronization
1 parent 1a280e9 commit c52c2aa

3 files changed

Lines changed: 45 additions & 8 deletions

File tree

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

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -428,6 +428,9 @@ private PendingStream(PickSubchannelArgs args, ClientStreamTracer[] tracers,
428428
* structured transition event is appended to the active span without span re-creation.
429429
*/
430430
synchronized void updateDelay(@Nullable String newType, @Nullable String newReason) {
431+
if (getRealStream() != null) {
432+
return;
433+
}
431434
if (!Objects.equals(activeDelayType, newType)) {
432435
// Delay type changed (e.g., from RLS lookup to connecting). End the previous delay.
433436
if (activeDelayType != null) {
@@ -466,7 +469,12 @@ synchronized void endDelay() {
466469
}
467470

468471
Runnable setStreamAndEndDelay(ClientStream stream) {
469-
endDelay();
472+
synchronized (this) {
473+
if (getRealStream() != null) {
474+
return null;
475+
}
476+
endDelay();
477+
}
470478
return setStream(stream);
471479
}
472480

@@ -493,7 +501,6 @@ private Runnable createRealStream(ClientTransport transport, String authorityOve
493501

494502
@Override
495503
public void cancel(Status reason) {
496-
endDelay();
497504
super.cancel(reason);
498505
synchronized (lock) {
499506
if (reportTransportTerminated != null) {
@@ -512,6 +519,7 @@ public void cancel(Status reason) {
512519

513520
@Override
514521
protected void onEarlyCancellation(Status reason) {
522+
endDelay();
515523
for (ClientStreamTracer tracer : tracers) {
516524
tracer.streamClosed(reason);
517525
}

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

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -939,6 +939,39 @@ public PickResult pickSubchannel(PickSubchannelArgs args) {
939939
};
940940
}
941941

942+
@Test
943+
public void streamDelayMetrics_updateDelayAfterCancellation_isNoOp() {
944+
final FakeStreamTracer fakeTracer = new FakeStreamTracer();
945+
ClientStreamTracer[] customTracers = new ClientStreamTracer[] { fakeTracer };
946+
947+
// Create a pending stream
948+
delayedTransport.reprocess(fakePicker(
949+
PickResult.withNoResult("connecting", "pick_first: attempting to connect")));
950+
final ClientStream stream = delayedTransport.newStream(
951+
method, headers, callOptions, customTracers);
952+
stream.start(streamListener);
953+
954+
assertEquals(Collections.singletonList("connecting"), fakeTracer.startedDelayTypes);
955+
assertEquals(0, fakeTracer.delayEndedCount);
956+
957+
// Reprocess with a custom picker that cancels the stream during the pick.
958+
// This exactly simulates the concurrent timing race.
959+
SubchannelPicker racePicker = new SubchannelPicker() {
960+
@Override
961+
public PickResult pickSubchannel(PickSubchannelArgs args) {
962+
stream.cancel(Status.CANCELLED);
963+
return PickResult.withNoResult("connecting", "retry-backoff");
964+
}
965+
};
966+
967+
delayedTransport.reprocess(racePicker);
968+
969+
// VERIFICATION:
970+
// The updateDelay called after the cancellation must be ignored because of our check.
971+
assertEquals(Collections.singletonList("connecting"), fakeTracer.startedDelayTypes);
972+
assertEquals(1, fakeTracer.delayEndedCount);
973+
}
974+
942975
private static TransportProvider newTransportProvider(final ClientTransport transport) {
943976
return new TransportProvider() {
944977
@Override

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

Lines changed: 2 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -44,16 +44,12 @@ final class RoundRobinLoadBalancer extends MultiChildLoadBalancer {
4444
private static final PickResult CONNECTING_RESULT = PickResult.withNoResult("connecting",
4545
"round_robin connecting: TCP/TLS handshake in progress to child balancers");
4646
private final AtomicInteger sequence = new AtomicInteger(new Random().nextInt());
47-
private SubchannelPicker currentPicker = new FixedResultPicker(connectingResult());
47+
private SubchannelPicker currentPicker = new FixedResultPicker(CONNECTING_RESULT);
4848

4949
public RoundRobinLoadBalancer(Helper helper) {
5050
super(helper);
5151
}
5252

53-
private PickResult connectingResult() {
54-
return CONNECTING_RESULT;
55-
}
56-
5753
/**
5854
* Updates picker with the list of active subchannels (state == READY).
5955
*/
@@ -74,7 +70,7 @@ protected void updateOverallBalancingState() {
7470
}
7571

7672
if (isConnecting) {
77-
updateBalancingState(CONNECTING, new FixedResultPicker(connectingResult()));
73+
updateBalancingState(CONNECTING, new FixedResultPicker(CONNECTING_RESULT));
7874
} else {
7975
updateBalancingState(TRANSIENT_FAILURE, createReadyPicker(getChildLbStates()));
8076
}

0 commit comments

Comments
 (0)