• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

grpc / grpc-java / #20494

30 Sep 2026 08:29AM UTC coverage: 89.351% (+0.05%) from 89.306%
#20494

push

github

web-flow
core: Implement [A121](https://github.com/grpc/proposal/pull/556) (#12893)

39251 of 43929 relevant lines covered (89.35%)

0.89 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

94.23
/../opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryTracingModule.java
1
/*
2
 * Copyright 2024 The gRPC Authors
3
 *
4
 * Licensed under the Apache License, Version 2.0 (the "License");
5
 * you may not use this file except in compliance with the License.
6
 * You may obtain a copy of the License at
7
 *
8
 *     http://www.apache.org/licenses/LICENSE-2.0
9
 *
10
 * Unless required by applicable law or agreed to in writing, software
11
 * distributed under the License is distributed on an "AS IS" BASIS,
12
 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13
 * See the License for the specific language governing permissions and
14
 * limitations under the License.
15
 */
16

17
package io.grpc.opentelemetry;
18

19
import static com.google.common.base.Preconditions.checkNotNull;
20
import static io.grpc.ClientStreamTracer.NAME_RESOLUTION_DELAYED;
21
import static io.grpc.internal.GrpcUtil.IMPLEMENTATION_VERSION;
22
import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.BAGGAGE_KEY;
23

24
import com.google.common.annotations.VisibleForTesting;
25
import com.google.errorprone.annotations.concurrent.GuardedBy;
26
import io.grpc.CallOptions;
27
import io.grpc.Channel;
28
import io.grpc.ClientCall;
29
import io.grpc.ClientInterceptor;
30
import io.grpc.ClientStreamTracer;
31
import io.grpc.ForwardingClientCall.SimpleForwardingClientCall;
32
import io.grpc.ForwardingClientCallListener.SimpleForwardingClientCallListener;
33
import io.grpc.ForwardingServerCallListener;
34
import io.grpc.Metadata;
35
import io.grpc.MethodDescriptor;
36
import io.grpc.ServerCall;
37
import io.grpc.ServerCallHandler;
38
import io.grpc.ServerInterceptor;
39
import io.grpc.ServerStreamTracer;
40
import io.grpc.Status;
41
import io.grpc.internal.GrpcUtil;
42
import io.grpc.opentelemetry.internal.OpenTelemetryConstants;
43
import io.opentelemetry.api.OpenTelemetry;
44
import io.opentelemetry.api.baggage.Baggage;
45
import io.opentelemetry.api.common.AttributeKey;
46
import io.opentelemetry.api.common.Attributes;
47
import io.opentelemetry.api.common.AttributesBuilder;
48
import io.opentelemetry.api.trace.Span;
49
import io.opentelemetry.api.trace.StatusCode;
50
import io.opentelemetry.api.trace.Tracer;
51
import io.opentelemetry.context.Context;
52
import io.opentelemetry.context.Scope;
53
import io.opentelemetry.context.propagation.ContextPropagators;
54
import java.util.concurrent.atomic.AtomicIntegerFieldUpdater;
55
import java.util.logging.Level;
56
import java.util.logging.Logger;
57
import javax.annotation.Nullable;
58

59
/**
60
 * Provides factories for {@link io.grpc.StreamTracer} that records tracing to OpenTelemetry.
61
 */
62
final class OpenTelemetryTracingModule {
63
  private static final Logger logger = Logger.getLogger(OpenTelemetryTracingModule.class.getName());
1 ✔
64

65
  @VisibleForTesting
1 ✔
66
  final io.grpc.Context.Key<Span> otelSpan = io.grpc.Context.key("opentelemetry-span-key");
1 ✔
67

68
  @Nullable
69
  private static final AtomicIntegerFieldUpdater<CallAttemptsTracerFactory> callEndedUpdater;
70
  @Nullable
71
  private static final AtomicIntegerFieldUpdater<ServerTracer> streamClosedUpdater;
72

73
  /*
74
   * When using Atomic*FieldUpdater, some Samsung Android 5.0.x devices encounter a bug in their JDK
75
   * reflection API that triggers a NoSuchFieldException. When this occurs, we fallback to
76
   * (potentially racy) direct updates of the volatile variables.
77
   */
78
  static {
79
    AtomicIntegerFieldUpdater<CallAttemptsTracerFactory> tmpCallEndedUpdater;
80
    AtomicIntegerFieldUpdater<ServerTracer> tmpStreamClosedUpdater;
81
    try {
82
      tmpCallEndedUpdater =
1 ✔
83
          AtomicIntegerFieldUpdater.newUpdater(CallAttemptsTracerFactory.class, "callEnded");
1 ✔
84
      tmpStreamClosedUpdater =
1 ✔
85
          AtomicIntegerFieldUpdater.newUpdater(ServerTracer.class, "streamClosed");
1 ✔
86
    } catch (Throwable t) {
×
87
      logger.log(Level.SEVERE, "Creating atomic field updaters failed", t);
×
88
      tmpCallEndedUpdater = null;
×
89
      tmpStreamClosedUpdater = null;
×
90
    }
1 ✔
91
    callEndedUpdater = tmpCallEndedUpdater;
1 ✔
92
    streamClosedUpdater = tmpStreamClosedUpdater;
1 ✔
93
  }
1 ✔
94

95
  private final Tracer otelTracer;
96
  private final ContextPropagators contextPropagators;
97
  private final MetadataGetter metadataGetter = MetadataGetter.getInstance();
1 ✔
98
  private final MetadataSetter metadataSetter = MetadataSetter.getInstance();
1 ✔
99
  private final TracingClientInterceptor clientInterceptor = new TracingClientInterceptor();
1 ✔
100
  private final ServerInterceptor serverSpanPropagationInterceptor =
1 ✔
101
      new TracingServerSpanPropagationInterceptor();
102
  private final ServerTracerFactory serverTracerFactory = new ServerTracerFactory();
1 ✔
103

104
  OpenTelemetryTracingModule(OpenTelemetry openTelemetry) {
1 ✔
105
    this.otelTracer = checkNotNull(openTelemetry.getTracerProvider(), "tracerProvider")
1 ✔
106
        .tracerBuilder(OpenTelemetryConstants.INSTRUMENTATION_SCOPE)
1 ✔
107
        .setInstrumentationVersion(IMPLEMENTATION_VERSION)
1 ✔
108
        .build();
1 ✔
109
    this.contextPropagators = checkNotNull(openTelemetry.getPropagators(), "contextPropagators");
1 ✔
110
  }
1 ✔
111

112
  @VisibleForTesting
113
  Tracer getTracer() {
114
    return otelTracer;
1 ✔
115
  }
116

117
  /**
118
   * Creates a {@link CallAttemptsTracerFactory} for a new call.
119
   */
120
  @VisibleForTesting
121
  CallAttemptsTracerFactory newClientCallTracer(Span clientSpan, MethodDescriptor<?, ?> method) {
122
    return new CallAttemptsTracerFactory(clientSpan, method);
1 ✔
123
  }
124

125
  /**
126
   * Returns the server tracer factory.
127
   */
128
  ServerStreamTracer.Factory getServerTracerFactory() {
129
    return serverTracerFactory;
1 ✔
130
  }
131

132
  /**
133
   * Returns the client interceptor that facilitates otel tracing reporting.
134
   */
135
  ClientInterceptor getClientInterceptor() {
136
    return clientInterceptor;
1 ✔
137
  }
138

139
  ServerInterceptor getServerSpanPropagationInterceptor() {
140
    return serverSpanPropagationInterceptor;
1 ✔
141
  }
142

143
  @VisibleForTesting
144
  final class CallAttemptsTracerFactory extends ClientStreamTracer.Factory {
145
    volatile int callEnded;
146
    private final Span clientSpan;
147
    private final String fullMethodName;
148
    @GuardedBy("this")
149
    @Nullable private Span activeCallDelaySpan;
150
    @GuardedBy("this")
151
    @Nullable private String activeCallDelayType;
152

153
    CallAttemptsTracerFactory(Span clientSpan, MethodDescriptor<?, ?> method) {
1 ✔
154
      checkNotNull(method, "method");
1 ✔
155
      this.fullMethodName = checkNotNull(method.getFullMethodName(), "fullMethodName");
1 ✔
156
      this.clientSpan = checkNotNull(clientSpan, "clientSpan");
1 ✔
157
    }
1 ✔
158

159
    @Override
160
    public ClientStreamTracer newClientStreamTracer(
161
        ClientStreamTracer.StreamInfo info, Metadata headers) {
162
      Span attemptSpan = otelTracer.spanBuilder(
1 ✔
163
              "Attempt." + fullMethodName.replace('/', '.'))
1 ✔
164
          .setParent(Context.current().with(clientSpan))
1 ✔
165
          .startSpan();
1 ✔
166
      attemptSpan.setAttribute(
1 ✔
167
          "previous-rpc-attempts", info.getPreviousAttempts());
1 ✔
168
      attemptSpan.setAttribute(
1 ✔
169
          "transparent-retry",info.isTransparentRetry());
1 ✔
170
      if (info.getCallOptions().getOption(NAME_RESOLUTION_DELAYED) != null) {
1 ✔
171
        clientSpan.addEvent("Delayed name resolution complete");
1 ✔
172
      }
173
      return new ClientTracer(attemptSpan, clientSpan);
1 ✔
174
    }
175

176
    private boolean isCallEnded() {
177
      if (callEndedUpdater != null) {
1 ✔
178
        return callEndedUpdater.get(this) != 0;
1 ✔
179
      }
180
      return callEnded != 0;
×
181
    }
182

183
    /**
184
     * Record a finished call and mark the current time as the end time.
185
     *
186
     * <p>Can be called from any thread without synchronization.  Calling it the second time or more
187
     * is a no-op.
188
     */
189
    void callEnded(io.grpc.Status status) {
190
      if (callEndedUpdater != null) {
1 ✔
191
        if (callEndedUpdater.getAndSet(this, 1) != 0) {
1 ✔
192
          return;
×
193
        }
194
      } else {
195
        if (callEnded != 0) {
×
196
          return;
×
197
        }
198
        callEnded = 1;
×
199
      }
200
      synchronized (this) {
1 ✔
201
        endActiveDelaySpan();
1 ✔
202
      }
1 ✔
203
      endSpanWithStatus(clientSpan, status);
1 ✔
204
    }
1 ✔
205

206
    @Override
207
    public void recordDelayStart(String delayType, String delayReason) {
208
      if (isCallEnded()) {
1 ✔
209
        return;
1 ✔
210
      }
211
      synchronized (this) {
1 ✔
212
        if (isCallEnded()) {
1 ✔
213
          return;
×
214
        }
215
        if (activeCallDelaySpan != null && delayType.equals(activeCallDelayType)) {
1 ✔
216
          addDelayEvent(activeCallDelaySpan, delayReason);
1 ✔
217
          return;
1 ✔
218
        }
219
        endActiveDelaySpan();
1 ✔
220
        activeCallDelayType = delayType;
1 ✔
221
        activeCallDelaySpan = otelTracer.spanBuilder("Delay")
1 ✔
222
            .setParent(Context.current().with(clientSpan))
1 ✔
223
            .setAttribute("grpc.delay_type", delayType)
1 ✔
224
            .startSpan();
1 ✔
225
        addDelayEvent(activeCallDelaySpan, delayReason);
1 ✔
226
      }
1 ✔
227
    }
1 ✔
228

229
    @Override
230
    public void recordDelayReasonChanged(String delayType, String delayReason) {
231
      if (isCallEnded()) {
1 ✔
232
        return;
1 ✔
233
      }
234
      synchronized (this) {
1 ✔
235
        if (isCallEnded() || activeCallDelaySpan == null) {
1 ✔
236
          return;
1 ✔
237
        }
238
        addDelayEvent(activeCallDelaySpan, delayReason);
1 ✔
239
      }
1 ✔
240
    }
1 ✔
241

242
    @Override
243
    public synchronized void recordDelayEnd(String delayType) {
244
      endActiveDelaySpan();
1 ✔
245
    }
1 ✔
246

247
    @GuardedBy("this")
248
    private void endActiveDelaySpan() {
249
      if (activeCallDelaySpan != null) {
1 ✔
250
        activeCallDelaySpan.end();
1 ✔
251
        activeCallDelaySpan = null;
1 ✔
252
        activeCallDelayType = null;
1 ✔
253
      }
254
    }
1 ✔
255
  }
256

257
  private final class ClientTracer extends ClientStreamTracer {
258
    private final Span span;
259
    private final Span parentSpan;
260
    volatile int seqNo;
261
    boolean isPendingStream;
262
    @GuardedBy("this")
263
    @Nullable private Span activeDelaySpan;
264
    @GuardedBy("this")
265
    @Nullable private String activeDelayType;
266
    @GuardedBy("this")
267
    private boolean streamClosed;
268

269
    ClientTracer(Span span, Span parentSpan) {
1 ✔
270
      this.span = checkNotNull(span, "span");
1 ✔
271
      this.parentSpan = checkNotNull(parentSpan, "parent span");
1 ✔
272
    }
1 ✔
273

274
    @Override
275
    public void streamCreated(io.grpc.Attributes transportAtts, Metadata headers) {
276
      synchronized (this) {
1 ✔
277
        endActiveDelaySpan();
1 ✔
278
      }
1 ✔
279
      contextPropagators.getTextMapPropagator().inject(Context.current().with(span), headers,
1 ✔
280
          metadataSetter);
1 ✔
281
      if (isPendingStream) {
1 ✔
282
        span.addEvent("Delayed LB pick complete");
1 ✔
283
      }
284
    }
1 ✔
285

286
    @Override
287
    public void createPendingStream() {
288
      isPendingStream = true;
1 ✔
289
    }
1 ✔
290

291
    @Override
292
    public synchronized void recordDelayStart(String delayType, String delayReason) {
293
      if (streamClosed) {
1 ✔
294
        return;
1 ✔
295
      }
296
      if (activeDelaySpan != null && delayType.equals(activeDelayType)) {
1 ✔
297
        // Do not recreate the span if the delay type is unchanged (e.g., priority failover).
298
        addDelayEvent(activeDelaySpan, delayReason);
1 ✔
299
        return;
1 ✔
300
      }
301
      endActiveDelaySpan();
1 ✔
302
      activeDelayType = delayType;
1 ✔
303
      activeDelaySpan = otelTracer.spanBuilder("Delay")
1 ✔
304
          .setParent(Context.current().with(span))
1 ✔
305
          .setAttribute("grpc.delay_type", delayType)
1 ✔
306
          .startSpan();
1 ✔
307
      addDelayEvent(activeDelaySpan, delayReason);
1 ✔
308
    }
1 ✔
309

310
    @Override
311
    public synchronized void recordDelayReasonChanged(String delayType, String delayReason) {
312
      if (streamClosed || activeDelaySpan == null) {
1 ✔
313
        return;
1 ✔
314
      }
315
      addDelayEvent(activeDelaySpan, delayReason);
1 ✔
316
    }
1 ✔
317

318
    @Override
319
    public synchronized void recordDelayEnd(String delayType) {
320
      endActiveDelaySpan();
1 ✔
321
    }
1 ✔
322

323
    @GuardedBy("this")
324
    private void endActiveDelaySpan() {
325
      if (activeDelaySpan != null) {
1 ✔
326
        activeDelaySpan.end();
1 ✔
327
        activeDelaySpan = null;
1 ✔
328
        activeDelayType = null;
1 ✔
329
      }
330
    }
1 ✔
331

332
    @Override
333
    public void outboundMessageSent(
334
        int seqNo, long optionalWireSize, long optionalUncompressedSize) {
335
      recordOutboundMessageSentEvent(span, seqNo, optionalWireSize, optionalUncompressedSize);
1 ✔
336
    }
1 ✔
337

338
    @Override
339
    public void inboundMessageRead(
340
        int seqNo, long optionalWireSize, long optionalUncompressedSize) {
341
      if (optionalWireSize != optionalUncompressedSize) {
1 ✔
342
        recordInboundCompressedMessage(span, seqNo, optionalWireSize);
1 ✔
343
      }
344
    }
1 ✔
345

346
    @Override
347
    public void inboundMessage(int seqNo) {
348
      this.seqNo = seqNo;
1 ✔
349
    }
1 ✔
350

351
    @Override
352
    public void inboundUncompressedSize(long bytes) {
353
      recordInboundMessageSize(parentSpan, seqNo, bytes);
1 ✔
354
    }
1 ✔
355

356
    @Override
357
    public synchronized void streamClosed(Status status) {
358
      if (streamClosed) {
1 ✔
359
        return;
×
360
      }
361
      streamClosed = true;
1 ✔
362
      endActiveDelaySpan();
1 ✔
363
      endSpanWithStatus(span, status);
1 ✔
364
    }
1 ✔
365
  }
366

367
  private final class ServerTracer extends ServerStreamTracer {
368
    private final Span span;
369
    volatile int streamClosed;
370
    private int seqNo;
371
    private Baggage baggage;
372

373
    ServerTracer(String fullMethodName, @Nullable Span remoteSpan, Baggage baggage) {
1 ✔
374
      checkNotNull(fullMethodName, "fullMethodName");
1 ✔
375
      this.span =
1 ✔
376
          otelTracer.spanBuilder(generateTraceSpanName(true, fullMethodName))
1 ✔
377
              .setParent(remoteSpan == null ? null : Context.current().with(remoteSpan))
1 ✔
378
              .startSpan();
1 ✔
379
      this.baggage = baggage;
1 ✔
380
    }
1 ✔
381

382
    /**
383
     * Record a finished stream and mark the current time as the end time.
384
     *
385
     * <p>Can be called from any thread without synchronization.  Calling it the second time or more
386
     * is a no-op.
387
     */
388
    @Override
389
    public void streamClosed(io.grpc.Status status) {
390
      if (streamClosedUpdater != null) {
1 ✔
391
        if (streamClosedUpdater.getAndSet(this, 1) != 0) {
1 ✔
392
          return;
×
393
        }
394
      } else {
395
        if (streamClosed != 0) {
×
396
          return;
×
397
        }
398
        streamClosed = 1;
×
399
      }
400
      endSpanWithStatus(span, status);
1 ✔
401
    }
1 ✔
402

403
    @Override
404
    public io.grpc.Context filterContext(io.grpc.Context context) {
405
      return context
1 ✔
406
          .withValue(otelSpan, span)
1 ✔
407
          .withValue(BAGGAGE_KEY, baggage);
1 ✔
408
    }
409

410
    @Override
411
    public void outboundMessageSent(
412
        int seqNo, long optionalWireSize, long optionalUncompressedSize) {
413
      recordOutboundMessageSentEvent(span, seqNo, optionalWireSize, optionalUncompressedSize);
1 ✔
414
    }
1 ✔
415

416
    @Override
417
    public void inboundMessageRead(
418
        int seqNo, long optionalWireSize, long optionalUncompressedSize) {
419
      if (optionalWireSize != optionalUncompressedSize) {
1 ✔
420
        recordInboundCompressedMessage(span, seqNo, optionalWireSize);
1 ✔
421
      }
422
    }
1 ✔
423

424
    @Override
425
    public void inboundMessage(int seqNo) {
426
      this.seqNo = seqNo;
1 ✔
427
    }
1 ✔
428

429
    @Override
430
    public void inboundUncompressedSize(long bytes) {
431
      recordInboundMessageSize(span, seqNo, bytes);
1 ✔
432
    }
1 ✔
433
  }
434

435
  @VisibleForTesting
436
  final class ServerTracerFactory extends ServerStreamTracer.Factory {
1 ✔
437
    @SuppressWarnings("ReferenceEquality")
438
    @Override
439
    public ServerStreamTracer newServerStreamTracer(String fullMethodName, Metadata headers) {
440
      Context context = contextPropagators.getTextMapPropagator().extract(
1 ✔
441
          Context.current(), headers, metadataGetter
1 ✔
442
      );
443
      Span remoteSpan = Span.fromContext(context);
1 ✔
444
      if (remoteSpan == Span.getInvalid()) {
1 ✔
445
        remoteSpan = null;
1 ✔
446
      }
447
      Baggage baggage = Baggage.fromContext(context);
1 ✔
448
      return new ServerTracer(fullMethodName, remoteSpan, baggage);
1 ✔
449
    }
450
  }
451

452
  @VisibleForTesting
453
  final class TracingServerSpanPropagationInterceptor implements ServerInterceptor {
1 ✔
454
    @Override
455
    public <ReqT, RespT> ServerCall.Listener<ReqT> interceptCall(ServerCall<ReqT, RespT> call,
456
        Metadata headers, ServerCallHandler<ReqT, RespT> next) {
457
      Span span = otelSpan.get(io.grpc.Context.current());
1 ✔
458
      if (span == null) {
1 ✔
459
        logger.log(Level.FINE, "Server span not found. ServerTracerFactory for server "
1 ✔
460
            + "tracing must be set.");
461
        return next.startCall(call, headers);
1 ✔
462
      }
463
      Context serverCallContext = Context.current();
1 ✔
464
      serverCallContext = serverCallContext.with(span);
1 ✔
465
      Baggage baggage = BAGGAGE_KEY.get();
1 ✔
466
      if (baggage != null) {
1 ✔
467
        serverCallContext = serverCallContext.with(baggage);
1 ✔
468
      } else {
469
        logger.log(Level.WARNING, "Server baggage not found which is unexpected, "
1 ✔
470
            + "as it is being added unconditionally in filterContext().");
471
      }
472
      try (Scope scope = serverCallContext.makeCurrent()) {
1 ✔
473
        return new ContextServerCallListener<>(next.startCall(call, headers), serverCallContext);
1 ✔
474
      }
475
    }
476
  }
477

478
  private static class ContextServerCallListener<ReqT> extends
479
      ForwardingServerCallListener.SimpleForwardingServerCallListener<ReqT> {
480
    private final Context context;
481

482
    protected ContextServerCallListener(ServerCall.Listener<ReqT> delegate, Context context) {
483
      super(delegate);
1 ✔
484
      this.context = checkNotNull(context, "context");
1 ✔
485
    }
1 ✔
486

487
    @Override
488
    public void onMessage(ReqT message) {
489
      try (Scope scope = context.makeCurrent()) {
1 ✔
490
        delegate().onMessage(message);
1 ✔
491
      }
492
    }
1 ✔
493

494
    @Override
495
    public void onHalfClose() {
496
      try (Scope scope = context.makeCurrent()) {
1 ✔
497
        delegate().onHalfClose();
1 ✔
498
      }
499
    }
1 ✔
500

501
    @Override
502
    public void onCancel() {
503
      try (Scope scope = context.makeCurrent()) {
1 ✔
504
        delegate().onCancel();
1 ✔
505
      }
506
    }
1 ✔
507

508
    @Override
509
    public void onComplete() {
510
      try (Scope scope = context.makeCurrent()) {
1 ✔
511
        delegate().onComplete();
1 ✔
512
      }
513
    }
1 ✔
514

515
    @Override
516
    public void onReady() {
517
      try (Scope scope = context.makeCurrent()) {
1 ✔
518
        delegate().onReady();
1 ✔
519
      }
520
    }
1 ✔
521

522
    @Override
523
    public void onEvent(Object event) {
524
      try (Scope scope = context.makeCurrent()) {
1 ✔
525
        delegate().onEvent(event);
1 ✔
526
      }
527
    }
1 ✔
528
  }
529

530
  @VisibleForTesting
531
  final class TracingClientInterceptor implements ClientInterceptor {
1 ✔
532

533
    @Override
534
    public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(
535
        MethodDescriptor<ReqT, RespT> method, CallOptions callOptions, Channel next) {
536
      Span clientSpan = otelTracer.spanBuilder(
1 ✔
537
          generateTraceSpanName(false, method.getFullMethodName()))
1 ✔
538
          .startSpan();
1 ✔
539

540
      final CallAttemptsTracerFactory tracerFactory = newClientCallTracer(clientSpan, method);
1 ✔
541
      ClientCall<ReqT, RespT> call =
1 ✔
542
          next.newCall(
1 ✔
543
              method,
544
              callOptions.withStreamTracerFactory(tracerFactory));
1 ✔
545
      return new SimpleForwardingClientCall<ReqT, RespT>(call) {
1 ✔
546
        @Override
547
        public void start(Listener<RespT> responseListener, Metadata headers) {
548
          delegate().start(
1 ✔
549
              new SimpleForwardingClientCallListener<RespT>(responseListener) {
1 ✔
550
                @Override
551
                public void onClose(io.grpc.Status status, Metadata trailers) {
552
                  tracerFactory.callEnded(status);
1 ✔
553
                  super.onClose(status, trailers);
1 ✔
554
                }
1 ✔
555
              },
556
              headers);
557
        }
1 ✔
558
      };
559
    }
560
  }
561

562
  // Attribute named "message-size" always means the message size the application sees.
563
  // If there was compression, additional event reports "message-size-compressed".
564
  //
565
  // An example trace with message compression:
566
  //
567
  // Sending:
568
  // |-- Event 'Outbound message sent', attributes('sequence-numer' = 0, 'message-size' = 7854,
569
  //                                               'message-size-compressed' = 5493) ----|
570
  //
571
  // Receiving:
572
  // |-- Event 'Inbound compressed message', attributes('sequence-numer' = 0,
573
  //                                                    'message-size-compressed' = 5493 ) ----|
574
  // |-- Event 'Inbound message received', attributes('sequence-numer' = 0,
575
  //                                                  'message-size' = 7854) ----|
576
  //
577
  // An example trace with no message compression:
578
  //
579
  // Sending:
580
  // |-- Event 'Outbound message sent', attributes('sequence-numer' = 0, 'message-size' = 7854) ---|
581
  //
582
  // Receiving:
583
  // |-- Event 'Inbound message received', attributes('sequence-numer' = 0,
584
  //                                                  'message-size' = 7854) ----|
585
  private static void addDelayEvent(Span delaySpan, String delayReason) {
586
    delaySpan.addEvent(
1 ✔
587
        "Delay triggered",
588
        Attributes.of(AttributeKey.stringKey("grpc.delay_reason"), delayReason));
1 ✔
589
  }
1 ✔
590

591
  private void recordOutboundMessageSentEvent(Span span,
592
      int seqNo, long optionalWireSize, long optionalUncompressedSize) {
593
    AttributesBuilder attributesBuilder = Attributes.builder();
1 ✔
594
    attributesBuilder.put("sequence-number", seqNo);
1 ✔
595
    if (optionalUncompressedSize != -1) {
1 ✔
596
      attributesBuilder.put("message-size", optionalUncompressedSize);
1 ✔
597
    }
598
    if (optionalWireSize != -1 && optionalWireSize != optionalUncompressedSize) {
1 ✔
599
      attributesBuilder.put("message-size-compressed", optionalWireSize);
1 ✔
600
    }
601
    span.addEvent("Outbound message", attributesBuilder.build());
1 ✔
602
  }
1 ✔
603

604
  private void recordInboundCompressedMessage(Span span, int seqNo, long optionalWireSize) {
605
    AttributesBuilder attributesBuilder = Attributes.builder();
1 ✔
606
    attributesBuilder.put("sequence-number", seqNo);
1 ✔
607
    attributesBuilder.put("message-size-compressed", optionalWireSize);
1 ✔
608
    span.addEvent("Inbound compressed message", attributesBuilder.build());
1 ✔
609
  }
1 ✔
610

611
  private void recordInboundMessageSize(Span span, int seqNo, long bytes) {
612
    AttributesBuilder attributesBuilder = Attributes.builder();
1 ✔
613
    attributesBuilder.put("sequence-number", seqNo);
1 ✔
614
    attributesBuilder.put("message-size", bytes);
1 ✔
615
    span.addEvent("Inbound message", attributesBuilder.build());
1 ✔
616
  }
1 ✔
617

618
  private void endSpanWithStatus(Span span, io.grpc.Status status) {
619
    if (status.isOk()) {
1 ✔
620
      span.setStatus(StatusCode.OK);
1 ✔
621
    } else {
622
      span.setStatus(StatusCode.ERROR, GrpcUtil.statusToPrettyString(status));
1 ✔
623
    }
624
    span.end();
1 ✔
625
  }
1 ✔
626

627
  /**
628
   * Convert a full method name to a tracing span name.
629
   *
630
   * @param isServer {@code false} if the span is on the client-side, {@code true} if on the
631
   *                 server-side
632
   * @param fullMethodName the method name as returned by
633
   *        {@link MethodDescriptor#getFullMethodName}.
634
   */
635
  @VisibleForTesting
636
  static String generateTraceSpanName(boolean isServer, String fullMethodName) {
637
    String prefix = isServer ? "Recv" : "Sent";
1 ✔
638
    return prefix + "." + fullMethodName.replace('/', '.');
1 ✔
639
  }
640
}
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc