Skip to content

Commit d25d064

Browse files
committed
Implement Name Resolution and unified RPC Delay Observability specification
1 parent 9a4f21f commit d25d064

17 files changed

Lines changed: 519 additions & 91 deletions

File tree

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

Lines changed: 59 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -58,29 +58,42 @@ public void createPendingStream() {
5858
}
5959

6060
/**
61-
* A delay segment started with a canonical root cause.
61+
* Called when an attempt-level delay segment (such as waiting for a load balancing pick or
62+
* connection establishment) starts.
6263
*
63-
* @param delayType the canonical root cause label (e.g., "connecting", "client_channel_init")
64+
* <p>This method is invoked synchronously on the attempt thread. Implementations should start
65+
* internal timers or child tracing spans (named strictly {@code "Attempt Delay"}) carrying the
66+
* canonical {@code grpc.delay_type} attribute.
67+
*
68+
* @param delayType canonical low-cardinality label categorizing the delay (e.g., "connecting")
69+
* @param delayReason high-cardinality diagnostic string describing granular runtime conditions
6470
* @since 1.82.0
6571
*/
66-
public void delayTypeStarted(String delayType) {
72+
public void recordAttemptDelayStart(String delayType, String delayReason) {
6773
}
6874

6975
/**
70-
* High-cardinality diagnostic context attached to the active delay span.
76+
* Called when an attempt-level delay reason changes while the overall delay type remains
77+
* constant (for example, when a priority load balancing policy fails over between tiers).
78+
*
79+
* <p>Implementations should record structured events (such as {@code "Delay state transition"})
80+
* on the active delay span without recreating the span or resetting cumulative timers.
7181
*
72-
* @param delayReason verbose diagnostic description of the delay
82+
* @param delayReason updated high-cardinality diagnostic string describing new conditions
7383
* @since 1.82.0
7484
*/
75-
public void delayReasonAttached(String delayReason) {
85+
public void recordAttemptDelayReasonChanged(String delayReason) {
7686
}
7787

7888
/**
79-
* The current delay segment ended.
89+
* Called when an attempt-level delay segment ends upon successful pick or stream creation.
90+
*
91+
* <p>Implementations should simultaneously close active child tracing spans and record elapsed
92+
* duration to the {@code grpc.client.attempt.delay.duration} histogram.
8093
*
8194
* @since 1.82.0
8295
*/
83-
public void delayEnded() {
96+
public void recordAttemptDelayEnd() {
8497
}
8598

8699
/**
@@ -143,6 +156,44 @@ public abstract static class Factory {
143156
public ClientStreamTracer newClientStreamTracer(StreamInfo info, Metadata headers) {
144157
throw new UnsupportedOperationException("Not implemented");
145158
}
159+
160+
/**
161+
* Called when a call-level delay segment (such as waiting for name resolution or service
162+
* configuration parsing) starts before any individual RPC attempt is created.
163+
*
164+
* <p>Implementations should start logical timers and create child tracing spans (named strictly
165+
* {@code "Call Delay"}) carrying the canonical {@code grpc.delay_type} attribute.
166+
*
167+
* @param delayType canonical low-cardinality label categorizing the delay (e.g., "resolving")
168+
* @param delayReason high-cardinality diagnostic string describing granular runtime conditions
169+
* @since 1.82.0
170+
*/
171+
public void recordCallDelayStart(String delayType, String delayReason) {
172+
}
173+
174+
/**
175+
* Called when a call-level delay reason changes while the active delay segment continues.
176+
*
177+
* <p>Implementations should emit structured events (such as {@code "Delay state transition"})
178+
* on the active call delay span without recreating the span or resetting timers.
179+
*
180+
* @param delayReason updated high-cardinality diagnostic string describing new conditions
181+
* @since 1.82.0
182+
*/
183+
public void recordCallDelayReasonChanged(String delayReason) {
184+
}
185+
186+
/**
187+
* Called when a call-level delay segment ends upon successful name resolution or when an RPC
188+
* is cancelled before resolution completes.
189+
*
190+
* <p>Implementations should close active call delay spans and record elapsed duration to the
191+
* {@code grpc.client.call.delay.duration} histogram.
192+
*
193+
* @since 1.82.0
194+
*/
195+
public void recordCallDelayEnd() {
196+
}
146197
}
147198

148199
/**

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

Lines changed: 19 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -413,43 +413,52 @@ private PendingStream(PickSubchannelArgs args, ClientStreamTracer[] tracers,
413413
this.activeDelayReason = initialReason;
414414
if (initialType != null) {
415415
for (ClientStreamTracer tracer : tracers) {
416-
tracer.delayTypeStarted(initialType);
417-
}
418-
}
419-
if (initialReason != null) {
420-
for (ClientStreamTracer tracer : tracers) {
421-
tracer.delayReasonAttached(initialReason);
416+
tracer.recordAttemptDelayStart(initialType, initialReason != null ? initialReason : "");
422417
}
423418
}
424419
}
425420

421+
/**
422+
* Updates active attempt delay telemetry state upon load balancing state transitions.
423+
*
424+
* <p>If {@code newType} differs from the active delay type, active segment timers and child
425+
* spans are ended and a new segment is initiated. If only {@code newReason} changes, a
426+
* structured transition event is appended to the active span without span re-creation.
427+
*/
426428
void updateDelay(@Nullable String newType, @Nullable String newReason) {
427429
if (!Objects.equals(activeDelayType, newType)) {
430+
// Delay categorization changed (e.g., from RLS lookup to TCP connecting).
431+
// Close prior active segment across all tracers before starting new canonical segment.
428432
if (activeDelayType != null) {
429433
for (ClientStreamTracer tracer : tracers) {
430-
tracer.delayEnded();
434+
tracer.recordAttemptDelayEnd();
431435
}
432436
}
433437
activeDelayType = newType;
434438
activeDelayReason = null;
435439
if (newType != null) {
436440
for (ClientStreamTracer tracer : tracers) {
437-
tracer.delayTypeStarted(newType);
441+
tracer.recordAttemptDelayStart(newType, newReason != null ? newReason : "");
438442
}
439443
}
440444
}
441445
if (newType != null && newReason != null && !Objects.equals(activeDelayReason, newReason)) {
446+
// Categorization remained constant, but granular runtime diagnostics updated
447+
// (e.g., priority policy failover between tiers). Emit transition event.
442448
activeDelayReason = newReason;
443449
for (ClientStreamTracer tracer : tracers) {
444-
tracer.delayReasonAttached(newReason);
450+
tracer.recordAttemptDelayReasonChanged(newReason);
445451
}
446452
}
447453
}
448454

455+
/**
456+
* Ends active attempt delay segment telemetry upon stream creation or stream cancellation.
457+
*/
449458
void endDelay() {
450459
if (activeDelayType != null) {
451460
for (ClientStreamTracer tracer : tracers) {
452-
tracer.delayEnded();
461+
tracer.recordAttemptDelayEnd();
453462
}
454463
activeDelayType = null;
455464
activeDelayReason = null;

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

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

4242
@Override
43-
public void delayTypeStarted(String delayType) {
44-
delegate().delayTypeStarted(delayType);
43+
public void recordAttemptDelayStart(String delayType, String delayReason) {
44+
delegate().recordAttemptDelayStart(delayType, delayReason);
4545
}
4646

4747
@Override
48-
public void delayReasonAttached(String delayReason) {
49-
delegate().delayReasonAttached(delayReason);
48+
public void recordAttemptDelayReasonChanged(String delayReason) {
49+
delegate().recordAttemptDelayReasonChanged(delayReason);
5050
}
5151

5252
@Override
53-
public void delayEnded() {
54-
delegate().delayEnded();
53+
public void recordAttemptDelayEnd() {
54+
delegate().recordAttemptDelayEnd();
5555
}
5656

5757
@Override

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

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -993,10 +993,20 @@ private final class PendingCall<ReqT, RespT> extends DelayedClientCall<ReqT, Res
993993
this.method = method;
994994
this.callOptions = callOptions;
995995
this.callCreationTime = ticker.nanoTime();
996+
// Category A (Resolver delay): Notify all registered tracer factories that this RPC
997+
// is queued waiting for initial name resolution or service configuration parsing.
998+
for (ClientStreamTracer.Factory factory : callOptions.getStreamTracerFactories()) {
999+
factory.recordCallDelayStart(
1000+
"resolving", "waiting for name resolution or service config");
1001+
}
9961002
}
9971003

9981004
/** Called when it's ready to create a real call and reprocess the pending call. */
9991005
void reprocess() {
1006+
// Name resolution succeeded; end Call-Level delay segment before launching attempts.
1007+
for (ClientStreamTracer.Factory factory : callOptions.getStreamTracerFactories()) {
1008+
factory.recordCallDelayEnd();
1009+
}
10001010
ClientCall<ReqT, RespT> realCall;
10011011
Context previous = context.attach();
10021012
try {
@@ -1022,6 +1032,9 @@ public void run() {
10221032

10231033
@Override
10241034
protected void callCancelled() {
1035+
for (ClientStreamTracer.Factory factory : callOptions.getStreamTracerFactories()) {
1036+
factory.recordCallDelayEnd();
1037+
}
10251038
super.callCancelled();
10261039
syncContext.execute(new PendingCallRemoval());
10271040
}

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

Lines changed: 22 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -785,22 +785,22 @@ public void streamDelayMetrics() {
785785
delayedTransport.newStream(method, headers, callOptions, customTracers);
786786

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

791791
SubchannelPicker customDelayPicker = mock(SubchannelPicker.class);
792792
when(customDelayPicker.pickSubchannel(any(PickSubchannelArgs.class)))
793793
.thenReturn(PickResult.withNoResult("rls_lookup_pending", "RLS request pending."));
794794

795795
delayedTransport.reprocess(customDelayPicker);
796796

797-
inOrder.verify(mockTracer).delayEnded();
798-
inOrder.verify(mockTracer).delayTypeStarted("rls_lookup_pending");
799-
inOrder.verify(mockTracer).delayReasonAttached("RLS request pending.");
797+
inOrder.verify(mockTracer).recordAttemptDelayEnd();
798+
inOrder.verify(mockTracer).recordAttemptDelayStart(
799+
"rls_lookup_pending", "RLS request pending.");
800800

801801
delayedTransport.reprocess(mockPicker);
802802

803-
inOrder.verify(mockTracer).delayEnded();
803+
inOrder.verify(mockTracer).recordAttemptDelayEnd();
804804
}
805805

806806
@Test
@@ -816,12 +816,12 @@ public void streamDelayMetrics_cancelled() {
816816
ClientStream stream = delayedTransport.newStream(method, headers, callOptions, customTracers);
817817
stream.start(streamListener);
818818

819-
verify(mockTracer).delayTypeStarted("connecting");
820-
verify(mockTracer).delayReasonAttached("pick_first: attempting to connect");
819+
verify(mockTracer).recordAttemptDelayStart(
820+
"connecting", "pick_first: attempting to connect");
821821

822822
stream.cancel(Status.CANCELLED);
823823

824-
verify(mockTracer).delayEnded();
824+
verify(mockTracer).recordAttemptDelayEnd();
825825
}
826826

827827
@Test
@@ -837,12 +837,12 @@ public void streamDelayMetrics_shutdownNow() {
837837
ClientStream stream = delayedTransport.newStream(method, headers, callOptions, customTracers);
838838
stream.start(streamListener);
839839

840-
verify(mockTracer).delayTypeStarted("connecting");
841-
verify(mockTracer).delayReasonAttached("pick_first: attempting to connect");
840+
verify(mockTracer).recordAttemptDelayStart(
841+
"connecting", "pick_first: attempting to connect");
842842

843843
delayedTransport.shutdownNow(Status.UNAVAILABLE);
844844

845-
verify(mockTracer).delayEnded();
845+
verify(mockTracer).recordAttemptDelayEnd();
846846
}
847847

848848
@Test
@@ -857,18 +857,17 @@ public void streamDelayMetrics_cadenceReasonUpdate_doesNotStartNewTypeSegment()
857857
delayedTransport.reprocess(picker1);
858858
delayedTransport.newStream(method, headers, callOptions, customTracers);
859859

860-
verify(mockTracer, times(1)).delayTypeStarted("connecting");
861-
verify(mockTracer).delayReasonAttached("attempt 1");
860+
verify(mockTracer, times(1)).recordAttemptDelayStart("connecting", "attempt 1");
862861

863862
SubchannelPicker picker2 = mock(SubchannelPicker.class);
864863
when(picker2.pickSubchannel(any(PickSubchannelArgs.class)))
865864
.thenReturn(PickResult.withNoResult("connecting", "attempt 2"));
866865

867866
delayedTransport.reprocess(picker2);
868867

869-
verify(mockTracer, times(1)).delayTypeStarted("connecting");
870-
verify(mockTracer).delayReasonAttached("attempt 2");
871-
verify(mockTracer, never()).delayEnded();
868+
verify(mockTracer, times(1)).recordAttemptDelayStart("connecting", "attempt 1");
869+
verify(mockTracer).recordAttemptDelayReasonChanged("attempt 2");
870+
verify(mockTracer, never()).recordAttemptDelayEnd();
872871
}
873872

874873
@Test
@@ -879,8 +878,8 @@ public void streamDelayMetrics_channelFallback_clientChannelInit() {
879878
// No picker reprocessed yet (lastPicker == null)
880879
delayedTransport.newStream(method, headers, callOptions, customTracers);
881880

882-
verify(mockTracer).delayTypeStarted("client_channel_init");
883-
verify(mockTracer).delayReasonAttached("client channel: created LB policy.");
881+
verify(mockTracer).recordAttemptDelayStart(
882+
"client_channel_init", "client channel: created LB policy.");
884883
}
885884

886885
@Test
@@ -900,8 +899,8 @@ public void streamDelayMetrics_channelFallback_subchannelStateMismatch() {
900899
delayedTransport.reprocess(stalePicker);
901900
delayedTransport.newStream(method, headers, callOptions, customTracers);
902901

903-
verify(mockTracer).delayTypeStarted("subchannel_state_mismatch");
904-
verify(mockTracer).delayReasonAttached(
902+
verify(mockTracer).recordAttemptDelayStart(
903+
"subchannel_state_mismatch",
905904
"subchannel returned by LB picker has no connected subchannel");
906905
}
907906

@@ -918,8 +917,8 @@ public void streamDelayMetrics_channelFallback_waitForReadyFailed() {
918917
CallOptions wfrOptions = callOptions.withWaitForReady();
919918
delayedTransport.newStream(method, headers, wfrOptions, customTracers);
920919

921-
verify(mockTracer).delayTypeStarted("wait_for_ready_failed");
922-
verify(mockTracer).delayReasonAttached(
920+
verify(mockTracer).recordAttemptDelayStart(
921+
"wait_for_ready_failed",
923922
"wait_for_ready RPC failed with status: " + Status.UNAVAILABLE);
924923
}
925924

opentelemetry/src/main/java/io/grpc/opentelemetry/GrpcOpenTelemetry.java

Lines changed: 26 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -232,12 +232,24 @@ static OpenTelemetryMetricsResource createMetricInstruments(Meter meter,
232232
.build());
233233
}
234234

235-
if (isMetricEnabled("grpc.client.attempt.delay", enableMetrics, disableDefault)) {
235+
if (isDelayObservabilityEnabled()
236+
&& isMetricEnabled("grpc.client.attempt.delay.duration", enableMetrics, disableDefault)) {
236237
builder.clientAttemptDelayCounter(
237238
meter.histogramBuilder(
238-
"grpc.client.attempt.delay")
239+
"grpc.client.attempt.delay.duration")
239240
.setUnit("s")
240-
.setDescription("Time taken to complete a client call attempt delay")
241+
.setDescription("Time taken before a client call attempt starts")
242+
.setExplicitBucketBoundariesAdvice(LATENCY_BUCKETS)
243+
.build());
244+
}
245+
246+
if (isDelayObservabilityEnabled()
247+
&& isMetricEnabled("grpc.client.call.delay.duration", enableMetrics, disableDefault)) {
248+
builder.clientCallDelayCounter(
249+
meter.histogramBuilder(
250+
"grpc.client.call.delay.duration")
251+
.setUnit("s")
252+
.setDescription("Time taken before a client call starts")
241253
.setExplicitBucketBoundariesAdvice(LATENCY_BUCKETS)
242254
.build());
243255
}
@@ -359,6 +371,17 @@ static OpenTelemetryMetricsResource createMetricInstruments(Meter meter,
359371
return builder.build();
360372
}
361373

374+
/**
375+
* Checks whether experimental client attempt and call delay observability is globally enabled.
376+
*
377+
* <p>Guarded strictly by the {@code GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY} environment
378+
* variable or JVM system property (defaults to {@code false}). When disabled, delay spans and
379+
* duration histograms are suppressed to avoid runtime overhead.
380+
*/
381+
static boolean isDelayObservabilityEnabled() {
382+
return GrpcUtil.getFlag("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY", false);
383+
}
384+
362385
static boolean isMetricEnabled(String metricName, Map<String, Boolean> enableMetrics,
363386
boolean disableDefault) {
364387
Boolean explicitlyEnabled = enableMetrics.get(metricName);

0 commit comments

Comments
 (0)