3939import io .grpc .InsecureChannelCredentials ;
4040import io .grpc .InsecureServerCredentials ;
4141import io .grpc .ManagedChannel ;
42+ import io .grpc .ManagedChannelBuilder ;
4243import io .grpc .Metadata ;
4344import io .grpc .MethodDescriptor ;
4445import io .grpc .Server ;
6061import io .grpc .testing .integration .Messages .SimpleRequest ;
6162import io .grpc .testing .integration .Messages .SimpleResponse ;
6263import io .grpc .xds .XdsChannelCredentials ;
64+ import io .opentelemetry .sdk .OpenTelemetrySdk ;
6365import io .opentelemetry .sdk .autoconfigure .AutoConfiguredOpenTelemetrySdk ;
6466import java .util .ArrayList ;
6567import java .util .Collections ;
@@ -104,6 +106,7 @@ public final class XdsTestClient {
104106 private long currentRequestId ;
105107 private ListeningScheduledExecutorService exec ;
106108 private CsmObservability csmObservability ;
109+ private OpenTelemetrySdk openTelemetrySdk ;
107110
108111 /**
109112 * The main application allowing this client to be launched from the command line.
@@ -265,14 +268,23 @@ private static RpcType parseRpc(String rpc) {
265268 @ IgnoreJRERequirement // OpenTelemetry uses Java 8+ APIs
266269 private void run () {
267270 if (enableCsmObservability ) {
271+ Map <String , String > props = new HashMap <>();
272+ props .put ("otel.logs.exporter" , "none" );
273+ props .put ("otel.metrics.exporter" , "otlp" );
274+ String tracesExporter = System .getenv ("OTEL_TRACES_EXPORTER" );
275+ if (tracesExporter != null ) {
276+ props .put ("otel.traces.exporter" , tracesExporter );
277+ } else {
278+ props .put ("otel.traces.exporter" , "none" );
279+ }
280+
281+ AutoConfiguredOpenTelemetrySdk autoSdk = AutoConfiguredOpenTelemetrySdk .builder ()
282+ .addPropertiesSupplier (() -> props )
283+ .build ();
284+ openTelemetrySdk = autoSdk .getOpenTelemetrySdk ();
268285 csmObservability = CsmObservability .newBuilder ()
269- .sdk (AutoConfiguredOpenTelemetrySdk .builder ()
270- .addPropertiesSupplier (() -> ImmutableMap .of (
271- "otel.logs.exporter" , "none" ,
272- "otel.metrics.exporter" , "prometheus" ,
273- "otel.traces.exporter" , "none" ))
274- .build ()
275- .getOpenTelemetrySdk ())
286+ .sdk (openTelemetrySdk )
287+ .enableTracing (!"none" .equals (props .get ("otel.traces.exporter" )))
276288 .build ();
277289 csmObservability .registerGlobal ();
278290 }
@@ -289,14 +301,16 @@ private void run() {
289301 try {
290302 statsServer .start ();
291303 for (int i = 0 ; i < numChannels ; i ++) {
292- channels .add (
293- Grpc .newChannelBuilder (
304+ ManagedChannelBuilder <?> builder = Grpc .newChannelBuilder (
294305 server ,
295306 secureMode
296307 ? XdsChannelCredentials .create (InsecureChannelCredentials .create ())
297308 : InsecureChannelCredentials .create ())
298- .enableRetry ()
299- .build ());
309+ .enableRetry ();
310+ if (enableCsmObservability ) {
311+ csmObservability .configureChannelBuilder (builder );
312+ }
313+ channels .add (builder .build ());
300314 }
301315 exec = MoreExecutors .listeningDecorator (Executors .newSingleThreadScheduledExecutor ());
302316 Payload requestPayload = Payload .newBuilder ()
@@ -325,6 +339,9 @@ private void stop() throws InterruptedException {
325339 if (csmObservability != null ) {
326340 csmObservability .close ();
327341 }
342+ if (openTelemetrySdk != null ) {
343+ openTelemetrySdk .close ();
344+ }
328345 }
329346
330347
@@ -373,6 +390,13 @@ public void start(Listener<RespT> responseListener, Metadata headers) {
373390 @ Override
374391 public void onHeaders (Metadata headers ) {
375392 hostnameRef .set (headers .get (XdsTestServer .HOSTNAME_KEY ));
393+ io .opentelemetry .api .trace .Span currentSpan = io .opentelemetry .api .trace .Span .current ();
394+ for (String key : config .metadata .keys ()) {
395+ String value = config .metadata .get (Metadata .Key .of (key , Metadata .ASCII_STRING_MARSHALLER ));
396+ if (value != null ) {
397+ currentSpan .setAttribute ("custom.metadata." + key , value );
398+ }
399+ }
376400 super .onHeaders (headers );
377401 }
378402 },
@@ -406,44 +430,56 @@ public void onNext(EmptyProtos.Empty response) {}
406430 .setPayload (requestPayload )
407431 .setResponseSize (responseSize )
408432 .build ();
409- stub .unaryCall (
410- request ,
411- new StreamObserver <SimpleResponse >() {
412- @ Override
413- public void onCompleted () {
414- handleRpcCompleted (requestId , config .rpcType , hostnameRef .get (), savedWatchers );
415- }
416433
417- @ Override
418- public void onError (Throwable t ) {
419- if (printResponse ) {
420- logger .log (Level .WARNING , "Rpc failed" , t );
434+ io .opentelemetry .api .baggage .BaggageBuilder baggageBuilder = io .opentelemetry .api .baggage .Baggage .builder ();
435+ for (String key : config .metadata .keys ()) {
436+ String value = config .metadata .get (Metadata .Key .of (key , Metadata .ASCII_STRING_MARSHALLER ));
437+ if (value != null ) {
438+ baggageBuilder .put (key , value );
439+ }
440+ }
441+ io .opentelemetry .api .baggage .Baggage baggage = baggageBuilder .build ();
442+
443+ try (io .opentelemetry .context .Scope scope = io .opentelemetry .context .Context .current ().with (baggage ).makeCurrent ()) {
444+ stub .unaryCall (
445+ request ,
446+ new StreamObserver <SimpleResponse >() {
447+ @ Override
448+ public void onCompleted () {
449+ handleRpcCompleted (requestId , config .rpcType , hostnameRef .get (), savedWatchers );
421450 }
422- handleRpcError (requestId , config .rpcType , Status .fromThrowable (t ),
423- savedWatchers );
424- }
425451
426- @ Override
427- public void onNext (SimpleResponse response ) {
428- // TODO(ericgribkoff) Currently some test environments cannot access the stats RPC
429- // service and rely on parsing stdout.
430- if (printResponse ) {
431- System .out .println (
432- "Greeting: Hello world, this is "
433- + response .getHostname ()
434- + ", from "
435- + clientCallRef
436- .get ()
437- .getAttributes ()
438- .get (Grpc .TRANSPORT_ATTR_REMOTE_ADDR ));
452+ @ Override
453+ public void onError (Throwable t ) {
454+ if (printResponse ) {
455+ logger .log (Level .WARNING , "Rpc failed" , t );
456+ }
457+ handleRpcError (requestId , config .rpcType , Status .fromThrowable (t ),
458+ savedWatchers );
439459 }
440- // Use the hostname from the response if not present in the metadata.
441- // TODO(ericgribkoff) Delete when server is deployed that sets metadata value.
442- if (hostnameRef .get () == null ) {
443- hostnameRef .set (response .getHostname ());
460+
461+ @ Override
462+ public void onNext (SimpleResponse response ) {
463+ // TODO(ericgribkoff) Currently some test environments cannot access the stats RPC
464+ // service and rely on parsing stdout.
465+ if (printResponse ) {
466+ System .out .println (
467+ "Greeting: Hello world, this is "
468+ + response .getHostname ()
469+ + ", from "
470+ + clientCallRef
471+ .get ()
472+ .getAttributes ()
473+ .get (Grpc .TRANSPORT_ATTR_REMOTE_ADDR ));
474+ }
475+ // Use the hostname from the response if not present in the metadata.
476+ // TODO(ericgribkoff) Delete when server is deployed that sets metadata value.
477+ if (hostnameRef .get () == null ) {
478+ hostnameRef .set (response .getHostname ());
479+ }
444480 }
445- }
446- });
481+ });
482+ }
447483 } else {
448484 throw new AssertionError ("Unknown RPC type: " + config .rpcType );
449485 }
0 commit comments