Skip to content

Commit b50ca84

Browse files
committed
core,api,xds: Implement load balancing policy delay plumbing
This commit implements the plumbing required to propagate delay reason tokens from load balancing policies up to the transport layer and tracers, as specified in the LB policy delay design.
1 parent c62cdef commit b50ca84

16 files changed

Lines changed: 256 additions & 22 deletions

File tree

api/src/main/java/io/grpc/ClientStreamTracer.java

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -57,6 +57,23 @@ public void streamCreated(@Grpc.TransportAttr Attributes transportAttrs, Metadat
5757
public void createPendingStream() {
5858
}
5959

60+
/**
61+
* A delay segment started with a specific reason during load balancing.
62+
*
63+
* @param reasonToken the reason for the delay, e.g., "pick_first:connecting"
64+
* @since 1.82.0
65+
*/
66+
public void delayStarted(String reasonToken) {
67+
}
68+
69+
/**
70+
* The current delay segment ended.
71+
*
72+
* @since 1.82.0
73+
*/
74+
public void delayEnded() {
75+
}
76+
6077
/**
6178
* Headers has been sent to the socket.
6279
*/

api/src/main/java/io/grpc/LoadBalancer.java

Lines changed: 26 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -549,25 +549,30 @@ public static final class PickResult {
549549
// True if the result is created by withDrop()
550550
private final boolean drop;
551551
@Nullable private final String authorityOverride;
552+
@Nullable private final String delayReasonToken;
552553

553554
private PickResult(
554555
@Nullable Subchannel subchannel, @Nullable ClientStreamTracer.Factory streamTracerFactory,
555556
Status status, boolean drop) {
556-
this.subchannel = subchannel;
557-
this.streamTracerFactory = streamTracerFactory;
558-
this.status = checkNotNull(status, "status");
559-
this.drop = drop;
560-
this.authorityOverride = null;
557+
this(subchannel, streamTracerFactory, status, drop, null, null);
561558
}
562559

563560
private PickResult(
564561
@Nullable Subchannel subchannel, @Nullable ClientStreamTracer.Factory streamTracerFactory,
565562
Status status, boolean drop, @Nullable String authorityOverride) {
563+
this(subchannel, streamTracerFactory, status, drop, authorityOverride, null);
564+
}
565+
566+
private PickResult(
567+
@Nullable Subchannel subchannel, @Nullable ClientStreamTracer.Factory streamTracerFactory,
568+
Status status, boolean drop, @Nullable String authorityOverride,
569+
@Nullable String delayReasonToken) {
566570
this.subchannel = subchannel;
567571
this.streamTracerFactory = streamTracerFactory;
568572
this.status = checkNotNull(status, "status");
569573
this.drop = drop;
570574
this.authorityOverride = authorityOverride;
575+
this.delayReasonToken = delayReasonToken;
571576
}
572577

573578
/**
@@ -727,6 +732,22 @@ public static PickResult withNoResult() {
727732
return NO_RESULT;
728733
}
729734

735+
/**
736+
* No decision could be made. The RPC will stay buffered with a specific reason.
737+
*
738+
* @since 1.82.0
739+
*/
740+
public static PickResult withNoResult(String delayReasonToken) {
741+
Preconditions.checkNotNull(delayReasonToken, "delayReasonToken");
742+
return new PickResult(null, null, Status.OK, false, null, delayReasonToken);
743+
}
744+
745+
/** Returns the delay reason token if any. */
746+
@Nullable
747+
public String getDelayReasonToken() {
748+
return delayReasonToken;
749+
}
750+
730751
/** Returns the authority override if any. */
731752
@ExperimentalApi("https://github.com/grpc/grpc-java/issues/11656")
732753
@Nullable

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

Lines changed: 42 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -157,7 +157,8 @@ public final ClientStream newStream(
157157
synchronized (lock) {
158158
PickerState newerState = pickerState;
159159
if (state == newerState) {
160-
return createPendingStream(args, tracers, pickResult);
160+
String token = pickResult != null ? pickResult.getDelayReasonToken() : null;
161+
return createPendingStream(args, tracers, pickResult, token);
161162
}
162163
state = newerState;
163164
}
@@ -173,8 +174,8 @@ public final ClientStream newStream(
173174
*/
174175
@GuardedBy("lock")
175176
private PendingStream createPendingStream(PickSubchannelArgs args, ClientStreamTracer[] tracers,
176-
PickResult pickResult) {
177-
PendingStream pendingStream = new PendingStream(args, tracers);
177+
PickResult pickResult, @Nullable String delayReasonToken) {
178+
PendingStream pendingStream = new PendingStream(args, tracers, delayReasonToken);
178179
if (args.getCallOptions().isWaitForReady() && pickResult != null && pickResult.hasResult()) {
179180
pendingStream.lastPickStatus = pickResult.getStatus();
180181
}
@@ -303,6 +304,7 @@ final void reprocess(@Nullable SubchannelPicker picker) {
303304
final ClientTransport transport = GrpcUtil.getTransportFromPickResult(pickResult,
304305
callOptions.isWaitForReady());
305306
if (transport != null) {
307+
stream.endDelay();
306308
Executor executor = defaultAppExecutor;
307309
// createRealStream may be expensive. It will start real streams on the transport. If
308310
// there are pending requests, they will be serialized too, which may be expensive. Since
@@ -315,7 +317,9 @@ final void reprocess(@Nullable SubchannelPicker picker) {
315317
executor.execute(runnable);
316318
}
317319
toRemove.add(stream);
318-
} // else: stay pending
320+
} else { // stay pending
321+
stream.updateDelayReason(pickResult.getDelayReasonToken());
322+
}
319323
}
320324

321325
synchronized (lock) {
@@ -361,11 +365,43 @@ private class PendingStream extends DelayedStream {
361365
private final Context context = Context.current();
362366
private final ClientStreamTracer[] tracers;
363367
private volatile Status lastPickStatus;
368+
@Nullable private String delayReasonToken;
364369

365-
private PendingStream(PickSubchannelArgs args, ClientStreamTracer[] tracers) {
370+
private PendingStream(PickSubchannelArgs args, ClientStreamTracer[] tracers, @Nullable String initialToken) {
366371
super("connecting_and_lb");
367372
this.args = args;
368373
this.tracers = tracers;
374+
this.delayReasonToken = initialToken;
375+
if (initialToken != null) {
376+
for (ClientStreamTracer tracer : tracers) {
377+
tracer.delayStarted(initialToken);
378+
}
379+
}
380+
}
381+
382+
void updateDelayReason(String newToken) {
383+
if (!java.util.Objects.equals(delayReasonToken, newToken)) {
384+
if (delayReasonToken != null) {
385+
for (ClientStreamTracer tracer : tracers) {
386+
tracer.delayEnded();
387+
}
388+
}
389+
delayReasonToken = newToken;
390+
if (newToken != null) {
391+
for (ClientStreamTracer tracer : tracers) {
392+
tracer.delayStarted(newToken);
393+
}
394+
}
395+
}
396+
}
397+
398+
void endDelay() {
399+
if (delayReasonToken != null) {
400+
for (ClientStreamTracer tracer : tracers) {
401+
tracer.delayEnded();
402+
}
403+
delayReasonToken = null;
404+
}
369405
}
370406

371407
/** Runnable may be null. */
@@ -391,6 +427,7 @@ private Runnable createRealStream(ClientTransport transport, String authorityOve
391427

392428
@Override
393429
public void cancel(Status reason) {
430+
endDelay();
394431
super.cancel(reason);
395432
synchronized (lock) {
396433
if (reportTransportTerminated != null) {

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -38,6 +38,7 @@
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");
4142
private final Helper helper;
4243
private Subchannel subchannel;
4344
private ConnectivityState currentState = IDLE;
@@ -83,7 +84,7 @@ public void onSubchannelState(ConnectivityStateInfo stateInfo) {
8384

8485
// The channel state does not get updated when doing name resolving today, so for the moment
8586
// let LB report CONNECTION and call subchannel.requestConnection() immediately.
86-
updateBalancingState(CONNECTING, new FixedResultPicker(PickResult.withNoResult()));
87+
updateBalancingState(CONNECTING, new FixedResultPicker(CONNECTING_RESULT));
8788
subchannel.requestConnection();
8889
} else {
8990
subchannel.updateAddresses(servers);
@@ -135,7 +136,7 @@ private void processSubchannelState(Subchannel subchannel, ConnectivityStateInfo
135136
case CONNECTING:
136137
// It's safe to use RequestConnectionPicker here, so when coming from IDLE we could leave
137138
// the current picker in-place. But ignoring the potential optimization is simpler.
138-
picker = new FixedResultPicker(PickResult.withNoResult());
139+
picker = new FixedResultPicker(CONNECTING_RESULT);
139140
break;
140141
case READY:
141142
picker = new FixedResultPicker(PickResult.withSubchannel(subchannel));

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

Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -772,6 +772,54 @@ public void pendingStream_appendTimeoutInsight_waitForReady_withLastPickFailure(
772772
+ " connecting_and_lb_delay=[0-9]+ns, was_still_waiting]");
773773
}
774774

775+
@Test
776+
public void streamDelayMetrics() {
777+
ClientStreamTracer mockTracer = mock(ClientStreamTracer.class);
778+
ClientStreamTracer[] customTracers = new ClientStreamTracer[] { mockTracer };
779+
780+
SubchannelPicker connectingPicker = mock(SubchannelPicker.class);
781+
when(connectingPicker.pickSubchannel(any(PickSubchannelArgs.class)))
782+
.thenReturn(PickResult.withNoResult("pick_first:connecting"));
783+
784+
delayedTransport.reprocess(connectingPicker);
785+
ClientStream stream = delayedTransport.newStream(method, headers, callOptions, customTracers);
786+
787+
InOrder inOrder = inOrder(mockTracer);
788+
inOrder.verify(mockTracer).delayStarted("pick_first:connecting");
789+
790+
SubchannelPicker customDelayPicker = mock(SubchannelPicker.class);
791+
when(customDelayPicker.pickSubchannel(any(PickSubchannelArgs.class)))
792+
.thenReturn(PickResult.withNoResult("rls:lookup_pending"));
793+
794+
delayedTransport.reprocess(customDelayPicker);
795+
796+
inOrder.verify(mockTracer).delayEnded();
797+
inOrder.verify(mockTracer).delayStarted("rls:lookup_pending");
798+
799+
delayedTransport.reprocess(mockPicker);
800+
801+
inOrder.verify(mockTracer).delayEnded();
802+
}
803+
804+
@Test
805+
public void streamDelayMetrics_cancelled() {
806+
ClientStreamTracer mockTracer = mock(ClientStreamTracer.class);
807+
ClientStreamTracer[] customTracers = new ClientStreamTracer[] { mockTracer };
808+
809+
SubchannelPicker connectingPicker = mock(SubchannelPicker.class);
810+
when(connectingPicker.pickSubchannel(any(PickSubchannelArgs.class)))
811+
.thenReturn(PickResult.withNoResult("pick_first:connecting"));
812+
813+
delayedTransport.reprocess(connectingPicker);
814+
ClientStream stream = delayedTransport.newStream(method, headers, callOptions, customTracers);
815+
816+
verify(mockTracer).delayStarted("pick_first:connecting");
817+
818+
stream.cancel(Status.CANCELLED);
819+
820+
verify(mockTracer).delayEnded();
821+
}
822+
775823
private static TransportProvider newTransportProvider(final ClientTransport transport) {
776824
return new TransportProvider() {
777825
@Override

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -147,8 +147,9 @@ public void pickAfterResolved() throws Exception {
147147
verify(mockSubchannel).requestConnection();
148148

149149
// Calling pickSubchannel() twice gave the same result
150-
assertEquals(pickerCaptor.getValue().pickSubchannel(mockArgs),
151-
pickerCaptor.getValue().pickSubchannel(mockArgs));
150+
PickResult result = pickerCaptor.getValue().pickSubchannel(mockArgs);
151+
assertThat(result.getDelayReasonToken()).isEqualTo("pick_first:connecting");
152+
assertEquals(result, pickerCaptor.getValue().pickSubchannel(mockArgs));
152153

153154
verifyNoMoreInteractions(mockHelper);
154155
}

rls/src/main/java/io/grpc/rls/CachingRlsLbClient.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1050,7 +1050,7 @@ public PickResult pickSubchannel(PickSubchannelArgs args) {
10501050
convertRlsServerStatus(response.getStatus(),
10511051
lbPolicyConfig.getRouteLookupConfig().lookupService()));
10521052
} else {
1053-
return PickResult.withNoResult();
1053+
return PickResult.withNoResult("rls:lookup_pending");
10541054
}
10551055
}
10561056

rls/src/test/java/io/grpc/rls/RlsLoadBalancerTest.java

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -262,6 +262,7 @@ public void lb_working_withDefaultTarget_rlsResponding() throws Exception {
262262
PickResult res = picker.pickSubchannel(searchSubchannelArgs);
263263
assertThat(res.getStatus().isOk()).isTrue();
264264
assertThat(res.getSubchannel()).isNull();
265+
assertThat(res.getDelayReasonToken()).isEqualTo("rls:lookup_pending");
265266
// Cache is warm, but still unconnected
266267
res = picker.pickSubchannel(searchSubchannelArgs);
267268
inOrder.verify(helper).createSubchannel(any(CreateSubchannelArgs.class));
@@ -493,6 +494,7 @@ public void lb_working_withoutDefaultTarget() throws Exception {
493494
PickResult res = picker.pickSubchannel(searchSubchannelArgs);
494495
assertThat(res.getStatus().isOk()).isTrue();
495496
assertThat(res.getSubchannel()).isNull();
497+
assertThat(res.getDelayReasonToken()).isEqualTo("rls:lookup_pending");
496498
// Cache is warm, but still unconnected
497499
res = picker.pickSubchannel(searchSubchannelArgs);
498500
inOrder.verify(helper).createSubchannel(any(CreateSubchannelArgs.class));

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

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -41,8 +41,9 @@
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");
4445
private final AtomicInteger sequence = new AtomicInteger(new Random().nextInt());
45-
private SubchannelPicker currentPicker = new FixedResultPicker(PickResult.withNoResult());
46+
private SubchannelPicker currentPicker = new FixedResultPicker(CONNECTING_RESULT);
4647

4748
public RoundRobinLoadBalancer(Helper helper) {
4849
super(helper);
@@ -68,7 +69,7 @@ protected void updateOverallBalancingState() {
6869
}
6970

7071
if (isConnecting) {
71-
updateBalancingState(CONNECTING, new FixedResultPicker(PickResult.withNoResult()));
72+
updateBalancingState(CONNECTING, new FixedResultPicker(CONNECTING_RESULT));
7273
} else {
7374
updateBalancingState(TRANSIENT_FAILURE, createReadyPicker(getChildLbStates()));
7475
}

util/src/test/java/io/grpc/util/RoundRobinLoadBalancerTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -86,7 +86,7 @@
8686
public class RoundRobinLoadBalancerTest {
8787
private static final Attributes.Key<String> MAJOR_KEY = Attributes.Key.create("major-key");
8888
private static final SubchannelPicker EMPTY_PICKER =
89-
new FixedResultPicker(PickResult.withNoResult());
89+
new FixedResultPicker(PickResult.withNoResult("round_robin:connecting"));
9090

9191
@Rule public final MockitoRule mocks = MockitoJUnit.rule();
9292

0 commit comments

Comments
 (0)