Skip to content

Commit 389b96f

Browse files
committed
core,api,rls,util,xds: Implement dual Load Balancer delay APIs and cadence invariants
- Refactor ClientStreamTracer to expose delayTypeStarted(String) and delayReasonAttached(String) - Enhance PickResult with separate delayType and delayReason diagnostic fields - Implement Mark Roth's hybrid telemetry cadence model in DelayedClientTransport.PendingStream - Support channel fallback delay states (client_channel_init, subchannel_state_mismatch, wait_for_ready_failed) - Simplify leaf and container LB policies to emit canonical unified connecting metric labels
1 parent 6a55ff2 commit 389b96f

18 files changed

Lines changed: 169 additions & 68 deletions

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

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -58,12 +58,21 @@ public void createPendingStream() {
5858
}
5959

6060
/**
61-
* A delay segment started with a specific reason during load balancing.
61+
* A delay segment started with a canonical root cause.
6262
*
63-
* @param reasonToken the reason for the delay, e.g., "pick_first:connecting"
63+
* @param delayType the canonical root cause label (e.g., "connecting", "client_channel_init")
6464
* @since 1.82.0
6565
*/
66-
public void delayStarted(String reasonToken) {
66+
public void delayTypeStarted(String delayType) {
67+
}
68+
69+
/**
70+
* High-cardinality diagnostic context attached to the active delay span.
71+
*
72+
* @param delayReason verbose diagnostic description of the delay
73+
* @since 1.82.0
74+
*/
75+
public void delayReasonAttached(String delayReason) {
6776
}
6877

6978
/**

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

Lines changed: 25 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -549,30 +549,32 @@ 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;
552+
@Nullable private final String delayType;
553+
@Nullable private final String delayReason;
553554

554555
private PickResult(
555556
@Nullable Subchannel subchannel, @Nullable ClientStreamTracer.Factory streamTracerFactory,
556557
Status status, boolean drop) {
557-
this(subchannel, streamTracerFactory, status, drop, null, null);
558+
this(subchannel, streamTracerFactory, status, drop, null, null, null);
558559
}
559560

560561
private PickResult(
561562
@Nullable Subchannel subchannel, @Nullable ClientStreamTracer.Factory streamTracerFactory,
562563
Status status, boolean drop, @Nullable String authorityOverride) {
563-
this(subchannel, streamTracerFactory, status, drop, authorityOverride, null);
564+
this(subchannel, streamTracerFactory, status, drop, authorityOverride, null, null);
564565
}
565566

566567
private PickResult(
567568
@Nullable Subchannel subchannel, @Nullable ClientStreamTracer.Factory streamTracerFactory,
568569
Status status, boolean drop, @Nullable String authorityOverride,
569-
@Nullable String delayReasonToken) {
570+
@Nullable String delayType, @Nullable String delayReason) {
570571
this.subchannel = subchannel;
571572
this.streamTracerFactory = streamTracerFactory;
572573
this.status = checkNotNull(status, "status");
573574
this.drop = drop;
574575
this.authorityOverride = authorityOverride;
575-
this.delayReasonToken = delayReasonToken;
576+
this.delayType = delayType;
577+
this.delayReason = delayReason;
576578
}
577579

578580
/**
@@ -684,7 +686,7 @@ public static PickResult withSubchannel(Subchannel subchannel) {
684686
*/
685687
public PickResult copyWithSubchannel(Subchannel subchannel) {
686688
return new PickResult(checkNotNull(subchannel, "subchannel"), streamTracerFactory,
687-
status, drop, authorityOverride);
689+
status, drop, authorityOverride, delayType, delayReason);
688690
}
689691

690692
/**
@@ -695,7 +697,7 @@ public PickResult copyWithSubchannel(Subchannel subchannel) {
695697
*/
696698
public PickResult copyWithStreamTracerFactory(
697699
@Nullable ClientStreamTracer.Factory streamTracerFactory) {
698-
return new PickResult(subchannel, streamTracerFactory, status, drop, authorityOverride);
700+
return new PickResult(subchannel, streamTracerFactory, status, drop, authorityOverride, delayType, delayReason);
699701
}
700702

701703
/**
@@ -733,19 +735,28 @@ public static PickResult withNoResult() {
733735
}
734736

735737
/**
736-
* No decision could be made. The RPC will stay buffered with a specific reason.
738+
* No decision could be made. The RPC will stay buffered with a specific delay type and reason.
737739
*
740+
* @param delayType low-cardinality root cause label (e.g., "connecting")
741+
* @param delayReason high-cardinality diagnostic string for trace events
738742
* @since 1.82.0
739743
*/
740-
public static PickResult withNoResult(String delayReasonToken) {
741-
Preconditions.checkNotNull(delayReasonToken, "delayReasonToken");
742-
return new PickResult(null, null, Status.OK, false, null, delayReasonToken);
744+
public static PickResult withNoResult(String delayType, String delayReason) {
745+
Preconditions.checkNotNull(delayType, "delayType");
746+
Preconditions.checkNotNull(delayReason, "delayReason");
747+
return new PickResult(null, null, Status.OK, false, null, delayType, delayReason);
743748
}
744749

745-
/** Returns the delay reason token if any. */
750+
/** Returns the delay type label if any. */
746751
@Nullable
747-
public String getDelayReasonToken() {
748-
return delayReasonToken;
752+
public String getDelayType() {
753+
return delayType;
754+
}
755+
756+
/** Returns the diagnostic delay reason if any. */
757+
@Nullable
758+
public String getDelayReason() {
759+
return delayReason;
749760
}
750761

751762
/** Returns the authority override if any. */

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

Lines changed: 72 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -158,8 +158,9 @@ public final ClientStream newStream(
158158
synchronized (lock) {
159159
PickerState newerState = pickerState;
160160
if (state == newerState) {
161-
String token = pickResult != null ? pickResult.getDelayReasonToken() : null;
162-
return createPendingStream(args, tracers, pickResult, token);
161+
String delayType = determineQueuingDelayType(pickResult, callOptions.isWaitForReady());
162+
String delayReason = determineQueuingDelayReason(pickResult, callOptions.isWaitForReady());
163+
return createPendingStream(args, tracers, pickResult, delayType, delayReason);
163164
}
164165
state = newerState;
165166
}
@@ -175,8 +176,8 @@ public final ClientStream newStream(
175176
*/
176177
@GuardedBy("lock")
177178
private PendingStream createPendingStream(PickSubchannelArgs args, ClientStreamTracer[] tracers,
178-
PickResult pickResult, @Nullable String delayReasonToken) {
179-
PendingStream pendingStream = new PendingStream(args, tracers, delayReasonToken);
179+
PickResult pickResult, @Nullable String delayType, @Nullable String delayReason) {
180+
PendingStream pendingStream = new PendingStream(args, tracers, delayType, delayReason);
180181
if (args.getCallOptions().isWaitForReady() && pickResult != null && pickResult.hasResult()) {
181182
pendingStream.lastPickStatus = pickResult.getStatus();
182183
}
@@ -319,7 +320,11 @@ final void reprocess(@Nullable SubchannelPicker picker) {
319320
}
320321
toRemove.add(stream);
321322
} else { // stay pending
322-
stream.updateDelayReason(pickResult.getDelayReasonToken());
323+
String delayType = determineQueuingDelayType(
324+
pickResult, stream.args.getCallOptions().isWaitForReady());
325+
String delayReason = determineQueuingDelayReason(
326+
pickResult, stream.args.getCallOptions().isWaitForReady());
327+
stream.updateDelay(delayType, delayReason);
323328
}
324329
}
325330

@@ -361,48 +366,97 @@ public InternalLogId getLogId() {
361366
return logId;
362367
}
363368

369+
private static String determineQueuingDelayType(
370+
@Nullable PickResult pickResult, boolean isWaitForReady) {
371+
if (pickResult == null) {
372+
return "client_channel_init";
373+
}
374+
if (pickResult.getSubchannel() != null) {
375+
return "subchannel_state_mismatch";
376+
}
377+
if (!pickResult.getStatus().isOk()) {
378+
return "wait_for_ready_failed";
379+
}
380+
if (pickResult.getDelayType() != null) {
381+
return pickResult.getDelayType();
382+
}
383+
return "client_channel_init";
384+
}
385+
386+
private static String determineQueuingDelayReason(
387+
@Nullable PickResult pickResult, boolean isWaitForReady) {
388+
if (pickResult == null) {
389+
return "client channel: created LB policy.";
390+
}
391+
if (pickResult.getSubchannel() != null) {
392+
return "subchannel returned by LB picker has no connected subchannel";
393+
}
394+
if (!pickResult.getStatus().isOk()) {
395+
return "wait_for_ready RPC failed with status: " + pickResult.getStatus();
396+
}
397+
if (pickResult.getDelayReason() != null) {
398+
return pickResult.getDelayReason();
399+
}
400+
return "client channel: waiting for picker";
401+
}
402+
364403
private class PendingStream extends DelayedStream {
365404
private final PickSubchannelArgs args;
366405
private final Context context = Context.current();
367406
private final ClientStreamTracer[] tracers;
368407
private volatile Status lastPickStatus;
369-
@Nullable private String delayReasonToken;
408+
@Nullable private String activeDelayType;
409+
@Nullable private String activeDelayReason;
370410

371411
private PendingStream(PickSubchannelArgs args, ClientStreamTracer[] tracers,
372-
@Nullable String initialToken) {
412+
@Nullable String initialType, @Nullable String initialReason) {
373413
super("connecting_and_lb");
374414
this.args = args;
375415
this.tracers = tracers;
376-
this.delayReasonToken = initialToken;
377-
if (initialToken != null) {
416+
this.activeDelayType = initialType;
417+
this.activeDelayReason = initialReason;
418+
if (initialType != null) {
378419
for (ClientStreamTracer tracer : tracers) {
379-
tracer.delayStarted(initialToken);
420+
tracer.delayTypeStarted(initialType);
421+
}
422+
}
423+
if (initialReason != null) {
424+
for (ClientStreamTracer tracer : tracers) {
425+
tracer.delayReasonAttached(initialReason);
380426
}
381427
}
382428
}
383429

384-
void updateDelayReason(String newToken) {
385-
if (!Objects.equals(delayReasonToken, newToken)) {
386-
if (delayReasonToken != null) {
430+
void updateDelay(@Nullable String newType, @Nullable String newReason) {
431+
if (!Objects.equals(activeDelayType, newType)) {
432+
if (activeDelayType != null) {
387433
for (ClientStreamTracer tracer : tracers) {
388434
tracer.delayEnded();
389435
}
390436
}
391-
delayReasonToken = newToken;
392-
if (newToken != null) {
437+
activeDelayType = newType;
438+
activeDelayReason = null;
439+
if (newType != null) {
393440
for (ClientStreamTracer tracer : tracers) {
394-
tracer.delayStarted(newToken);
441+
tracer.delayTypeStarted(newType);
395442
}
396443
}
397444
}
445+
if (newType != null && newReason != null && !Objects.equals(activeDelayReason, newReason)) {
446+
activeDelayReason = newReason;
447+
for (ClientStreamTracer tracer : tracers) {
448+
tracer.delayReasonAttached(newReason);
449+
}
450+
}
398451
}
399452

400453
void endDelay() {
401-
if (delayReasonToken != null) {
454+
if (activeDelayType != null) {
402455
for (ClientStreamTracer tracer : tracers) {
403456
tracer.delayEnded();
404457
}
405-
delayReasonToken = null;
458+
activeDelayType = null;
459+
activeDelayReason = null;
406460
}
407461
}
408462

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

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -40,8 +40,13 @@ public void createPendingStream() {
4040
}
4141

4242
@Override
43-
public void delayStarted(String reasonToken) {
44-
delegate().delayStarted(reasonToken);
43+
public void delayTypeStarted(String delayType) {
44+
delegate().delayTypeStarted(delayType);
45+
}
46+
47+
@Override
48+
public void delayReasonAttached(String delayReason) {
49+
delegate().delayReasonAttached(delayReason);
4550
}
4651

4752
@Override

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -38,8 +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 =
42-
PickResult.withNoResult("pick_first:connecting");
41+
private static final PickResult CONNECTING_RESULT =
42+
PickResult.withNoResult("connecting", "pick_first: attempting to connect");
4343
private final Helper helper;
4444
private Subchannel subchannel;
4545
private ConnectivityState currentState = IDLE;

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

Lines changed: 12 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -779,22 +779,24 @@ public void streamDelayMetrics() {
779779

780780
SubchannelPicker connectingPicker = mock(SubchannelPicker.class);
781781
when(connectingPicker.pickSubchannel(any(PickSubchannelArgs.class)))
782-
.thenReturn(PickResult.withNoResult("pick_first:connecting"));
782+
.thenReturn(PickResult.withNoResult("connecting", "pick_first: attempting to connect"));
783783

784784
delayedTransport.reprocess(connectingPicker);
785785
delayedTransport.newStream(method, headers, callOptions, customTracers);
786786

787787
InOrder inOrder = inOrder(mockTracer);
788-
inOrder.verify(mockTracer).delayStarted("pick_first:connecting");
788+
inOrder.verify(mockTracer).delayTypeStarted("connecting");
789+
inOrder.verify(mockTracer).delayReasonAttached("pick_first: attempting to connect");
789790

790791
SubchannelPicker customDelayPicker = mock(SubchannelPicker.class);
791792
when(customDelayPicker.pickSubchannel(any(PickSubchannelArgs.class)))
792-
.thenReturn(PickResult.withNoResult("rls:lookup_pending"));
793+
.thenReturn(PickResult.withNoResult("rls_lookup_pending", "RLS request pending."));
793794

794795
delayedTransport.reprocess(customDelayPicker);
795796

796797
inOrder.verify(mockTracer).delayEnded();
797-
inOrder.verify(mockTracer).delayStarted("rls:lookup_pending");
798+
inOrder.verify(mockTracer).delayTypeStarted("rls_lookup_pending");
799+
inOrder.verify(mockTracer).delayReasonAttached("RLS request pending.");
798800

799801
delayedTransport.reprocess(mockPicker);
800802

@@ -808,13 +810,14 @@ public void streamDelayMetrics_cancelled() {
808810

809811
SubchannelPicker connectingPicker = mock(SubchannelPicker.class);
810812
when(connectingPicker.pickSubchannel(any(PickSubchannelArgs.class)))
811-
.thenReturn(PickResult.withNoResult("pick_first:connecting"));
813+
.thenReturn(PickResult.withNoResult("connecting", "pick_first: attempting to connect"));
812814

813815
delayedTransport.reprocess(connectingPicker);
814816
ClientStream stream = delayedTransport.newStream(method, headers, callOptions, customTracers);
815817
stream.start(streamListener);
816818

817-
verify(mockTracer).delayStarted("pick_first:connecting");
819+
verify(mockTracer).delayTypeStarted("connecting");
820+
verify(mockTracer).delayReasonAttached("pick_first: attempting to connect");
818821

819822
stream.cancel(Status.CANCELLED);
820823

@@ -828,13 +831,14 @@ public void streamDelayMetrics_shutdownNow() {
828831

829832
SubchannelPicker connectingPicker = mock(SubchannelPicker.class);
830833
when(connectingPicker.pickSubchannel(any(PickSubchannelArgs.class)))
831-
.thenReturn(PickResult.withNoResult("pick_first:connecting"));
834+
.thenReturn(PickResult.withNoResult("connecting", "pick_first: attempting to connect"));
832835

833836
delayedTransport.reprocess(connectingPicker);
834837
ClientStream stream = delayedTransport.newStream(method, headers, callOptions, customTracers);
835838
stream.start(streamListener);
836839

837-
verify(mockTracer).delayStarted("pick_first:connecting");
840+
verify(mockTracer).delayTypeStarted("connecting");
841+
verify(mockTracer).delayReasonAttached("pick_first: attempting to connect");
838842

839843
delayedTransport.shutdownNow(Status.UNAVAILABLE);
840844

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -148,7 +148,8 @@ public void pickAfterResolved() throws Exception {
148148

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

154155
verifyNoMoreInteractions(mockHelper);

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("rls:lookup_pending");
1053+
return PickResult.withNoResult("rls_lookup_pending", "RLS request pending.");
10541054
}
10551055
}
10561056

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

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -262,7 +262,8 @@ 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");
265+
assertThat(res.getDelayType()).isEqualTo("rls_lookup_pending");
266+
assertThat(res.getDelayReason()).isEqualTo("RLS request pending.");
266267
// Cache is warm, but still unconnected
267268
res = picker.pickSubchannel(searchSubchannelArgs);
268269
inOrder.verify(helper).createSubchannel(any(CreateSubchannelArgs.class));
@@ -494,7 +495,8 @@ public void lb_working_withoutDefaultTarget() throws Exception {
494495
PickResult res = picker.pickSubchannel(searchSubchannelArgs);
495496
assertThat(res.getStatus().isOk()).isTrue();
496497
assertThat(res.getSubchannel()).isNull();
497-
assertThat(res.getDelayReasonToken()).isEqualTo("rls:lookup_pending");
498+
assertThat(res.getDelayType()).isEqualTo("rls_lookup_pending");
499+
assertThat(res.getDelayReason()).isEqualTo("RLS request pending.");
498500
// Cache is warm, but still unconnected
499501
res = picker.pickSubchannel(searchSubchannelArgs);
500502
inOrder.verify(helper).createSubchannel(any(CreateSubchannelArgs.class));

0 commit comments

Comments
 (0)