diff --git a/api/src/main/java/io/grpc/ClientStreamTracer.java b/api/src/main/java/io/grpc/ClientStreamTracer.java
index d654fd33358..285a02cbc49 100644
--- a/api/src/main/java/io/grpc/ClientStreamTracer.java
+++ b/api/src/main/java/io/grpc/ClientStreamTracer.java
@@ -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.
*
- *
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.
+ *
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).
*
- *
Implementations should record structured events (such as {@code "Delay state transition"})
- * on the active delay span without recreating the span or resetting cumulative timers.
+ *
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.
*
- *
Implementations should simultaneously close active child tracing spans and record elapsed
- * duration to the {@code grpc.client.attempt.delay.duration} histogram.
+ *
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) {
}
/**
@@ -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.
+ *
+ *
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.
+ *
+ *
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.
+ *
+ *
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) {
+ }
}
/**
diff --git a/api/src/main/java/io/grpc/LoadBalancer.java b/api/src/main/java/io/grpc/LoadBalancer.java
index e5c3d053ee7..3d38806adfd 100644
--- a/api/src/main/java/io/grpc/LoadBalancer.java
+++ b/api/src/main/java/io/grpc/LoadBalancer.java
@@ -739,7 +739,7 @@ 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");
@@ -747,13 +747,21 @@ public static PickResult withNoResult(String delayType, String 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;
diff --git a/api/src/test/java/io/grpc/ClientStreamTracerTest.java b/api/src/test/java/io/grpc/ClientStreamTracerTest.java
index 5ddee77f5c0..7f059539e97 100644
--- a/api/src/test/java/io/grpc/ClientStreamTracerTest.java
+++ b/api/src/test/java/io/grpc/ClientStreamTracerTest.java
@@ -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");
+ }
}
diff --git a/api/src/test/java/io/grpc/LoadBalancerTest.java b/api/src/test/java/io/grpc/LoadBalancerTest.java
index 22fdc220081..02bde5cc4bc 100644
--- a/api/src/test/java/io/grpc/LoadBalancerTest.java
+++ b/api/src/test/java/io/grpc/LoadBalancerTest.java
@@ -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
@@ -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);
@@ -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);
diff --git a/core/src/main/java/io/grpc/internal/DelayedClientTransport.java b/core/src/main/java/io/grpc/internal/DelayedClientTransport.java
index aa1d820a570..7fc6efe1b4c 100644
--- a/core/src/main/java/io/grpc/internal/DelayedClientTransport.java
+++ b/core/src/main/java/io/grpc/internal/DelayedClientTransport.java
@@ -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;
@@ -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();
@@ -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();
@@ -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);
}
}
@@ -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);
}
}
}
@@ -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);
+ }
}
}
diff --git a/core/src/main/java/io/grpc/internal/ForwardingClientStreamTracer.java b/core/src/main/java/io/grpc/internal/ForwardingClientStreamTracer.java
index 3ecdcdfbcaa..d197af3aa38 100644
--- a/core/src/main/java/io/grpc/internal/ForwardingClientStreamTracer.java
+++ b/core/src/main/java/io/grpc/internal/ForwardingClientStreamTracer.java
@@ -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
diff --git a/core/src/main/java/io/grpc/internal/ManagedChannelImpl.java b/core/src/main/java/io/grpc/internal/ManagedChannelImpl.java
index fcc82a526af..34ff8db2275 100644
--- a/core/src/main/java/io/grpc/internal/ManagedChannelImpl.java
+++ b/core/src/main/java/io/grpc/internal/ManagedChannelImpl.java
@@ -923,6 +923,7 @@ public void run() {
inUseStateAggregator.updateObjectInUse(pendingCallsInUseObject, true);
}
pendingCalls.add(pendingCall);
+ pendingCall.notifyQueuedForNameResolution();
} else {
pendingCall.reprocess();
}
@@ -1003,6 +1004,10 @@ private final class PendingCall extends DelayedClientCall method;
final CallOptions callOptions;
private final long callCreationTime;
+ @GuardedBy("this")
+ private boolean queuedForResolution;
+ @GuardedBy("this")
+ private boolean delayEnded;
PendingCall(Context context, MethodDescriptor method, CallOptions callOptions) {
super(
@@ -1016,8 +1021,35 @@ private final class PendingCall extends DelayedClientCall realCall;
Context previous = context.attach();
try {
@@ -1043,6 +1075,7 @@ public void run() {
@Override
protected void callCancelled() {
+ endDelayIfNeeded();
super.callCancelled();
syncContext.execute(new PendingCallRemoval());
}
diff --git a/core/src/test/java/io/grpc/internal/DelayedClientTransportTest.java b/core/src/test/java/io/grpc/internal/DelayedClientTransportTest.java
index afcce806f2b..8f1e4f86e91 100644
--- a/core/src/test/java/io/grpc/internal/DelayedClientTransportTest.java
+++ b/core/src/test/java/io/grpc/internal/DelayedClientTransportTest.java
@@ -42,6 +42,7 @@
import io.grpc.IntegerMarshaller;
import io.grpc.LoadBalancer.PickResult;
import io.grpc.LoadBalancer.PickSubchannelArgs;
+import io.grpc.LoadBalancer.Subchannel;
import io.grpc.LoadBalancer.SubchannelPicker;
import io.grpc.Metadata;
import io.grpc.MethodDescriptor;
@@ -792,7 +793,7 @@ public void streamDelayMetrics() {
delayedTransport.reprocess(fakePicker(
PickResult.withNoResult("rls_lookup_pending", "RLS request pending.")));
- assertEquals(1, fakeTracer.delayEndedCount);
+ assertEquals(Collections.singletonList("connecting"), fakeTracer.endedDelayTypes);
assertEquals(Arrays.asList("connecting", "rls_lookup_pending"),
fakeTracer.startedDelayTypes);
assertEquals(Arrays.asList("pick_first: attempting to connect", "RLS request pending."),
@@ -800,7 +801,7 @@ public void streamDelayMetrics() {
delayedTransport.reprocess(mockPicker);
- assertEquals(2, fakeTracer.delayEndedCount);
+ assertEquals(Arrays.asList("connecting", "rls_lookup_pending"), fakeTracer.endedDelayTypes);
}
@Test
@@ -876,8 +877,8 @@ public void streamDelayMetrics_channelFallback_subchannelStateMismatch() {
FakeStreamTracer fakeTracer = new FakeStreamTracer();
ClientStreamTracer[] customTracers = new ClientStreamTracer[] { fakeTracer };
- io.grpc.LoadBalancer.Subchannel disconnectedSubchannel =
- mock(io.grpc.LoadBalancer.Subchannel.class);
+ Subchannel disconnectedSubchannel =
+ mock(Subchannel.class);
when(disconnectedSubchannel.getInternalSubchannel())
.thenReturn(newTransportProvider(null));
@@ -896,14 +897,15 @@ public void streamDelayMetrics_channelFallback_waitForReadyFailed() {
FakeStreamTracer fakeTracer = new FakeStreamTracer();
ClientStreamTracer[] customTracers = new ClientStreamTracer[] { fakeTracer };
- delayedTransport.reprocess(fakePicker(PickResult.withError(Status.UNAVAILABLE)));
+ delayedTransport.reprocess(fakePicker(PickResult.withError(
+ Status.UNAVAILABLE.withDescription("io exception").withCause(new Exception("boom")))));
CallOptions wfrOptions = callOptions.withWaitForReady();
delayedTransport.newStream(method, headers, wfrOptions, customTracers);
assertEquals(Collections.singletonList("picker_failing_with_wait_for_ready"),
fakeTracer.startedDelayTypes);
assertEquals(Collections.singletonList(
- "wait_for_ready RPC failed with status: " + Status.UNAVAILABLE),
+ "wait_for_ready RPC failed with status: UNAVAILABLE: io exception"),
fakeTracer.startedDelayReasons);
}
@@ -911,21 +913,23 @@ private static final class FakeStreamTracer extends ClientStreamTracer {
final List startedDelayTypes = new ArrayList<>();
final List startedDelayReasons = new ArrayList<>();
final List changedDelayReasons = new ArrayList<>();
+ final List endedDelayTypes = new ArrayList<>();
int delayEndedCount = 0;
@Override
- public void recordAttemptDelayStart(String delayType, String delayReason) {
+ public void recordDelayStart(String delayType, String delayReason) {
startedDelayTypes.add(delayType);
startedDelayReasons.add(delayReason);
}
@Override
- public void recordAttemptDelayReasonChanged(String delayReason) {
+ public void recordDelayReasonChanged(String delayType, String delayReason) {
changedDelayReasons.add(delayReason);
}
@Override
- public void recordAttemptDelayEnd() {
+ public void recordDelayEnd(String delayType) {
+ endedDelayTypes.add(delayType);
delayEndedCount++;
}
}
diff --git a/core/src/test/java/io/grpc/internal/ForwardingClientStreamTracerTest.java b/core/src/test/java/io/grpc/internal/ForwardingClientStreamTracerTest.java
index 5eb5b49fa19..0355800274f 100644
--- a/core/src/test/java/io/grpc/internal/ForwardingClientStreamTracerTest.java
+++ b/core/src/test/java/io/grpc/internal/ForwardingClientStreamTracerTest.java
@@ -40,6 +40,7 @@ public void allMethodsForwarded() throws Exception {
Collections.emptyList());
}
+
private final class TestClientStreamTracer extends ForwardingClientStreamTracer {
@Override
protected ClientStreamTracer delegate() {
diff --git a/core/src/test/java/io/grpc/internal/ManagedChannelImplTest.java b/core/src/test/java/io/grpc/internal/ManagedChannelImplTest.java
index d052cce137b..66c722baf6c 100644
--- a/core/src/test/java/io/grpc/internal/ManagedChannelImplTest.java
+++ b/core/src/test/java/io/grpc/internal/ManagedChannelImplTest.java
@@ -40,6 +40,7 @@
import static org.junit.Assert.fail;
import static org.mockito.AdditionalAnswers.delegatesTo;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.ArgumentMatchers.isA;
import static org.mockito.ArgumentMatchers.same;
@@ -81,6 +82,7 @@
import io.grpc.ConnectivityState;
import io.grpc.ConnectivityStateInfo;
import io.grpc.Context;
+import io.grpc.Deadline;
import io.grpc.EquivalentAddressGroup;
import io.grpc.InsecureChannelCredentials;
import io.grpc.IntegerMarshaller;
@@ -4875,4 +4877,210 @@ private static ManagedChannelServiceConfig createManagedChannelServiceConfig(
return ManagedChannelServiceConfig
.fromServiceConfig(rawServiceConfig, true, 3, 4, policySelection);
}
+
+ @Test
+ public void callDelay_normalDeferredResolution() {
+ FakeNameResolverFactory nsFactory = new FakeNameResolverFactory.Builder(expectedUri)
+ .setResolvedAtStart(false).build();
+ channelBuilder.nameResolverFactory(nsFactory);
+ createChannel();
+
+ ClientStreamTracer.Factory mockTracerFactory = mock(ClientStreamTracer.Factory.class);
+ when(mockTracerFactory.newClientStreamTracer(any(StreamInfo.class), any(Metadata.class)))
+ .thenReturn(new ClientStreamTracer() {});
+ CallOptions callOptions = CallOptions.DEFAULT.withStreamTracerFactory(mockTracerFactory);
+ ClientCall call = channel.newCall(method, callOptions);
+ call.start(mockCallListener, new Metadata());
+
+ verify(mockTracerFactory).recordDelayStart(
+ eq("resolving"), eq("waiting for name resolution or service config"));
+ verify(mockTracerFactory, never()).recordDelayEnd(anyString());
+
+ nsFactory.allResolved();
+
+ verify(mockTracerFactory).recordDelayEnd("resolving");
+ executor.runDueTasks();
+ }
+
+ @Test
+ public void callDelay_immediateResolution() {
+ FakeNameResolverFactory nsFactory = new FakeNameResolverFactory.Builder(expectedUri)
+ .setResolvedAtStart(true).build();
+ channelBuilder.nameResolverFactory(nsFactory);
+ createChannel();
+
+ ClientStreamTracer.Factory mockTracerFactory = mock(ClientStreamTracer.Factory.class);
+ when(mockTracerFactory.newClientStreamTracer(any(StreamInfo.class), any(Metadata.class)))
+ .thenReturn(new ClientStreamTracer() {});
+ CallOptions callOptions = CallOptions.DEFAULT.withStreamTracerFactory(mockTracerFactory);
+ ClientCall call = channel.newCall(method, callOptions);
+ call.start(mockCallListener, new Metadata());
+
+ channel.syncContext.execute(new Runnable() {
+ @Override
+ public void run() {}
+ });
+
+ verify(mockTracerFactory, never()).recordDelayStart(anyString(), anyString());
+ verify(mockTracerFactory, never()).recordDelayEnd(anyString());
+ executor.runDueTasks();
+ }
+
+ @Test
+ public void callDelay_cancellationWhileQueued() {
+ FakeNameResolverFactory nsFactory = new FakeNameResolverFactory.Builder(expectedUri)
+ .setResolvedAtStart(false).build();
+ channelBuilder.nameResolverFactory(nsFactory);
+ createChannel();
+
+ ClientStreamTracer.Factory mockTracerFactory = mock(ClientStreamTracer.Factory.class);
+ when(mockTracerFactory.newClientStreamTracer(any(StreamInfo.class), any(Metadata.class)))
+ .thenReturn(new ClientStreamTracer() {});
+ CallOptions callOptions = CallOptions.DEFAULT.withStreamTracerFactory(mockTracerFactory);
+ ClientCall call = channel.newCall(method, callOptions);
+ call.start(mockCallListener, new Metadata());
+
+ verify(mockTracerFactory).recordDelayStart(
+ eq("resolving"), eq("waiting for name resolution or service config"));
+ verify(mockTracerFactory, never()).recordDelayEnd(anyString());
+
+ call.cancel("Cancelled while queued", null);
+
+ verify(mockTracerFactory).recordDelayEnd("resolving");
+ executor.runDueTasks();
+ }
+
+ @Test
+ public void callDelay_deadlineExpirationWhileQueued() {
+ FakeNameResolverFactory nsFactory = new FakeNameResolverFactory.Builder(expectedUri)
+ .setResolvedAtStart(false).build();
+ channelBuilder.nameResolverFactory(nsFactory);
+ createChannel();
+
+ ClientStreamTracer.Factory mockTracerFactory = mock(ClientStreamTracer.Factory.class);
+ when(mockTracerFactory.newClientStreamTracer(any(StreamInfo.class), any(Metadata.class)))
+ .thenReturn(new ClientStreamTracer() {});
+ CallOptions callOptions = CallOptions.DEFAULT
+ .withStreamTracerFactory(mockTracerFactory)
+ .withDeadline(Deadline.after(100, TimeUnit.MILLISECONDS, timer.getDeadlineTicker()));
+ ClientCall call = channel.newCall(method, callOptions);
+ call.start(mockCallListener, new Metadata());
+
+ verify(mockTracerFactory).recordDelayStart(
+ eq("resolving"), eq("waiting for name resolution or service config"));
+ verify(mockTracerFactory, never()).recordDelayEnd(anyString());
+
+ timer.forwardTime(101, TimeUnit.MILLISECONDS);
+
+ verify(mockTracerFactory).recordDelayEnd("resolving");
+ executor.runDueTasks();
+ }
+
+ @Test
+ public void callDelay_resolutionFailure() {
+ FakeNameResolverFactory nsFactory = new FakeNameResolverFactory.Builder(expectedUri)
+ .setResolvedAtStart(false)
+ .setError(Status.UNAVAILABLE.withDescription("Simulated resolver failure"))
+ .build();
+ channelBuilder.nameResolverFactory(nsFactory);
+ createChannel();
+
+ ClientStreamTracer.Factory mockTracerFactory = mock(ClientStreamTracer.Factory.class);
+ when(mockTracerFactory.newClientStreamTracer(any(StreamInfo.class), any(Metadata.class)))
+ .thenReturn(new ClientStreamTracer() {});
+ CallOptions callOptions = CallOptions.DEFAULT.withStreamTracerFactory(mockTracerFactory);
+ ClientCall call = channel.newCall(method, callOptions);
+ call.start(mockCallListener, new Metadata());
+
+ verify(mockTracerFactory).recordDelayStart(
+ eq("resolving"), eq("waiting for name resolution or service config"));
+ verify(mockTracerFactory, never()).recordDelayEnd(anyString());
+
+ nsFactory.allResolved();
+
+ verify(mockTracerFactory).recordDelayEnd("resolving");
+ executor.runDueTasks();
+ }
+
+ @Test
+ public void callDelay_forcefulShutdown() {
+ FakeNameResolverFactory nsFactory = new FakeNameResolverFactory.Builder(expectedUri)
+ .setResolvedAtStart(false).build();
+ channelBuilder.nameResolverFactory(nsFactory);
+ createChannel();
+
+ ClientStreamTracer.Factory mockTracerFactory = mock(ClientStreamTracer.Factory.class);
+ when(mockTracerFactory.newClientStreamTracer(any(StreamInfo.class), any(Metadata.class)))
+ .thenReturn(new ClientStreamTracer() {});
+ CallOptions callOptions = CallOptions.DEFAULT.withStreamTracerFactory(mockTracerFactory);
+ ClientCall call = channel.newCall(method, callOptions);
+ call.start(mockCallListener, new Metadata());
+
+ verify(mockTracerFactory).recordDelayStart(
+ eq("resolving"), eq("waiting for name resolution or service config"));
+ verify(mockTracerFactory, never()).recordDelayEnd(anyString());
+
+ channel.shutdownNow();
+
+ verify(mockTracerFactory).recordDelayEnd("resolving");
+ executor.runDueTasks();
+ }
+
+ @Test
+ public void callDelay_cancelledBeforeQueuedOnSyncContext() {
+ FakeNameResolverFactory nsFactory = new FakeNameResolverFactory.Builder(expectedUri)
+ .setResolvedAtStart(false).build();
+ channelBuilder.nameResolverFactory(nsFactory);
+ createChannel();
+
+ ClientStreamTracer.Factory mockTracerFactory = mock(ClientStreamTracer.Factory.class);
+ when(mockTracerFactory.newClientStreamTracer(any(StreamInfo.class), any(Metadata.class)))
+ .thenReturn(new ClientStreamTracer() {});
+ CallOptions callOptions = CallOptions.DEFAULT.withStreamTracerFactory(mockTracerFactory);
+
+ channel.syncContext.execute(() -> {
+ ClientCall call = channel.newCall(method, callOptions);
+ call.cancel("Cancelled before syncContext drains", null);
+ });
+
+ verify(mockTracerFactory, never()).recordDelayStart(anyString(), anyString());
+ verify(mockTracerFactory, never()).recordDelayEnd(anyString());
+ executor.runDueTasks();
+ }
+
+ @Test
+ public void callDelay_callCancelledDuringTracerIteration_abortsLoop() {
+ FakeNameResolverFactory nsFactory = new FakeNameResolverFactory.Builder(expectedUri)
+ .setResolvedAtStart(false).build();
+ channelBuilder.nameResolverFactory(nsFactory);
+ createChannel();
+
+ final AtomicReference> callRef = new AtomicReference<>();
+ ClientStreamTracer.Factory tracer1 = new ClientStreamTracer.Factory() {
+ @Override
+ public ClientStreamTracer newClientStreamTracer(StreamInfo info, Metadata headers) {
+ return new ClientStreamTracer() {};
+ }
+
+ @Override
+ public void recordDelayStart(String delayType, String delayReason) {
+ callRef.get().cancel("Cancel inside first tracer start", null);
+ }
+ };
+ ClientStreamTracer.Factory tracer2 = mock(ClientStreamTracer.Factory.class);
+ when(tracer2.newClientStreamTracer(any(StreamInfo.class), any(Metadata.class)))
+ .thenReturn(new ClientStreamTracer() {});
+
+ CallOptions callOptions = CallOptions.DEFAULT
+ .withStreamTracerFactory(tracer1)
+ .withStreamTracerFactory(tracer2);
+
+ channel.syncContext.execute(() -> {
+ ClientCall call = channel.newCall(method, callOptions);
+ callRef.set(call);
+ });
+
+ verify(tracer2, never()).recordDelayStart(anyString(), anyString());
+ executor.runDueTasks();
+ }
}
diff --git a/opentelemetry/src/main/java/io/grpc/opentelemetry/GrpcOpenTelemetry.java b/opentelemetry/src/main/java/io/grpc/opentelemetry/GrpcOpenTelemetry.java
index 1243e0fff59..31482a6dc80 100644
--- a/opentelemetry/src/main/java/io/grpc/opentelemetry/GrpcOpenTelemetry.java
+++ b/opentelemetry/src/main/java/io/grpc/opentelemetry/GrpcOpenTelemetry.java
@@ -232,17 +232,29 @@ static OpenTelemetryMetricsResource createMetricInstruments(Meter meter,
.build());
}
- if (isDelayObservabilityEnabled()
- && isMetricEnabled("grpc.client.attempt.delay.duration", enableMetrics, disableDefault)) {
+ if (isMetricEnabled("grpc.client.attempt.delay.duration", enableMetrics, disableDefault)) {
builder.clientAttemptDelayCounter(
meter.histogramBuilder(
"grpc.client.attempt.delay.duration")
.setUnit("s")
- .setDescription("Time taken before a client call attempt starts")
+ .setDescription(
+ "EXPERIMENTAL. Time an RPC attempt spent waiting for a load balancing pick"
+ + " or connection establishment.")
.setExplicitBucketBoundariesAdvice(LATENCY_BUCKETS)
.build());
}
+ if (isMetricEnabled("grpc.client.call.delay.duration", enableMetrics, disableDefault)) {
+ builder.clientCallDelayCounter(
+ meter.histogramBuilder(
+ "grpc.client.call.delay.duration")
+ .setUnit("s")
+ .setDescription(
+ "EXPERIMENTAL. Time an RPC spent waiting at the call level before an attempt was"
+ + " initiated, such as waiting for name resolution.")
+ .setExplicitBucketBoundariesAdvice(LATENCY_BUCKETS)
+ .build());
+ }
if (isMetricEnabled("grpc.client.attempt.sent_total_compressed_message_size", enableMetrics,
disableDefault)) {
builder.clientTotalSentCompressedMessageSizeCounter(
@@ -360,17 +372,6 @@ && isMetricEnabled("grpc.client.attempt.delay.duration", enableMetrics, disableD
return builder.build();
}
- /**
- * Checks whether experimental client attempt and call delay observability is globally enabled.
- *
- * Guarded strictly by the {@code GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY} environment
- * variable (defaults to {@code false}). When disabled, delay spans and
- * duration histograms are suppressed to avoid runtime overhead.
- */
- static boolean isDelayObservabilityEnabled() {
- return GrpcUtil.getFlag("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY", false);
- }
-
static boolean isMetricEnabled(String metricName, Map enableMetrics,
boolean disableDefault) {
Boolean explicitlyEnabled = enableMetrics.get(metricName);
diff --git a/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsModule.java b/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsModule.java
index 008a9754b29..7ea9e1d7b3c 100644
--- a/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsModule.java
+++ b/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsModule.java
@@ -57,7 +57,6 @@
import java.util.Collection;
import java.util.Collections;
import java.util.List;
-import java.util.Objects;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicIntegerFieldUpdater;
import java.util.concurrent.atomic.AtomicLong;
@@ -205,8 +204,12 @@ private static final class ClientTracer extends ClientStreamTracer {
volatile String backendService;
long attemptNanos;
Code statusCode;
- @Nullable private volatile Stopwatch activeDelayStopwatch;
- @Nullable private volatile String activeDelayType;
+ @GuardedBy("this")
+ @Nullable private Stopwatch activeDelayStopwatch;
+ @GuardedBy("this")
+ @Nullable private String activeDelayType;
+ @GuardedBy("this")
+ private boolean streamClosed;
ClientTracer(CallAttemptsTracerFactory attemptsState, OpenTelemetryMetricsModule module,
StreamInfo info, String target, String fullMethodName,
@@ -221,51 +224,59 @@ private static final class ClientTracer extends ClientStreamTracer {
}
@Override
- public void streamCreated(io.grpc.Attributes transportAtts, Metadata headers) {
- recordAttemptDelayEnd();
+ public synchronized void streamCreated(io.grpc.Attributes transportAtts, Metadata headers) {
+ endOpenDelay();
}
@Override
- public void recordAttemptDelayStart(String delayType, String delayReason) {
- if (!GrpcOpenTelemetry.isDelayObservabilityEnabled()
- || (activeDelayStopwatch != null && Objects.equals(activeDelayType, delayType))) {
+ public synchronized void recordDelayStart(String delayType, String delayReason) {
+ if (streamClosed) {
+ return;
+ }
+ if (activeDelayStopwatch != null && delayType.equals(activeDelayType)) {
// Do not reset the stopwatch if the delay type is unchanged.
return;
}
- recordAttemptDelayEnd();
+ endOpenDelay();
activeDelayType = delayType;
activeDelayStopwatch = module.stopwatchSupplier.get().start();
}
@Override
- public void recordAttemptDelayReasonChanged(String delayReason) {
+ public void recordDelayReasonChanged(String delayType, String delayReason) {
// Reason strings are high-cardinality diagnostics intended for tracing spans.
}
@Override
- public void recordAttemptDelayEnd() {
- Stopwatch delayStopwatch = activeDelayStopwatch;
- String delayType = activeDelayType;
- if (delayStopwatch != null && delayType != null) {
- delayStopwatch.stop();
- long delayNanos = delayStopwatch.elapsed(TimeUnit.NANOSECONDS);
+ public synchronized void recordDelayEnd(String delayType) {
+ endOpenDelay();
+ }
+
+ /** Ends the delay that is still open, if any, under the type it was started with. */
+ @GuardedBy("this")
+ private void endOpenDelay() {
+ if (activeDelayStopwatch != null) {
+ recordDelay(activeDelayStopwatch.stop().elapsed(TimeUnit.NANOSECONDS), activeDelayType);
activeDelayStopwatch = null;
activeDelayType = null;
- if (module.resource.clientAttemptDelayCounter() != null) {
- AttributesBuilder builder = Attributes.builder()
- .put(METHOD_KEY, fullMethodName)
- .put(TARGET_KEY, target)
- .put("grpc.delay_type", delayType);
- if (module.customLabelEnabled) {
- builder.put(
- CUSTOM_LABEL_KEY, info.getCallOptions().getOption(Grpc.CALL_OPTION_CUSTOM_LABEL));
- }
- for (OpenTelemetryPlugin.ClientStreamPlugin plugin : streamPlugins) {
- plugin.addLabels(builder);
- }
- module.resource.clientAttemptDelayCounter()
- .record(delayNanos * SECONDS_PER_NANO, builder.build(), attemptsState.otelContext);
+ }
+ }
+
+ private void recordDelay(long delayNanos, String delayType) {
+ if (module.resource.clientAttemptDelayCounter() != null) {
+ AttributesBuilder builder = Attributes.builder()
+ .put(METHOD_KEY, fullMethodName)
+ .put(TARGET_KEY, target)
+ .put("grpc.delay_type", delayType);
+ if (module.customLabelEnabled) {
+ builder.put(
+ CUSTOM_LABEL_KEY, info.getCallOptions().getOption(Grpc.CALL_OPTION_CUSTOM_LABEL));
+ }
+ for (OpenTelemetryPlugin.ClientStreamPlugin plugin : streamPlugins) {
+ plugin.addLabels(builder);
}
+ module.resource.clientAttemptDelayCounter()
+ .record(delayNanos * SECONDS_PER_NANO, builder.build(), attemptsState.otelContext);
}
}
@@ -315,7 +326,13 @@ public void inboundTrailers(Metadata trailers) {
@Override
public void streamClosed(Status status) {
- recordAttemptDelayEnd();
+ synchronized (this) {
+ if (streamClosed) {
+ return;
+ }
+ streamClosed = true;
+ endOpenDelay();
+ }
stopwatch.stop();
attemptNanos = stopwatch.elapsed(TimeUnit.NANOSECONDS);
Deadline deadline = info.getCallOptions().getDeadline();
@@ -387,6 +404,11 @@ static final class CallAttemptsTracerFactory extends ClientStreamTracer.Factory
private final List callPlugins;
private final Context otelContext;
private Status status;
+ @GuardedBy("lock")
+ @Nullable private Stopwatch activeCallDelayStopwatch;
+ @GuardedBy("lock")
+ @Nullable private String activeCallDelayType;
+ private final Attributes callLevelBaseAttributes;
private long retryDelayNanos;
private long callLatencyNanos;
private final Object lock = new Object();
@@ -412,18 +434,18 @@ static final class CallAttemptsTracerFactory extends ClientStreamTracer.Factory
this.attemptDelayStopwatch = module.stopwatchSupplier.get();
this.callStopWatch = module.stopwatchSupplier.get().start();
- AttributesBuilder builder = io.opentelemetry.api.common.Attributes.builder()
+ AttributesBuilder builder = Attributes.builder()
.put(METHOD_KEY, fullMethodName)
.put(TARGET_KEY, target);
if (module.customLabelEnabled) {
builder.put(
CUSTOM_LABEL_KEY, callOptions.getOption(Grpc.CALL_OPTION_CUSTOM_LABEL));
}
- io.opentelemetry.api.common.Attributes attribute = builder.build();
+ this.callLevelBaseAttributes = builder.build();
- // Record here in case mewClientStreamTracer() would never be called.
+ // Record here in case newClientStreamTracer() would never be called.
if (module.resource.clientAttemptCountCounter() != null) {
- module.resource.clientAttemptCountCounter().add(1, attribute, otelContext);
+ module.resource.clientAttemptCountCounter().add(1, callLevelBaseAttributes, otelContext);
}
}
@@ -504,6 +526,7 @@ void callEnded(Status status, CallOptions callOptions) {
return;
}
callEnded = true;
+ endOpenDelay();
if (activeStreams == 0 && !finishedCallToBeRecorded) {
shouldRecordFinishedCall = true;
finishedCallToBeRecorded = true;
@@ -579,6 +602,52 @@ void recordFinishedCall(CallOptions callOptions) {
);
}
}
+
+ @Override
+ public void recordDelayStart(String delayType, String delayReason) {
+ synchronized (lock) {
+ if (callEnded) {
+ return;
+ }
+ if (activeCallDelayStopwatch != null && delayType.equals(activeCallDelayType)) {
+ return;
+ }
+ endOpenDelay();
+ activeCallDelayType = delayType;
+ activeCallDelayStopwatch = module.stopwatchSupplier.get().start();
+ }
+ }
+
+ @Override
+ public void recordDelayReasonChanged(String delayType, String delayReason) {
+ // Reason strings are high-cardinality diagnostics intended for tracing spans.
+ }
+
+ @Override
+ public void recordDelayEnd(String delayType) {
+ synchronized (lock) {
+ endOpenDelay();
+ }
+ }
+
+ @GuardedBy("lock")
+ private void endOpenDelay() {
+ if (activeCallDelayStopwatch != null) {
+ recordDelay(
+ activeCallDelayStopwatch.stop().elapsed(TimeUnit.NANOSECONDS), activeCallDelayType);
+ activeCallDelayStopwatch = null;
+ activeCallDelayType = null;
+ }
+ }
+
+ private void recordDelay(long delayNanos, String delayType) {
+ if (module.resource.clientCallDelayCounter() != null) {
+ module.resource.clientCallDelayCounter().record(
+ delayNanos * SECONDS_PER_NANO,
+ callLevelBaseAttributes.toBuilder().put("grpc.delay_type", delayType).build(),
+ otelContext);
+ }
+ }
}
private static final class ServerTracer extends ServerStreamTracer
diff --git a/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsResource.java b/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsResource.java
index 085498d746e..2771b3eb9af 100644
--- a/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsResource.java
+++ b/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsResource.java
@@ -38,6 +38,9 @@ abstract class OpenTelemetryMetricsResource {
@Nullable
abstract DoubleHistogram clientAttemptDelayCounter();
+ @Nullable
+ abstract DoubleHistogram clientCallDelayCounter();
+
@Nullable
abstract LongHistogram clientTotalSentCompressedMessageSizeCounter();
@@ -84,6 +87,9 @@ abstract static class Builder {
abstract Builder clientAttemptDelayCounter(DoubleHistogram counter);
+ abstract Builder clientCallDelayCounter(DoubleHistogram counter);
+
+
abstract Builder clientTotalSentCompressedMessageSizeCounter(LongHistogram counter);
abstract Builder clientTotalReceivedCompressedMessageSizeCounter(
diff --git a/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryTracingModule.java b/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryTracingModule.java
index 088d6dc9845..a7afa563f5c 100644
--- a/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryTracingModule.java
+++ b/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryTracingModule.java
@@ -22,7 +22,7 @@
import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.BAGGAGE_KEY;
import com.google.common.annotations.VisibleForTesting;
-import io.grpc.Attributes;
+import com.google.errorprone.annotations.concurrent.GuardedBy;
import io.grpc.CallOptions;
import io.grpc.Channel;
import io.grpc.ClientCall;
@@ -37,11 +37,13 @@
import io.grpc.ServerCallHandler;
import io.grpc.ServerInterceptor;
import io.grpc.ServerStreamTracer;
+import io.grpc.Status;
import io.grpc.internal.GrpcUtil;
import io.grpc.opentelemetry.internal.OpenTelemetryConstants;
import io.opentelemetry.api.OpenTelemetry;
import io.opentelemetry.api.baggage.Baggage;
import io.opentelemetry.api.common.AttributeKey;
+import io.opentelemetry.api.common.Attributes;
import io.opentelemetry.api.common.AttributesBuilder;
import io.opentelemetry.api.trace.Span;
import io.opentelemetry.api.trace.StatusCode;
@@ -49,7 +51,6 @@
import io.opentelemetry.context.Context;
import io.opentelemetry.context.Scope;
import io.opentelemetry.context.propagation.ContextPropagators;
-import java.util.Objects;
import java.util.concurrent.atomic.AtomicIntegerFieldUpdater;
import java.util.logging.Level;
import java.util.logging.Logger;
@@ -144,6 +145,10 @@ final class CallAttemptsTracerFactory extends ClientStreamTracer.Factory {
volatile int callEnded;
private final Span clientSpan;
private final String fullMethodName;
+ @GuardedBy("this")
+ @Nullable private Span activeCallDelaySpan;
+ @GuardedBy("this")
+ @Nullable private String activeCallDelayType;
CallAttemptsTracerFactory(Span clientSpan, MethodDescriptor, ?> method) {
checkNotNull(method, "method");
@@ -168,6 +173,13 @@ public ClientStreamTracer newClientStreamTracer(
return new ClientTracer(attemptSpan, clientSpan);
}
+ private boolean isCallEnded() {
+ if (callEndedUpdater != null) {
+ return callEndedUpdater.get(this) != 0;
+ }
+ return callEnded != 0;
+ }
+
/**
* Record a finished call and mark the current time as the end time.
*
@@ -185,8 +197,61 @@ void callEnded(io.grpc.Status status) {
}
callEnded = 1;
}
+ synchronized (this) {
+ endActiveDelaySpan();
+ }
endSpanWithStatus(clientSpan, status);
}
+
+ @Override
+ public void recordDelayStart(String delayType, String delayReason) {
+ if (isCallEnded()) {
+ return;
+ }
+ synchronized (this) {
+ if (isCallEnded()) {
+ return;
+ }
+ if (activeCallDelaySpan != null && delayType.equals(activeCallDelayType)) {
+ addDelayEvent(activeCallDelaySpan, delayReason);
+ return;
+ }
+ endActiveDelaySpan();
+ activeCallDelayType = delayType;
+ activeCallDelaySpan = otelTracer.spanBuilder("Delay")
+ .setParent(Context.current().with(clientSpan))
+ .setAttribute("grpc.delay_type", delayType)
+ .startSpan();
+ addDelayEvent(activeCallDelaySpan, delayReason);
+ }
+ }
+
+ @Override
+ public void recordDelayReasonChanged(String delayType, String delayReason) {
+ if (isCallEnded()) {
+ return;
+ }
+ synchronized (this) {
+ if (isCallEnded() || activeCallDelaySpan == null) {
+ return;
+ }
+ addDelayEvent(activeCallDelaySpan, delayReason);
+ }
+ }
+
+ @Override
+ public synchronized void recordDelayEnd(String delayType) {
+ endActiveDelaySpan();
+ }
+
+ @GuardedBy("this")
+ private void endActiveDelaySpan() {
+ if (activeCallDelaySpan != null) {
+ activeCallDelaySpan.end();
+ activeCallDelaySpan = null;
+ activeCallDelayType = null;
+ }
+ }
}
private final class ClientTracer extends ClientStreamTracer {
@@ -194,8 +259,12 @@ private final class ClientTracer extends ClientStreamTracer {
private final Span parentSpan;
volatile int seqNo;
boolean isPendingStream;
- @Nullable private volatile Span activeDelaySpan;
- @Nullable private volatile String activeDelayType;
+ @GuardedBy("this")
+ @Nullable private Span activeDelaySpan;
+ @GuardedBy("this")
+ @Nullable private String activeDelayType;
+ @GuardedBy("this")
+ private boolean streamClosed;
ClientTracer(Span span, Span parentSpan) {
this.span = checkNotNull(span, "span");
@@ -203,8 +272,10 @@ private final class ClientTracer extends ClientStreamTracer {
}
@Override
- public void streamCreated(Attributes transportAtts, Metadata headers) {
- recordAttemptDelayEnd();
+ public void streamCreated(io.grpc.Attributes transportAtts, Metadata headers) {
+ synchronized (this) {
+ endActiveDelaySpan();
+ }
contextPropagators.getTextMapPropagator().inject(Context.current().with(span), headers,
metadataSetter);
if (isPendingStream) {
@@ -218,50 +289,41 @@ public void createPendingStream() {
}
@Override
- public void recordAttemptDelayStart(String delayType, String delayReason) {
- if (!GrpcOpenTelemetry.isDelayObservabilityEnabled()) {
+ public synchronized void recordDelayStart(String delayType, String delayReason) {
+ if (streamClosed) {
return;
}
- if (activeDelaySpan != null && Objects.equals(activeDelayType, delayType)) {
+ if (activeDelaySpan != null && delayType.equals(activeDelayType)) {
// Do not recreate the span if the delay type is unchanged (e.g., priority failover).
- recordAttemptDelayReasonChanged(delayReason);
+ addDelayEvent(activeDelaySpan, delayReason);
return;
}
- // Close any previous delay segment before starting a new canonical segment.
- recordAttemptDelayEnd();
+ endActiveDelaySpan();
activeDelayType = delayType;
- // All attempt queuing segments use the strict child span name "Attempt Delay".
- Span delaySpan = otelTracer.spanBuilder("Attempt Delay")
+ activeDelaySpan = otelTracer.spanBuilder("Delay")
.setParent(Context.current().with(span))
.setAttribute("grpc.delay_type", delayType)
.startSpan();
- activeDelaySpan = delaySpan;
- delaySpan.addEvent(
- "Delay state transition",
- io.opentelemetry.api.common.Attributes.of(
- AttributeKey.stringKey("grpc.delay_type"), delayType,
- AttributeKey.stringKey("grpc.delay_reason"), delayReason));
+ addDelayEvent(activeDelaySpan, delayReason);
}
@Override
- public void recordAttemptDelayReasonChanged(String delayReason) {
- if (!GrpcOpenTelemetry.isDelayObservabilityEnabled() || activeDelaySpan == null) {
+ public synchronized void recordDelayReasonChanged(String delayType, String delayReason) {
+ if (streamClosed || activeDelaySpan == null) {
return;
}
- String type = activeDelayType;
- activeDelaySpan.addEvent(
- "Delay state transition",
- io.opentelemetry.api.common.Attributes.of(
- AttributeKey.stringKey("grpc.delay_type"), type != null ? type : "",
- AttributeKey.stringKey("grpc.delay_reason"), delayReason));
+ addDelayEvent(activeDelaySpan, delayReason);
}
@Override
- public void recordAttemptDelayEnd() {
- Span delaySpan = activeDelaySpan;
- if (delaySpan != null) {
- // End active child span upon pick completion or transport cancellation.
- delaySpan.end();
+ public synchronized void recordDelayEnd(String delayType) {
+ endActiveDelaySpan();
+ }
+
+ @GuardedBy("this")
+ private void endActiveDelaySpan() {
+ if (activeDelaySpan != null) {
+ activeDelaySpan.end();
activeDelaySpan = null;
activeDelayType = null;
}
@@ -292,8 +354,12 @@ public void inboundUncompressedSize(long bytes) {
}
@Override
- public void streamClosed(io.grpc.Status status) {
- recordAttemptDelayEnd();
+ public synchronized void streamClosed(Status status) {
+ if (streamClosed) {
+ return;
+ }
+ streamClosed = true;
+ endActiveDelaySpan();
endSpanWithStatus(span, status);
}
}
@@ -516,9 +582,15 @@ public void onClose(io.grpc.Status status, Metadata trailers) {
// Receiving:
// |-- Event 'Inbound message received', attributes('sequence-numer' = 0,
// 'message-size' = 7854) ----|
+ private static void addDelayEvent(Span delaySpan, String delayReason) {
+ delaySpan.addEvent(
+ "Delay triggered",
+ Attributes.of(AttributeKey.stringKey("grpc.delay_reason"), delayReason));
+ }
+
private void recordOutboundMessageSentEvent(Span span,
int seqNo, long optionalWireSize, long optionalUncompressedSize) {
- AttributesBuilder attributesBuilder = io.opentelemetry.api.common.Attributes.builder();
+ AttributesBuilder attributesBuilder = Attributes.builder();
attributesBuilder.put("sequence-number", seqNo);
if (optionalUncompressedSize != -1) {
attributesBuilder.put("message-size", optionalUncompressedSize);
@@ -530,14 +602,14 @@ private void recordOutboundMessageSentEvent(Span span,
}
private void recordInboundCompressedMessage(Span span, int seqNo, long optionalWireSize) {
- AttributesBuilder attributesBuilder = io.opentelemetry.api.common.Attributes.builder();
+ AttributesBuilder attributesBuilder = Attributes.builder();
attributesBuilder.put("sequence-number", seqNo);
attributesBuilder.put("message-size-compressed", optionalWireSize);
span.addEvent("Inbound compressed message", attributesBuilder.build());
}
private void recordInboundMessageSize(Span span, int seqNo, long bytes) {
- AttributesBuilder attributesBuilder = io.opentelemetry.api.common.Attributes.builder();
+ AttributesBuilder attributesBuilder = Attributes.builder();
attributesBuilder.put("sequence-number", seqNo);
attributesBuilder.put("message-size", bytes);
span.addEvent("Inbound message", attributesBuilder.build());
diff --git a/opentelemetry/src/test/java/io/grpc/opentelemetry/GrpcOpenTelemetryTest.java b/opentelemetry/src/test/java/io/grpc/opentelemetry/GrpcOpenTelemetryTest.java
index 77eadf9ebbb..68c1cacb4f7 100644
--- a/opentelemetry/src/test/java/io/grpc/opentelemetry/GrpcOpenTelemetryTest.java
+++ b/opentelemetry/src/test/java/io/grpc/opentelemetry/GrpcOpenTelemetryTest.java
@@ -17,6 +17,8 @@
package io.grpc.opentelemetry;
import static com.google.common.truth.Truth.assertThat;
+import static java.nio.charset.StandardCharsets.UTF_8;
+import static java.util.Collections.emptyList;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.times;
@@ -24,27 +26,104 @@
import static org.mockito.Mockito.verifyNoMoreInteractions;
import com.google.common.collect.ImmutableList;
+import com.google.common.collect.ImmutableMap;
+import com.google.common.io.ByteStreams;
+import io.grpc.CallOptions;
+import io.grpc.ClientCall;
import io.grpc.ClientInterceptor;
+import io.grpc.ClientStreamTracer;
+import io.grpc.ConnectivityState;
+import io.grpc.EquivalentAddressGroup;
import io.grpc.ForwardingChannelBuilder2;
+import io.grpc.LoadBalancer;
+import io.grpc.LoadBalancerProvider;
+import io.grpc.LoadBalancerRegistry;
+import io.grpc.ManagedChannel;
import io.grpc.ManagedChannelBuilder;
+import io.grpc.Metadata;
+import io.grpc.MethodDescriptor;
import io.grpc.MetricSink;
+import io.grpc.NameResolver;
+import io.grpc.NameResolverProvider;
+import io.grpc.NameResolverRegistry;
import io.grpc.ServerBuilder;
+import io.grpc.ServerCall;
+import io.grpc.ServerCallHandler;
+import io.grpc.ServerServiceDefinition;
+import io.grpc.ServiceDescriptor;
+import io.grpc.Status;
+import io.grpc.StatusOr;
+import io.grpc.SynchronizationContext;
+import io.grpc.inprocess.InProcessChannelBuilder;
+import io.grpc.inprocess.InProcessServerBuilder;
+import io.grpc.inprocess.InProcessSocketAddress;
+import io.grpc.internal.FakeClock;
import io.grpc.internal.GrpcUtil;
import io.grpc.opentelemetry.GrpcOpenTelemetry.TargetFilter;
+import io.grpc.testing.GrpcCleanupRule;
import io.opentelemetry.api.OpenTelemetry;
+import io.opentelemetry.api.common.AttributeKey;
import io.opentelemetry.sdk.OpenTelemetrySdk;
import io.opentelemetry.sdk.metrics.SdkMeterProvider;
+import io.opentelemetry.sdk.metrics.data.MetricData;
+import io.opentelemetry.sdk.testing.assertj.OpenTelemetryAssertions;
import io.opentelemetry.sdk.testing.exporter.InMemoryMetricReader;
+import io.opentelemetry.sdk.testing.junit4.OpenTelemetryRule;
import io.opentelemetry.sdk.trace.SdkTracerProvider;
+import io.opentelemetry.sdk.trace.data.EventData;
+import io.opentelemetry.sdk.trace.data.SpanData;
+import java.io.ByteArrayInputStream;
+import java.io.IOException;
+import java.io.InputStream;
+import java.net.SocketAddress;
+import java.net.URI;
import java.util.Arrays;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.stream.Collectors;
import org.junit.After;
import org.junit.Before;
+import org.junit.Rule;
import org.junit.Test;
import org.junit.runner.RunWith;
import org.junit.runners.JUnit4;
@RunWith(JUnit4.class)
public class GrpcOpenTelemetryTest {
+ @Rule
+ public final OpenTelemetryRule openTelemetryRule = OpenTelemetryRule.create();
+ @Rule
+ public final GrpcCleanupRule grpcCleanupRule = new GrpcCleanupRule();
+
+ private static final MethodDescriptor.Marshaller MARSHALLER =
+ new MethodDescriptor.Marshaller() {
+ @Override
+ public InputStream stream(String value) {
+ return new ByteArrayInputStream(value.getBytes(UTF_8));
+ }
+
+ @Override
+ public String parse(InputStream stream) {
+ try {
+ return new String(ByteStreams.toByteArray(stream), UTF_8);
+ } catch (IOException ex) {
+ throw new RuntimeException(ex);
+ }
+ }
+ };
+
+ private final MethodDescriptor method =
+ MethodDescriptor.newBuilder()
+ .setType(MethodDescriptor.MethodType.UNARY)
+ .setRequestMarshaller(MARSHALLER)
+ .setResponseMarshaller(MARSHALLER)
+ .setFullMethodName("test.service/method")
+ .build();
+
private final InMemoryMetricReader inMemoryMetricReader = InMemoryMetricReader.create();
private final SdkMeterProvider meterProvider =
SdkMeterProvider.builder().registerMetricReader(inMemoryMetricReader).build();
@@ -179,6 +258,425 @@ public void configureChannelBuilder_registersMetricSink() {
assertThat(testBuilder.interceptorFactory).isNotNull();
}
+ @Test
+ public void delayHistograms_optedIn_recordedWithSpecAttributes() {
+ // gRFC A121 fixes the delay histogram label set to grpc.target, grpc.method and
+ // grpc.delay_type, and both histograms are opt-in. Drive the real metrics module against the
+ // real OpenTelemetry SDK and assert the emitted instruments match the spec.
+ OpenTelemetrySdk sdk = (OpenTelemetrySdk) openTelemetryRule.getOpenTelemetry();
+ OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(
+ sdk.getMeterProvider().get("grpc-java"),
+ ImmutableMap.of(
+ "grpc.client.attempt.delay.duration", true,
+ "grpc.client.call.delay.duration", true),
+ false);
+ assertThat(resource.clientAttemptDelayCounter()).isNotNull();
+ assertThat(resource.clientCallDelayCounter()).isNotNull();
+
+ OpenTelemetryMetricsModule module = new OpenTelemetryMetricsModule(
+ new FakeClock().getStopwatchSupplier(), resource, emptyList(), emptyList());
+ OpenTelemetryMetricsModule.CallAttemptsTracerFactory factory =
+ new OpenTelemetryMetricsModule.CallAttemptsTracerFactory(
+ module, "target:///", CallOptions.DEFAULT, method.getFullMethodName(),
+ emptyList(), io.opentelemetry.context.Context.root());
+
+ ClientStreamTracer delayTracer = factory.newClientStreamTracer(
+ ClientStreamTracer.StreamInfo.newBuilder().setCallOptions(CallOptions.DEFAULT).build(),
+ new Metadata());
+ delayTracer.recordDelayStart("connecting", "DNS server unreachable temporarily");
+ delayTracer.recordDelayEnd("connecting");
+ factory.recordDelayStart("resolving", "DNS resolution pending");
+ factory.recordDelayEnd("resolving");
+
+ OpenTelemetryAssertions.assertThat(openTelemetryRule.getMetrics())
+ .anySatisfy(
+ metric -> OpenTelemetryAssertions.assertThat(metric)
+ .hasName("grpc.client.attempt.delay.duration")
+ .hasDescription(
+ "EXPERIMENTAL. Time an RPC attempt spent waiting for a load balancing pick"
+ + " or connection establishment.")
+ .hasUnit("s")
+ .hasHistogramSatisfying(
+ histogram -> histogram.hasPointsSatisfying(
+ point -> point
+ .hasAttribute(AttributeKey.stringKey("grpc.target"), "target:///")
+ .hasAttribute(
+ AttributeKey.stringKey("grpc.method"), method.getFullMethodName())
+ .hasAttribute(
+ AttributeKey.stringKey("grpc.delay_type"), "connecting"))));
+ OpenTelemetryAssertions.assertThat(openTelemetryRule.getMetrics())
+ .anySatisfy(
+ metric -> OpenTelemetryAssertions.assertThat(metric)
+ .hasName("grpc.client.call.delay.duration")
+ .hasDescription(
+ "EXPERIMENTAL. Time an RPC spent waiting at the call level before an attempt"
+ + " was initiated, such as waiting for name resolution.")
+ .hasUnit("s")
+ .hasHistogramSatisfying(
+ histogram -> histogram.hasPointsSatisfying(
+ point -> point
+ .hasAttribute(AttributeKey.stringKey("grpc.target"), "target:///")
+ .hasAttribute(
+ AttributeKey.stringKey("grpc.method"), method.getFullMethodName())
+ .hasAttribute(
+ AttributeKey.stringKey("grpc.delay_type"), "resolving"))));
+ }
+
+ @Test
+ public void delayMetrics_notOptedIn_noInstrumentsAndNoMetrics() {
+ OpenTelemetrySdk sdk = (OpenTelemetrySdk) openTelemetryRule.getOpenTelemetry();
+ OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(
+ sdk.getMeterProvider().get("grpc-java"),
+ ImmutableMap.of(),
+ false);
+
+ assertThat(resource.clientAttemptDelayCounter()).isNull();
+ assertThat(resource.clientCallDelayCounter()).isNull();
+
+ OpenTelemetryMetricsModule module = new OpenTelemetryMetricsModule(
+ new FakeClock().getStopwatchSupplier(), resource, emptyList(), emptyList());
+ OpenTelemetryMetricsModule.CallAttemptsTracerFactory factory =
+ new OpenTelemetryMetricsModule.CallAttemptsTracerFactory(
+ module, "target:///", CallOptions.DEFAULT, method.getFullMethodName(),
+ emptyList(), io.opentelemetry.context.Context.root());
+
+ // Verify call delay methods execute cleanly when counters are null
+ factory.recordDelayStart("resolving", "resolving name");
+ factory.recordDelayEnd("resolving");
+ factory.callEnded(Status.OK, CallOptions.DEFAULT);
+
+ // Verify attempt delay methods execute cleanly when counters are null
+ ClientStreamTracer delayTracer = factory.newClientStreamTracer(
+ ClientStreamTracer.StreamInfo.newBuilder().setCallOptions(CallOptions.DEFAULT).build(),
+ new Metadata());
+ delayTracer.recordDelayStart("connecting", "connecting to backend");
+ delayTracer.recordDelayReasonChanged("connecting", "still connecting");
+ delayTracer.recordDelayEnd("connecting");
+ delayTracer.streamClosed(Status.OK);
+
+ for (MetricData m : openTelemetryRule.getMetrics()) {
+ assertThat(m.getName()).isNotIn(
+ ImmutableList.of(
+ "grpc.client.attempt.delay.duration", "grpc.client.call.delay.duration"));
+ }
+ }
+
+ @Test
+ public void delayObservability_endToEnd_fullCallAndAttemptLifecycle() throws Exception {
+ MethodDescriptor method =
+ this.method.toBuilder().setSampledToLocalTracing(true).build();
+ String serverName = InProcessServerBuilder.generateName();
+ String target = "testdelaye2e:///" + serverName;
+ ServerCallHandler handler = new ServerCallHandler() {
+ @Override
+ public ServerCall.Listener startCall(
+ ServerCall call, Metadata headers) {
+ call.sendHeaders(new Metadata());
+ call.sendMessage("response");
+ call.close(Status.OK, new Metadata());
+ return new ServerCall.Listener() {};
+ }
+ };
+ grpcCleanupRule.register(InProcessServerBuilder.forName(serverName)
+ .directExecutor()
+ .addService(ServerServiceDefinition.builder(
+ ServiceDescriptor.newBuilder("test.service").addMethod(method).build())
+ .addMethod(method, handler)
+ .build())
+ .build()
+ .start());
+
+ final AtomicReference listenerRef = new AtomicReference<>();
+ final AtomicReference syncContextRef = new AtomicReference<>();
+ final AtomicReference lbHelperRef = new AtomicReference<>();
+ final AtomicReference> addressesRef = new AtomicReference<>();
+ final AtomicReference subchannelRef = new AtomicReference<>();
+
+ NameResolverProvider resolverProvider = new NameResolverProvider() {
+ @Override
+ protected boolean isAvailable() {
+ return true;
+ }
+
+ @Override
+ protected int priority() {
+ return 5;
+ }
+
+ @Override
+ public String getDefaultScheme() {
+ return "testdelaye2e";
+ }
+
+ @Override
+ public Collection> getProducedSocketAddressTypes() {
+ return Collections.singleton(InProcessSocketAddress.class);
+ }
+
+ @Override
+ public NameResolver newNameResolver(URI targetUri, NameResolver.Args args) {
+ syncContextRef.set(args.getSynchronizationContext());
+ return new NameResolver() {
+ @Override
+ public String getServiceAuthority() {
+ return "localhost";
+ }
+
+ @Override
+ public void start(Listener2 listener) {
+ listenerRef.set(listener);
+ }
+
+ @Override
+ public void shutdown() {}
+ };
+ }
+ };
+
+ LoadBalancerProvider lbProvider = new LoadBalancerProvider() {
+ @Override
+ public boolean isAvailable() {
+ return true;
+ }
+
+ @Override
+ public int getPriority() {
+ return 5;
+ }
+
+ @Override
+ public String getPolicyName() {
+ return "test_delay_e2e_lb";
+ }
+
+ @Override
+ public LoadBalancer newLoadBalancer(LoadBalancer.Helper helper) {
+ lbHelperRef.set(helper);
+ return new LoadBalancer() {
+ @Override
+ public Status acceptResolvedAddresses(ResolvedAddresses resolvedAddresses) {
+ addressesRef.set(resolvedAddresses.getAddresses());
+ helper.updateBalancingState(
+ ConnectivityState.CONNECTING,
+ new FixedResultPicker(
+ PickResult.withNoResult(
+ "0:connecting", "waiting on priority group 0 (connecting)")));
+ return Status.OK;
+ }
+
+ @Override
+ public void handleNameResolutionError(Status error) {
+ helper.updateBalancingState(
+ ConnectivityState.TRANSIENT_FAILURE,
+ new FixedResultPicker(PickResult.withError(error)));
+ }
+
+ @Override
+ public void shutdown() {
+ if (subchannelRef.get() != null) {
+ subchannelRef.get().shutdown();
+ }
+ }
+ };
+ }
+ };
+
+ NameResolverRegistry.getDefaultRegistry().register(resolverProvider);
+ LoadBalancerRegistry.getDefaultRegistry().register(lbProvider);
+ try {
+ GrpcOpenTelemetry grpcOpenTelemetry = GrpcOpenTelemetry.newBuilder()
+ .sdk(openTelemetryRule.getOpenTelemetry())
+ .enableTracing(true)
+ .enableMetrics(ImmutableList.of(
+ "grpc.client.call.delay.duration", "grpc.client.attempt.delay.duration"))
+ .build();
+ InProcessChannelBuilder channelBuilder =
+ InProcessChannelBuilder.forTarget(target)
+ .defaultLoadBalancingPolicy("test_delay_e2e_lb")
+ .directExecutor();
+ grpcOpenTelemetry.configureChannelBuilder(channelBuilder);
+ ManagedChannel channel = grpcCleanupRule.register(channelBuilder.build());
+
+ final CountDownLatch closeLatch = new CountDownLatch(1);
+ final AtomicReference closeStatus = new AtomicReference<>();
+ final AtomicReference responseRef = new AtomicReference<>();
+ ClientCall call =
+ channel.newCall(method, CallOptions.DEFAULT.withWaitForReady());
+ call.start(new ClientCall.Listener() {
+ @Override
+ public void onMessage(String message) {
+ responseRef.set(message);
+ }
+
+ @Override
+ public void onClose(Status status, Metadata trailers) {
+ closeStatus.set(status);
+ closeLatch.countDown();
+ }
+ }, new Metadata());
+ call.sendMessage("request");
+ call.halfClose();
+ call.request(1);
+
+ // 1. Complete name resolution -> ends call-level "resolving" delay and starts attempt-level
+ // "0:connecting" delay.
+ assertThat(listenerRef.get()).isNotNull();
+ syncContextRef.get().execute(() -> listenerRef.get().onResult2(
+ NameResolver.ResolutionResult.newBuilder()
+ .setAddressesOrError(StatusOr.fromValue(ImmutableList.of(
+ new EquivalentAddressGroup(new InProcessSocketAddress(serverName)))))
+ .build()));
+
+ // 2. Same delayType ("0:connecting") with updated delayReason -> triggers
+ // recordDelayReasonChanged without closing the span or resetting the metric stopwatch.
+ syncContextRef.get().execute(() -> lbHelperRef.get().updateBalancingState(
+ ConnectivityState.CONNECTING,
+ new LoadBalancer.FixedResultPicker(
+ LoadBalancer.PickResult.withNoResult(
+ "0:connecting", "waiting on priority group 0 (subchannel connecting)"))));
+
+ // 3. Transition to TRANSIENT_FAILURE while wait-for-ready -> transitions attempt delay to
+ // "picker_failing_with_wait_for_ready".
+ syncContextRef.get().execute(() -> lbHelperRef.get().updateBalancingState(
+ ConnectivityState.TRANSIENT_FAILURE,
+ new LoadBalancer.FixedResultPicker(
+ LoadBalancer.PickResult.withError(
+ Status.UNAVAILABLE.withDescription("backend down")))));
+
+ // 4. Connect a real InProcess subchannel and transition to READY -> ends attempt delay and
+ // completes the RPC with Status.OK.
+ syncContextRef.get().execute(() -> {
+ LoadBalancer.Helper helper = lbHelperRef.get();
+ final LoadBalancer.Subchannel subchannel = helper.createSubchannel(
+ LoadBalancer.CreateSubchannelArgs.newBuilder()
+ .setAddresses(addressesRef.get())
+ .build());
+ subchannelRef.set(subchannel);
+ subchannel.start(stateInfo -> {
+ if (stateInfo.getState() == ConnectivityState.READY) {
+ helper.updateBalancingState(
+ ConnectivityState.READY,
+ new LoadBalancer.FixedResultPicker(
+ LoadBalancer.PickResult.withSubchannel(subchannel)));
+ }
+ });
+ subchannel.requestConnection();
+ });
+
+ assertThat(closeLatch.await(5, TimeUnit.SECONDS)).isTrue();
+ assertThat(closeStatus.get().getCode()).isEqualTo(Status.Code.OK);
+ assertThat(responseRef.get()).isEqualTo("response");
+
+ // Verify call-level delay histogram
+ OpenTelemetryAssertions.assertThat(openTelemetryRule.getMetrics())
+ .anySatisfy(
+ metric -> OpenTelemetryAssertions.assertThat(metric)
+ .hasName("grpc.client.call.delay.duration")
+ .hasHistogramSatisfying(
+ histogram -> histogram.hasPointsSatisfying(
+ point -> point
+ .hasCount(1)
+ .hasAttribute(AttributeKey.stringKey("grpc.target"), target)
+ .hasAttribute(
+ AttributeKey.stringKey("grpc.method"), method.getFullMethodName())
+ .hasAttribute(
+ AttributeKey.stringKey("grpc.delay_type"), "resolving"))));
+
+ // Verify attempt-level delay histogram points across all 3 delay types
+ OpenTelemetryAssertions.assertThat(openTelemetryRule.getMetrics())
+ .anySatisfy(
+ metric -> OpenTelemetryAssertions.assertThat(metric)
+ .hasName("grpc.client.attempt.delay.duration")
+ .hasHistogramSatisfying(
+ histogram -> histogram.hasPointsSatisfying(
+ point -> point
+ .hasCount(1)
+ .hasAttribute(AttributeKey.stringKey("grpc.target"), target)
+ .hasAttribute(
+ AttributeKey.stringKey("grpc.method"), method.getFullMethodName())
+ .hasAttribute(
+ AttributeKey.stringKey("grpc.delay_type"), "0:connecting"),
+ point -> point
+ .hasCount(1)
+ .hasAttribute(AttributeKey.stringKey("grpc.target"), target)
+ .hasAttribute(
+ AttributeKey.stringKey("grpc.method"), method.getFullMethodName())
+ .hasAttribute(
+ AttributeKey.stringKey("grpc.delay_type"),
+ "picker_failing_with_wait_for_ready"),
+ point -> point
+ .hasCount(1)
+ .hasAttribute(AttributeKey.stringKey("grpc.target"), target)
+ .hasAttribute(
+ AttributeKey.stringKey("grpc.method"), method.getFullMethodName())
+ .hasAttribute(
+ AttributeKey.stringKey("grpc.delay_type"), "connecting"))));
+
+ // Verify trace spans and events
+ List spans = openTelemetryRule.getSpans();
+ SpanData callSpan = spans.stream()
+ .filter(s -> s.getName().equals("Sent.test.service.method"))
+ .findFirst()
+ .orElseThrow(AssertionError::new);
+ SpanData attemptSpan = spans.stream()
+ .filter(s -> s.getName().equals("Attempt.test.service.method"))
+ .findFirst()
+ .orElseThrow(AssertionError::new);
+ List delaySpans = spans.stream()
+ .filter(s -> s.getName().equals("Delay"))
+ .collect(Collectors.toList());
+ assertThat(delaySpans).hasSize(4);
+
+ SpanData resolvingSpan = delaySpans.get(0);
+ assertThat(resolvingSpan.getParentSpanId()).isEqualTo(callSpan.getSpanId());
+ assertThat(resolvingSpan.getAttributes().get(AttributeKey.stringKey("grpc.delay_type")))
+ .isEqualTo("resolving");
+ assertThat(resolvingSpan.getEvents().stream()
+ .map(e -> e.getAttributes().get(AttributeKey.stringKey("grpc.delay_reason")))
+ .collect(Collectors.toList()))
+ .containsExactly("waiting for name resolution or service config");
+
+ SpanData initialConnectingSpan = delaySpans.get(1);
+ assertThat(initialConnectingSpan.getParentSpanId()).isEqualTo(attemptSpan.getSpanId());
+ assertThat(
+ initialConnectingSpan.getAttributes().get(AttributeKey.stringKey("grpc.delay_type")))
+ .isEqualTo("connecting");
+ assertThat(initialConnectingSpan.getEvents().stream()
+ .map(e -> e.getAttributes().get(AttributeKey.stringKey("grpc.delay_reason")))
+ .collect(Collectors.toList()))
+ .containsExactly("client channel: waiting for picker");
+
+ SpanData priorityConnectingSpan = delaySpans.get(2);
+ assertThat(priorityConnectingSpan.getParentSpanId()).isEqualTo(attemptSpan.getSpanId());
+ assertThat(
+ priorityConnectingSpan.getAttributes().get(AttributeKey.stringKey("grpc.delay_type")))
+ .isEqualTo("0:connecting");
+ assertThat(priorityConnectingSpan.getEvents().stream()
+ .map(EventData::getName)
+ .collect(Collectors.toList()))
+ .containsExactly("Delay triggered", "Delay triggered");
+ assertThat(priorityConnectingSpan.getEvents().stream()
+ .map(e -> e.getAttributes().get(AttributeKey.stringKey("grpc.delay_reason")))
+ .collect(Collectors.toList()))
+ .containsExactly(
+ "waiting on priority group 0 (connecting)",
+ "waiting on priority group 0 (subchannel connecting)")
+ .inOrder();
+
+ SpanData waitForReadySpan = delaySpans.get(3);
+ assertThat(waitForReadySpan.getParentSpanId()).isEqualTo(attemptSpan.getSpanId());
+ assertThat(waitForReadySpan.getAttributes().get(AttributeKey.stringKey("grpc.delay_type")))
+ .isEqualTo("picker_failing_with_wait_for_ready");
+ assertThat(waitForReadySpan.getEvents().stream()
+ .map(e -> e.getAttributes().get(AttributeKey.stringKey("grpc.delay_reason")))
+ .collect(Collectors.toList()))
+ .containsExactly("wait_for_ready RPC failed with status: UNAVAILABLE: backend down");
+ } finally {
+ LoadBalancerRegistry.getDefaultRegistry().deregister(lbProvider);
+ NameResolverRegistry.getDefaultRegistry().deregister(resolverProvider);
+ }
+ }
+
private static class TestChannelBuilder extends ForwardingChannelBuilder2 {
Object interceptorFactory;
MetricSink metricSink;
diff --git a/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryMetricsModuleTest.java b/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryMetricsModuleTest.java
index bd613888f94..ccff3210d88 100644
--- a/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryMetricsModuleTest.java
+++ b/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryMetricsModuleTest.java
@@ -45,6 +45,9 @@
import io.grpc.Grpc;
import io.grpc.KnownLength;
import io.grpc.LoadBalancer;
+import io.grpc.LoadBalancer.PickResult;
+import io.grpc.LoadBalancer.PickSubchannelArgs;
+import io.grpc.LoadBalancer.SubchannelPicker;
import io.grpc.LoadBalancerProvider;
import io.grpc.LoadBalancerRegistry;
import io.grpc.ManagedChannel;
@@ -219,7 +222,6 @@ public String parse(InputStream stream) {
@Before
public void setUp() throws Exception {
- System.setProperty("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY", "true");
testMeter = openTelemetryTesting.getOpenTelemetry()
.getMeter(OpenTelemetryConstants.INSTRUMENTATION_SCOPE);
@@ -227,7 +229,6 @@ public void setUp() throws Exception {
@After
public void tearDown() {
- System.clearProperty("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY");
if (channel != null) {
channel.shutdownNow();
}
@@ -1644,9 +1645,9 @@ public void clientAttemptDelayDuration_recorded() {
ClientStreamTracer tracer =
callAttemptsTracerFactory.newClientStreamTracer(STREAM_INFO, new Metadata());
- tracer.recordAttemptDelayStart("connecting", "connecting reason");
+ tracer.recordDelayStart("connecting", "connecting reason");
fakeClock.forwardTime(250, TimeUnit.MILLISECONDS);
- tracer.recordAttemptDelayEnd();
+ tracer.recordDelayEnd("connecting");
assertThat(openTelemetryTesting.getMetrics())
.anySatisfy(
@@ -1663,6 +1664,179 @@ public void clientAttemptDelayDuration_recorded() {
})));
}
+
+ @Test
+ public void clientCallDelayDuration_recorded() {
+ Map enabledMetrics = ImmutableMap.of(
+ "grpc.client.call.delay.duration", true
+ );
+ OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(
+ testMeter, enabledMetrics, disableDefaultMetrics);
+ OpenTelemetryMetricsModule module = new OpenTelemetryMetricsModule(
+ fakeClock.getStopwatchSupplier(), resource, emptyList(), emptyList());
+ CallAttemptsTracerFactory callAttemptsTracerFactory =
+ new CallAttemptsTracerFactory(
+ module, "target:///", STREAM_INFO.getCallOptions(), method.getFullMethodName(),
+ emptyList(), Context.root());
+
+ callAttemptsTracerFactory.recordDelayStart("resolving", "dns resolution pending");
+ fakeClock.forwardTime(500, TimeUnit.MILLISECONDS);
+ callAttemptsTracerFactory.recordDelayEnd("resolving");
+
+ assertThat(openTelemetryTesting.getMetrics())
+ .anySatisfy(
+ metric -> assertThat(metric)
+ .hasName("grpc.client.call.delay.duration")
+ .hasHistogramSatisfying(
+ histogram -> histogram.hasPointsSatisfying(
+ point -> {
+ point.hasSum(0.5);
+ point.hasAttribute(
+ AttributeKey.stringKey("grpc.delay_type"), "resolving");
+ })));
+ }
+
+ @Test
+ public void clientCallDelayDuration_defensiveCleanupOnAbruptCallEnded() {
+ Map enabledMetrics = ImmutableMap.of(
+ "grpc.client.call.delay.duration", true
+ );
+ OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(
+ testMeter, enabledMetrics, disableDefaultMetrics);
+ OpenTelemetryMetricsModule module = new OpenTelemetryMetricsModule(
+ fakeClock.getStopwatchSupplier(), resource, emptyList(), emptyList());
+ CallAttemptsTracerFactory callAttemptsTracerFactory =
+ new CallAttemptsTracerFactory(
+ module, "target:///", STREAM_INFO.getCallOptions(), method.getFullMethodName(),
+ emptyList(), Context.root());
+
+ callAttemptsTracerFactory.recordDelayStart("resolving", "dns resolution pending");
+ fakeClock.forwardTime(500, TimeUnit.MILLISECONDS);
+ // Call ends abruptly without prior recordDelayEnd("resolving")
+ callAttemptsTracerFactory.callEnded(
+ Status.CANCELLED.withDescription("abrupt cancellation"), STREAM_INFO.getCallOptions());
+
+ assertThat(openTelemetryTesting.getMetrics())
+ .anySatisfy(
+ metric -> assertThat(metric)
+ .hasName("grpc.client.call.delay.duration")
+ .hasHistogramSatisfying(
+ histogram -> histogram.hasPointsSatisfying(
+ point -> {
+ point.hasSum(0.5);
+ point.hasAttribute(METHOD_KEY, method.getFullMethodName());
+ point.hasAttribute(TARGET_KEY, "target:///");
+ point.hasAttribute(
+ AttributeKey.stringKey("grpc.delay_type"), "resolving");
+ })));
+
+ // Ensure subsequent calls to recordDelayEnd are safe no-ops
+ callAttemptsTracerFactory.recordDelayEnd("resolving");
+ // Ensure subsequent calls to recordDelayStart are rejected after callEnded
+ callAttemptsTracerFactory.recordDelayStart("resolving", "late start attempt");
+ fakeClock.forwardTime(200, TimeUnit.MILLISECONDS);
+ callAttemptsTracerFactory.recordDelayEnd("resolving");
+
+ // Metric sum remains 0.5, no additional recordings
+ assertThat(openTelemetryTesting.getMetrics())
+ .anySatisfy(
+ metric -> assertThat(metric)
+ .hasName("grpc.client.call.delay.duration")
+ .hasHistogramSatisfying(
+ histogram -> histogram.hasPointsSatisfying(
+ point -> point.hasSum(0.5))));
+ }
+
+ @Test
+ public void clientCallDelayDuration_endToEnd_nameResolutionDelay() throws Exception {
+ final CountDownLatch resolutionLatch = new CountDownLatch(1);
+ final AtomicReference capturedListener = new AtomicReference<>();
+
+ NameResolverProvider slowResolverProvider = new NameResolverProvider() {
+ @Override
+ protected boolean isAvailable() {
+ return true;
+ }
+
+ @Override
+ protected int priority() {
+ return 5;
+ }
+
+ @Override
+ public String getDefaultScheme() {
+ return "slowresmetric";
+ }
+
+ @Override
+ public Collection> getProducedSocketAddressTypes() {
+ return Collections.singleton(InProcessSocketAddress.class);
+ }
+
+ @Override
+ public NameResolver newNameResolver(URI targetUri, NameResolver.Args args) {
+ return new NameResolver() {
+ @Override
+ public String getServiceAuthority() {
+ return "slowresmetric";
+ }
+
+ @Override
+ public void start(Listener2 listener) {
+ capturedListener.set(listener);
+ resolutionLatch.countDown();
+ }
+
+ @Override
+ public void shutdown() {}
+ };
+ }
+ };
+ NameResolverRegistry.getDefaultRegistry().register(slowResolverProvider);
+
+ GrpcOpenTelemetry grpcOpenTelemetry = GrpcOpenTelemetry.newBuilder()
+ .sdk(openTelemetryTesting.getOpenTelemetry())
+ .enableMetrics(Collections.singleton("grpc.client.call.delay.duration"))
+ .build();
+
+ InProcessChannelBuilder channelBuilder =
+ InProcessChannelBuilder.forTarget("slowresmetric:///test-metric-service")
+ .defaultLoadBalancingPolicy("pick_first");
+ grpcOpenTelemetry.configureChannelBuilder(channelBuilder);
+ ManagedChannel channel = channelBuilder.build();
+ try {
+ ClientCall call = channel.newCall(method, CallOptions.DEFAULT);
+ call.start(new ClientCall.Listener() {}, new Metadata());
+ call.request(1);
+
+ resolutionLatch.await(5, TimeUnit.SECONDS);
+
+ // Complete name resolution
+ capturedListener.get().onResult(NameResolver.ResolutionResult.newBuilder()
+ .setAddressesOrError(StatusOr.fromValue(Collections.singletonList(
+ new EquivalentAddressGroup(new InProcessSocketAddress("test-slow-metric")))))
+ .build());
+
+ call.cancel("End test", null);
+ } finally {
+ channel.shutdownNow();
+ channel.awaitTermination(5, TimeUnit.SECONDS);
+ NameResolverRegistry.getDefaultRegistry().deregister(slowResolverProvider);
+ }
+
+ assertThat(openTelemetryTesting.getMetrics())
+ .anySatisfy(
+ metric -> assertThat(metric)
+ .hasName("grpc.client.call.delay.duration")
+ .hasHistogramSatisfying(
+ histogram -> histogram.hasPointsSatisfying(
+ point -> {
+ point.hasAttribute(METHOD_KEY, method.getFullMethodName());
+ point.hasAttribute(
+ AttributeKey.stringKey("grpc.delay_type"), "resolving");
+ })));
+ }
+
@Test
public void clientAttemptDelayDuration_endToEnd_inProcessTransport() throws Exception {
final CountDownLatch latch = new CountDownLatch(1);
@@ -1791,36 +1965,6 @@ public void shutdown() {}
})));
}
- @Test
- public void clientAttemptDelayStart_featureFlagDisabled_zeroMetrics() {
- System.setProperty("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY", "false");
- try {
- Map enabledMetrics = ImmutableMap.of(
- "grpc.client.attempt.delay.duration", true
- );
- OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(
- testMeter, enabledMetrics, disableDefaultMetrics);
- OpenTelemetryMetricsModule module = new OpenTelemetryMetricsModule(
- fakeClock.getStopwatchSupplier(), resource, emptyList(), emptyList());
- CallAttemptsTracerFactory callAttemptsTracerFactory =
- new CallAttemptsTracerFactory(
- module, "target:///", STREAM_INFO.getCallOptions(), method.getFullMethodName(),
- emptyList(), Context.root());
-
- ClientStreamTracer tracer =
- callAttemptsTracerFactory.newClientStreamTracer(STREAM_INFO, new Metadata());
- tracer.recordAttemptDelayStart("connecting", "connecting reason");
- fakeClock.forwardTime(250, TimeUnit.MILLISECONDS);
- tracer.recordAttemptDelayEnd();
-
- assertThat(openTelemetryTesting.getMetrics())
- .extracting("name")
- .doesNotContain("grpc.client.attempt.delay.duration");
- } finally {
- System.setProperty("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY", "true");
- }
- }
-
@Test
public void serverBasicMetrics() {
OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(testMeter,
@@ -2364,7 +2508,183 @@ public ClientCall interceptCall(
assertEquals(
"baggage-val-1", capturedBaggage.getEntryValue("baggage-key-1"));
}
-
+
+ @Test
+ public void clientCallDelayDuration_sameDelayType_doesNotResetStopwatch() {
+ Map enabledMetrics = ImmutableMap.of(
+ "grpc.client.call.delay.duration", true
+ );
+ OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(testMeter,
+ enabledMetrics, disableDefaultMetrics);
+ OpenTelemetryMetricsModule module = newOpenTelemetryMetricsModule(resource);
+ OpenTelemetryMetricsModule.CallAttemptsTracerFactory callAttemptsTracerFactory =
+ new CallAttemptsTracerFactory(module, "target:///", CALL_OPTIONS,
+ method.getFullMethodName(), emptyList(), Context.root());
+
+ callAttemptsTracerFactory.recordDelayStart("resolving", "reason1");
+ fakeClock.forwardTime(100, TimeUnit.MILLISECONDS);
+ // Same delay type: the running stopwatch is kept.
+ callAttemptsTracerFactory.recordDelayStart("resolving", "reason2");
+ fakeClock.forwardTime(100, TimeUnit.MILLISECONDS);
+ // Different delay type: the current delay is recorded and a new one is started.
+ callAttemptsTracerFactory.recordDelayStart("connecting", "transition to connecting");
+ fakeClock.forwardTime(50, TimeUnit.MILLISECONDS);
+ callAttemptsTracerFactory.recordDelayEnd("connecting");
+
+ assertThat(openTelemetryTesting.getMetrics())
+ .anySatisfy(
+ metric -> assertThat(metric)
+ .hasName("grpc.client.call.delay.duration")
+ .hasHistogramSatisfying(
+ histogram -> histogram.hasPointsSatisfying(
+ point -> {
+ point.hasSum(0.2);
+ point.hasAttribute(
+ AttributeKey.stringKey("grpc.delay_type"), "resolving");
+ },
+ point -> {
+ point.hasSum(0.05);
+ point.hasAttribute(
+ AttributeKey.stringKey("grpc.delay_type"), "connecting");
+ })));
+ }
+
+ @Test
+ public void clientAttemptDelayDuration_sameDelayType_doesNotResetStopwatch() {
+ Map enabledMetrics = ImmutableMap.of(
+ "grpc.client.attempt.delay.duration", true
+ );
+ OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(testMeter,
+ enabledMetrics, disableDefaultMetrics);
+ OpenTelemetryMetricsModule module = newOpenTelemetryMetricsModule(resource);
+ OpenTelemetryMetricsModule.CallAttemptsTracerFactory callAttemptsTracerFactory =
+ new CallAttemptsTracerFactory(module, "target:///", CALL_OPTIONS,
+ method.getFullMethodName(), emptyList(), Context.root());
+ ClientStreamTracer tracer = callAttemptsTracerFactory.newClientStreamTracer(
+ ClientStreamTracer.StreamInfo.newBuilder().build(), new Metadata());
+
+ tracer.recordDelayStart("connecting", "reason1");
+ fakeClock.forwardTime(100, TimeUnit.MILLISECONDS);
+ // Same delay type: the running stopwatch is kept.
+ tracer.recordDelayStart("connecting", "reason2");
+ fakeClock.forwardTime(100, TimeUnit.MILLISECONDS);
+ // Different delay type: the current delay is recorded and a new one is started.
+ tracer.recordDelayStart("0:connecting", "transition to different delay type");
+ fakeClock.forwardTime(50, TimeUnit.MILLISECONDS);
+ tracer.recordDelayEnd("0:connecting");
+
+ assertThat(openTelemetryTesting.getMetrics())
+ .anySatisfy(
+ metric -> assertThat(metric)
+ .hasName("grpc.client.attempt.delay.duration")
+ .hasHistogramSatisfying(
+ histogram -> histogram.hasPointsSatisfying(
+ point -> {
+ point.hasSum(0.2);
+ point.hasAttribute(
+ AttributeKey.stringKey("grpc.delay_type"), "connecting");
+ },
+ point -> {
+ point.hasSum(0.05);
+ point.hasAttribute(
+ AttributeKey.stringKey("grpc.delay_type"), "0:connecting");
+ })));
+ }
+
+ @Test
+ public void clientAttemptDelayDuration_excludesOptionalLabels() {
+ Map enabledMetrics = ImmutableMap.of(
+ "grpc.client.attempt.delay.duration", true
+ );
+ OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(
+ testMeter, enabledMetrics, disableDefaultMetrics);
+ OpenTelemetryMetricsModule module = new OpenTelemetryMetricsModule(
+ fakeClock.getStopwatchSupplier(),
+ resource,
+ Arrays.asList("grpc.lb.locality", "grpc.lb.backend_service"),
+ emptyList());
+ CallAttemptsTracerFactory callAttemptsTracerFactory =
+ new CallAttemptsTracerFactory(
+ module, "target:///", STREAM_INFO.getCallOptions(), method.getFullMethodName(),
+ emptyList(), Context.root());
+
+ ClientStreamTracer tracer =
+ callAttemptsTracerFactory.newClientStreamTracer(STREAM_INFO, new Metadata());
+ tracer.addOptionalLabel("grpc.lb.locality", "us-east1-a");
+ tracer.addOptionalLabel("grpc.lb.backend_service", "backend-service-1");
+ tracer.recordDelayStart("connecting", "connecting reason");
+ fakeClock.forwardTime(250, TimeUnit.MILLISECONDS);
+ tracer.recordDelayEnd("connecting");
+
+ // Per gRFC A121 the delay histogram carries only target, method and delay_type, even when
+ // the per-attempt optional labels are enabled.
+ assertThat(openTelemetryTesting.getMetrics())
+ .anySatisfy(
+ metric -> assertThat(metric)
+ .hasName("grpc.client.attempt.delay.duration")
+ .hasHistogramSatisfying(
+ histogram -> histogram.hasPointsSatisfying(
+ point -> {
+ point.hasSum(0.25);
+ point.hasAttributes(
+ io.opentelemetry.api.common.Attributes.of(
+ METHOD_KEY, method.getFullMethodName(),
+ TARGET_KEY, "target:///",
+ AttributeKey.stringKey("grpc.delay_type"), "connecting"));
+ })));
+ }
+
+ @Test
+ public void clientAttemptDelayStart_afterStreamClosed_noOp() {
+ Map enabledMetrics = ImmutableMap.of(
+ "grpc.client.attempt.delay.duration", true
+ );
+ OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(
+ testMeter, enabledMetrics, disableDefaultMetrics);
+ OpenTelemetryMetricsModule module = new OpenTelemetryMetricsModule(
+ fakeClock.getStopwatchSupplier(), resource, emptyList(), emptyList());
+ CallAttemptsTracerFactory callAttemptsTracerFactory =
+ new CallAttemptsTracerFactory(
+ module, "target:///", STREAM_INFO.getCallOptions(), method.getFullMethodName(),
+ emptyList(), Context.root());
+ ClientStreamTracer tracer =
+ callAttemptsTracerFactory.newClientStreamTracer(STREAM_INFO, new Metadata());
+
+ tracer.streamClosed(Status.OK);
+ tracer.recordDelayStart("connecting", "post-close");
+ tracer.recordDelayReasonChanged("connecting", "changed");
+ fakeClock.forwardTime(250, TimeUnit.MILLISECONDS);
+ tracer.recordDelayEnd("connecting");
+ callAttemptsTracerFactory.callEnded(Status.OK, STREAM_INFO.getCallOptions());
+
+ assertThat(openTelemetryTesting.getMetrics())
+ .noneSatisfy(metric -> assertThat(metric).hasName("grpc.client.attempt.delay.duration"));
+ }
+
+ @Test
+ public void clientCallDelayStart_afterCallEnded_noOp() {
+ Map enabledMetrics = ImmutableMap.of(
+ "grpc.client.call.delay.duration", true
+ );
+ OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(
+ testMeter, enabledMetrics, disableDefaultMetrics);
+ OpenTelemetryMetricsModule module = new OpenTelemetryMetricsModule(
+ fakeClock.getStopwatchSupplier(), resource, emptyList(), emptyList());
+ CallAttemptsTracerFactory callAttemptsTracerFactory =
+ new CallAttemptsTracerFactory(
+ module, "target:///", STREAM_INFO.getCallOptions(), method.getFullMethodName(),
+ emptyList(), Context.root());
+
+ callAttemptsTracerFactory.callEnded(Status.OK, STREAM_INFO.getCallOptions());
+ callAttemptsTracerFactory.recordDelayStart("resolving", "post-close");
+ callAttemptsTracerFactory.recordDelayReasonChanged("resolving", "changed");
+ fakeClock.forwardTime(250, TimeUnit.MILLISECONDS);
+ callAttemptsTracerFactory.recordDelayEnd("resolving");
+
+ assertThat(openTelemetryTesting.getMetrics())
+ .noneSatisfy(metric -> assertThat(metric).hasName("grpc.client.call.delay.duration"));
+ }
+
private static List sortByName(List metrics) {
metrics.sort((m1, m2) -> m1.getName().compareTo(m2.getName()));
return metrics;
diff --git a/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryTracingModuleTest.java b/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryTracingModuleTest.java
index ee7e86e05cc..e5b4971f7bd 100644
--- a/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryTracingModuleTest.java
+++ b/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryTracingModuleTest.java
@@ -20,6 +20,7 @@
import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.BAGGAGE_KEY;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
+import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertSame;
import static org.junit.Assert.assertTrue;
import static org.mockito.ArgumentMatchers.any;
@@ -46,6 +47,9 @@
import io.grpc.EquivalentAddressGroup;
import io.grpc.KnownLength;
import io.grpc.LoadBalancer;
+import io.grpc.LoadBalancer.PickResult;
+import io.grpc.LoadBalancer.PickSubchannelArgs;
+import io.grpc.LoadBalancer.SubchannelPicker;
import io.grpc.LoadBalancerProvider;
import io.grpc.LoadBalancerRegistry;
import io.grpc.ManagedChannel;
@@ -202,7 +206,6 @@ public String parse(InputStream stream) {
@Before
public void setUp() {
- System.setProperty("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY", "true");
tracerRule = openTelemetryRule.getOpenTelemetry().getTracer(
OpenTelemetryConstants.INSTRUMENTATION_SCOPE);
TracerProvider mockTracerProvider = mock(TracerProvider.class);
@@ -220,7 +223,6 @@ public void setUp() {
@After
public void tearDown() {
- System.clearProperty("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY");
}
// Use mock instead of OpenTelemetryRule to verify inOrder and propagator.
@@ -298,6 +300,294 @@ public void clientBasicTracingMocking() {
inOrder.verifyNoMoreInteractions();
}
+ @Test
+ public void clientDelayTracingMocking() {
+ Span mockDelaySpan = mock(Span.class);
+ when(mockSpanBuilder.setAttribute(
+ org.mockito.ArgumentMatchers.anyString(),
+ org.mockito.ArgumentMatchers.anyString()))
+ .thenReturn(mockSpanBuilder);
+ when(mockSpanBuilder.startSpan()).thenReturn(mockAttemptSpan, mockDelaySpan);
+
+ OpenTelemetryTracingModule tracingModule = new OpenTelemetryTracingModule(mockOpenTelemetry);
+ CallAttemptsTracerFactory callTracer =
+ tracingModule.newClientCallTracer(mockClientSpan, method);
+ ClientStreamTracer clientStreamTracer =
+ callTracer.newClientStreamTracer(STREAM_INFO, new Metadata());
+
+ clientStreamTracer.recordDelayStart("connecting", "pick_first: attempting to connect");
+ clientStreamTracer.recordDelayEnd("connecting");
+
+ verify(mockTracer).spanBuilder(eq("Delay"));
+ verify(mockSpanBuilder).setAttribute(eq("grpc.delay_type"), eq("connecting"));
+ verify(mockDelaySpan).addEvent(
+ eq("Delay triggered"),
+ org.mockito.ArgumentMatchers.any());
+ verify(mockDelaySpan).end();
+ }
+
+ @Test
+ public void clientCallDelayTracingMocking() {
+ Span mockDelaySpan = mock(Span.class);
+ when(mockSpanBuilder.setAttribute(
+ org.mockito.ArgumentMatchers.anyString(),
+ org.mockito.ArgumentMatchers.anyString()))
+ .thenReturn(mockSpanBuilder);
+ when(mockSpanBuilder.startSpan()).thenReturn(mockDelaySpan);
+
+ OpenTelemetryTracingModule tracingModule = new OpenTelemetryTracingModule(mockOpenTelemetry);
+ CallAttemptsTracerFactory callTracer =
+ tracingModule.newClientCallTracer(mockClientSpan, method);
+
+ callTracer.recordDelayStart("resolving", "waiting for DNS query");
+ callTracer.recordDelayEnd("resolving");
+
+ verify(mockTracer).spanBuilder(eq("Delay"));
+ verify(mockSpanBuilder).setAttribute(eq("grpc.delay_type"), eq("resolving"));
+ verify(mockDelaySpan).addEvent(
+ eq("Delay triggered"),
+ org.mockito.ArgumentMatchers.any());
+ verify(mockDelaySpan).end();
+ }
+
+ @Test
+ public void clientCallDelayTracing_endToEnd_nameResolutionDelay() throws Exception {
+ final CountDownLatch resolutionLatch = new CountDownLatch(1);
+ final AtomicReference capturedListener = new AtomicReference<>();
+
+ NameResolverProvider slowResolverProvider = new NameResolverProvider() {
+ @Override
+ protected boolean isAvailable() {
+ return true;
+ }
+
+ @Override
+ protected int priority() {
+ return 5;
+ }
+
+ @Override
+ public String getDefaultScheme() {
+ return "slowres";
+ }
+
+ @Override
+ public Collection> getProducedSocketAddressTypes() {
+ return Collections.singleton(InProcessSocketAddress.class);
+ }
+
+ @Override
+ public NameResolver newNameResolver(URI targetUri, NameResolver.Args args) {
+ return new NameResolver() {
+ @Override
+ public String getServiceAuthority() {
+ return "slowres";
+ }
+
+ @Override
+ public void start(Listener2 listener) {
+ capturedListener.set(listener);
+ resolutionLatch.countDown();
+ }
+
+ @Override
+ public void shutdown() {}
+ };
+ }
+ };
+ NameResolverRegistry.getDefaultRegistry().register(slowResolverProvider);
+
+ GrpcOpenTelemetry grpcOpenTelemetry = GrpcOpenTelemetry.newBuilder()
+ .sdk(openTelemetryRule.getOpenTelemetry())
+ .enableTracing(true)
+ .build();
+
+ InProcessChannelBuilder channelBuilder =
+ InProcessChannelBuilder.forTarget("slowres:///test-service")
+ .defaultLoadBalancingPolicy("pick_first");
+ grpcOpenTelemetry.configureChannelBuilder(channelBuilder);
+ ManagedChannel channel = channelBuilder.build();
+ try {
+ ClientCall call = channel.newCall(method, CallOptions.DEFAULT);
+ call.start(new ClientCall.Listener() {}, new Metadata());
+ call.request(1);
+
+ resolutionLatch.await(5, TimeUnit.SECONDS);
+
+ // Now complete name resolution
+ capturedListener.get().onResult(NameResolver.ResolutionResult.newBuilder()
+ .setAddressesOrError(StatusOr.fromValue(Collections.singletonList(
+ new EquivalentAddressGroup(new InProcessSocketAddress("test-slow-res")))))
+ .build());
+
+ call.cancel("End test", null);
+ } finally {
+ channel.shutdownNow();
+ channel.awaitTermination(5, TimeUnit.SECONDS);
+ NameResolverRegistry.getDefaultRegistry().deregister(slowResolverProvider);
+ }
+
+ List spans = openTelemetryRule.getSpans();
+ SpanData callDelaySpan = null;
+ for (SpanData s : spans) {
+ if ("Delay".equals(s.getName())) {
+ callDelaySpan = s;
+ break;
+ }
+ }
+ assertNotNull(callDelaySpan);
+ assertEquals("resolving",
+ callDelaySpan.getAttributes().get(AttributeKey.stringKey("grpc.delay_type")));
+
+ boolean foundTransition = false;
+ for (EventData event : callDelaySpan.getEvents()) {
+ if ("Delay triggered".equals(event.getName())
+ && "waiting for name resolution or service config".equals(
+ event.getAttributes().get(AttributeKey.stringKey("grpc.delay_reason")))) {
+ foundTransition = true;
+ break;
+ }
+ }
+ assertTrue(foundTransition);
+ }
+
+ @Test
+ public void clientAttemptDelayTracing_endToEnd_inProcessTransport() throws Exception {
+ final CountDownLatch latch = new CountDownLatch(1);
+ LoadBalancerProvider slowLbProvider = new LoadBalancerProvider() {
+ @Override
+ public boolean isAvailable() {
+ return true;
+ }
+
+ @Override
+ public int getPriority() {
+ return 5;
+ }
+
+ @Override
+ public String getPolicyName() {
+ return "slow_connecting_policy";
+ }
+
+ @Override
+ public LoadBalancer newLoadBalancer(LoadBalancer.Helper helper) {
+ return new LoadBalancer() {
+ @Override
+ public Status acceptResolvedAddresses(LoadBalancer.ResolvedAddresses resolvedAddresses) {
+ helper.updateBalancingState(ConnectivityState.CONNECTING, new SubchannelPicker() {
+ @Override
+ public PickResult pickSubchannel(PickSubchannelArgs args) {
+ latch.countDown();
+ return PickResult.withNoResult("connecting",
+ "Simulated slow TLS handshake with backend");
+ }
+ });
+ return Status.OK;
+ }
+
+ @Override
+ public void handleNameResolutionError(Status error) {}
+
+ @Override
+ public void shutdown() {}
+ };
+ }
+ };
+ LoadBalancerRegistry.getDefaultRegistry().register(slowLbProvider);
+
+ NameResolverProvider customResolverProvider = new NameResolverProvider() {
+ @Override
+ protected boolean isAvailable() {
+ return true;
+ }
+
+ @Override
+ protected int priority() {
+ return 5;
+ }
+
+ @Override
+ public String getDefaultScheme() {
+ return "inproce2e";
+ }
+
+ @Override
+ public Collection> getProducedSocketAddressTypes() {
+ return Collections.singleton(InProcessSocketAddress.class);
+ }
+
+ @Override
+ public NameResolver newNameResolver(URI targetUri, NameResolver.Args args) {
+ return new NameResolver() {
+ @Override
+ public String getServiceAuthority() {
+ return "inproce2e";
+ }
+
+ @Override
+ public void start(Listener2 listener) {
+ listener.onResult(NameResolver.ResolutionResult.newBuilder()
+ .setAddressesOrError(StatusOr.fromValue(Collections.singletonList(
+ new EquivalentAddressGroup(new InProcessSocketAddress("test-e2e")))))
+ .build());
+ }
+
+ @Override
+ public void shutdown() {}
+ };
+ }
+ };
+ NameResolverRegistry.getDefaultRegistry().register(customResolverProvider);
+
+ GrpcOpenTelemetry grpcOpenTelemetry = GrpcOpenTelemetry.newBuilder()
+ .sdk(openTelemetryRule.getOpenTelemetry())
+ .enableTracing(true)
+ .build();
+
+ InProcessChannelBuilder channelBuilder =
+ InProcessChannelBuilder.forTarget("inproce2e:///test-e2e")
+ .defaultLoadBalancingPolicy("slow_connecting_policy");
+ grpcOpenTelemetry.configureChannelBuilder(channelBuilder);
+ ManagedChannel channel = channelBuilder.build();
+ try {
+ ClientCall call = channel.newCall(method, CallOptions.DEFAULT);
+ call.start(new ClientCall.Listener() {}, new Metadata());
+ call.request(1);
+
+ latch.await(5, TimeUnit.SECONDS);
+ call.cancel("End test delay segment", null);
+ } finally {
+ channel.shutdownNow();
+ channel.awaitTermination(5, TimeUnit.SECONDS);
+ LoadBalancerRegistry.getDefaultRegistry().deregister(slowLbProvider);
+ NameResolverRegistry.getDefaultRegistry().deregister(customResolverProvider);
+ }
+
+ List spans = openTelemetryRule.getSpans();
+ SpanData delaySpanData = null;
+ for (SpanData s : spans) {
+ if ("Delay".equals(s.getName())) {
+ delaySpanData = s;
+ break;
+ }
+ }
+ assertNotNull(delaySpanData);
+ assertEquals("connecting",
+ delaySpanData.getAttributes().get(AttributeKey.stringKey("grpc.delay_type")));
+
+ boolean foundTransition = false;
+ for (EventData event : delaySpanData.getEvents()) {
+ if ("Delay triggered".equals(event.getName())
+ && "Simulated slow TLS handshake with backend".equals(
+ event.getAttributes().get(AttributeKey.stringKey("grpc.delay_reason")))) {
+ foundTransition = true;
+ break;
+ }
+ }
+ assertTrue(foundTransition);
+ }
+
@Test
public void clientBasicTracingRule() {
OpenTelemetryTracingModule tracingModule = new OpenTelemetryTracingModule(
@@ -415,10 +705,10 @@ public void clientAttemptDelayTracing_reasonChangedInvariant() {
ClientStreamTracer clientStreamTracer =
callTracer.newClientStreamTracer(STREAM_INFO, new Metadata());
- clientStreamTracer.recordAttemptDelayStart("connecting", "reason1");
- clientStreamTracer.recordAttemptDelayReasonChanged("reason2");
- clientStreamTracer.recordAttemptDelayStart("connecting", "reason3");
- clientStreamTracer.recordAttemptDelayEnd();
+ clientStreamTracer.recordDelayStart("connecting", "reason1");
+ clientStreamTracer.recordDelayReasonChanged("connecting", "reason2");
+ clientStreamTracer.recordDelayStart("connecting", "reason3");
+ clientStreamTracer.recordDelayEnd("connecting");
clientStreamTracer.streamClosed(Status.OK);
callTracer.callEnded(Status.OK);
clientSpan.end();
@@ -427,196 +717,172 @@ public void clientAttemptDelayTracing_reasonChangedInvariant() {
assertEquals(3, spans.size());
SpanData delaySpanData = spans.get(0);
- assertEquals("Attempt Delay", delaySpanData.getName());
+ assertEquals("Delay", delaySpanData.getName());
assertEquals("connecting", delaySpanData.getAttributes().get(
AttributeKey.stringKey("grpc.delay_type")));
assertEquals(3, delaySpanData.getEvents().size());
EventData event1 = delaySpanData.getEvents().get(0);
- assertEquals("Delay state transition", event1.getName());
- assertEquals("connecting", event1.getAttributes().get(
- AttributeKey.stringKey("grpc.delay_type")));
+ assertEquals("Delay triggered", event1.getName());
assertEquals("reason1", event1.getAttributes().get(
AttributeKey.stringKey("grpc.delay_reason")));
+ assertNull(event1.getAttributes().get(AttributeKey.stringKey("grpc.delay_type")));
EventData event2 = delaySpanData.getEvents().get(1);
- assertEquals("Delay state transition", event2.getName());
- assertEquals("connecting", event2.getAttributes().get(
- AttributeKey.stringKey("grpc.delay_type")));
+ assertEquals("Delay triggered", event2.getName());
assertEquals("reason2", event2.getAttributes().get(
AttributeKey.stringKey("grpc.delay_reason")));
+ assertNull(event2.getAttributes().get(AttributeKey.stringKey("grpc.delay_type")));
EventData event3 = delaySpanData.getEvents().get(2);
- assertEquals("Delay state transition", event3.getName());
- assertEquals("connecting", event3.getAttributes().get(
- AttributeKey.stringKey("grpc.delay_type")));
+ assertEquals("Delay triggered", event3.getName());
assertEquals("reason3", event3.getAttributes().get(
AttributeKey.stringKey("grpc.delay_reason")));
+ assertNull(event3.getAttributes().get(AttributeKey.stringKey("grpc.delay_type")));
}
@Test
- public void clientAttemptDelayTracing_endToEnd_inProcessTransport() throws Exception {
- final CountDownLatch latch = new CountDownLatch(1);
- LoadBalancerProvider slowLbProvider = new LoadBalancerProvider() {
- @Override
- public boolean isAvailable() {
- return true;
- }
-
- @Override
- public int getPriority() {
- return 5;
- }
+ public void clientCallDelayTracing_reasonChangedInvariant() {
+ OpenTelemetryTracingModule tracingModule = new OpenTelemetryTracingModule(
+ openTelemetryRule.getOpenTelemetry());
+ Span clientSpan = tracerRule.spanBuilder("test-client-span").startSpan();
+ CallAttemptsTracerFactory callTracer =
+ tracingModule.newClientCallTracer(clientSpan, method);
- @Override
- public String getPolicyName() {
- return "slow_connecting_policy";
- }
+ callTracer.recordDelayStart("resolving", "reason1");
+ callTracer.recordDelayReasonChanged("resolving", "reason2");
+ callTracer.recordDelayStart("resolving", "reason3");
+ callTracer.recordDelayEnd("resolving");
+ callTracer.callEnded(Status.OK);
+ clientSpan.end();
- @Override
- public LoadBalancer newLoadBalancer(LoadBalancer.Helper helper) {
- return new LoadBalancer() {
- @Override
- public Status acceptResolvedAddresses(LoadBalancer.ResolvedAddresses resolvedAddresses) {
- helper.updateBalancingState(ConnectivityState.CONNECTING, new SubchannelPicker() {
- @Override
- public PickResult pickSubchannel(PickSubchannelArgs args) {
- latch.countDown();
- return PickResult.withNoResult("connecting",
- "Simulated slow TLS handshake with backend");
- }
- });
- return Status.OK;
- }
+ List spans = openTelemetryRule.getSpans();
+ assertEquals(2, spans.size());
+ SpanData callDelaySpan = spans.stream()
+ .filter(s -> "Delay".equals(s.getName()))
+ .findFirst()
+ .orElseThrow(() -> new AssertionError("Expected 'Delay' span not found"));
- @Override
- public void handleNameResolutionError(Status error) {}
+ assertEquals("resolving", callDelaySpan.getAttributes().get(
+ AttributeKey.stringKey("grpc.delay_type")));
+ assertEquals(3, callDelaySpan.getEvents().size());
+
+ EventData event1 = callDelaySpan.getEvents().get(0);
+ assertEquals("Delay triggered", event1.getName());
+ assertEquals("reason1",
+ event1.getAttributes().get(AttributeKey.stringKey("grpc.delay_reason")));
+ assertNull(event1.getAttributes().get(AttributeKey.stringKey("grpc.delay_type")));
+
+ EventData event2 = callDelaySpan.getEvents().get(1);
+ assertEquals("Delay triggered", event2.getName());
+ assertEquals("reason2",
+ event2.getAttributes().get(AttributeKey.stringKey("grpc.delay_reason")));
+ assertNull(event2.getAttributes().get(AttributeKey.stringKey("grpc.delay_type")));
+
+ EventData event3 = callDelaySpan.getEvents().get(2);
+ assertEquals("Delay triggered", event3.getName());
+ assertEquals("reason3",
+ event3.getAttributes().get(AttributeKey.stringKey("grpc.delay_reason")));
+ assertNull(event3.getAttributes().get(AttributeKey.stringKey("grpc.delay_type")));
+ }
- @Override
- public void shutdown() {}
- };
- }
- };
- LoadBalancerRegistry.getDefaultRegistry().register(slowLbProvider);
+ @Test
+ public void clientCallDelayStart_afterCallEnded_noSpansRecorded() {
+ OpenTelemetryTracingModule tracingModule = new OpenTelemetryTracingModule(
+ openTelemetryRule.getOpenTelemetry());
+ Span clientSpan = tracerRule.spanBuilder("test-client-span").startSpan();
+ CallAttemptsTracerFactory callTracer =
+ tracingModule.newClientCallTracer(clientSpan, method);
- NameResolverProvider customResolverProvider = new NameResolverProvider() {
- @Override
- protected boolean isAvailable() {
- return true;
- }
+ callTracer.callEnded(Status.OK);
+ callTracer.recordDelayStart("resolving", "reasonAfterCallEnded");
+ callTracer.recordDelayReasonChanged("resolving", "reason2");
+ callTracer.recordDelayEnd("resolving");
+ clientSpan.end();
- @Override
- protected int priority() {
- return 5;
- }
+ List spans = openTelemetryRule.getSpans();
+ for (SpanData span : spans) {
+ assertTrue(!span.getName().equals("Delay"));
+ }
+ }
- @Override
- public String getDefaultScheme() {
- return "inproce2e";
- }
+ @Test
+ public void clientAttemptDelayStart_afterStreamClosed_noSpansRecorded() {
+ OpenTelemetryTracingModule tracingModule = new OpenTelemetryTracingModule(
+ openTelemetryRule.getOpenTelemetry());
+ Span clientSpan = tracerRule.spanBuilder("test-client-span").startSpan();
+ CallAttemptsTracerFactory callTracer =
+ tracingModule.newClientCallTracer(clientSpan, method);
+ ClientStreamTracer clientStreamTracer =
+ callTracer.newClientStreamTracer(STREAM_INFO, new Metadata());
- @Override
- public Collection> getProducedSocketAddressTypes() {
- return Collections.singleton(InProcessSocketAddress.class);
- }
+ clientStreamTracer.streamClosed(Status.OK);
+ clientStreamTracer.recordDelayStart("connecting", "reasonAfterStreamClosed");
+ clientStreamTracer.recordDelayReasonChanged("connecting", "reason2");
+ clientStreamTracer.recordDelayEnd("connecting");
+ callTracer.callEnded(Status.OK);
+ clientSpan.end();
- @Override
- public NameResolver newNameResolver(URI targetUri, NameResolver.Args args) {
- return new NameResolver() {
- @Override
- public String getServiceAuthority() {
- return "inproce2e";
- }
+ List spans = openTelemetryRule.getSpans();
+ for (SpanData span : spans) {
+ assertTrue(!span.getName().equals("Delay"));
+ }
+ }
- @Override
- public void start(Listener2 listener) {
- listener.onResult(ResolutionResult.newBuilder()
- .setAddressesOrError(StatusOr.fromValue(Collections.singletonList(
- new EquivalentAddressGroup(new InProcessSocketAddress("test-e2e")))))
- .build());
- }
+ @Test
+ public void clientCallDelayStart_delayTypeTransition_closesPreviousSpan() {
+ OpenTelemetryTracingModule tracingModule = new OpenTelemetryTracingModule(
+ openTelemetryRule.getOpenTelemetry());
+ Span clientSpan = tracerRule.spanBuilder("test-client-span").startSpan();
+ CallAttemptsTracerFactory callTracer =
+ tracingModule.newClientCallTracer(clientSpan, method);
- @Override
- public void shutdown() {}
- };
- }
- };
- NameResolverRegistry.getDefaultRegistry().register(customResolverProvider);
+ callTracer.recordDelayStart("resolving", "dns lookup");
+ callTracer.recordDelayStart("connecting", "pick first connect");
+ callTracer.recordDelayEnd("connecting");
+ callTracer.callEnded(Status.OK);
+ clientSpan.end();
- GrpcOpenTelemetry grpcOpenTelemetry = GrpcOpenTelemetry.newBuilder()
- .sdk(openTelemetryRule.getOpenTelemetry())
- .enableTracing(true)
- .build();
+ List spans = openTelemetryRule.getSpans();
+ long callDelaySpanCount = spans.stream().filter(s -> "Delay".equals(s.getName())).count();
+ assertEquals(2L, callDelaySpanCount);
+ }
- InProcessChannelBuilder channelBuilder =
- InProcessChannelBuilder.forTarget("inproce2e:///test-e2e")
- .defaultLoadBalancingPolicy("slow_connecting_policy");
- grpcOpenTelemetry.configureChannelBuilder(channelBuilder);
- ManagedChannel channel = channelBuilder.build();
- try {
- ClientCall call = channel.newCall(method, CallOptions.DEFAULT);
- call.start(new ClientCall.Listener() {}, new Metadata());
- call.request(1);
+ @Test
+ public void clientCallDelayReasonChanged_noActiveSpan_noOp() {
+ OpenTelemetryTracingModule tracingModule = new OpenTelemetryTracingModule(
+ openTelemetryRule.getOpenTelemetry());
+ Span clientSpan = tracerRule.spanBuilder("test-client-span").startSpan();
+ CallAttemptsTracerFactory callTracer =
+ tracingModule.newClientCallTracer(clientSpan, method);
- latch.await(5, TimeUnit.SECONDS);
- call.cancel("End test delay segment", null);
- } finally {
- channel.shutdownNow();
- channel.awaitTermination(5, TimeUnit.SECONDS);
- LoadBalancerRegistry.getDefaultRegistry().deregister(slowLbProvider);
- NameResolverRegistry.getDefaultRegistry().deregister(customResolverProvider);
- }
+ callTracer.recordDelayReasonChanged("resolving", "reasonWithoutSpan");
+ callTracer.recordDelayEnd("resolving");
+ callTracer.callEnded(Status.OK);
+ clientSpan.end();
List spans = openTelemetryRule.getSpans();
- SpanData delaySpanData = null;
- for (SpanData s : spans) {
- if ("Attempt Delay".equals(s.getName())) {
- delaySpanData = s;
- break;
- }
- }
- assertNotNull(delaySpanData);
- assertEquals("connecting",
- delaySpanData.getAttributes().get(AttributeKey.stringKey("grpc.delay_type")));
-
- boolean foundTransition = false;
- for (EventData event : delaySpanData.getEvents()) {
- if ("Delay state transition".equals(event.getName())
- && "Simulated slow TLS handshake with backend".equals(
- event.getAttributes().get(AttributeKey.stringKey("grpc.delay_reason")))) {
- foundTransition = true;
- break;
- }
- }
- assertTrue(foundTransition);
+ assertEquals(0L, spans.stream().filter(s -> "Delay".equals(s.getName())).count());
}
@Test
- public void clientAttemptDelayStart_featureFlagDisabled_zeroChildSpans() {
- System.setProperty("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY", "false");
- try {
- OpenTelemetryTracingModule tracingModule = new OpenTelemetryTracingModule(
- openTelemetryRule.getOpenTelemetry());
- Span clientSpan = tracerRule.spanBuilder("test-client-span").startSpan();
- CallAttemptsTracerFactory callTracer =
- tracingModule.newClientCallTracer(clientSpan, method);
- ClientStreamTracer clientStreamTracer =
- callTracer.newClientStreamTracer(STREAM_INFO, new Metadata());
-
- clientStreamTracer.recordAttemptDelayStart("connecting", "attempting to connect");
- clientStreamTracer.recordAttemptDelayEnd();
- clientStreamTracer.streamClosed(Status.OK);
- callTracer.callEnded(Status.OK);
- clientSpan.end();
-
- List spans = openTelemetryRule.getSpans();
- assertEquals(2, spans.size());
- for (SpanData span : spans) {
- assertTrue(!span.getName().equals("Attempt Delay"));
- }
- } finally {
- System.setProperty("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY", "true");
- }
+ public void clientAttemptDelayReasonChanged_noActiveSpan_noOp() {
+ OpenTelemetryTracingModule tracingModule = new OpenTelemetryTracingModule(
+ openTelemetryRule.getOpenTelemetry());
+ Span clientSpan = tracerRule.spanBuilder("test-client-span").startSpan();
+ CallAttemptsTracerFactory callTracer =
+ tracingModule.newClientCallTracer(clientSpan, method);
+ ClientStreamTracer clientStreamTracer =
+ callTracer.newClientStreamTracer(STREAM_INFO, new Metadata());
+
+ clientStreamTracer.recordDelayReasonChanged("connecting", "reasonWithoutSpan");
+ clientStreamTracer.recordDelayEnd("connecting");
+ clientStreamTracer.streamClosed(Status.OK);
+ callTracer.callEnded(Status.OK);
+ clientSpan.end();
+
+ List spans = openTelemetryRule.getSpans();
+ assertEquals(0L, spans.stream().filter(s -> "Delay".equals(s.getName())).count());
}
@Test
diff --git a/util/src/main/java/io/grpc/util/ForwardingClientStreamTracer.java b/util/src/main/java/io/grpc/util/ForwardingClientStreamTracer.java
index 78f234a8a30..3d0988eda8d 100644
--- a/util/src/main/java/io/grpc/util/ForwardingClientStreamTracer.java
+++ b/util/src/main/java/io/grpc/util/ForwardingClientStreamTracer.java
@@ -39,18 +39,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
diff --git a/util/src/test/java/io/grpc/util/ForwardingClientStreamTracerTest.java b/util/src/test/java/io/grpc/util/ForwardingClientStreamTracerTest.java
index dbd7e99b29a..e6b16101fa4 100644
--- a/util/src/test/java/io/grpc/util/ForwardingClientStreamTracerTest.java
+++ b/util/src/test/java/io/grpc/util/ForwardingClientStreamTracerTest.java
@@ -40,6 +40,7 @@ public void allMethodsForwarded() throws Exception {
Collections.emptyList());
}
+
@SuppressWarnings("deprecation")
private final class TestClientStreamTracer extends ForwardingClientStreamTracer {
@Override
diff --git a/xds/src/main/java/io/grpc/xds/PriorityLoadBalancer.java b/xds/src/main/java/io/grpc/xds/PriorityLoadBalancer.java
index e06d819f66f..ddc43ef9696 100644
--- a/xds/src/main/java/io/grpc/xds/PriorityLoadBalancer.java
+++ b/xds/src/main/java/io/grpc/xds/PriorityLoadBalancer.java
@@ -331,8 +331,9 @@ public void updateBalancingState(final ConnectivityState newState,
}
ConnectivityState oldState = connectivityState;
connectivityState = newState;
- if (newState == CONNECTING || newState == IDLE) {
- picker = new PriorityPicker(newPicker, priority);
+ int priorityIndex = priorityNames.indexOf(priority);
+ if (priorityIndex >= 0 && (newState == CONNECTING || newState == IDLE)) {
+ picker = new PriorityPicker(newPicker, String.valueOf(priorityIndex));
} else {
picker = newPicker;
}
diff --git a/xds/src/test/java/io/grpc/xds/PriorityLoadBalancerTest.java b/xds/src/test/java/io/grpc/xds/PriorityLoadBalancerTest.java
index b0fd3b9c087..09836ab7758 100644
--- a/xds/src/test/java/io/grpc/xds/PriorityLoadBalancerTest.java
+++ b/xds/src/test/java/io/grpc/xds/PriorityLoadBalancerTest.java
@@ -1027,7 +1027,7 @@ public void noDuplicateOverallBalancingStateUpdate() {
}
@Test
- public void priorityPicker_prependsToken() throws Exception {
+ public void priorityPicker_prependsNumericPriority() throws Exception {
PriorityChildConfig priorityChildConfig0 =
new PriorityChildConfig(newChildConfig(fooLbProvider, new Object()), true);
PriorityLbConfig priorityLbConfig =
@@ -1055,13 +1055,13 @@ public PickResult pickSubchannel(PickSubchannelArgs args) {
SubchannelPicker priorityPicker = pickerCaptor.getValue();
PickResult result = priorityPicker.pickSubchannel(mock(PickSubchannelArgs.class));
- assertThat(result.getDelayType()).isEqualTo("p0:connecting");
+ assertThat(result.getDelayType()).isEqualTo("0:connecting");
assertThat(result.getDelayReason())
- .isEqualTo("waiting on priority group p0 (child_reason)");
+ .isEqualTo("waiting on priority group 0 (child_reason)");
}
@Test
- public void priorityPicker_nestedPriorities_composesTokens() throws Exception {
+ public void priorityPicker_nestedPriorities_composesNumericPrefixes() throws Exception {
PriorityChildConfig priorityChildConfig0 =
new PriorityChildConfig(newChildConfig(fooLbProvider, new Object()), true);
PriorityLbConfig priorityLbConfig =
@@ -1078,8 +1078,8 @@ public void priorityPicker_nestedPriorities_composesTokens() throws Exception {
SubchannelPicker nestedChildPicker = new SubchannelPicker() {
@Override
public PickResult pickSubchannel(PickSubchannelArgs args) {
- return PickResult.withNoResult("p1:connecting",
- "waiting on priority group p1 (child_reason)");
+ return PickResult.withNoResult("1:connecting",
+ "waiting on priority group 1 (child_reason)");
}
};
helper0.updateBalancingState(CONNECTING, nestedChildPicker);
@@ -1090,9 +1090,9 @@ public PickResult pickSubchannel(PickSubchannelArgs args) {
SubchannelPicker priorityPicker = pickerCaptor.getValue();
PickResult result = priorityPicker.pickSubchannel(mock(PickSubchannelArgs.class));
- assertThat(result.getDelayType()).isEqualTo("p0:p1:connecting");
+ assertThat(result.getDelayType()).isEqualTo("0:1:connecting");
assertThat(result.getDelayReason()).isEqualTo(
- "waiting on priority group p0 (waiting on priority group p1 (child_reason))");
+ "waiting on priority group 0 (waiting on priority group 1 (child_reason))");
}
@Test
@@ -1117,6 +1117,175 @@ public void initialChildPicker_returnsAnnotatedDelayAttributes() throws Exceptio
"priority child state uninitialized");
}
+ @Test
+ public void priorityPicker_idleState_prependsNumericPriority() throws Exception {
+ PriorityChildConfig priorityChildConfig0 =
+ new PriorityChildConfig(newChildConfig(fooLbProvider, new Object()), true);
+ PriorityLbConfig priorityLbConfig =
+ new PriorityLbConfig(ImmutableMap.of("p0", priorityChildConfig0), ImmutableList.of("p0"));
+ priorityLb.acceptResolvedAddresses(
+ ResolvedAddresses.newBuilder()
+ .setAddresses(ImmutableList.of())
+ .setLoadBalancingPolicyConfig(priorityLbConfig)
+ .build());
+
+ Helper helper0 = Iterables.getOnlyElement(fooHelpers);
+ SubchannelPicker idleChildPicker = new SubchannelPicker() {
+ @Override
+ public PickResult pickSubchannel(PickSubchannelArgs args) {
+ return PickResult.withNoResult("connecting", "idle: connection requested");
+ }
+ };
+ helper0.updateBalancingState(IDLE, idleChildPicker);
+
+ assertLatestConnectivityState(IDLE);
+ PickResult result = pickerCaptor.getValue().pickSubchannel(mock(PickSubchannelArgs.class));
+ assertThat(result.getDelayType()).isEqualTo("0:connecting");
+ assertThat(result.getDelayReason())
+ .isEqualTo("waiting on priority group 0 (idle: connection requested)");
+ }
+
+ @Test
+ public void priorityPicker_passthroughWhenNoDelayTypeOrHasResult() throws Exception {
+ PriorityChildConfig priorityChildConfig0 =
+ new PriorityChildConfig(newChildConfig(fooLbProvider, new Object()), true);
+ PriorityLbConfig priorityLbConfig =
+ new PriorityLbConfig(ImmutableMap.of("p0", priorityChildConfig0), ImmutableList.of("p0"));
+ priorityLb.acceptResolvedAddresses(
+ ResolvedAddresses.newBuilder()
+ .setAddresses(ImmutableList.of())
+ .setLoadBalancingPolicyConfig(priorityLbConfig)
+ .build());
+
+ Helper helper0 = Iterables.getOnlyElement(fooHelpers);
+ Subchannel subchannel = mock(Subchannel.class);
+ final PickResult[] nextResult = new PickResult[] {PickResult.withNoResult()};
+ SubchannelPicker dynamicPicker = new SubchannelPicker() {
+ @Override
+ public PickResult pickSubchannel(PickSubchannelArgs args) {
+ return nextResult[0];
+ }
+ };
+ helper0.updateBalancingState(CONNECTING, dynamicPicker);
+ verify(helper, atLeastOnce()).updateBalancingState(eq(CONNECTING), pickerCaptor.capture());
+ SubchannelPicker priorityPicker = pickerCaptor.getValue();
+
+ // No delayType -> passed through unchanged
+ PickResult noDelayTypeResult = priorityPicker.pickSubchannel(mock(PickSubchannelArgs.class));
+ assertThat(noDelayTypeResult).isEqualTo(PickResult.withNoResult());
+ assertThat(noDelayTypeResult.getDelayType()).isNull();
+
+ // Has result (subchannel) -> passed through unchanged
+ nextResult[0] = PickResult.withSubchannel(subchannel);
+ PickResult subchannelResult = priorityPicker.pickSubchannel(mock(PickSubchannelArgs.class));
+ assertThat(subchannelResult).isEqualTo(PickResult.withSubchannel(subchannel));
+ assertThat(subchannelResult.getDelayType()).isNull();
+ }
+
+ @Test
+ public void priorityPicker_deactivatedChild_doesNotPrependNegativePriorityIndex() {
+ PriorityLoadBalancer.enablePriorityLbChildPolicyCache = true;
+ try {
+ PriorityChildConfig priorityChildConfig0 =
+ new PriorityChildConfig(newChildConfig(fooLbProvider, new Object()), true);
+ PriorityChildConfig priorityChildConfig1 =
+ new PriorityChildConfig(newChildConfig(fooLbProvider, new Object()), true);
+ PriorityLbConfig initialConfig =
+ new PriorityLbConfig(
+ ImmutableMap.of("p0", priorityChildConfig0, "p1", priorityChildConfig1),
+ ImmutableList.of("p0", "p1"));
+ priorityLb.acceptResolvedAddresses(
+ ResolvedAddresses.newBuilder()
+ .setAddresses(ImmutableList.of())
+ .setLoadBalancingPolicyConfig(initialConfig)
+ .build());
+
+ // Trigger creation of p1 by failing p0
+ Helper helper0 = fooHelpers.get(0);
+ helper0.updateBalancingState(
+ TRANSIENT_FAILURE,
+ new FixedResultPicker(PickResult.withError(Status.UNAVAILABLE)));
+ Helper helper1 = fooHelpers.get(1);
+
+ // Remove p1 from priorityNames so p1 is deactivated (retained in cache for 15m)
+ PriorityLbConfig onlyP0Config =
+ new PriorityLbConfig(
+ ImmutableMap.of("p0", priorityChildConfig0), ImmutableList.of("p0"));
+ priorityLb.acceptResolvedAddresses(
+ ResolvedAddresses.newBuilder()
+ .setAddresses(ImmutableList.of())
+ .setLoadBalancingPolicyConfig(onlyP0Config)
+ .build());
+
+ // Deactivated p1 (priorityIndex == -1) updates state to CONNECTING
+ helper1.updateBalancingState(
+ CONNECTING,
+ new FixedResultPicker(PickResult.withNoResult("connecting", "deactivated_child")));
+ assertLatestConnectivityState(TRANSIENT_FAILURE);
+ } finally {
+ PriorityLoadBalancer.enablePriorityLbChildPolicyCache = false;
+ }
+ }
+
+ @Test
+ public void priorityPicker_equalsAndHashCode() {
+ PriorityChildConfig priorityChildConfig0 =
+ new PriorityChildConfig(newChildConfig(fooLbProvider, new Object()), true);
+ PriorityChildConfig priorityChildConfig1 =
+ new PriorityChildConfig(newChildConfig(fooLbProvider, new Object()), true);
+ PriorityLbConfig priorityLbConfig =
+ new PriorityLbConfig(
+ ImmutableMap.of("p0", priorityChildConfig0, "p1", priorityChildConfig1),
+ ImmutableList.of("p0", "p1"));
+ priorityLb.acceptResolvedAddresses(
+ ResolvedAddresses.newBuilder()
+ .setAddresses(ImmutableList.of())
+ .setLoadBalancingPolicyConfig(priorityLbConfig)
+ .build());
+
+ Helper helper0 = fooHelpers.get(0);
+ SubchannelPicker childPickerA = new SubchannelPicker() {
+ @Override
+ public PickResult pickSubchannel(PickSubchannelArgs args) {
+ return PickResult.withNoResult("connecting", "a");
+ }
+ };
+ SubchannelPicker childPickerB = new SubchannelPicker() {
+ @Override
+ public PickResult pickSubchannel(PickSubchannelArgs args) {
+ return PickResult.withNoResult("connecting", "b");
+ }
+ };
+
+ helper0.updateBalancingState(CONNECTING, childPickerA);
+ verify(helper, atLeastOnce()).updateBalancingState(eq(CONNECTING), pickerCaptor.capture());
+ SubchannelPicker picker0A1 = pickerCaptor.getValue();
+
+ helper0.updateBalancingState(CONNECTING, childPickerB);
+ verify(helper, atLeastOnce()).updateBalancingState(eq(CONNECTING), pickerCaptor.capture());
+ SubchannelPicker picker0B = pickerCaptor.getValue();
+
+ helper0.updateBalancingState(CONNECTING, childPickerA);
+ verify(helper, atLeastOnce()).updateBalancingState(eq(CONNECTING), pickerCaptor.capture());
+ SubchannelPicker picker0A2 = pickerCaptor.getValue();
+
+ // Fail over to p1 (priority index 1) with same delegate childPickerA
+ helper0.updateBalancingState(
+ TRANSIENT_FAILURE, new FixedResultPicker(PickResult.withError(Status.UNAVAILABLE)));
+ Helper helper1 = fooHelpers.get(1);
+ helper1.updateBalancingState(CONNECTING, childPickerA);
+ verify(helper, atLeastOnce()).updateBalancingState(eq(CONNECTING), pickerCaptor.capture());
+ SubchannelPicker picker1A = pickerCaptor.getValue();
+
+ assertThat(picker0A1.equals(picker0A1)).isTrue();
+ assertThat(picker0A1).isEqualTo(picker0A2);
+ assertThat(picker0A1.hashCode()).isEqualTo(picker0A2.hashCode());
+ assertThat(picker0A1).isNotEqualTo(picker0B);
+ assertThat(picker0A1).isNotEqualTo(picker1A);
+ assertThat(picker0A1.equals(null)).isFalse();
+ assertThat(picker0A1.equals(new Object())).isFalse();
+ }
+
private void assertLatestConnectivityState(ConnectivityState expectedState) {
verify(helper, atLeastOnce())
.updateBalancingState(connectivityStateCaptor.capture(), pickerCaptor.capture());