|
28 | 28 | import io.grpc.TlsServerCredentials; |
29 | 29 | import io.grpc.alts.AltsServerCredentials; |
30 | 30 | import io.grpc.netty.NettyServerBuilder; |
| 31 | +import io.grpc.opentelemetry.GrpcOpenTelemetry; |
| 32 | +import io.grpc.opentelemetry.GrpcTraceBinContextPropagator; |
| 33 | +import io.grpc.opentelemetry.InternalGrpcOpenTelemetry; |
31 | 34 | import io.grpc.services.MetricRecorder; |
32 | 35 | import io.grpc.testing.TlsTesting; |
33 | 36 | import io.grpc.xds.orca.OrcaMetricReportingServerInterceptor; |
34 | 37 | import io.grpc.xds.orca.OrcaServiceImpl; |
| 38 | +import io.opentelemetry.context.propagation.TextMapPropagator; |
| 39 | +import io.opentelemetry.sdk.OpenTelemetrySdk; |
| 40 | +import io.opentelemetry.sdk.autoconfigure.AutoConfiguredOpenTelemetrySdk; |
35 | 41 | import java.net.InetSocketAddress; |
36 | 42 | import java.net.SocketAddress; |
37 | 43 | import java.util.List; |
@@ -76,6 +82,8 @@ public void run() { |
76 | 82 | private boolean useTls = true; |
77 | 83 | private boolean useAlts = false; |
78 | 84 | private int mcsLimit = -1; |
| 85 | + private boolean enableOpentelemetry = false; |
| 86 | + private OpenTelemetrySdk openTelemetrySdk; |
79 | 87 |
|
80 | 88 | private ScheduledExecutorService executor; |
81 | 89 | private Server server; |
@@ -123,6 +131,8 @@ void parseArgs(String[] args) { |
123 | 131 | mcsLimit = Integer.parseInt(value); |
124 | 132 | // TODO: Make Netty server builder usable for IPV6 as well (not limited to MCS handling) |
125 | 133 | addressType = Util.AddressType.IPV4; // To use NettyServerBuilder |
| 134 | + } else if ("enable_opentelemetry".equals(key)) { |
| 135 | + enableOpentelemetry = Boolean.parseBoolean(value); |
126 | 136 | } else { |
127 | 137 | System.err.println("Unknown argument: " + key); |
128 | 138 | usage = true; |
@@ -156,6 +166,20 @@ void parseArgs(String[] args) { |
156 | 166 | @SuppressWarnings("AddressSelection") |
157 | 167 | @VisibleForTesting |
158 | 168 | void start() throws Exception { |
| 169 | + if (enableOpentelemetry) { |
| 170 | + AutoConfiguredOpenTelemetrySdk autoSdk = AutoConfiguredOpenTelemetrySdk.builder() |
| 171 | + .addPropagatorCustomizer( |
| 172 | + (previous, config) -> |
| 173 | + TextMapPropagator.composite( |
| 174 | + previous, GrpcTraceBinContextPropagator.defaultInstance())) |
| 175 | + .build(); |
| 176 | + this.openTelemetrySdk = autoSdk.getOpenTelemetrySdk(); |
| 177 | + GrpcOpenTelemetry.Builder grpcOpentelemetryBuilder = GrpcOpenTelemetry.newBuilder() |
| 178 | + .sdk(openTelemetrySdk); |
| 179 | + InternalGrpcOpenTelemetry.enableTracing(grpcOpentelemetryBuilder, true); |
| 180 | + GrpcOpenTelemetry grpcOpenTelemetry = grpcOpentelemetryBuilder.build(); |
| 181 | + grpcOpenTelemetry.registerGlobal(); |
| 182 | + } |
159 | 183 | executor = Executors.newSingleThreadScheduledExecutor(); |
160 | 184 | ServerCredentials serverCreds; |
161 | 185 | if (useAlts) { |
@@ -224,11 +248,17 @@ void start() throws Exception { |
224 | 248 |
|
225 | 249 | @VisibleForTesting |
226 | 250 | void stop() throws Exception { |
227 | | - server.shutdownNow(); |
228 | | - if (!server.awaitTermination(5, TimeUnit.SECONDS)) { |
229 | | - System.err.println("Timed out waiting for server shutdown"); |
| 251 | + try { |
| 252 | + server.shutdownNow(); |
| 253 | + if (!server.awaitTermination(5, TimeUnit.SECONDS)) { |
| 254 | + System.err.println("Timed out waiting for server shutdown"); |
| 255 | + } |
| 256 | + MoreExecutors.shutdownAndAwaitTermination(executor, 5, TimeUnit.SECONDS); |
| 257 | + } finally { |
| 258 | + if (openTelemetrySdk != null) { |
| 259 | + openTelemetrySdk.close(); |
| 260 | + } |
230 | 261 | } |
231 | | - MoreExecutors.shutdownAndAwaitTermination(executor, 5, TimeUnit.SECONDS); |
232 | 262 | } |
233 | 263 |
|
234 | 264 | @VisibleForTesting |
|
0 commit comments