Skip to content
1 change: 1 addition & 0 deletions gradle/libs.versions.toml
Original file line number Diff line number Diff line change
Expand Up @@ -122,6 +122,7 @@ opencensus-impl = { module = "io.opencensus:opencensus-impl", version.ref = "ope
opentelemetry-api = "io.opentelemetry:opentelemetry-api:1.63.0"
opentelemetry-exporter-prometheus = "io.opentelemetry:opentelemetry-exporter-prometheus:1.63.0-alpha"
opentelemetry-gcp-resources = "io.opentelemetry.contrib:opentelemetry-gcp-resources:1.57.0-alpha"
opentelemetry-exporter-otlp = "io.opentelemetry:opentelemetry-exporter-otlp:1.63.0"
opentelemetry-sdk-extension-autoconfigure = "io.opentelemetry:opentelemetry-sdk-extension-autoconfigure:1.63.0"
opentelemetry-sdk-testing = "io.opentelemetry:opentelemetry-sdk-testing:1.63.0"
perfmark-api = "io.perfmark:perfmark-api:0.27.0"
Expand Down
1 change: 1 addition & 0 deletions interop-testing/build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ dependencies {
libraries.opencensus.contrib.grpc.metrics,
libraries.google.auth.oauth2Http,
libraries.opentelemetry.sdk.extension.autoconfigure,
libraries.opentelemetry.exporter.otlp,
libraries.guava.jre // Fix checkUpperBoundDeps using -android
api project(':grpc-api'),
project(':grpc-stub'),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,9 @@
import io.grpc.netty.NettyChannelBuilder;
import io.grpc.okhttp.InternalOkHttpChannelBuilder;
import io.grpc.okhttp.OkHttpChannelBuilder;
import io.grpc.opentelemetry.GrpcOpenTelemetry;
import io.grpc.opentelemetry.GrpcTraceBinContextPropagator;
import io.grpc.opentelemetry.InternalGrpcOpenTelemetry;
import io.grpc.stub.ClientCalls;
import io.grpc.stub.MetadataUtils;
import io.grpc.stub.StreamObserver;
Expand All @@ -68,6 +71,9 @@
import io.grpc.testing.integration.Messages.StreamingOutputCallRequest;
import io.grpc.testing.integration.Messages.StreamingOutputCallResponse;
import io.grpc.testing.integration.Messages.TestOrcaReport;
import io.opentelemetry.context.propagation.TextMapPropagator;
import io.opentelemetry.sdk.OpenTelemetrySdk;
import io.opentelemetry.sdk.autoconfigure.AutoConfiguredOpenTelemetrySdk;
import java.io.File;
import java.io.FileInputStream;
import java.io.InputStream;
Expand All @@ -79,6 +85,7 @@
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import javax.annotation.Nullable;
import org.codehaus.mojo.animal_sniffer.IgnoreJRERequirement;

/**
* Application that starts a client for the {@link TestServiceGrpc.TestServiceImplBase} and runs
Expand All @@ -101,9 +108,9 @@ public static void main(String[] args) throws Exception {
client.parseArgs(args);
customBackendMetricsLoadBalancerProvider = new CustomBackendMetricsLoadBalancerProvider();
LoadBalancerRegistry.getDefaultRegistry().register(customBackendMetricsLoadBalancerProvider);
client.setUp();

try {
client.setUp();
client.run();
} finally {
client.tearDown();
Expand All @@ -118,6 +125,8 @@ public static void main(String[] args) throws Exception {
private boolean useTls = true;
private boolean useAlts = false;
private boolean useH2cUpgrade = false;
private boolean enableOpentelemetry = false;
private OpenTelemetrySdk openTelemetrySdk;
private String customCredentialsType;
private boolean useTestCa;
private boolean useOkHttp;
Expand Down Expand Up @@ -219,6 +228,8 @@ void parseArgs(String[] args) throws Exception {
numThreads = Integer.parseInt(value);
} else if ("additional_metadata".equals(key)) {
additionalMetadata = value;
} else if ("enable_opentelemetry".equals(key)) {
enableOpentelemetry = Boolean.parseBoolean(value);
} else {
System.err.println("Unknown argument: " + key);
usage = true;
Expand Down Expand Up @@ -306,21 +317,36 @@ void parseArgs(String[] args) throws Exception {
}

@VisibleForTesting
@IgnoreJRERequirement // OpenTelemetry uses Java 8+ APIs
void setUp() {
if (enableOpentelemetry) {
AutoConfiguredOpenTelemetrySdk autoSdk = AutoConfiguredOpenTelemetrySdk.builder()
.addPropagatorCustomizer(
(previous, config) ->
TextMapPropagator.composite(
previous, GrpcTraceBinContextPropagator.defaultInstance()))
.build();
this.openTelemetrySdk = autoSdk.getOpenTelemetrySdk();
GrpcOpenTelemetry.Builder grpcOpentelemetryBuilder = GrpcOpenTelemetry.newBuilder()
.sdk(openTelemetrySdk);
InternalGrpcOpenTelemetry.enableTracing(grpcOpentelemetryBuilder, true);
GrpcOpenTelemetry grpcOpenTelemetry = grpcOpentelemetryBuilder.build();
grpcOpenTelemetry.registerGlobal();
}
tester.setUp();
}

private synchronized void tearDown() {
try {
tester.tearDown();
} finally {
if (customBackendMetricsLoadBalancerProvider != null) {
LoadBalancerRegistry.getDefaultRegistry()
.deregister(customBackendMetricsLoadBalancerProvider);
}
} catch (RuntimeException ex) {
throw ex;
} catch (Exception ex) {
throw new RuntimeException(ex);
if (openTelemetrySdk != null) {
openTelemetrySdk.close();
}
}
}

Expand Down Expand Up @@ -424,28 +450,36 @@ private void runTest(TestCases testCase) throws Exception {

case SERVICE_ACCOUNT_CREDS: {
String jsonKey = Files.asCharSource(new File(serviceAccountKeyFile), UTF_8).read();
FileInputStream credentialsStream = new FileInputStream(new File(serviceAccountKeyFile));
tester.serviceAccountCreds(jsonKey, credentialsStream, oauthScope);
try (FileInputStream credentialsStream =
new FileInputStream(new File(serviceAccountKeyFile))) {
tester.serviceAccountCreds(jsonKey, credentialsStream, oauthScope);
}
break;
}

case JWT_TOKEN_CREDS: {
FileInputStream credentialsStream = new FileInputStream(new File(serviceAccountKeyFile));
tester.jwtTokenCreds(credentialsStream);
try (FileInputStream credentialsStream =
new FileInputStream(new File(serviceAccountKeyFile))) {
tester.jwtTokenCreds(credentialsStream);
}
break;
}

case OAUTH2_AUTH_TOKEN: {
String jsonKey = Files.asCharSource(new File(serviceAccountKeyFile), UTF_8).read();
FileInputStream credentialsStream = new FileInputStream(new File(serviceAccountKeyFile));
tester.oauth2AuthToken(jsonKey, credentialsStream, oauthScope);
try (FileInputStream credentialsStream =
new FileInputStream(new File(serviceAccountKeyFile))) {
tester.oauth2AuthToken(jsonKey, credentialsStream, oauthScope);
}
break;
}

case PER_RPC_CREDS: {
String jsonKey = Files.asCharSource(new File(serviceAccountKeyFile), UTF_8).read();
FileInputStream credentialsStream = new FileInputStream(new File(serviceAccountKeyFile));
tester.perRpcCreds(jsonKey, credentialsStream, oauthScope);
try (FileInputStream credentialsStream =
new FileInputStream(new File(serviceAccountKeyFile))) {
tester.perRpcCreds(jsonKey, credentialsStream, oauthScope);
}
break;
}

Expand Down Expand Up @@ -677,7 +711,8 @@ protected ManagedChannelBuilder<?> createChannelBuilder() {
if (serverPort == 0) {
nettyBuilder = NettyChannelBuilder.forTarget(serverHost, channelCredentials);
} else {
nettyBuilder = NettyChannelBuilder.forAddress(serverHost, serverPort, channelCredentials);
nettyBuilder =
NettyChannelBuilder.forAddress(serverHost, serverPort, channelCredentials);
}
nettyBuilder.flowControlWindow(AbstractInteropTest.TEST_FLOW_CONTROL_WINDOW);
if (serverHostOverride != null) {
Expand Down Expand Up @@ -795,8 +830,8 @@ public void cacheableUnary() {
}

/** Sends a large unary rpc with service account credentials. */
public void serviceAccountCreds(String jsonKey, InputStream credentialsStream, String authScope)
throws Exception {
public void serviceAccountCreds(
String jsonKey, InputStream credentialsStream, String authScope) throws Exception {
// cast to ServiceAccountCredentials to double-check the right type of object was created.
GoogleCredentials credentials =
ServiceAccountCredentials.class.cast(GoogleCredentials.fromStream(credentialsStream));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,17 +28,24 @@
import io.grpc.TlsServerCredentials;
import io.grpc.alts.AltsServerCredentials;
import io.grpc.netty.NettyServerBuilder;
import io.grpc.opentelemetry.GrpcOpenTelemetry;
import io.grpc.opentelemetry.GrpcTraceBinContextPropagator;
import io.grpc.opentelemetry.InternalGrpcOpenTelemetry;
import io.grpc.services.MetricRecorder;
import io.grpc.testing.TlsTesting;
import io.grpc.xds.orca.OrcaMetricReportingServerInterceptor;
import io.grpc.xds.orca.OrcaServiceImpl;
import io.opentelemetry.context.propagation.TextMapPropagator;
import io.opentelemetry.sdk.OpenTelemetrySdk;
import io.opentelemetry.sdk.autoconfigure.AutoConfiguredOpenTelemetrySdk;
import java.net.InetSocketAddress;
import java.net.SocketAddress;
import java.util.List;
import java.util.Locale;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import org.codehaus.mojo.animal_sniffer.IgnoreJRERequirement;

/** Server that manages startup/shutdown of a single {@code TestService}. */
public class TestServiceServer {
Expand Down Expand Up @@ -76,6 +83,8 @@ public void run() {
private boolean useTls = true;
private boolean useAlts = false;
private int mcsLimit = -1;
private boolean enableOpentelemetry = false;
private OpenTelemetrySdk openTelemetrySdk;

private ScheduledExecutorService executor;
private Server server;
Expand Down Expand Up @@ -123,6 +132,8 @@ void parseArgs(String[] args) {
mcsLimit = Integer.parseInt(value);
// TODO: Make Netty server builder usable for IPV6 as well (not limited to MCS handling)
addressType = Util.AddressType.IPV4; // To use NettyServerBuilder
} else if ("enable_opentelemetry".equals(key)) {
enableOpentelemetry = Boolean.parseBoolean(value);
} else {
System.err.println("Unknown argument: " + key);
usage = true;
Expand Down Expand Up @@ -155,80 +166,110 @@ void parseArgs(String[] args) {

@SuppressWarnings("AddressSelection")
@VisibleForTesting
@IgnoreJRERequirement // OpenTelemetry uses Java 8+ APIs
void start() throws Exception {
executor = Executors.newSingleThreadScheduledExecutor();
ServerCredentials serverCreds;
if (useAlts) {
if (localHandshakerPort > -1) {
serverCreds = AltsServerCredentials.newBuilder()
.enableUntrustedAltsForTesting()
.setHandshakerAddressForTesting("localhost:" + localHandshakerPort).build();
try {
if (enableOpentelemetry) {
AutoConfiguredOpenTelemetrySdk autoSdk = AutoConfiguredOpenTelemetrySdk.builder()
.addPropagatorCustomizer(
(previous, config) ->
TextMapPropagator.composite(
previous, GrpcTraceBinContextPropagator.defaultInstance()))
.build();
this.openTelemetrySdk = autoSdk.getOpenTelemetrySdk();
GrpcOpenTelemetry.Builder grpcOpentelemetryBuilder = GrpcOpenTelemetry.newBuilder()
.sdk(openTelemetrySdk);
InternalGrpcOpenTelemetry.enableTracing(grpcOpentelemetryBuilder, true);
GrpcOpenTelemetry grpcOpenTelemetry = grpcOpentelemetryBuilder.build();
grpcOpenTelemetry.registerGlobal();
}
executor = Executors.newSingleThreadScheduledExecutor();
ServerCredentials serverCreds;
if (useAlts) {
if (localHandshakerPort > -1) {
serverCreds = AltsServerCredentials.newBuilder()
.enableUntrustedAltsForTesting()
.setHandshakerAddressForTesting("localhost:" + localHandshakerPort).build();
} else {
serverCreds = AltsServerCredentials.create();
}
} else if (useTls) {
serverCreds = TlsServerCredentials.create(
TlsTesting.loadCert("server1.pem"), TlsTesting.loadCert("server1.key"));
} else {
serverCreds = AltsServerCredentials.create();
serverCreds = InsecureServerCredentials.create();
}
} else if (useTls) {
serverCreds = TlsServerCredentials.create(
TlsTesting.loadCert("server1.pem"), TlsTesting.loadCert("server1.key"));
} else {
serverCreds = InsecureServerCredentials.create();
}
MetricRecorder metricRecorder = MetricRecorder.newInstance();
BindableService orcaOobService =
OrcaServiceImpl.createService(executor, metricRecorder, 1, TimeUnit.SECONDS);
MetricRecorder metricRecorder = MetricRecorder.newInstance();
BindableService orcaOobService =
OrcaServiceImpl.createService(executor, metricRecorder, 1, TimeUnit.SECONDS);

// Create ServerBuilder with appropriate addresses
// - IPV4_IPV6: bind to wildcard which covers all addresses on all interfaces of both families
// - IPV4: bind to v4 address for local hostname + v4 localhost
// - IPV6: bind to all v6 addresses for local hostname + v6 localhost
ServerBuilder<?> serverBuilder;
switch (addressType) {
case IPV4_IPV6:
serverBuilder = Grpc.newServerBuilderForPort(port, serverCreds);
break;
case IPV4:
SocketAddress v4Address = Util.getV4Address(port);
InetSocketAddress localV4Address = new InetSocketAddress("127.0.0.1", port);
serverBuilder =
NettyServerBuilder.forAddress(localV4Address, serverCreds);
if (v4Address != null && !v4Address.equals(localV4Address)) {
((NettyServerBuilder) serverBuilder).addListenAddress(v4Address);
}
if (mcsLimit != -1) {
((NettyServerBuilder) serverBuilder).maxConcurrentCallsPerConnection(mcsLimit);
}
break;
case IPV6:
List<SocketAddress> v6Addresses = Util.getV6Addresses(port);
InetSocketAddress localV6Address = new InetSocketAddress("::1", port);
serverBuilder =
NettyServerBuilder.forAddress(localV6Address, serverCreds);
for (SocketAddress address : v6Addresses) {
if (!address.equals(localV6Address)) {
((NettyServerBuilder) serverBuilder).addListenAddress(address);
// Create ServerBuilder with appropriate addresses
// - IPV4_IPV6: bind to wildcard which covers all addresses on all interfaces of both families
// - IPV4: bind to v4 address for local hostname + v4 localhost
// - IPV6: bind to all v6 addresses for local hostname + v6 localhost
ServerBuilder<?> serverBuilder;
switch (addressType) {
case IPV4_IPV6:
serverBuilder = Grpc.newServerBuilderForPort(port, serverCreds);
break;
case IPV4:
SocketAddress v4Address = Util.getV4Address(port);
InetSocketAddress localV4Address = new InetSocketAddress("127.0.0.1", port);
serverBuilder =
NettyServerBuilder.forAddress(localV4Address, serverCreds);
if (v4Address != null && !v4Address.equals(localV4Address)) {
((NettyServerBuilder) serverBuilder).addListenAddress(v4Address);
}
}
break;
default:
throw new AssertionError("Unknown address type: " + addressType);
if (mcsLimit != -1) {
((NettyServerBuilder) serverBuilder).maxConcurrentCallsPerConnection(mcsLimit);
}
break;
case IPV6:
List<SocketAddress> v6Addresses = Util.getV6Addresses(port);
InetSocketAddress localV6Address = new InetSocketAddress("::1", port);
serverBuilder =
NettyServerBuilder.forAddress(localV6Address, serverCreds);
for (SocketAddress address : v6Addresses) {
if (!address.equals(localV6Address)) {
((NettyServerBuilder) serverBuilder).addListenAddress(address);
}
}
break;
default:
throw new AssertionError("Unknown address type: " + addressType);
}
server = serverBuilder
.maxInboundMessageSize(AbstractInteropTest.MAX_MESSAGE_SIZE)
.addService(
ServerInterceptors.intercept(
new TestServiceImpl(executor, metricRecorder), TestServiceImpl.interceptors()))
.addService(orcaOobService)
.intercept(OrcaMetricReportingServerInterceptor.create(metricRecorder))
.build()
.start();
} catch (Throwable t) {
stop();
throw t;
}
server = serverBuilder
.maxInboundMessageSize(AbstractInteropTest.MAX_MESSAGE_SIZE)
.addService(
ServerInterceptors.intercept(
new TestServiceImpl(executor, metricRecorder), TestServiceImpl.interceptors()))
.addService(orcaOobService)
.intercept(OrcaMetricReportingServerInterceptor.create(metricRecorder))
.build()
.start();
}

@VisibleForTesting
void stop() throws Exception {
server.shutdownNow();
if (!server.awaitTermination(5, TimeUnit.SECONDS)) {
System.err.println("Timed out waiting for server shutdown");
try {
if (server != null) {
server.shutdownNow();
if (!server.awaitTermination(5, TimeUnit.SECONDS)) {
System.err.println("Timed out waiting for server shutdown");
}
}
if (executor != null) {
MoreExecutors.shutdownAndAwaitTermination(executor, 5, TimeUnit.SECONDS);
}
} finally {
if (openTelemetrySdk != null) {
openTelemetrySdk.close();
}
}
MoreExecutors.shutdownAndAwaitTermination(executor, 5, TimeUnit.SECONDS);
}

@VisibleForTesting
Expand Down
Loading