Skip to content

Commit 0ca843d

Browse files
committed
opentelemetry: Add real channel E2E delay tracing & metrics tests via InProcessChannelBuilder (Proposal A121)
1 parent 5245f15 commit 0ca843d

2 files changed

Lines changed: 158 additions & 11 deletions

File tree

opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryMetricsModuleTest.java

Lines changed: 142 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,11 +40,19 @@
4040
import io.grpc.ClientInterceptor;
4141
import io.grpc.ClientInterceptors;
4242
import io.grpc.ClientStreamTracer;
43+
import io.grpc.ConnectivityState;
44+
import io.grpc.EquivalentAddressGroup;
4345
import io.grpc.Grpc;
4446
import io.grpc.KnownLength;
47+
import io.grpc.LoadBalancer;
48+
import io.grpc.LoadBalancerProvider;
49+
import io.grpc.LoadBalancerRegistry;
4550
import io.grpc.ManagedChannel;
4651
import io.grpc.Metadata;
4752
import io.grpc.MethodDescriptor;
53+
import io.grpc.NameResolver;
54+
import io.grpc.NameResolverProvider;
55+
import io.grpc.NameResolverRegistry;
4856
import io.grpc.Server;
4957
import io.grpc.ServerCall;
5058
import io.grpc.ServerCallHandler;
@@ -56,6 +64,7 @@
5664
import io.grpc.Status.Code;
5765
import io.grpc.inprocess.InProcessChannelBuilder;
5866
import io.grpc.inprocess.InProcessServerBuilder;
67+
import io.grpc.inprocess.InProcessSocketAddress;
5968
import io.grpc.internal.FakeClock;
6069
import io.grpc.internal.StatsTraceContext.ServerCallMethodListener;
6170
import io.grpc.opentelemetry.GrpcOpenTelemetry.TargetFilter;
@@ -79,10 +88,15 @@
7988
import io.opentelemetry.sdk.testing.junit4.OpenTelemetryRule;
8089
import java.io.IOException;
8190
import java.io.InputStream;
91+
import java.net.SocketAddress;
92+
import java.net.URI;
8293
import java.util.Arrays;
94+
import java.util.Collection;
95+
import java.util.Collections;
8396
import java.util.List;
8497
import java.util.Map;
8598
import java.util.Optional;
99+
import java.util.concurrent.CountDownLatch;
86100
import java.util.concurrent.TimeUnit;
87101
import java.util.concurrent.atomic.AtomicReference;
88102
import javax.annotation.Nullable;
@@ -1648,6 +1662,134 @@ public void clientAttemptDelayDuration_recorded() {
16481662
})));
16491663
}
16501664

1665+
@Test
1666+
public void clientAttemptDelayDuration_endToEnd_inProcessTransport() throws Exception {
1667+
final CountDownLatch latch = new CountDownLatch(1);
1668+
LoadBalancerProvider slowLbProvider = new LoadBalancerProvider() {
1669+
@Override
1670+
public boolean isAvailable() {
1671+
return true;
1672+
}
1673+
1674+
@Override
1675+
public int getPriority() {
1676+
return 5;
1677+
}
1678+
1679+
@Override
1680+
public String getPolicyName() {
1681+
return "slow_metrics_connecting_policy";
1682+
}
1683+
1684+
@Override
1685+
public LoadBalancer newLoadBalancer(LoadBalancer.Helper helper) {
1686+
return new LoadBalancer() {
1687+
@Override
1688+
public Status acceptResolvedAddresses(LoadBalancer.ResolvedAddresses resolvedAddresses) {
1689+
helper.updateBalancingState(ConnectivityState.CONNECTING, new SubchannelPicker() {
1690+
@Override
1691+
public PickResult pickSubchannel(PickSubchannelArgs args) {
1692+
return PickResult.withNoResult("connecting",
1693+
"Simulated slow TLS handshake with backend");
1694+
}
1695+
});
1696+
latch.countDown();
1697+
return Status.OK;
1698+
}
1699+
1700+
@Override
1701+
public void handleNameResolutionError(Status error) {}
1702+
1703+
@Override
1704+
public void shutdown() {}
1705+
};
1706+
}
1707+
};
1708+
LoadBalancerRegistry.getDefaultRegistry().register(slowLbProvider);
1709+
1710+
NameResolverProvider customResolverProvider = new NameResolverProvider() {
1711+
@Override
1712+
protected boolean isAvailable() {
1713+
return true;
1714+
}
1715+
1716+
@Override
1717+
protected int priority() {
1718+
return 5;
1719+
}
1720+
1721+
@Override
1722+
public String getDefaultScheme() {
1723+
return "inprocmetricse2e";
1724+
}
1725+
1726+
@Override
1727+
public Collection<Class<? extends SocketAddress>> getProducedSocketAddressTypes() {
1728+
return Collections.singleton(InProcessSocketAddress.class);
1729+
}
1730+
1731+
@Override
1732+
public NameResolver newNameResolver(URI targetUri, NameResolver.Args args) {
1733+
return new NameResolver() {
1734+
@Override
1735+
public String getServiceAuthority() {
1736+
return "inprocmetricse2e";
1737+
}
1738+
1739+
@Override
1740+
public void start(Listener2 listener) {
1741+
listener.onResult(ResolutionResult.newBuilder()
1742+
.setAddresses(Collections.singletonList(new EquivalentAddressGroup(
1743+
new InProcessSocketAddress("test-metrics-e2e"))))
1744+
.build());
1745+
}
1746+
1747+
@Override
1748+
public void shutdown() {}
1749+
};
1750+
}
1751+
};
1752+
NameResolverRegistry.getDefaultRegistry().register(customResolverProvider);
1753+
1754+
GrpcOpenTelemetry grpcOpenTelemetry = GrpcOpenTelemetry.newBuilder()
1755+
.sdk(openTelemetryTesting.getOpenTelemetry())
1756+
.enableMetrics(Collections.singleton("grpc.client.attempt.delay.duration"))
1757+
.build();
1758+
1759+
InProcessChannelBuilder channelBuilder =
1760+
InProcessChannelBuilder.forTarget("inprocmetricse2e:///test-metrics-e2e")
1761+
.defaultLoadBalancingPolicy("slow_metrics_connecting_policy");
1762+
grpcOpenTelemetry.configureChannelBuilder(channelBuilder);
1763+
ManagedChannel channel = channelBuilder.build();
1764+
try {
1765+
ClientCall<String, String> call = channel.newCall(method, CallOptions.DEFAULT);
1766+
call.start(new ClientCall.Listener<String>() {}, new Metadata());
1767+
call.request(1);
1768+
1769+
latch.await(5, TimeUnit.SECONDS);
1770+
Thread.sleep(50);
1771+
call.cancel("End test delay segment", null);
1772+
} finally {
1773+
channel.shutdownNow();
1774+
channel.awaitTermination(5, TimeUnit.SECONDS);
1775+
LoadBalancerRegistry.getDefaultRegistry().deregister(slowLbProvider);
1776+
NameResolverRegistry.getDefaultRegistry().deregister(customResolverProvider);
1777+
}
1778+
1779+
assertThat(openTelemetryTesting.getMetrics())
1780+
.anySatisfy(
1781+
metric -> assertThat(metric)
1782+
.hasName("grpc.client.attempt.delay.duration")
1783+
.hasHistogramSatisfying(
1784+
histogram -> histogram.hasPointsSatisfying(
1785+
point -> {
1786+
point.hasAttribute(METHOD_KEY, method.getFullMethodName());
1787+
point.hasAttribute(TARGET_KEY, "inprocmetricse2e:///test-metrics-e2e");
1788+
point.hasAttribute(
1789+
AttributeKey.stringKey("grpc.delay_type"), "connecting");
1790+
})));
1791+
}
1792+
16511793
@Test
16521794
public void clientAttemptDelayStart_featureFlagDisabled_zeroMetrics() {
16531795
System.setProperty("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY", "false");

opentelemetry/src/test/java/io/grpc/opentelemetry/OpenTelemetryTracingModuleTest.java

Lines changed: 16 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,7 @@
6565
import io.grpc.Status;
6666
import io.grpc.inprocess.InProcessChannelBuilder;
6767
import io.grpc.inprocess.InProcessServerBuilder;
68+
import io.grpc.inprocess.InProcessSocketAddress;
6869
import io.grpc.opentelemetry.OpenTelemetryTracingModule.CallAttemptsTracerFactory;
6970
import io.grpc.opentelemetry.internal.OpenTelemetryConstants;
7071
import io.grpc.testing.GrpcCleanupRule;
@@ -92,8 +93,14 @@
9293
import io.opentelemetry.sdk.trace.data.SpanData;
9394
import java.io.IOException;
9495
import java.io.InputStream;
96+
import java.net.SocketAddress;
97+
import java.net.URI;
9598
import java.util.Arrays;
99+
import java.util.Collection;
100+
import java.util.Collections;
96101
import java.util.List;
102+
import java.util.concurrent.CountDownLatch;
103+
import java.util.concurrent.TimeUnit;
97104
import java.util.concurrent.atomic.AtomicReference;
98105
import org.junit.After;
99106
import org.junit.Before;
@@ -448,8 +455,7 @@ public void clientAttemptDelayTracing_reasonChangedInvariant() {
448455

449456
@Test
450457
public void clientAttemptDelayTracing_endToEnd_inProcessTransport() throws Exception {
451-
final java.util.concurrent.CountDownLatch latch =
452-
new java.util.concurrent.CountDownLatch(1);
458+
final CountDownLatch latch = new CountDownLatch(1);
453459
LoadBalancerProvider slowLbProvider = new LoadBalancerProvider() {
454460
@Override
455461
public boolean isAvailable() {
@@ -509,13 +515,12 @@ public String getDefaultScheme() {
509515
}
510516

511517
@Override
512-
public java.util.Collection<Class<? extends java.net.SocketAddress>>
513-
getProducedSocketAddressTypes() {
514-
return java.util.Collections.singleton(io.grpc.inprocess.InProcessSocketAddress.class);
518+
public Collection<Class<? extends SocketAddress>> getProducedSocketAddressTypes() {
519+
return Collections.singleton(InProcessSocketAddress.class);
515520
}
516521

517522
@Override
518-
public NameResolver newNameResolver(java.net.URI targetUri, NameResolver.Args args) {
523+
public NameResolver newNameResolver(URI targetUri, NameResolver.Args args) {
519524
return new NameResolver() {
520525
@Override
521526
public String getServiceAuthority() {
@@ -525,8 +530,8 @@ public String getServiceAuthority() {
525530
@Override
526531
public void start(Listener2 listener) {
527532
listener.onResult(ResolutionResult.newBuilder()
528-
.setAddresses(java.util.Collections.singletonList(new EquivalentAddressGroup(
529-
new io.grpc.inprocess.InProcessSocketAddress("test-e2e"))))
533+
.setAddresses(Collections.singletonList(new EquivalentAddressGroup(
534+
new InProcessSocketAddress("test-e2e"))))
530535
.build());
531536
}
532537

@@ -552,12 +557,12 @@ public void shutdown() {}
552557
call.start(new ClientCall.Listener<String>() {}, new Metadata());
553558
call.request(1);
554559

555-
latch.await(5, java.util.concurrent.TimeUnit.SECONDS);
560+
latch.await(5, TimeUnit.SECONDS);
556561
Thread.sleep(50);
557562
call.cancel("End test delay segment", null);
558563
} finally {
559564
channel.shutdownNow();
560-
channel.awaitTermination(5, java.util.concurrent.TimeUnit.SECONDS);
565+
channel.awaitTermination(5, TimeUnit.SECONDS);
561566
LoadBalancerRegistry.getDefaultRegistry().deregister(slowLbProvider);
562567
NameResolverRegistry.getDefaultRegistry().deregister(customResolverProvider);
563568
}
@@ -575,7 +580,7 @@ public void shutdown() {}
575580
delaySpanData.getAttributes().get(AttributeKey.stringKey("grpc.delay_type")));
576581

577582
boolean foundTransition = false;
578-
for (io.opentelemetry.sdk.trace.data.EventData event : delaySpanData.getEvents()) {
583+
for (EventData event : delaySpanData.getEvents()) {
579584
if ("Delay state transition".equals(event.getName())
580585
&& "Simulated slow TLS handshake with backend".equals(
581586
event.getAttributes().get(AttributeKey.stringKey("grpc.delay_reason")))) {

0 commit comments

Comments
 (0)