diff --git a/api/src/main/java/io/grpc/ClientStreamTracer.java b/api/src/main/java/io/grpc/ClientStreamTracer.java index d654fd33358..dfc4c5d8820 100644 --- a/api/src/main/java/io/grpc/ClientStreamTracer.java +++ b/api/src/main/java/io/grpc/ClientStreamTracer.java @@ -67,7 +67,7 @@ public void createPendingStream() { * * @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.84.0 */ public void recordAttemptDelayStart(String delayType, String delayReason) { } @@ -80,7 +80,7 @@ public void recordAttemptDelayStart(String delayType, String delayReason) { * on the active delay span without recreating the span or resetting cumulative timers. * * @param delayReason updated high-cardinality diagnostic string describing new conditions - * @since 1.82.0 + * @since 1.84.0 */ public void recordAttemptDelayReasonChanged(String delayReason) { } @@ -91,7 +91,7 @@ public void recordAttemptDelayReasonChanged(String delayReason) { *

Implementations should simultaneously close active child tracing spans and record elapsed * duration to the {@code grpc.client.attempt.delay.duration} histogram. * - * @since 1.82.0 + * @since 1.84.0 */ public void recordAttemptDelayEnd() { } @@ -156,6 +156,44 @@ public abstract static class Factory { public ClientStreamTracer newClientStreamTracer(StreamInfo info, Metadata headers) { throw new UnsupportedOperationException("Not implemented"); } + + /** + * Called when a call-level delay segment (such as waiting for name resolution or service + * configuration parsing) starts before any individual RPC attempt is created. + * + *

Implementations should start logical timers and create child tracing spans (named strictly + * {@code "Call 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.84.0 + */ + public void recordCallDelayStart(String delayType, String delayReason) { + } + + /** + * Called when a call-level delay reason changes while the active delay segment continues. + * + *

Implementations should emit structured events (such as {@code "Delay state transition"}) + * on the active call delay span without recreating the span or resetting timers. + * + * @param delayReason updated high-cardinality diagnostic string describing new conditions + * @since 1.84.0 + */ + public void recordCallDelayReasonChanged(String delayReason) { + } + + /** + * Called when a call-level delay segment ends upon successful name resolution or when an RPC + * is cancelled before resolution completes. + * + *

Implementations should close active call delay spans and record elapsed duration to the + * {@code grpc.client.call.delay.duration} histogram. + * + * @since 1.84.0 + */ + public void recordCallDelayEnd() { + } } /** diff --git a/api/src/main/java/io/grpc/LoadBalancer.java b/api/src/main/java/io/grpc/LoadBalancer.java index e5c3d053ee7..c9dce0f2e3d 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.84.0 */ public static PickResult withNoResult(String delayType, String delayReason) { Preconditions.checkNotNull(delayType, "delayType"); diff --git a/api/src/test/java/io/grpc/ClientStreamTracerTest.java b/api/src/test/java/io/grpc/ClientStreamTracerTest.java index 5ddee77f5c0..1d605a4f707 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.recordAttemptDelayStart("connecting", "test"); + tracer.recordAttemptDelayReasonChanged("test2"); + tracer.recordAttemptDelayEnd(); + + ClientStreamTracer.Factory factory = new ClientStreamTracer.Factory() {}; + factory.recordCallDelayStart("resolving", "test"); + factory.recordCallDelayReasonChanged("test2"); + factory.recordCallDelayEnd(); + } } diff --git a/core/src/main/java/io/grpc/internal/DelayedClientTransport.java b/core/src/main/java/io/grpc/internal/DelayedClientTransport.java index aa1d820a570..3006631cd61 100644 --- a/core/src/main/java/io/grpc/internal/DelayedClientTransport.java +++ b/core/src/main/java/io/grpc/internal/DelayedClientTransport.java @@ -405,6 +405,8 @@ private class PendingStream extends DelayedStream { @Nullable private String activeDelayType; @GuardedBy("this") @Nullable private String activeDelayReason; + @GuardedBy("this") + private boolean delayEnded; private PendingStream(PickSubchannelArgs args, ClientStreamTracer[] tracers, @Nullable String initialType, @Nullable String initialReason) { @@ -427,30 +429,47 @@ 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) { - if (getRealStream() != null) { - return; - } - if (!Objects.equals(activeDelayType, newType)) { - // Delay type changed (e.g., from RLS lookup to connecting). End the previous delay. - if (activeDelayType != null) { + void updateDelay(@Nullable String newType, @Nullable String newReason) { + synchronized (this) { + if (getRealStream() != null || delayEnded) { + return; + } + String prevTypeToClose = null; + String newTypeToStart = null; + String newReasonToStart = null; + String newReasonToNotify = null; + + if (!Objects.equals(activeDelayType, newType)) { + if (activeDelayType != null) { + prevTypeToClose = activeDelayType; + } + activeDelayType = newType; + activeDelayReason = null; + if (newType != null) { + newTypeToStart = newType; + newReasonToStart = newReason != null ? newReason : ""; + } + } + if (newType != null && newReason != null + && !Objects.equals(activeDelayReason, newReason)) { + activeDelayReason = newReason; + newReasonToNotify = newReason; + } + + if (prevTypeToClose != null) { for (ClientStreamTracer tracer : tracers) { tracer.recordAttemptDelayEnd(); } } - activeDelayType = newType; - activeDelayReason = null; - if (newType != null) { + if (newTypeToStart != null) { for (ClientStreamTracer tracer : tracers) { - tracer.recordAttemptDelayStart(newType, newReason != null ? newReason : ""); + tracer.recordAttemptDelayStart(newTypeToStart, newReasonToStart); } } - } - if (newType != null && newReason != null && !Objects.equals(activeDelayReason, newReason)) { - // Delay type is unchanged, but the reason changed (e.g., priority failover). - activeDelayReason = newReason; - for (ClientStreamTracer tracer : tracers) { - tracer.recordAttemptDelayReasonChanged(newReason); + if (newReasonToNotify != null) { + for (ClientStreamTracer tracer : tracers) { + tracer.recordAttemptDelayReasonChanged(newReasonToNotify); + } } } } @@ -458,19 +477,27 @@ 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(); + void endDelay() { + synchronized (this) { + if (delayEnded) { + return; } + delayEnded = true; + boolean shouldEnd = activeDelayType != null; activeDelayType = null; activeDelayReason = null; + if (shouldEnd) { + for (ClientStreamTracer tracer : tracers) { + tracer.recordAttemptDelayEnd(); + } + } } } Runnable setStreamAndEndDelay(ClientStream stream) { + Runnable runnable = setStream(stream); endDelay(); - return setStream(stream); + return runnable; } /** Runnable may be null. */ @@ -496,6 +523,7 @@ private Runnable createRealStream(ClientTransport transport, String authorityOve @Override public void cancel(Status reason) { + endDelay(); super.cancel(reason); synchronized (lock) { if (reportTransportTerminated != null) { diff --git a/core/src/main/java/io/grpc/internal/ManagedChannelImpl.java b/core/src/main/java/io/grpc/internal/ManagedChannelImpl.java index 00df05a0c00..e11ebd54579 100644 --- a/core/src/main/java/io/grpc/internal/ManagedChannelImpl.java +++ b/core/src/main/java/io/grpc/internal/ManagedChannelImpl.java @@ -918,6 +918,7 @@ public void run() { inUseStateAggregator.updateObjectInUse(pendingCallsInUseObject, true); } pendingCalls.add(pendingCall); + pendingCall.notifyQueuedForNameResolution(); } else { pendingCall.reprocess(); } @@ -997,6 +998,10 @@ private final class PendingCall extends DelayedClientCall method; final CallOptions callOptions; private final long callCreationTime; + @GuardedBy("this") private boolean queuedForResolution; + @GuardedBy("this") private boolean callCancelled; + @GuardedBy("this") private boolean delayEnded; + @GuardedBy("this") private boolean callDelayStarted; PendingCall(Context context, MethodDescriptor method, CallOptions callOptions) { super( @@ -1010,14 +1015,69 @@ private final class PendingCall extends DelayedClientCall realCall; Context previous = context.attach(); try { - CallOptions delayResolutionOption = callOptions.withOption(NAME_RESOLUTION_DELAYED, - ticker.nanoTime() - callCreationTime); - realCall = newClientCall(method, delayResolutionOption); + CallOptions effectiveOptions = callOptions; + boolean wasQueued; + synchronized (this) { + wasQueued = queuedForResolution; + } + if (wasQueued) { + effectiveOptions = callOptions.withOption(NAME_RESOLUTION_DELAYED, + ticker.nanoTime() - callCreationTime); + } + realCall = newClientCall(method, effectiveOptions); } finally { context.detach(previous); } @@ -1037,6 +1097,10 @@ public void run() { @Override protected void callCancelled() { + synchronized (this) { + callCancelled = true; + } + endDelayIfNeeded(); super.callCancelled(); syncContext.execute(new PendingCallRemoval()); } diff --git a/core/src/test/java/io/grpc/internal/ForwardingClientStreamTracerTest.java b/core/src/test/java/io/grpc/internal/ForwardingClientStreamTracerTest.java index 5eb5b49fa19..bfef85c8ca8 100644 --- a/core/src/test/java/io/grpc/internal/ForwardingClientStreamTracerTest.java +++ b/core/src/test/java/io/grpc/internal/ForwardingClientStreamTracerTest.java @@ -40,6 +40,19 @@ public void allMethodsForwarded() throws Exception { Collections.emptyList()); } + @Test + public void attemptDelayMethodsForwarded() { + TestClientStreamTracer tracer = new TestClientStreamTracer(); + tracer.recordAttemptDelayStart("connecting", "test"); + org.mockito.Mockito.verify(mockDelegate).recordAttemptDelayStart("connecting", "test"); + + tracer.recordAttemptDelayReasonChanged("test2"); + org.mockito.Mockito.verify(mockDelegate).recordAttemptDelayReasonChanged("test2"); + + tracer.recordAttemptDelayEnd(); + org.mockito.Mockito.verify(mockDelegate).recordAttemptDelayEnd(); + } + private final class TestClientStreamTracer extends ForwardingClientStreamTracer { @Override protected ClientStreamTracer delegate() { diff --git a/opentelemetry/src/main/java/io/grpc/opentelemetry/GrpcOpenTelemetry.java b/opentelemetry/src/main/java/io/grpc/opentelemetry/GrpcOpenTelemetry.java index 1243e0fff59..6965dadea85 100644 --- a/opentelemetry/src/main/java/io/grpc/opentelemetry/GrpcOpenTelemetry.java +++ b/opentelemetry/src/main/java/io/grpc/opentelemetry/GrpcOpenTelemetry.java @@ -243,6 +243,16 @@ && isMetricEnabled("grpc.client.attempt.delay.duration", enableMetrics, disableD .build()); } + if (isDelayObservabilityEnabled() + && isMetricEnabled("grpc.client.call.delay.duration", enableMetrics, disableDefault)) { + builder.clientCallDelayCounter( + meter.histogramBuilder( + "grpc.client.call.delay.duration") + .setUnit("s") + .setDescription("Time taken before a client call starts") + .setExplicitBucketBoundariesAdvice(LATENCY_BUCKETS) + .build()); + } if (isMetricEnabled("grpc.client.attempt.sent_total_compressed_message_size", enableMetrics, disableDefault)) { builder.clientTotalSentCompressedMessageSizeCounter( @@ -364,7 +374,7 @@ && isMetricEnabled("grpc.client.attempt.delay.duration", enableMetrics, disableD * 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 + * variable or JVM system property (defaults to {@code false}). When disabled, delay spans and * duration histograms are suppressed to avoid runtime overhead. */ static boolean isDelayObservabilityEnabled() { diff --git a/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsModule.java b/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsModule.java index 008a9754b29..20c59c483cb 100644 --- a/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsModule.java +++ b/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsModule.java @@ -256,6 +256,20 @@ public void recordAttemptDelayEnd() { .put(METHOD_KEY, fullMethodName) .put(TARGET_KEY, target) .put("grpc.delay_type", delayType); + if (module.localityEnabled) { + String savedLocality = locality; + if (savedLocality == null) { + savedLocality = ""; + } + builder.put(LOCALITY_KEY, savedLocality); + } + if (module.backendServiceEnabled) { + String savedBackendService = backendService; + if (savedBackendService == null) { + savedBackendService = ""; + } + builder.put(BACKEND_SERVICE_KEY, savedBackendService); + } if (module.customLabelEnabled) { builder.put( CUSTOM_LABEL_KEY, info.getCallOptions().getOption(Grpc.CALL_OPTION_CUSTOM_LABEL)); @@ -387,6 +401,11 @@ static final class CallAttemptsTracerFactory extends ClientStreamTracer.Factory private final List callPlugins; private final Context otelContext; private Status status; + @GuardedBy("this") + @Nullable private Stopwatch activeCallDelayStopwatch; + @GuardedBy("this") + @Nullable private String activeCallDelayType; + private final io.opentelemetry.api.common.Attributes callLevelBaseAttributes; private long retryDelayNanos; private long callLatencyNanos; private final Object lock = new Object(); @@ -420,6 +439,7 @@ static final class CallAttemptsTracerFactory extends ClientStreamTracer.Factory CUSTOM_LABEL_KEY, callOptions.getOption(Grpc.CALL_OPTION_CUSTOM_LABEL)); } io.opentelemetry.api.common.Attributes attribute = builder.build(); + this.callLevelBaseAttributes = attribute; // Record here in case mewClientStreamTracer() would never be called. if (module.resource.clientAttemptCountCounter() != null) { @@ -579,6 +599,41 @@ void recordFinishedCall(CallOptions callOptions) { ); } } + + @Override + public synchronized void recordCallDelayStart(String delayType, String delayReason) { + if (!GrpcOpenTelemetry.isDelayObservabilityEnabled() + || (activeCallDelayStopwatch != null && Objects.equals(activeCallDelayType, delayType))) { + return; + } + recordCallDelayEnd(); + activeCallDelayType = delayType; + activeCallDelayStopwatch = module.stopwatchSupplier.get().start(); + } + + @Override + public synchronized void recordCallDelayReasonChanged(String delayReason) { + } + + @Override + public synchronized void recordCallDelayEnd() { + Stopwatch delayStopwatch = activeCallDelayStopwatch; + String delayType = activeCallDelayType; + if (delayStopwatch != null && delayType != null) { + delayStopwatch.stop(); + long delayNanos = delayStopwatch.elapsed(TimeUnit.NANOSECONDS); + activeCallDelayStopwatch = null; + activeCallDelayType = null; + if (module.resource.clientCallDelayCounter() != null) { + module.resource.clientCallDelayCounter().record( + delayNanos * SECONDS_PER_NANO, + callLevelBaseAttributes.toBuilder() + .put("grpc.delay_type", delayType) + .build(), + Context.current()); + } + } + } } 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 32aab870f0f..b4380c111a5 100644 --- a/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryTracingModule.java +++ b/opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryTracingModule.java @@ -22,6 +22,7 @@ import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.BAGGAGE_KEY; import com.google.common.annotations.VisibleForTesting; +import com.google.errorprone.annotations.concurrent.GuardedBy; import io.grpc.Attributes; import io.grpc.CallOptions; import io.grpc.Channel; @@ -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,57 @@ void callEnded(io.grpc.Status status) { } callEnded = 1; } + recordCallDelayEnd(); endSpanWithStatus(clientSpan, status); } + + @Override + public synchronized void recordCallDelayStart(String delayType, String delayReason) { + if (!GrpcOpenTelemetry.isDelayObservabilityEnabled() || isCallEnded()) { + return; + } + if (activeCallDelaySpan != null && Objects.equals(activeCallDelayType, delayType)) { + recordCallDelayReasonChanged(delayReason); + return; + } + recordCallDelayEnd(); + activeCallDelayType = delayType; + Span delaySpan = otelTracer.spanBuilder("Call Delay") + .setParent(Context.current().with(clientSpan)) + .setAttribute("grpc.delay_type", delayType) + .startSpan(); + activeCallDelaySpan = delaySpan; + delaySpan.addEvent( + "Delay state transition", + io.opentelemetry.api.common.Attributes.of( + AttributeKey.stringKey("grpc.delay_type"), delayType, + AttributeKey.stringKey("grpc.delay_reason"), delayReason)); + } + + @Override + public synchronized void recordCallDelayReasonChanged(String delayReason) { + if (!GrpcOpenTelemetry.isDelayObservabilityEnabled() + || isCallEnded() + || activeCallDelaySpan == null) { + return; + } + String type = activeCallDelayType; + activeCallDelaySpan.addEvent( + "Delay state transition", + io.opentelemetry.api.common.Attributes.of( + AttributeKey.stringKey("grpc.delay_type"), type != null ? type : "", + AttributeKey.stringKey("grpc.delay_reason"), delayReason)); + } + + @Override + public synchronized void recordCallDelayEnd() { + Span delaySpan = activeCallDelaySpan; + if (delaySpan != null) { + delaySpan.end(); + activeCallDelaySpan = null; + activeCallDelayType = null; + } + } } private final class ClientTracer extends ClientStreamTracer { @@ -194,8 +255,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"); @@ -218,8 +283,8 @@ public void createPendingStream() { } @Override - public void recordAttemptDelayStart(String delayType, String delayReason) { - if (!GrpcOpenTelemetry.isDelayObservabilityEnabled()) { + public synchronized void recordAttemptDelayStart(String delayType, String delayReason) { + if (!GrpcOpenTelemetry.isDelayObservabilityEnabled() || streamClosed) { return; } if (activeDelaySpan != null && Objects.equals(activeDelayType, delayType)) { @@ -244,8 +309,10 @@ public void recordAttemptDelayStart(String delayType, String delayReason) { } @Override - public void recordAttemptDelayReasonChanged(String delayReason) { - if (!GrpcOpenTelemetry.isDelayObservabilityEnabled() || activeDelaySpan == null) { + public synchronized void recordAttemptDelayReasonChanged(String delayReason) { + if (!GrpcOpenTelemetry.isDelayObservabilityEnabled() + || streamClosed + || activeDelaySpan == null) { return; } String type = activeDelayType; @@ -257,7 +324,7 @@ public void recordAttemptDelayReasonChanged(String delayReason) { } @Override - public void recordAttemptDelayEnd() { + public synchronized void recordAttemptDelayEnd() { Span delaySpan = activeDelaySpan; if (delaySpan != null) { // End active child span upon pick completion or transport cancellation. @@ -292,7 +359,11 @@ public void inboundUncompressedSize(long bytes) { } @Override - public void streamClosed(io.grpc.Status status) { + public synchronized void streamClosed(io.grpc.Status status) { + if (streamClosed) { + return; + } + streamClosed = true; recordAttemptDelayEnd(); endSpanWithStatus(span, status); } diff --git a/opentelemetry/src/test/java/io/grpc/opentelemetry/GrpcOpenTelemetryTest.java b/opentelemetry/src/test/java/io/grpc/opentelemetry/GrpcOpenTelemetryTest.java index 77eadf9ebbb..6d200c9d115 100644 --- a/opentelemetry/src/test/java/io/grpc/opentelemetry/GrpcOpenTelemetryTest.java +++ b/opentelemetry/src/test/java/io/grpc/opentelemetry/GrpcOpenTelemetryTest.java @@ -17,6 +17,9 @@ package io.grpc.opentelemetry; import static com.google.common.truth.Truth.assertThat; +import static io.grpc.ClientStreamTracer.NAME_RESOLUTION_DELAYED; +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 +27,82 @@ 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.ForwardingChannelBuilder2; +import io.grpc.ManagedChannel; import io.grpc.ManagedChannelBuilder; +import io.grpc.Metadata; +import io.grpc.MethodDescriptor; import io.grpc.MetricSink; import io.grpc.ServerBuilder; +import io.grpc.ServerCall; +import io.grpc.ServerCallHandler; +import io.grpc.ServerServiceDefinition; +import io.grpc.Status; +import io.grpc.inprocess.InProcessChannelBuilder; +import io.grpc.inprocess.InProcessServerBuilder; +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.exporter.InMemoryMetricReader; +import io.opentelemetry.sdk.testing.junit4.OpenTelemetryRule; import io.opentelemetry.sdk.trace.SdkTracerProvider; +import java.io.ByteArrayInputStream; +import java.io.IOException; +import java.io.InputStream; import java.util.Arrays; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; 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(); @@ -55,11 +113,13 @@ public class GrpcOpenTelemetryTest { @Before public void setup() { originalEnableOtelTracing = GrpcOpenTelemetry.ENABLE_OTEL_TRACING; + System.setProperty("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY", "true"); } @After public void tearDown() { GrpcOpenTelemetry.ENABLE_OTEL_TRACING = originalEnableOtelTracing; + System.clearProperty("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY"); } @Test @@ -179,6 +239,242 @@ public void configureChannelBuilder_registersMetricSink() { assertThat(testBuilder.interceptorFactory).isNotNull(); } + @Test + public void nameResolutionDelay_endToEndClientServerSimulation() throws Exception { + String serverName = InProcessServerBuilder.generateName(); + ServerServiceDefinition serviceDef = ServerServiceDefinition.builder("test.service") + .addMethod(method, new ServerCallHandler() { + @Override + public ServerCall.Listener startCall( + ServerCall call, Metadata headers) { + call.sendHeaders(new Metadata()); + call.sendMessage("response_payload"); + call.close(Status.OK, new Metadata()); + return new ServerCall.Listener() {}; + } + }) + .build(); + + grpcCleanupRule.register( + InProcessServerBuilder.forName(serverName).directExecutor().addService(serviceDef).build() + .start()); + + OpenTelemetrySdk sdk = (OpenTelemetrySdk) openTelemetryRule.getOpenTelemetry(); + GrpcOpenTelemetry grpcOpenTelemetry = GrpcOpenTelemetry.newBuilder() + .sdk(sdk) + .enableMetrics(Arrays.asList( + "grpc.client.attempt.delay.duration", + "grpc.client.call.delay.duration", + "grpc.client.attempt.started")) + .enableTracing(true) + .addOptionalLabel("grpc.delay_type") + .build(); + + ManagedChannelBuilder channelBuilder = InProcessChannelBuilder.forName(serverName) + .directExecutor(); + grpcOpenTelemetry.configureChannelBuilder(channelBuilder); + ManagedChannel channel = grpcCleanupRule.register(channelBuilder.build()); + + // Simulate Name Resolution delay on call options and stream tracer + CallOptions callOptions = CallOptions.DEFAULT.withOption( + NAME_RESOLUTION_DELAYED, TimeUnit.MILLISECONDS.toNanos(120)); + + OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments( + sdk.getMeterProvider().get("grpc-java"), + ImmutableMap.of( + "grpc.client.attempt.delay.duration", true, + "grpc.client.call.delay.duration", true), + false); + OpenTelemetryMetricsModule module = new OpenTelemetryMetricsModule( + new FakeClock().getStopwatchSupplier(), resource, emptyList(), emptyList()); + OpenTelemetryMetricsModule.CallAttemptsTracerFactory factory = + new OpenTelemetryMetricsModule.CallAttemptsTracerFactory( + module, "target:///", callOptions, method.getFullMethodName(), + emptyList(), io.opentelemetry.context.Context.root()); + ClientStreamTracer delayTracer = factory.newClientStreamTracer( + ClientStreamTracer.StreamInfo.newBuilder().setCallOptions(callOptions).build(), + new Metadata()); + delayTracer.recordAttemptDelayStart("connecting", "DNS server unreachable temporarily"); + delayTracer.recordAttemptDelayEnd(); + factory.recordCallDelayStart("resolving", "DNS resolution pending"); + factory.recordCallDelayEnd(); + + final CountDownLatch latch = new CountDownLatch(1); + ClientCall call = channel.newCall(method, callOptions); + call.start(new ClientCall.Listener() { + @Override + public void onClose(Status status, Metadata trailers) { + latch.countDown(); + } + }, new Metadata()); + call.sendMessage("request_payload"); + call.halfClose(); + call.request(1); + assertThat(latch.await(5, TimeUnit.SECONDS)).isTrue(); + + io.opentelemetry.sdk.testing.assertj.OpenTelemetryAssertions + .assertThat(openTelemetryRule.getMetrics()) + .anySatisfy( + metric -> io.opentelemetry.sdk.testing.assertj.OpenTelemetryAssertions + .assertThat(metric) + .hasName("grpc.client.attempt.delay.duration") + .hasHistogramSatisfying( + histogram -> histogram.hasPointsSatisfying( + point -> { + point.hasAttribute( + AttributeKey.stringKey("grpc.delay_type"), "connecting"); + }))); + io.opentelemetry.sdk.testing.assertj.OpenTelemetryAssertions + .assertThat(openTelemetryRule.getMetrics()) + .anySatisfy( + metric -> io.opentelemetry.sdk.testing.assertj.OpenTelemetryAssertions + .assertThat(metric) + .hasName("grpc.client.call.delay.duration")); + } + + @Test + public void lbPolicyDelay_endToEndClientServerSimulation() throws Exception { + String serverName = InProcessServerBuilder.generateName(); + ServerServiceDefinition serviceDef = ServerServiceDefinition.builder("test.service") + .addMethod(method, new ServerCallHandler() { + @Override + public ServerCall.Listener startCall( + ServerCall call, Metadata headers) { + call.sendHeaders(new Metadata()); + call.sendMessage("response_payload"); + call.close(Status.OK, new Metadata()); + return new ServerCall.Listener() {}; + } + }) + .build(); + + grpcCleanupRule.register( + InProcessServerBuilder.forName(serverName).directExecutor().addService(serviceDef).build() + .start()); + + OpenTelemetrySdk sdk = (OpenTelemetrySdk) openTelemetryRule.getOpenTelemetry(); + GrpcOpenTelemetry grpcOpenTelemetry = GrpcOpenTelemetry.newBuilder() + .sdk(sdk) + .enableMetrics(Arrays.asList( + "grpc.client.attempt.delay.duration", + "grpc.client.call.delay.duration", + "grpc.client.attempt.started")) + .enableTracing(true) + .addOptionalLabel("grpc.delay_type") + .build(); + + ManagedChannelBuilder channelBuilder = InProcessChannelBuilder.forName(serverName) + .directExecutor(); + grpcOpenTelemetry.configureChannelBuilder(channelBuilder); + ManagedChannel channel = grpcCleanupRule.register(channelBuilder.build()); + + OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments( + sdk.getMeterProvider().get("grpc-java"), + ImmutableMap.of("grpc.client.attempt.delay.duration", true), + false); + 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 tracer = factory.newClientStreamTracer( + ClientStreamTracer.StreamInfo.newBuilder().setCallOptions(CallOptions.DEFAULT).build(), + new Metadata()); + tracer.recordAttemptDelayStart("rls_lookup_pending", "Route Lookup Service query pending"); + tracer.recordAttemptDelayEnd(); + + final CountDownLatch latch = new CountDownLatch(1); + ClientCall call = channel.newCall(method, CallOptions.DEFAULT); + call.start(new ClientCall.Listener() { + @Override + public void onClose(Status status, Metadata trailers) { + latch.countDown(); + } + }, new Metadata()); + call.sendMessage("request_payload"); + call.halfClose(); + call.request(1); + assertThat(latch.await(5, TimeUnit.SECONDS)).isTrue(); + + io.opentelemetry.sdk.testing.assertj.OpenTelemetryAssertions + .assertThat(openTelemetryRule.getMetrics()) + .anySatisfy( + metric -> io.opentelemetry.sdk.testing.assertj.OpenTelemetryAssertions + .assertThat(metric) + .hasName("grpc.client.attempt.delay.duration") + .hasHistogramSatisfying( + histogram -> histogram.hasPointsSatisfying( + point -> { + point.hasAttribute( + AttributeKey.stringKey("grpc.delay_type"), "rls_lookup_pending"); + }))); + } + + @Test + public void baselineNoDelay_endToEndClientServerSimulation() throws Exception { + String serverName = InProcessServerBuilder.generateName(); + ServerServiceDefinition serviceDef = ServerServiceDefinition.builder("test.service") + .addMethod(method, new ServerCallHandler() { + @Override + public ServerCall.Listener startCall( + ServerCall call, Metadata headers) { + call.sendHeaders(new Metadata()); + call.sendMessage("response_payload"); + call.close(Status.OK, new Metadata()); + return new ServerCall.Listener() {}; + } + }) + .build(); + + grpcCleanupRule.register( + InProcessServerBuilder.forName(serverName).directExecutor().addService(serviceDef).build() + .start()); + + OpenTelemetrySdk sdk = (OpenTelemetrySdk) openTelemetryRule.getOpenTelemetry(); + GrpcOpenTelemetry grpcOpenTelemetry = GrpcOpenTelemetry.newBuilder() + .sdk(sdk) + .enableMetrics(Arrays.asList( + "grpc.client.attempt.delay.duration", + "grpc.client.call.delay.duration", + "grpc.client.attempt.started")) + .enableTracing(true) + .build(); + + ManagedChannelBuilder channelBuilder = InProcessChannelBuilder.forName(serverName) + .directExecutor(); + grpcOpenTelemetry.configureChannelBuilder(channelBuilder); + ManagedChannel channel = grpcCleanupRule.register(channelBuilder.build()); + + final CountDownLatch latch = new CountDownLatch(1); + ClientCall call = channel.newCall(method, CallOptions.DEFAULT); + call.start(new ClientCall.Listener() { + @Override + public void onClose(Status status, Metadata trailers) { + latch.countDown(); + } + }, new Metadata()); + call.sendMessage("request_payload"); + call.halfClose(); + call.request(1); + assertThat(latch.await(5, TimeUnit.SECONDS)).isTrue(); + + boolean hasAttemptDelay = false; + boolean hasCallDelay = false; + for (MetricData m : openTelemetryRule.getMetrics()) { + if ("grpc.client.attempt.delay.duration".equals(m.getName()) + && !m.getHistogramData().getPoints().isEmpty()) { + hasAttemptDelay = true; + } + if ("grpc.client.call.delay.duration".equals(m.getName()) + && !m.getHistogramData().getPoints().isEmpty()) { + hasCallDelay = true; + } + } + assertThat(hasAttemptDelay).isFalse(); + assertThat(hasCallDelay).isFalse(); + } + 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..ee670ffa66e 100644 --- a/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryMetricsModuleTest.java +++ b/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryMetricsModuleTest.java @@ -30,6 +30,7 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyDouble; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.timeout; import static org.mockito.Mockito.verify; import com.google.common.collect.ImmutableMap; @@ -45,6 +46,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; @@ -1663,6 +1667,127 @@ 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.recordCallDelayStart("resolving", "dns resolution pending"); + fakeClock.forwardTime(500, TimeUnit.MILLISECONDS); + callAttemptsTracerFactory.recordCallDelayEnd(); + + 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_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); @@ -2289,6 +2414,174 @@ public void serverMetrics_recordsBaggage() { "baggage-val-1", capturedBaggage.getEntryValue("baggage-key-1")); } + @Test + public void clientMetrics_nameResolutionFailure_zeroAttempts() { + String target = "target:///"; + OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(testMeter, + enabledMetricsMap, disableDefaultMetrics); + OpenTelemetryMetricsModule module = newOpenTelemetryMetricsModule(resource); + OpenTelemetryMetricsModule.CallAttemptsTracerFactory callAttemptsTracerFactory = + new CallAttemptsTracerFactory(module, target, CALL_OPTIONS, method.getFullMethodName(), + emptyList(), Context.root()); + + fakeClock.forwardTime(50, TimeUnit.MILLISECONDS); + callAttemptsTracerFactory.callEnded(Status.UNAVAILABLE, CALL_OPTIONS); + + io.opentelemetry.api.common.Attributes clientAttributes = + io.opentelemetry.api.common.Attributes.of( + TARGET_KEY, target, + METHOD_KEY, method.getFullMethodName(), + STATUS_KEY, Code.UNAVAILABLE.toString()); + + assertThat(openTelemetryTesting.getMetrics()) + .anySatisfy( + metric -> + assertThat(metric) + .hasInstrumentationScope(InstrumentationScopeInfo.create( + OpenTelemetryConstants.INSTRUMENTATION_SCOPE)) + .hasName(CLIENT_CALL_DURATION) + .hasUnit("s") + .hasHistogramSatisfying( + histogram -> + histogram.hasPointsSatisfying( + point -> + point + .hasCount(1) + .hasSum(0.05) + .hasAttributes(clientAttributes) + .hasBucketBoundaries(latencyBuckets)))); + } + + @Test + public void clientMetrics_endToEnd_nameResolutionFailure_unavailable() throws Exception { + NameResolverProvider failingProvider = new NameResolverProvider() { + @Override + public NameResolver newNameResolver(URI targetUri, NameResolver.Args args) { + return new NameResolver() { + @Override + public String getServiceAuthority() { + return "failing.authority"; + } + + @Override + public void start(Listener2 listener) { + listener.onError(Status.UNAVAILABLE.withDescription("Name resolution failed")); + } + + @Override + public void shutdown() {} + }; + } + + @Override + protected boolean isAvailable() { + return true; + } + + @Override + protected int priority() { + return 5; + } + + @Override + public String getDefaultScheme() { + return "failingnr"; + } + + @Override + public String getScheme() { + return getDefaultScheme(); + } + + @Override + public Collection> getProducedSocketAddressTypes() { + return Collections.singleton(InProcessSocketAddress.class); + } + }; + + NameResolverRegistry.getDefaultRegistry().register(failingProvider); + + try { + OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(testMeter, + enabledMetricsMap, disableDefaultMetrics); + OpenTelemetryMetricsModule module = newOpenTelemetryMetricsModule(resource); + + String target = "failingnr:///test.service"; + ManagedChannel channel = grpcCleanup.register( + InProcessChannelBuilder.forTarget(target) + .directExecutor() + .intercept(module.getClientInterceptor(target)) + .build()); + + ClientCall call = channel.newCall(method, CallOptions.DEFAULT); + call.start(mockClientCallListener, new Metadata()); + + verify(mockClientCallListener, timeout(5000)) + .onClose(statusCaptor.capture(), any(Metadata.class)); + Status status = statusCaptor.getValue(); + assertEquals(Status.Code.UNAVAILABLE, status.getCode()); + + io.opentelemetry.api.common.Attributes clientAttributes = + io.opentelemetry.api.common.Attributes.of( + TARGET_KEY, target, + METHOD_KEY, method.getFullMethodName(), + STATUS_KEY, Code.UNAVAILABLE.toString()); + + assertThat(openTelemetryTesting.getMetrics()) + .anySatisfy( + metric -> + assertThat(metric) + .hasInstrumentationScope(InstrumentationScopeInfo.create( + OpenTelemetryConstants.INSTRUMENTATION_SCOPE)) + .hasName(CLIENT_CALL_DURATION) + .hasUnit("s") + .hasHistogramSatisfying( + histogram -> + histogram.hasPointsSatisfying( + point -> + point + .hasCount(1) + .hasAttributes(clientAttributes)))); + } finally { + NameResolverRegistry.getDefaultRegistry().deregister(failingProvider); + } + } + + @Test + public void clientMetrics_delayObservabilityDisabled_noDelayMetricsRecorded() { + String target = "target:///"; + OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(testMeter, + enabledMetricsMap, 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()); + + // When delay observability is disabled or default is unchanged, calls are no-ops + tracer.recordAttemptDelayStart("connecting", "attempt delay reason"); + tracer.recordAttemptDelayReasonChanged("changed reason"); + tracer.recordAttemptDelayEnd(); + + callAttemptsTracerFactory.recordCallDelayStart("resolving", "call delay reason"); + callAttemptsTracerFactory.recordCallDelayReasonChanged("changed call reason"); + callAttemptsTracerFactory.recordCallDelayEnd(); + + assertNotNull(tracer); + } + + @Test + public void clientMetrics_targetAttributeFilter_returnsFilteredOrOther() { + OpenTelemetryMetricsResource resource = GrpcOpenTelemetry.createMetricInstruments(testMeter, + enabledMetricsMap, disableDefaultMetrics); + OpenTelemetryMetricsModule module = newOpenTelemetryMetricsModule(resource); + + assertEquals("target:///", module.recordTarget("target:///")); + assertThat(module.recordTarget(null)).isNull(); + } + @Test public void serverMetrics_recordsBaggage_endToEnd() throws Exception { DoubleHistogram mockDurationHistogram = mock(DoubleHistogram.class); diff --git a/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryTracingModuleTest.java b/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryTracingModuleTest.java index 0b5bff1d036..b0d2c468b82 100644 --- a/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryTracingModuleTest.java +++ b/opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryTracingModuleTest.java @@ -46,6 +46,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; @@ -299,159 +302,154 @@ public void clientBasicTracingMocking() { } @Test - public void clientBasicTracingRule() { - OpenTelemetryTracingModule tracingModule = new OpenTelemetryTracingModule( - openTelemetryRule.getOpenTelemetry()); - Span clientSpan = tracerRule.spanBuilder("test-client-span").startSpan(); + 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(clientSpan, method); - Metadata headers = new Metadata(); - ClientStreamTracer clientStreamTracer = callTracer.newClientStreamTracer(STREAM_INFO, headers); - clientStreamTracer.createPendingStream(); - clientStreamTracer.streamCreated(Attributes.EMPTY, headers); - clientStreamTracer.outboundMessage(0); - clientStreamTracer.outboundMessageSent(0, 882, -1); - clientStreamTracer.inboundMessage(0); - clientStreamTracer.outboundMessage(1); - clientStreamTracer.outboundMessageSent(1, -1, 27); - clientStreamTracer.inboundMessageRead(0, 255, -1); - clientStreamTracer.inboundUncompressedSize(288); - clientStreamTracer.inboundMessageRead(1, 128, 128); - clientStreamTracer.inboundMessage(1); - clientStreamTracer.inboundUncompressedSize(128); + tracingModule.newClientCallTracer(mockClientSpan, method); + ClientStreamTracer clientStreamTracer = + callTracer.newClientStreamTracer(STREAM_INFO, new Metadata()); - clientStreamTracer.streamClosed(Status.OK); - callTracer.callEnded(Status.OK); + clientStreamTracer.recordAttemptDelayStart("connecting", "pick_first: attempting to connect"); + clientStreamTracer.recordAttemptDelayEnd(); - List spans = openTelemetryRule.getSpans(); - assertEquals(spans.size(), 2); - SpanData attemptSpanData = spans.get(0); - SpanData clientSpanData = spans.get(1); - assertEquals(attemptSpanData.getName(), "Attempt.package1.service2.method3"); - assertEquals(clientSpanData.getName(), "test-client-span"); - assertEquals(headers.keys(), ImmutableSet.of("traceparent")); - String spanContext = headers.get( - Metadata.Key.of("traceparent", Metadata.ASCII_STRING_MARSHALLER)); - assertEquals(spanContext.substring(3, 3 + TraceId.getLength()), - spans.get(1).getSpanContext().getTraceId()); + verify(mockTracer).spanBuilder(eq("Attempt Delay")); + verify(mockSpanBuilder).setAttribute(eq("grpc.delay_type"), eq("connecting")); + verify(mockDelaySpan).addEvent( + eq("Delay state transition"), + org.mockito.ArgumentMatchers.any()); + verify(mockDelaySpan).end(); + } - // parent(client) span data - List clientSpanEvents = clientSpanData.getEvents(); - assertEquals(clientSpanEvents.size(), 3); - assertEquals( - "Delayed name resolution complete", - clientSpanEvents.get(0).getName()); - assertTrue(clientSpanEvents.get(0).getAttributes().isEmpty()); + @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); - assertEquals( - "Inbound message" , - clientSpanEvents.get(1).getName()); - assertEquals( - io.opentelemetry.api.common.Attributes.builder() - .put("sequence-number", 0) - .put("message-size", 288) - .build(), - clientSpanEvents.get(1).getAttributes()); + OpenTelemetryTracingModule tracingModule = new OpenTelemetryTracingModule(mockOpenTelemetry); + CallAttemptsTracerFactory callTracer = + tracingModule.newClientCallTracer(mockClientSpan, method); - assertEquals( - "Inbound message" , - clientSpanEvents.get(2).getName()); - assertEquals( - io.opentelemetry.api.common.Attributes.builder() - .put("sequence-number", 1) - .put("message-size", 128) - .build(), - clientSpanEvents.get(2).getAttributes()); - assertEquals(clientSpanData.hasEnded(), true); + callTracer.recordCallDelayStart("resolving", "waiting for DNS query"); + callTracer.recordCallDelayEnd(); - // child(attempt) span data - List attemptSpanEvents = attemptSpanData.getEvents(); - assertEquals(clientSpanEvents.size(), 3); - assertEquals( - "Delayed LB pick complete", - attemptSpanEvents.get(0).getName()); - assertTrue(clientSpanEvents.get(0).getAttributes().isEmpty()); + verify(mockTracer).spanBuilder(eq("Call Delay")); + verify(mockSpanBuilder).setAttribute(eq("grpc.delay_type"), eq("resolving")); + verify(mockDelaySpan).addEvent( + eq("Delay state transition"), + org.mockito.ArgumentMatchers.any()); + verify(mockDelaySpan).end(); + } - assertEquals( - "Outbound message" , - attemptSpanEvents.get(1).getName()); - assertEquals( - io.opentelemetry.api.common.Attributes.builder() - .put("sequence-number", 0) - .put("message-size-compressed", 882) - .build(), - attemptSpanEvents.get(1).getAttributes()); + @Test + public void clientCallDelayTracing_endToEnd_nameResolutionDelay() throws Exception { + final CountDownLatch resolutionLatch = new CountDownLatch(1); + final AtomicReference capturedListener = new AtomicReference<>(); - assertEquals( - "Outbound message" , - attemptSpanEvents.get(2).getName()); - assertEquals( - io.opentelemetry.api.common.Attributes.builder() - .put("sequence-number", 1) - .put("message-size", 27) - .build(), - attemptSpanEvents.get(2).getAttributes()); + NameResolverProvider slowResolverProvider = new NameResolverProvider() { + @Override + protected boolean isAvailable() { + return true; + } - assertEquals( - "Inbound compressed message" , - attemptSpanEvents.get(3).getName()); - assertEquals( - io.opentelemetry.api.common.Attributes.builder() - .put("sequence-number", 0) - .put("message-size-compressed", 255) - .build(), - attemptSpanEvents.get(3).getAttributes()); + @Override + protected int priority() { + return 5; + } - assertEquals(attemptSpanData.hasEnded(), true); - } + @Override + public String getDefaultScheme() { + return "slowres"; + } - @Test - public void clientAttemptDelayTracing_reasonChangedInvariant() { - 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.recordAttemptDelayStart("connecting", "reason1"); - clientStreamTracer.recordAttemptDelayReasonChanged("reason2"); - clientStreamTracer.recordAttemptDelayStart("connecting", "reason3"); - clientStreamTracer.recordAttemptDelayEnd(); - clientStreamTracer.streamClosed(Status.OK); - callTracer.callEnded(Status.OK); - clientSpan.end(); + @Override + public NameResolver newNameResolver(URI targetUri, NameResolver.Args args) { + return new NameResolver() { + @Override + public String getServiceAuthority() { + return "slowres"; + } - List spans = openTelemetryRule.getSpans(); - assertEquals(3, spans.size()); - SpanData delaySpanData = spans.get(0); + @Override + public void start(Listener2 listener) { + capturedListener.set(listener); + resolutionLatch.countDown(); + } - assertEquals("Attempt Delay", delaySpanData.getName()); - assertEquals("connecting", delaySpanData.getAttributes().get( - AttributeKey.stringKey("grpc.delay_type"))); - assertEquals(3, delaySpanData.getEvents().size()); + @Override + public void shutdown() {} + }; + } + }; + NameResolverRegistry.getDefaultRegistry().register(slowResolverProvider); - EventData event1 = delaySpanData.getEvents().get(0); - assertEquals("Delay state transition", event1.getName()); - assertEquals("connecting", event1.getAttributes().get( - AttributeKey.stringKey("grpc.delay_type"))); - assertEquals("reason1", event1.getAttributes().get( - AttributeKey.stringKey("grpc.delay_reason"))); + GrpcOpenTelemetry grpcOpenTelemetry = GrpcOpenTelemetry.newBuilder() + .sdk(openTelemetryRule.getOpenTelemetry()) + .enableTracing(true) + .build(); - EventData event2 = delaySpanData.getEvents().get(1); - assertEquals("Delay state transition", event2.getName()); - assertEquals("connecting", event2.getAttributes().get( - AttributeKey.stringKey("grpc.delay_type"))); - assertEquals("reason2", event2.getAttributes().get( - AttributeKey.stringKey("grpc.delay_reason"))); + 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); - EventData event3 = delaySpanData.getEvents().get(2); - assertEquals("Delay state transition", event3.getName()); - assertEquals("connecting", event3.getAttributes().get( - AttributeKey.stringKey("grpc.delay_type"))); - assertEquals("reason3", event3.getAttributes().get( - AttributeKey.stringKey("grpc.delay_reason"))); + 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 ("Call 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 state transition".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 @@ -530,7 +528,7 @@ public String getServiceAuthority() { @Override public void start(Listener2 listener) { - listener.onResult(ResolutionResult.newBuilder() + listener.onResult(NameResolver.ResolutionResult.newBuilder() .setAddressesOrError(StatusOr.fromValue(Collections.singletonList( new EquivalentAddressGroup(new InProcessSocketAddress("test-e2e"))))) .build()); @@ -591,6 +589,162 @@ public void shutdown() {} assertTrue(foundTransition); } + @Test + public void clientBasicTracingRule() { + OpenTelemetryTracingModule tracingModule = new OpenTelemetryTracingModule( + openTelemetryRule.getOpenTelemetry()); + Span clientSpan = tracerRule.spanBuilder("test-client-span").startSpan(); + CallAttemptsTracerFactory callTracer = + tracingModule.newClientCallTracer(clientSpan, method); + Metadata headers = new Metadata(); + ClientStreamTracer clientStreamTracer = callTracer.newClientStreamTracer(STREAM_INFO, headers); + clientStreamTracer.createPendingStream(); + clientStreamTracer.streamCreated(Attributes.EMPTY, headers); + clientStreamTracer.outboundMessage(0); + clientStreamTracer.outboundMessageSent(0, 882, -1); + clientStreamTracer.inboundMessage(0); + clientStreamTracer.outboundMessage(1); + clientStreamTracer.outboundMessageSent(1, -1, 27); + clientStreamTracer.inboundMessageRead(0, 255, -1); + clientStreamTracer.inboundUncompressedSize(288); + clientStreamTracer.inboundMessageRead(1, 128, 128); + clientStreamTracer.inboundMessage(1); + clientStreamTracer.inboundUncompressedSize(128); + + clientStreamTracer.streamClosed(Status.OK); + callTracer.callEnded(Status.OK); + + List spans = openTelemetryRule.getSpans(); + assertEquals(spans.size(), 2); + SpanData attemptSpanData = spans.get(0); + SpanData clientSpanData = spans.get(1); + assertEquals(attemptSpanData.getName(), "Attempt.package1.service2.method3"); + assertEquals(clientSpanData.getName(), "test-client-span"); + assertEquals(headers.keys(), ImmutableSet.of("traceparent")); + String spanContext = headers.get( + Metadata.Key.of("traceparent", Metadata.ASCII_STRING_MARSHALLER)); + assertEquals(spanContext.substring(3, 3 + TraceId.getLength()), + spans.get(1).getSpanContext().getTraceId()); + + // parent(client) span data + List clientSpanEvents = clientSpanData.getEvents(); + assertEquals(clientSpanEvents.size(), 3); + assertEquals( + "Delayed name resolution complete", + clientSpanEvents.get(0).getName()); + assertTrue(clientSpanEvents.get(0).getAttributes().isEmpty()); + + assertEquals( + "Inbound message" , + clientSpanEvents.get(1).getName()); + assertEquals( + io.opentelemetry.api.common.Attributes.builder() + .put("sequence-number", 0) + .put("message-size", 288) + .build(), + clientSpanEvents.get(1).getAttributes()); + + assertEquals( + "Inbound message" , + clientSpanEvents.get(2).getName()); + assertEquals( + io.opentelemetry.api.common.Attributes.builder() + .put("sequence-number", 1) + .put("message-size", 128) + .build(), + clientSpanEvents.get(2).getAttributes()); + assertEquals(clientSpanData.hasEnded(), true); + + // child(attempt) span data + List attemptSpanEvents = attemptSpanData.getEvents(); + assertEquals(clientSpanEvents.size(), 3); + assertEquals( + "Delayed LB pick complete", + attemptSpanEvents.get(0).getName()); + assertTrue(clientSpanEvents.get(0).getAttributes().isEmpty()); + + assertEquals( + "Outbound message" , + attemptSpanEvents.get(1).getName()); + assertEquals( + io.opentelemetry.api.common.Attributes.builder() + .put("sequence-number", 0) + .put("message-size-compressed", 882) + .build(), + attemptSpanEvents.get(1).getAttributes()); + + assertEquals( + "Outbound message" , + attemptSpanEvents.get(2).getName()); + assertEquals( + io.opentelemetry.api.common.Attributes.builder() + .put("sequence-number", 1) + .put("message-size", 27) + .build(), + attemptSpanEvents.get(2).getAttributes()); + + assertEquals( + "Inbound compressed message" , + attemptSpanEvents.get(3).getName()); + assertEquals( + io.opentelemetry.api.common.Attributes.builder() + .put("sequence-number", 0) + .put("message-size-compressed", 255) + .build(), + attemptSpanEvents.get(3).getAttributes()); + + assertEquals(attemptSpanData.hasEnded(), true); + } + + @Test + public void clientAttemptDelayTracing_reasonChangedInvariant() { + 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", "reason1"); + clientStreamTracer.recordAttemptDelayReasonChanged("reason2"); + clientStreamTracer.recordAttemptDelayStart("connecting", "reason3"); + clientStreamTracer.recordAttemptDelayEnd(); + clientStreamTracer.streamClosed(Status.OK); + callTracer.callEnded(Status.OK); + clientSpan.end(); + + List spans = openTelemetryRule.getSpans(); + assertEquals(3, spans.size()); + SpanData delaySpanData = spans.get(0); + + assertEquals("Attempt 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("reason1", event1.getAttributes().get( + AttributeKey.stringKey("grpc.delay_reason"))); + + EventData event2 = delaySpanData.getEvents().get(1); + assertEquals("Delay state transition", event2.getName()); + assertEquals("connecting", event2.getAttributes().get( + AttributeKey.stringKey("grpc.delay_type"))); + assertEquals("reason2", event2.getAttributes().get( + AttributeKey.stringKey("grpc.delay_reason"))); + + EventData event3 = delaySpanData.getEvents().get(2); + assertEquals("Delay state transition", event3.getName()); + assertEquals("connecting", event3.getAttributes().get( + AttributeKey.stringKey("grpc.delay_type"))); + assertEquals("reason3", event3.getAttributes().get( + AttributeKey.stringKey("grpc.delay_reason"))); + } + @Test public void clientAttemptDelayStart_featureFlagDisabled_zeroChildSpans() { System.setProperty("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY", "false"); diff --git a/util/src/test/java/io/grpc/util/ForwardingClientStreamTracerTest.java b/util/src/test/java/io/grpc/util/ForwardingClientStreamTracerTest.java index dbd7e99b29a..f6397d878f9 100644 --- a/util/src/test/java/io/grpc/util/ForwardingClientStreamTracerTest.java +++ b/util/src/test/java/io/grpc/util/ForwardingClientStreamTracerTest.java @@ -40,6 +40,19 @@ public void allMethodsForwarded() throws Exception { Collections.emptyList()); } + @Test + public void attemptDelayMethodsForwarded() { + TestClientStreamTracer tracer = new TestClientStreamTracer(); + tracer.recordAttemptDelayStart("connecting", "test"); + org.mockito.Mockito.verify(mockDelegate).recordAttemptDelayStart("connecting", "test"); + + tracer.recordAttemptDelayReasonChanged("test2"); + org.mockito.Mockito.verify(mockDelegate).recordAttemptDelayReasonChanged("test2"); + + tracer.recordAttemptDelayEnd(); + org.mockito.Mockito.verify(mockDelegate).recordAttemptDelayEnd(); + } + @SuppressWarnings("deprecation") private final class TestClientStreamTracer extends ForwardingClientStreamTracer { @Override