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());