Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
76 commits
Select commit Hold shift + click to select a range
b50ca84
core,api,xds: Implement load balancing policy delay plumbing
AgraVator May 13, 2026
c38ce1d
fix: tests
AgraVator May 14, 2026
a992bdf
fix: minor changes
AgraVator May 19, 2026
6a55ff2
add missing endDelay()
AgraVator Jun 8, 2026
389b96f
core,api,rls,util,xds: Implement dual Load Balancer delay APIs and ca…
AgraVator Jun 19, 2026
5e56f38
core: Add 100% test coverage for dual LB delay APIs and cadence rules
AgraVator Jun 19, 2026
6a3572b
opentelemetry: Implement dual Load Balancer delay spans and metrics
AgraVator Jun 22, 2026
9a4f21f
Merge remote-tracking branch 'upstream/master' into lb-policy-delay
AgraVator Jun 22, 2026
d25d064
Implement Name Resolution and unified RPC Delay Observability specifi…
AgraVator Jun 23, 2026
ffb485e
Ensure thread-safety and unit test coverage for Call-Level Delay APIs…
AgraVator Jul 6, 2026
42352f2
Add comprehensive End-to-End tests for Call-Level Name Resolution Del…
AgraVator Jul 6, 2026
a76996c
Ensure Call-Level delay recording only triggers when RPCs are queued …
AgraVator Jul 6, 2026
7c94743
Fix CdsLoadBalancer2Test atLeastOnce static import and assertions on …
AgraVator Jul 6, 2026
a7830ff
Fix PriorityLoadBalancerTest handleNameResolutionError assertion on n…
AgraVator Jul 6, 2026
b482c43
Fix checkstyle import ordering in OpenTelemetryTracingModuleTest
AgraVator Jul 6, 2026
fbbc1b9
Merge remote-tracking branch 'upstream/master' into name-resolution-d…
AgraVator Jul 27, 2026
df48bee
core, opentelemetry: Harden Name Resolution & LB delay state machines…
AgraVator Jul 27, 2026
fc9721d
api: update @since 1.82.0 to 1.84.0
AgraVator Jul 27, 2026
162caef
test: remove temporary stress test files prior to PR submission
AgraVator Jul 27, 2026
f3bdcd9
core: align PendingStream synchronized (this) blocks with ManagedChan…
AgraVator Jul 27, 2026
9cc310b
opentelemetry: add targeted unit tests to expand branch coverage for …
AgraVator Jul 27, 2026
0a7eb26
core: harden PendingStream synchronization and add multithreaded race…
AgraVator Jul 28, 2026
0bc1634
opentelemetry: add end-to-end Client/Server simulation tests for dela…
AgraVator Jul 29, 2026
20fc8d0
Fix PR #12893 CI failures: revert PickResult.withError default delay,…
AgraVator Jul 29, 2026
a7398d7
cleanup: remove extra stress and unit tests, keeping only essential C…
AgraVator Jul 29, 2026
a9ea9e7
test: restore unit tests from master and add test coverage for delay …
AgraVator Jul 29, 2026
35547a8
test: restore pickResult_withSubchannelReplacement and pickResult_wit…
AgraVator Jul 29, 2026
6571c2d
test: remove multi-threaded stress tests from DelayedClientTransportT…
AgraVator Jul 29, 2026
b7227a7
test: remove unused imports from ManagedChannelImplTest and DelayedCl…
AgraVator Jul 29, 2026
129b8e5
test(opentelemetry): add unit tests for delay metrics and tracing bra…
AgraVator Aug 1, 2026
ff1ed50
test(opentelemetry): remove A121DelayObservabilityWrapperTest
AgraVator Aug 1, 2026
d48258e
test(opentelemetry): add tests for callEnded/streamClosed guards and …
AgraVator Aug 1, 2026
24cd5c3
test(opentelemetry): add tests for delay reason-changed guards withou…
AgraVator Aug 1, 2026
e63b9ee
test: expand branch coverage for OobChannel, delay metrics, and traci…
AgraVator Aug 1, 2026
42e8686
refactor: simplify delay synchronization and telemetry style across P…
AgraVator Aug 3, 2026
f7216e3
core: restore state transition comments in DelayedClientTransport.upd…
AgraVator Aug 5, 2026
ffc15db
test, opentelemetry: address PR review feedback and add gRFC A121 nam…
AgraVator Aug 12, 2026
2f71b12
test, opentelemetry: add unit tests for missing code coverage branche…
AgraVator Aug 12, 2026
2416612
test, opentelemetry: fix unit test setups for checkstyle and null checks
AgraVator Aug 12, 2026
6f6edf2
test: add unit tests covering missing branch combinations in DelayedC…
AgraVator Aug 12, 2026
9c5406c
test: expand unit test coverage for targetAttributeFilter, clientAtte…
AgraVator Aug 12, 2026
e23ccfa
test: add coverage for initialType null, Baggage context, and stream …
AgraVator Aug 12, 2026
6d5feab
test: add unstarted delay end tests for OpenTelemetryMetricsModule
AgraVator Aug 12, 2026
3c67ba3
test: add emptyTracersAndNullInitialReason test for DelayedClientTran…
AgraVator Aug 12, 2026
0309b66
style: remove unnecessary fully qualified class names across PR files
AgraVator Aug 12, 2026
acf4c1e
test: resolve all review discussions on delay observability branch co…
AgraVator Aug 13, 2026
e9c1f65
core, opentelemetry: harden delay synchronization, address codecov fe…
AgraVator Sep 8, 2026
713d9d8
core, opentelemetry: optimize delay synchronization and zero-overhead…
AgraVator Sep 8, 2026
2b426b7
test: remove redundant and reflection-based tests, restore private en…
AgraVator Sep 8, 2026
eea9af5
api: update newly added delay observability APIs @since 1.84.0 to 1.85.0
AgraVator Sep 8, 2026
28a6d18
core, opentelemetry: optimize delay synchronization, remove redundant…
AgraVator Sep 8, 2026
234b8b2
api: update attempt-level delay APIs to @since 1.84.0
AgraVator Sep 8, 2026
3314902
opentelemetry, core, api: align delay observability with A121 and res…
AgraVator Sep 9, 2026
ece66db
opentelemetry: align delay metric descriptions with A121 and add test…
AgraVator Sep 10, 2026
d330fb9
core, opentelemetry, xds, api: complete gRFC A121 delay observability…
AgraVator Sep 15, 2026
9a051c4
core, util, opentelemetry, xds, rls, grpclb: close remaining A121 del…
AgraVator Sep 17, 2026
0be0d94
opentelemetry: add concurrency tests for the A121 delay span state ma…
AgraVator Sep 17, 2026
941c217
test: use the surrounding module's pick-args idiom in the new delay t…
AgraVator Sep 17, 2026
d55b434
core: cover the pending-call delay guards that a cancellation can reach
AgraVator Sep 17, 2026
8770f47
Merge upstream/master into the A121 delay observability branch
AgraVator Sep 17, 2026
73bcf15
A121: consolidate delay observability, drop PR-only machinery and rea…
AgraVator Sep 23, 2026
827b870
A121: replace re-entrancy deferral with start counters; simplify reas…
AgraVator Sep 23, 2026
4d9e74a
A121: realign PendingStream and PendingCall with the approved ece66db…
AgraVator Sep 23, 2026
72688d0
api: mark the A121 delay observability API as @since 1.85.0
AgraVator Sep 24, 2026
25cd25e
A121: re-baseline on master + the approved ece66db call-level delta
AgraVator Sep 24, 2026
bdd3d08
opentelemetry: keep A121 delay labels to the gRFC set; drop bogus tests
AgraVator Sep 24, 2026
8a180ee
core: end the call-level delay before the listener is closed; tidy de…
AgraVator Sep 24, 2026
828ece4
core: notify LB of name resolution error before releasing pending calls
AgraVator Sep 28, 2026
47c97d4
api: include delay type and reason in FixedResultPicker equality
AgraVator Sep 28, 2026
88787cb
xds: wrap priority child pickers in one place and prefix null-typed c…
AgraVator Sep 28, 2026
736ad07
opentelemetry,core,util: drop redundant delay tracer tests
AgraVator Sep 28, 2026
fbddfca
opentelemetry: add end-to-end test for resolver failure with wait-for…
AgraVator Sep 28, 2026
9dccd93
core, api, xds, opentelemetry: revert speculative delay changes and r…
AgraVator Sep 30, 2026
81d7dec
api: update @since to 1.86.0 for A121 delay observability APIs
AgraVator Sep 30, 2026
0071383
xds: add unit test coverage for PriorityLoadBalancer delay picker bra…
AgraVator Sep 30, 2026
8fa8c3e
opentelemetry: add full end-to-end test for A121 call and attempt del…
AgraVator Sep 30, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
80 changes: 62 additions & 18 deletions api/src/main/java/io/grpc/ClientStreamTracer.java
Original file line number Diff line number Diff line change
Expand Up @@ -58,42 +58,45 @@ public void createPendingStream() {
}

/**
* Called when an attempt-level delay segment (such as waiting for a load balancing pick or
* connection establishment) starts.
* Called when an attempt-level delay (such as waiting for a load balancing pick or connection
* establishment) starts, or when the {@code delayType} of the ongoing delay changes, in which
* case {@link #recordDelayEnd} is called for the previous delay first.
*
* <p>This method is invoked synchronously on the attempt thread. Implementations should start
* internal timers or child tracing spans (named strictly {@code "Attempt Delay"}) carrying the
* canonical {@code grpc.delay_type} attribute.
* <p>Implementations should start a timer and open a child tracing span (named strictly
* {@code "Delay"}) carrying the canonical {@code grpc.delay_type} attribute.
*
* @param delayType canonical low-cardinality label categorizing the delay (e.g., "connecting")
* @param delayReason high-cardinality diagnostic string describing granular runtime conditions
* @since 1.82.0
* @since 1.86.0
*/
public void recordAttemptDelayStart(String delayType, String delayReason) {
public void recordDelayStart(String delayType, String delayReason) {
}

/**
* Called when an attempt-level delay reason changes while the overall delay type remains
* constant (for example, when a priority load balancing policy fails over between tiers).
* Called when an attempt-level delay reason changes while the delay type remains constant (for
* example, when a load balancing policy updates its connection status detail).
*
* <p>Implementations should record structured events (such as {@code "Delay state transition"})
* on the active delay span without recreating the span or resetting cumulative timers.
* <p>Implementations should record a structured event (such as {@code "Delay triggered"}) on
* the active delay span without recreating the span or resetting cumulative timers.
*
* @param delayType canonical low-cardinality label of the ongoing delay
* @param delayReason updated high-cardinality diagnostic string describing new conditions
* @since 1.82.0
* @since 1.86.0
*/
public void recordAttemptDelayReasonChanged(String delayReason) {
public void recordDelayReasonChanged(String delayType, String delayReason) {
}

/**
* Called when an attempt-level delay segment ends upon successful pick or stream creation.
* Called when an attempt-level delay ends upon successful pick or stream creation, or when the
* attempt is cancelled or reaches its deadline while still waiting.
*
* <p>Implementations should simultaneously close active child tracing spans and record elapsed
* duration to the {@code grpc.client.attempt.delay.duration} histogram.
* <p>Implementations should close the active child tracing span and record the elapsed duration
* to the {@code grpc.client.attempt.delay.duration} histogram labeled with {@code delayType}.
*
* @since 1.82.0
* @param delayType canonical low-cardinality label of the delay being ended
* @since 1.86.0
*/
public void recordAttemptDelayEnd() {
public void recordDelayEnd(String delayType) {
}

/**
Expand Down Expand Up @@ -156,6 +159,47 @@ public abstract static class Factory {
public ClientStreamTracer newClientStreamTracer(StreamInfo info, Metadata headers) {
throw new UnsupportedOperationException("Not implemented");
}

/**
* Called when a call-level delay (such as waiting for name resolution) starts before any
* individual RPC attempt is created, or when the {@code delayType} of the ongoing delay
* changes, in which case {@link #recordDelayEnd} is called for the previous delay first.
*
* <p>Implementations should start a timer and open a child tracing span (named strictly
* {@code "Delay"}) carrying the canonical {@code grpc.delay_type} attribute.
*
* @param delayType canonical low-cardinality label categorizing the delay (e.g., "resolving")
* @param delayReason high-cardinality diagnostic string describing granular runtime conditions
* @since 1.86.0
*/
public void recordDelayStart(String delayType, String delayReason) {
}

/**
* Called when a call-level delay reason changes while the delay type remains constant.
*
* <p>Implementations should record a structured event (such as {@code "Delay triggered"}) on
* the active call delay span without recreating the span or resetting timers.
*
* @param delayType canonical low-cardinality label of the ongoing delay
* @param delayReason updated high-cardinality diagnostic string describing new conditions
* @since 1.86.0
*/
public void recordDelayReasonChanged(String delayType, String delayReason) {
}

/**
* Called when a call-level delay ends upon successful name resolution, or when the RPC is
* cancelled or reaches its deadline before resolution completes.
*
* <p>Implementations should close the active call delay span and record the elapsed duration
* to the {@code grpc.client.call.delay.duration} histogram labeled with {@code delayType}.
*
* @param delayType canonical low-cardinality label of the delay being ended
* @since 1.86.0
*/
public void recordDelayEnd(String delayType) {
}
}

/**
Expand Down
14 changes: 11 additions & 3 deletions api/src/main/java/io/grpc/LoadBalancer.java
Original file line number Diff line number Diff line change
Expand Up @@ -739,21 +739,29 @@ public static PickResult withNoResult() {
*
* @param delayType low-cardinality root cause label (e.g., "connecting")
* @param delayReason high-cardinality diagnostic string for trace events
* @since 1.82.0
* @since 1.86.0
*/
public static PickResult withNoResult(String delayType, String delayReason) {
Preconditions.checkNotNull(delayType, "delayType");
Preconditions.checkNotNull(delayReason, "delayReason");
return new PickResult(null, null, Status.OK, false, null, delayType, delayReason);
}

/** Returns the delay type label if any. */
/**
* Returns the delay type label if any.
*
* @since 1.86.0
*/
@Nullable
public String getDelayType() {
return delayType;
}

/** Returns the diagnostic delay reason if any. */
/**
* Returns the diagnostic delay reason if any.
*
* @since 1.86.0
*/
@Nullable
public String getDelayReason() {
return delayReason;
Expand Down
13 changes: 13 additions & 0 deletions api/src/test/java/io/grpc/ClientStreamTracerTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -57,4 +57,17 @@ public void streamInfo_toBuilder() {
StreamInfo info2 = info1.toBuilder().build();
assertThat(info2.getCallOptions()).isSameInstanceAs(callOptions);
}

@Test
public void defaultDelayMethodsNoOp() {
ClientStreamTracer tracer = new ClientStreamTracer() {};
tracer.recordDelayStart("connecting", "test");
tracer.recordDelayReasonChanged("connecting", "test2");
tracer.recordDelayEnd("connecting");

ClientStreamTracer.Factory factory = new ClientStreamTracer.Factory() {};
factory.recordDelayStart("resolving", "test");
factory.recordDelayReasonChanged("resolving", "test2");
factory.recordDelayEnd("resolving");
}
}
23 changes: 23 additions & 0 deletions api/src/test/java/io/grpc/LoadBalancerTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,19 @@ public void pickResult_withNoResult() {
assertThat(result.getStatus()).isSameInstanceAs(Status.OK);
assertThat(result.getStreamTracerFactory()).isNull();
assertThat(result.isDrop()).isFalse();
assertThat(result.getDelayType()).isNull();
assertThat(result.getDelayReason()).isNull();
}

@Test
public void pickResult_withNoResult_withDelay() {
PickResult result = PickResult.withNoResult("connecting", "trying backends");
assertThat(result.getSubchannel()).isNull();
assertThat(result.getStatus()).isSameInstanceAs(Status.OK);
assertThat(result.getStreamTracerFactory()).isNull();
assertThat(result.isDrop()).isFalse();
assertThat(result.getDelayType()).isEqualTo("connecting");
assertThat(result.getDelayReason()).isEqualTo("trying backends");
}

@Test
Expand Down Expand Up @@ -118,6 +131,10 @@ public void pickResult_equals() {
PickResult sc3 = PickResult.withSubchannel(subchannel, tracerFactory);
PickResult sc4 = PickResult.withSubchannel(subchannel2);
PickResult nr = PickResult.withNoResult();
PickResult nrDelay1 = PickResult.withNoResult("connecting", "trying 10.0.0.1");
PickResult nrDelay2 = PickResult.withNoResult("connecting", "trying 10.0.0.1");
PickResult nrDelayDiffReason = PickResult.withNoResult("connecting", "trying 10.0.0.2");
PickResult nrDelayDiffType = PickResult.withNoResult("rls_lookup_pending", "trying 10.0.0.1");
PickResult error1 = PickResult.withError(status);
PickResult error2 = PickResult.withError(status2);
PickResult error3 = PickResult.withError(status2);
Expand All @@ -132,6 +149,12 @@ public void pickResult_equals() {
assertThat(sc1).isNotEqualTo(sc3);
assertThat(sc1).isNotEqualTo(sc4);

assertThat(nr).isEqualTo(nrDelay1);
assertThat(nrDelay1).isEqualTo(nrDelay2);
assertThat(nrDelay1.hashCode()).isEqualTo(nrDelay2.hashCode());
assertThat(nrDelay1).isEqualTo(nrDelayDiffReason);
assertThat(nrDelay1).isEqualTo(nrDelayDiffType);

assertThat(error1).isNotEqualTo(error2);
assertThat(error2).isEqualTo(error3);

Expand Down
52 changes: 23 additions & 29 deletions core/src/main/java/io/grpc/internal/DelayedClientTransport.java
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,6 @@
import java.util.Collection;
import java.util.Collections;
import java.util.LinkedHashSet;
import java.util.Objects;
import java.util.concurrent.Executor;
import javax.annotation.Nonnull;
import javax.annotation.Nullable;
Expand Down Expand Up @@ -176,7 +175,7 @@ public final ClientStream newStream(
*/
@GuardedBy("lock")
private PendingStream createPendingStream(PickSubchannelArgs args, ClientStreamTracer[] tracers,
PickResult pickResult, @Nullable String delayType, @Nullable String delayReason) {
PickResult pickResult, String delayType, String delayReason) {
PendingStream pendingStream = new PendingStream(args, tracers, delayType, delayReason);
if (args.getCallOptions().isWaitForReady() && pickResult != null && pickResult.hasResult()) {
pendingStream.lastPickStatus = pickResult.getStatus();
Expand Down Expand Up @@ -388,7 +387,10 @@ private static String determineQueuingDelayReason(@Nullable PickResult pickResul
return "subchannel returned by LB picker has no connected subchannel";
}
if (!pickResult.getStatus().isOk()) {
return "wait_for_ready RPC failed with status: " + pickResult.getStatus();
Status status = pickResult.getStatus();
// Status.toString() would append the cause's stack trace.
return "wait_for_ready RPC failed with status: " + status.getCode()
+ (status.getDescription() == null ? "" : ": " + status.getDescription());
}
if (pickResult.getDelayReason() != null) {
return pickResult.getDelayReason();
Expand All @@ -407,16 +409,14 @@ private class PendingStream extends DelayedStream {
@Nullable private String activeDelayReason;

private PendingStream(PickSubchannelArgs args, ClientStreamTracer[] tracers,
@Nullable String initialType, @Nullable String initialReason) {
String delayType, String delayReason) {
super("connecting_and_lb");
this.args = args;
this.tracers = tracers;
this.activeDelayType = initialType;
this.activeDelayReason = initialReason;
if (initialType != null) {
for (ClientStreamTracer tracer : tracers) {
tracer.recordAttemptDelayStart(initialType, initialReason != null ? initialReason : "");
}
this.activeDelayType = delayType;
this.activeDelayReason = delayReason;
for (ClientStreamTracer tracer : tracers) {
tracer.recordDelayStart(delayType, delayReason);
}
}

Expand All @@ -427,30 +427,23 @@ private PendingStream(PickSubchannelArgs args, ClientStreamTracer[] tracers,
* spans are ended and a new segment is initiated. If only {@code newReason} changes, a
* structured transition event is appended to the active span without span re-creation.
*/
synchronized void updateDelay(@Nullable String newType, @Nullable String newReason) {
synchronized void updateDelay(String newType, String newReason) {
if (getRealStream() != null) {
return;
}
if (!Objects.equals(activeDelayType, newType)) {
if (!newType.equals(activeDelayType)) {
// Delay type changed (e.g., from RLS lookup to connecting). End the previous delay.
if (activeDelayType != null) {
for (ClientStreamTracer tracer : tracers) {
tracer.recordAttemptDelayEnd();
}
}
endDelay();
activeDelayType = newType;
activeDelayReason = null;
if (newType != null) {
for (ClientStreamTracer tracer : tracers) {
tracer.recordAttemptDelayStart(newType, newReason != null ? newReason : "");
}
activeDelayReason = newReason;
for (ClientStreamTracer tracer : tracers) {
tracer.recordDelayStart(newType, newReason);
}
}
if (newType != null && newReason != null && !Objects.equals(activeDelayReason, newReason)) {
} else if (!newReason.equals(activeDelayReason)) {
// Delay type is unchanged, but the reason changed (e.g., priority failover).
activeDelayReason = newReason;
for (ClientStreamTracer tracer : tracers) {
tracer.recordAttemptDelayReasonChanged(newReason);
tracer.recordDelayReasonChanged(newType, newReason);
}
}
}
Expand All @@ -459,12 +452,13 @@ synchronized void updateDelay(@Nullable String newType, @Nullable String newReas
* Ends active attempt delay segment telemetry upon stream creation or stream cancellation.
*/
synchronized void endDelay() {
if (activeDelayType != null) {
for (ClientStreamTracer tracer : tracers) {
tracer.recordAttemptDelayEnd();
}
String delayType = activeDelayType;
if (delayType != null) {
activeDelayType = null;
activeDelayReason = null;
for (ClientStreamTracer tracer : tracers) {
tracer.recordDelayEnd(delayType);
}
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -40,18 +40,18 @@ public void createPendingStream() {
}

@Override
public void recordAttemptDelayStart(String delayType, String delayReason) {
delegate().recordAttemptDelayStart(delayType, delayReason);
public void recordDelayStart(String delayType, String delayReason) {
delegate().recordDelayStart(delayType, delayReason);
}

@Override
public void recordAttemptDelayReasonChanged(String delayReason) {
delegate().recordAttemptDelayReasonChanged(delayReason);
public void recordDelayReasonChanged(String delayType, String delayReason) {
delegate().recordDelayReasonChanged(delayType, delayReason);
}

@Override
public void recordAttemptDelayEnd() {
delegate().recordAttemptDelayEnd();
public void recordDelayEnd(String delayType) {
delegate().recordDelayEnd(delayType);
}

@Override
Expand Down
Loading
Loading