• 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

93.53
/../opentelemetry/src/main/java/io/grpc/opentelemetry/OpenTelemetryMetricsModule.java
1
/*
2
 * Copyright 2023 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.opentelemetry.internal.OpenTelemetryConstants.BACKEND_SERVICE_KEY;
21
import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.BAGGAGE_KEY;
22
import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.CUSTOM_LABEL_KEY;
23
import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.LOCALITY_KEY;
24
import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.METHOD_KEY;
25
import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.STATUS_KEY;
26
import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.TARGET_KEY;
27

28
import com.google.common.annotations.VisibleForTesting;
29
import com.google.common.base.Stopwatch;
30
import com.google.common.base.Supplier;
31
import com.google.common.collect.ImmutableList;
32
import com.google.common.collect.ImmutableSet;
33
import com.google.errorprone.annotations.concurrent.GuardedBy;
34
import io.grpc.CallOptions;
35
import io.grpc.Channel;
36
import io.grpc.ClientCall;
37
import io.grpc.ClientInterceptor;
38
import io.grpc.ClientStreamTracer;
39
import io.grpc.ClientStreamTracer.StreamInfo;
40
import io.grpc.Deadline;
41
import io.grpc.ForwardingClientCall.SimpleForwardingClientCall;
42
import io.grpc.ForwardingClientCallListener.SimpleForwardingClientCallListener;
43
import io.grpc.Grpc;
44
import io.grpc.Metadata;
45
import io.grpc.MethodDescriptor;
46
import io.grpc.ServerStreamTracer;
47
import io.grpc.Status;
48
import io.grpc.Status.Code;
49
import io.grpc.StreamTracer;
50
import io.grpc.internal.StatsTraceContext.ServerCallMethodListener;
51
import io.grpc.opentelemetry.GrpcOpenTelemetry.TargetFilter;
52
import io.opentelemetry.api.baggage.Baggage;
53
import io.opentelemetry.api.common.Attributes;
54
import io.opentelemetry.api.common.AttributesBuilder;
55
import io.opentelemetry.context.Context;
56
import java.util.ArrayList;
57
import java.util.Collection;
58
import java.util.Collections;
59
import java.util.List;
60
import java.util.concurrent.TimeUnit;
61
import java.util.concurrent.atomic.AtomicIntegerFieldUpdater;
62
import java.util.concurrent.atomic.AtomicLong;
63
import java.util.concurrent.atomic.AtomicLongFieldUpdater;
64
import java.util.logging.Level;
65
import java.util.logging.Logger;
66
import javax.annotation.Nullable;
67

68
/**
69
 * Provides factories for {@link StreamTracer} that records metrics to OpenTelemetry.
70
 *
71
 * <p>On the client-side, a factory is created for each call, and the factory creates a stream
72
 * tracer for each attempt. If there is no stream created when the call is ended, we still create a
73
 * tracer. It's the tracer that reports per-attempt stats, and the factory that reports the stats
74
 * of the overall RPC, such as RETRIES_PER_CALL, to OpenTelemetry.
75
 *
76
 * <p>This module optionally applies a target attribute filter to limit the cardinality of
77
 * the {@code grpc.target} attribute in client-side metrics by mapping disallowed targets
78
 * to a stable placeholder value.
79
 *
80
 * <p>On the server-side, there is only one ServerStream per each ServerCall, and ServerStream
81
 * starts earlier than the ServerCall. Therefore, only one tracer is created per stream/call, and
82
 * it's the tracer that reports the summary to OpenTelemetry.
83
 */
84
final class OpenTelemetryMetricsModule {
85
  private static final Logger logger = Logger.getLogger(OpenTelemetryMetricsModule.class.getName());
1 ✔
86
  public static final ImmutableSet<String> DEFAULT_PER_CALL_METRICS_SET =
1 ✔
87
      ImmutableSet.of(
1 ✔
88
          "grpc.client.attempt.started",
89
          "grpc.client.attempt.duration",
90
          "grpc.client.attempt.sent_total_compressed_message_size",
91
          "grpc.client.attempt.rcvd_total_compressed_message_size",
92
          "grpc.client.call.duration",
93
          "grpc.server.call.started",
94
          "grpc.server.call.duration",
95
          "grpc.server.call.sent_total_compressed_message_size",
96
          "grpc.server.call.rcvd_total_compressed_message_size");
97

98
  // Using floating point because TimeUnit.NANOSECONDS.toSeconds would discard
99
  // fractional seconds.
100
  private static final double SECONDS_PER_NANO = 1e-9;
101

102
  private final OpenTelemetryMetricsResource resource;
103
  private final Supplier<Stopwatch> stopwatchSupplier;
104
  private final boolean localityEnabled;
105
  private final boolean backendServiceEnabled;
106
  private final boolean customLabelEnabled;
107
  private final ImmutableList<OpenTelemetryPlugin> plugins;
108
  @Nullable
109
  private final TargetFilter targetAttributeFilter;
110

111
  OpenTelemetryMetricsModule(Supplier<Stopwatch> stopwatchSupplier,
112
                             OpenTelemetryMetricsResource resource,
113
                             Collection<String> optionalLabels, List<OpenTelemetryPlugin> plugins) {
114
    this(stopwatchSupplier, resource, optionalLabels, plugins, null);
1 ✔
115
  }
1 ✔
116

117
  OpenTelemetryMetricsModule(Supplier<Stopwatch> stopwatchSupplier,
118
      OpenTelemetryMetricsResource resource,
119
      Collection<String> optionalLabels, List<OpenTelemetryPlugin> plugins,
120
      @Nullable TargetFilter targetAttributeFilter) {
1 ✔
121
    this.resource = checkNotNull(resource, "resource");
1 ✔
122
    this.stopwatchSupplier = checkNotNull(stopwatchSupplier, "stopwatchSupplier");
1 ✔
123
    this.localityEnabled = optionalLabels.contains(LOCALITY_KEY.getKey());
1 ✔
124
    this.backendServiceEnabled = optionalLabels.contains(BACKEND_SERVICE_KEY.getKey());
1 ✔
125
    this.customLabelEnabled = optionalLabels.contains(CUSTOM_LABEL_KEY.getKey());
1 ✔
126
    this.plugins = ImmutableList.copyOf(plugins);
1 ✔
127
    this.targetAttributeFilter = targetAttributeFilter;
1 ✔
128
  }
1 ✔
129

130
  @VisibleForTesting
131
  TargetFilter getTargetAttributeFilter() {
132
    return targetAttributeFilter;
1 ✔
133
  }
134

135
  /**
136
   * Returns the server tracer factory.
137
   */
138
  ServerStreamTracer.Factory getServerTracerFactory() {
139
    return new ServerTracerFactory();
1 ✔
140
  }
141

142
  /**
143
   * Returns the client interceptor that facilitates OpenTelemetry metrics reporting.
144
   */
145
  ClientInterceptor getClientInterceptor(String target) {
146
    ImmutableList.Builder<OpenTelemetryPlugin> pluginBuilder =
1 ✔
147
        ImmutableList.builderWithExpectedSize(plugins.size());
1 ✔
148
    for (OpenTelemetryPlugin plugin : plugins) {
1 ✔
149
      if (plugin.enablePluginForChannel(target)) {
1 ✔
150
        pluginBuilder.add(plugin);
1 ✔
151
      }
152
    }
1 ✔
153
    String filteredTarget = recordTarget(target);
1 ✔
154
    return new MetricsClientInterceptor(filteredTarget, pluginBuilder.build());
1 ✔
155
  }
156

157
  String recordTarget(String target) {
158
    if (targetAttributeFilter == null || target == null) {
1 ✔
159
      return target;
1 ✔
160
    }
161
    return targetAttributeFilter.test(target) ? target : "other";
1 ✔
162
  }
163

164
  static String recordMethodName(String fullMethodName, boolean isGeneratedMethod) {
165
    return isGeneratedMethod ? fullMethodName : "other";
1 ✔
166
  }
167

168
  private static final class ClientTracer extends ClientStreamTracer {
169
    @Nullable private static final AtomicLongFieldUpdater<ClientTracer> outboundWireSizeUpdater;
170
    @Nullable private static final AtomicLongFieldUpdater<ClientTracer> inboundWireSizeUpdater;
171

172
    /*
173
     * When using Atomic*FieldUpdater, some Samsung Android 5.0.x devices encounter a bug in their
174
     * JDK reflection API that triggers a NoSuchFieldException. When this occurs, we fall back to
175
     * (potentially racy) direct updates of the volatile variables.
176
     */
177
    static {
178
      AtomicLongFieldUpdater<ClientTracer> tmpOutboundWireSizeUpdater;
179
      AtomicLongFieldUpdater<ClientTracer> tmpInboundWireSizeUpdater;
180
      try {
181
        tmpOutboundWireSizeUpdater =
1 ✔
182
            AtomicLongFieldUpdater.newUpdater(ClientTracer.class, "outboundWireSize");
1 ✔
183
        tmpInboundWireSizeUpdater =
1 ✔
184
            AtomicLongFieldUpdater.newUpdater(ClientTracer.class, "inboundWireSize");
1 ✔
185
      } catch (Throwable t) {
×
186
        logger.log(Level.SEVERE, "Creating atomic field updaters failed", t);
×
187
        tmpOutboundWireSizeUpdater = null;
×
188
        tmpInboundWireSizeUpdater = null;
×
189
      }
1 ✔
190
      outboundWireSizeUpdater = tmpOutboundWireSizeUpdater;
1 ✔
191
      inboundWireSizeUpdater = tmpInboundWireSizeUpdater;
1 ✔
192
    }
1 ✔
193

194
    final Stopwatch stopwatch;
195
    final CallAttemptsTracerFactory attemptsState;
196
    final OpenTelemetryMetricsModule module;
197
    final StreamInfo info;
198
    final String target;
199
    final String fullMethodName;
200
    final List<OpenTelemetryPlugin.ClientStreamPlugin> streamPlugins;
201
    volatile long outboundWireSize;
202
    volatile long inboundWireSize;
203
    volatile String locality;
204
    volatile String backendService;
205
    long attemptNanos;
206
    Code statusCode;
207
    @GuardedBy("this")
208
    @Nullable private Stopwatch activeDelayStopwatch;
209
    @GuardedBy("this")
210
    @Nullable private String activeDelayType;
211
    @GuardedBy("this")
212
    private boolean streamClosed;
213

214
    ClientTracer(CallAttemptsTracerFactory attemptsState, OpenTelemetryMetricsModule module,
215
        StreamInfo info, String target, String fullMethodName,
216
        List<OpenTelemetryPlugin.ClientStreamPlugin> streamPlugins) {
1 ✔
217
      this.attemptsState = attemptsState;
1 ✔
218
      this.module = module;
1 ✔
219
      this.info = info;
1 ✔
220
      this.target = target;
1 ✔
221
      this.fullMethodName = fullMethodName;
1 ✔
222
      this.streamPlugins = streamPlugins;
1 ✔
223
      this.stopwatch = module.stopwatchSupplier.get().start();
1 ✔
224
    }
1 ✔
225

226
    @Override
227
    public synchronized void streamCreated(io.grpc.Attributes transportAtts, Metadata headers) {
228
      endOpenDelay();
1 ✔
229
    }
1 ✔
230

231
    @Override
232
    public synchronized void recordDelayStart(String delayType, String delayReason) {
233
      if (streamClosed) {
1 ✔
234
        return;
1 ✔
235
      }
236
      if (activeDelayStopwatch != null && delayType.equals(activeDelayType)) {
1 ✔
237
        // Do not reset the stopwatch if the delay type is unchanged.
238
        return;
1 ✔
239
      }
240
      endOpenDelay();
1 ✔
241
      activeDelayType = delayType;
1 ✔
242
      activeDelayStopwatch = module.stopwatchSupplier.get().start();
1 ✔
243
    }
1 ✔
244

245
    @Override
246
    public void recordDelayReasonChanged(String delayType, String delayReason) {
247
      // Reason strings are high-cardinality diagnostics intended for tracing spans.
248
    }
1 ✔
249

250
    @Override
251
    public synchronized void recordDelayEnd(String delayType) {
252
      endOpenDelay();
1 ✔
253
    }
1 ✔
254

255
    /** Ends the delay that is still open, if any, under the type it was started with. */
256
    @GuardedBy("this")
257
    private void endOpenDelay() {
258
      if (activeDelayStopwatch != null) {
1 ✔
259
        recordDelay(activeDelayStopwatch.stop().elapsed(TimeUnit.NANOSECONDS), activeDelayType);
1 ✔
260
        activeDelayStopwatch = null;
1 ✔
261
        activeDelayType = null;
1 ✔
262
      }
263
    }
1 ✔
264

265
    private void recordDelay(long delayNanos, String delayType) {
266
      if (module.resource.clientAttemptDelayCounter() != null) {
1 ✔
267
        AttributesBuilder builder = Attributes.builder()
1 ✔
268
            .put(METHOD_KEY, fullMethodName)
1 ✔
269
            .put(TARGET_KEY, target)
1 ✔
270
            .put("grpc.delay_type", delayType);
1 ✔
271
        if (module.customLabelEnabled) {
1 ✔
272
          builder.put(
×
273
              CUSTOM_LABEL_KEY, info.getCallOptions().getOption(Grpc.CALL_OPTION_CUSTOM_LABEL));
×
274
        }
275
        for (OpenTelemetryPlugin.ClientStreamPlugin plugin : streamPlugins) {
1 ✔
276
          plugin.addLabels(builder);
×
277
        }
×
278
        module.resource.clientAttemptDelayCounter()
1 ✔
279
            .record(delayNanos * SECONDS_PER_NANO, builder.build(), attemptsState.otelContext);
1 ✔
280
      }
281
    }
1 ✔
282

283
    @Override
284
    public void inboundHeaders(Metadata headers) {
285
      for (OpenTelemetryPlugin.ClientStreamPlugin plugin : streamPlugins) {
1 ✔
286
        plugin.inboundHeaders(headers);
1 ✔
287
      }
1 ✔
288
    }
1 ✔
289

290
    @Override
291
    @SuppressWarnings("NonAtomicVolatileUpdate")
292
    public void outboundWireSize(long bytes) {
293
      if (outboundWireSizeUpdater != null) {
1 ✔
294
        outboundWireSizeUpdater.getAndAdd(this, bytes);
1 ✔
295
      } else {
296
        outboundWireSize += bytes;
×
297
      }
298
    }
1 ✔
299

300
    @Override
301
    @SuppressWarnings("NonAtomicVolatileUpdate")
302
    public void inboundWireSize(long bytes) {
303
      if (inboundWireSizeUpdater != null) {
1 ✔
304
        inboundWireSizeUpdater.getAndAdd(this, bytes);
1 ✔
305
      } else {
306
        inboundWireSize += bytes;
×
307
      }
308
    }
1 ✔
309

310
    @Override
311
    public void addOptionalLabel(String key, String value) {
312
      if ("grpc.lb.locality".equals(key)) {
1 ✔
313
        locality = value;
1 ✔
314
      }
315
      if ("grpc.lb.backend_service".equals(key)) {
1 ✔
316
        backendService = value;
1 ✔
317
      }
318
    }
1 ✔
319

320
    @Override
321
    public void inboundTrailers(Metadata trailers) {
322
      for (OpenTelemetryPlugin.ClientStreamPlugin plugin : streamPlugins) {
1 ✔
323
        plugin.inboundTrailers(trailers);
1 ✔
324
      }
1 ✔
325
    }
1 ✔
326

327
    @Override
328
    public void streamClosed(Status status) {
329
      synchronized (this) {
1 ✔
330
        if (streamClosed) {
1 ✔
331
          return;
×
332
        }
333
        streamClosed = true;
1 ✔
334
        endOpenDelay();
1 ✔
335
      }
1 ✔
336
      stopwatch.stop();
1 ✔
337
      attemptNanos = stopwatch.elapsed(TimeUnit.NANOSECONDS);
1 ✔
338
      Deadline deadline = info.getCallOptions().getDeadline();
1 ✔
339
      statusCode = status.getCode();
1 ✔
340
      if (statusCode == Code.CANCELLED && deadline != null) {
1 ✔
341
        // When the server's deadline expires, it can only reset the stream with CANCEL and no
342
        // description. Since our timer may be delayed in firing, we double-check the deadline and
343
        // turn the failure into the likely more helpful DEADLINE_EXCEEDED status.
344
        if (deadline.isExpired()) {
1 ✔
345
          statusCode = Code.DEADLINE_EXCEEDED;
1 ✔
346
        }
347
      }
348
      attemptsState.attemptEnded(info.getCallOptions());
1 ✔
349
      recordFinishedAttempt();
1 ✔
350
    }
1 ✔
351

352
    void recordFinishedAttempt() {
353
      AttributesBuilder builder = io.opentelemetry.api.common.Attributes.builder()
1 ✔
354
          .put(METHOD_KEY, fullMethodName)
1 ✔
355
          .put(TARGET_KEY, target)
1 ✔
356
          .put(STATUS_KEY, statusCode.toString());
1 ✔
357
      if (module.localityEnabled) {
1 ✔
358
        String savedLocality = locality;
1 ✔
359
        if (savedLocality == null) {
1 ✔
360
          savedLocality = "";
1 ✔
361
        }
362
        builder.put(LOCALITY_KEY, savedLocality);
1 ✔
363
      }
364
      if (module.backendServiceEnabled) {
1 ✔
365
        String savedBackendService = backendService;
1 ✔
366
        if (savedBackendService == null) {
1 ✔
367
          savedBackendService = "";
1 ✔
368
        }
369
        builder.put(BACKEND_SERVICE_KEY, savedBackendService);
1 ✔
370
      }
371
      if (module.customLabelEnabled) {
1 ✔
372
        builder.put(
1 ✔
373
            CUSTOM_LABEL_KEY, info.getCallOptions().getOption(Grpc.CALL_OPTION_CUSTOM_LABEL));
1 ✔
374
      }
375
      for (OpenTelemetryPlugin.ClientStreamPlugin plugin : streamPlugins) {
1 ✔
376
        plugin.addLabels(builder);
1 ✔
377
      }
1 ✔
378
      io.opentelemetry.api.common.Attributes attribute = builder.build();
1 ✔
379

380
      if (module.resource.clientAttemptDurationCounter() != null ) {
1 ✔
381
        module.resource.clientAttemptDurationCounter()
1 ✔
382
            .record(attemptNanos * SECONDS_PER_NANO, attribute, attemptsState.otelContext);
1 ✔
383
      }
384
      if (module.resource.clientTotalSentCompressedMessageSizeCounter() != null) {
1 ✔
385
        module.resource.clientTotalSentCompressedMessageSizeCounter()
1 ✔
386
            .record(outboundWireSize, attribute, attemptsState.otelContext);
1 ✔
387
      }
388
      if (module.resource.clientTotalReceivedCompressedMessageSizeCounter() != null) {
1 ✔
389
        module.resource.clientTotalReceivedCompressedMessageSizeCounter()
1 ✔
390
            .record(inboundWireSize, attribute, attemptsState.otelContext);
1 ✔
391
      }
392
    }
1 ✔
393
  }
394

395
  @VisibleForTesting
396
  static final class CallAttemptsTracerFactory extends ClientStreamTracer.Factory {
397
    private final OpenTelemetryMetricsModule module;
398
    private final String target;
399
    private final Stopwatch attemptDelayStopwatch;
400
    private final Stopwatch callStopWatch;
401
    @GuardedBy("lock")
402
    private boolean callEnded;
403
    private final String fullMethodName;
404
    private final List<OpenTelemetryPlugin.ClientCallPlugin> callPlugins;
405
    private final Context otelContext;
406
    private Status status;
407
    @GuardedBy("lock")
408
    @Nullable private Stopwatch activeCallDelayStopwatch;
409
    @GuardedBy("lock")
410
    @Nullable private String activeCallDelayType;
411
    private final Attributes callLevelBaseAttributes;
412
    private long retryDelayNanos;
413
    private long callLatencyNanos;
414
    private final Object lock = new Object();
1 ✔
415
    private final AtomicLong attemptsPerCall = new AtomicLong();
1 ✔
416
    private final AtomicLong hedgedAttemptsPerCall = new AtomicLong();
1 ✔
417
    private final AtomicLong transparentRetriesPerCall = new AtomicLong();
1 ✔
418
    @GuardedBy("lock")
419
    private int activeStreams;
420
    @GuardedBy("lock")
421
    private boolean finishedCallToBeRecorded;
422

423
    CallAttemptsTracerFactory(
424
        OpenTelemetryMetricsModule module,
425
        String target,
426
        CallOptions callOptions,
427
        String fullMethodName,
428
        List<OpenTelemetryPlugin.ClientCallPlugin> callPlugins, Context otelContext) {
1 ✔
429
      this.module = checkNotNull(module, "module");
1 ✔
430
      this.target = checkNotNull(target, "target");
1 ✔
431
      this.fullMethodName = checkNotNull(fullMethodName, "fullMethodName");
1 ✔
432
      this.callPlugins = checkNotNull(callPlugins, "callPlugins");
1 ✔
433
      this.otelContext = checkNotNull(otelContext, "otelContext");
1 ✔
434
      this.attemptDelayStopwatch = module.stopwatchSupplier.get();
1 ✔
435
      this.callStopWatch = module.stopwatchSupplier.get().start();
1 ✔
436

437
      AttributesBuilder builder = Attributes.builder()
1 ✔
438
          .put(METHOD_KEY, fullMethodName)
1 ✔
439
          .put(TARGET_KEY, target);
1 ✔
440
      if (module.customLabelEnabled) {
1 ✔
441
        builder.put(
1 ✔
442
            CUSTOM_LABEL_KEY, callOptions.getOption(Grpc.CALL_OPTION_CUSTOM_LABEL));
1 ✔
443
      }
444
      this.callLevelBaseAttributes = builder.build();
1 ✔
445

446
      // Record here in case newClientStreamTracer() would never be called.
447
      if (module.resource.clientAttemptCountCounter() != null) {
1 ✔
448
        module.resource.clientAttemptCountCounter().add(1, callLevelBaseAttributes, otelContext);
1 ✔
449
      }
450
    }
1 ✔
451

452
    @Override
453
    public ClientStreamTracer newClientStreamTracer(StreamInfo info, Metadata metadata) {
454
      synchronized (lock) {
1 ✔
455
        if (finishedCallToBeRecorded) {
1 ✔
456
          // This can be the case when the call is cancelled but a retry attempt is created.
457
          return new ClientStreamTracer() {};
1 ✔
458
        }
459
        if (++activeStreams == 1 && attemptDelayStopwatch.isRunning()) {
1 ✔
460
          attemptDelayStopwatch.stop();
1 ✔
461
          retryDelayNanos = attemptDelayStopwatch.elapsed(TimeUnit.NANOSECONDS);
1 ✔
462
        }
463
      }
1 ✔
464
      // Skip recording for the first time, since it is already recorded in
465
      // CallAttemptsTracerFactory constructor. attemptsPerCall will be non-zero after the first
466
      // attempt, as first attempt cannot be a transparent retry.
467
      if (attemptsPerCall.get() > 0) {
1 ✔
468
        AttributesBuilder builder = io.opentelemetry.api.common.Attributes.builder()
1 ✔
469
            .put(METHOD_KEY, fullMethodName)
1 ✔
470
            .put(TARGET_KEY, target);
1 ✔
471
        if (module.customLabelEnabled) {
1 ✔
472
          builder.put(
1 ✔
473
              CUSTOM_LABEL_KEY, info.getCallOptions().getOption(Grpc.CALL_OPTION_CUSTOM_LABEL));
1 ✔
474
        }
475
        io.opentelemetry.api.common.Attributes attribute = builder.build();
1 ✔
476
        if (module.resource.clientAttemptCountCounter() != null) {
1 ✔
477
          module.resource.clientAttemptCountCounter().add(1, attribute, otelContext);
1 ✔
478
        }
479
      }
480
      if (info.isTransparentRetry()) {
1 ✔
481
        transparentRetriesPerCall.incrementAndGet();
1 ✔
482
      } else if (info.isHedging()) {
1 ✔
483
        hedgedAttemptsPerCall.incrementAndGet();
1 ✔
484
      } else {
485
        attemptsPerCall.incrementAndGet();
1 ✔
486
      }
487
      return newClientTracer(info);
1 ✔
488
    }
489

490
    private ClientTracer newClientTracer(StreamInfo info) {
491
      List<OpenTelemetryPlugin.ClientStreamPlugin> streamPlugins = Collections.emptyList();
1 ✔
492
      if (!callPlugins.isEmpty()) {
1 ✔
493
        streamPlugins = new ArrayList<>(callPlugins.size());
1 ✔
494
        for (OpenTelemetryPlugin.ClientCallPlugin plugin : callPlugins) {
1 ✔
495
          streamPlugins.add(plugin.newClientStreamPlugin());
1 ✔
496
        }
1 ✔
497
        streamPlugins = Collections.unmodifiableList(streamPlugins);
1 ✔
498
      }
499
      return new ClientTracer(this, module, info, target, fullMethodName, streamPlugins);
1 ✔
500
    }
501

502
    // Called whenever each attempt is ended.
503
    void attemptEnded(CallOptions callOptions) {
504
      boolean shouldRecordFinishedCall = false;
1 ✔
505
      synchronized (lock) {
1 ✔
506
        if (--activeStreams == 0) {
1 ✔
507
          attemptDelayStopwatch.start();
1 ✔
508
          if (callEnded && !finishedCallToBeRecorded) {
1 ✔
509
            shouldRecordFinishedCall = true;
×
510
            finishedCallToBeRecorded = true;
×
511
          }
512
        }
513
      }
1 ✔
514
      if (shouldRecordFinishedCall) {
1 ✔
515
        recordFinishedCall(callOptions);
×
516
      }
517
    }
1 ✔
518

519
    void callEnded(Status status, CallOptions callOptions) {
520
      callStopWatch.stop();
1 ✔
521
      this.status = status;
1 ✔
522
      boolean shouldRecordFinishedCall = false;
1 ✔
523
      synchronized (lock) {
1 ✔
524
        if (callEnded) {
1 ✔
525
          // TODO(https://github.com/grpc/grpc-java/issues/7921): this shouldn't happen
526
          return;
×
527
        }
528
        callEnded = true;
1 ✔
529
        endOpenDelay();
1 ✔
530
        if (activeStreams == 0 && !finishedCallToBeRecorded) {
1 ✔
531
          shouldRecordFinishedCall = true;
1 ✔
532
          finishedCallToBeRecorded = true;
1 ✔
533
        }
534
      }
1 ✔
535
      if (shouldRecordFinishedCall) {
1 ✔
536
        recordFinishedCall(callOptions);
1 ✔
537
      }
538
    }
1 ✔
539

540
    void recordFinishedCall(CallOptions callOptions) {
541
      if (attemptsPerCall.get() == 0) {
1 ✔
542
        ClientTracer tracer = newClientTracer(null);
1 ✔
543
        tracer.attemptNanos = attemptDelayStopwatch.elapsed(TimeUnit.NANOSECONDS);
1 ✔
544
        tracer.statusCode = status.getCode();
1 ✔
545
        tracer.recordFinishedAttempt();
1 ✔
546
      }
547
      callLatencyNanos = callStopWatch.elapsed(TimeUnit.NANOSECONDS);
1 ✔
548

549
      // Base attributes
550
      AttributesBuilder builder = io.opentelemetry.api.common.Attributes.builder()
1 ✔
551
          .put(METHOD_KEY, fullMethodName)
1 ✔
552
          .put(TARGET_KEY, target);
1 ✔
553
      if (module.customLabelEnabled) {
1 ✔
554
        builder.put(CUSTOM_LABEL_KEY, callOptions.getOption(Grpc.CALL_OPTION_CUSTOM_LABEL));
1 ✔
555
      }
556
      io.opentelemetry.api.common.Attributes baseAttributes = builder.build();
1 ✔
557

558
      // Duration
559
      if (module.resource.clientCallDurationCounter() != null) {
1 ✔
560
        module.resource.clientCallDurationCounter().record(
1 ✔
561
            callLatencyNanos * SECONDS_PER_NANO,
562
            baseAttributes.toBuilder()
1 ✔
563
                .put(STATUS_KEY, status.getCode().toString())
1 ✔
564
                .build(),
1 ✔
565
            otelContext
566
        );
567
      }
568

569
      // Retry counts
570
      if (module.resource.clientCallRetriesCounter() != null) {
1 ✔
571
        long retriesPerCall = Math.max(attemptsPerCall.get() - 1, 0);
1 ✔
572
        if (retriesPerCall > 0) {
1 ✔
573
          module.resource.clientCallRetriesCounter()
1 ✔
574
              .record(retriesPerCall, baseAttributes, otelContext);
1 ✔
575
        }
576
      }
577

578
      // Hedge counts
579
      if (module.resource.clientCallHedgesCounter() != null) {
1 ✔
580
        long hedges = hedgedAttemptsPerCall.get();
1 ✔
581
        if (hedges > 0) {
1 ✔
582
          module.resource.clientCallHedgesCounter()
1 ✔
583
              .record(hedges, baseAttributes, otelContext);
1 ✔
584
        }
585
      }
586

587
      // Transparent Retry counts
588
      if (module.resource.clientCallTransparentRetriesCounter() != null) {
1 ✔
589
        long transparentRetries = transparentRetriesPerCall.get();
1 ✔
590
        if (transparentRetries > 0) {
1 ✔
591
          module.resource.clientCallTransparentRetriesCounter()
1 ✔
592
              .record(transparentRetries, baseAttributes, otelContext);
1 ✔
593
        }
594
      }
595

596
      // Retry delay
597
      if (module.resource.clientCallRetryDelayCounter() != null) {
1 ✔
598
        module.resource.clientCallRetryDelayCounter().record(
1 ✔
599
            retryDelayNanos * SECONDS_PER_NANO,
600
            baseAttributes,
601
            otelContext
602
        );
603
      }
604
    }
1 ✔
605

606
    @Override
607
    public void recordDelayStart(String delayType, String delayReason) {
608
      synchronized (lock) {
1 ✔
609
        if (callEnded) {
1 ✔
610
          return;
1 ✔
611
        }
612
        if (activeCallDelayStopwatch != null && delayType.equals(activeCallDelayType)) {
1 ✔
613
          return;
1 ✔
614
        }
615
        endOpenDelay();
1 ✔
616
        activeCallDelayType = delayType;
1 ✔
617
        activeCallDelayStopwatch = module.stopwatchSupplier.get().start();
1 ✔
618
      }
1 ✔
619
    }
1 ✔
620

621
    @Override
622
    public void recordDelayReasonChanged(String delayType, String delayReason) {
623
      // Reason strings are high-cardinality diagnostics intended for tracing spans.
624
    }
1 ✔
625

626
    @Override
627
    public void recordDelayEnd(String delayType) {
628
      synchronized (lock) {
1 ✔
629
        endOpenDelay();
1 ✔
630
      }
1 ✔
631
    }
1 ✔
632

633
    @GuardedBy("lock")
634
    private void endOpenDelay() {
635
      if (activeCallDelayStopwatch != null) {
1 ✔
636
        recordDelay(
1 ✔
637
            activeCallDelayStopwatch.stop().elapsed(TimeUnit.NANOSECONDS), activeCallDelayType);
1 ✔
638
        activeCallDelayStopwatch = null;
1 ✔
639
        activeCallDelayType = null;
1 ✔
640
      }
641
    }
1 ✔
642

643
    private void recordDelay(long delayNanos, String delayType) {
644
      if (module.resource.clientCallDelayCounter() != null) {
1 ✔
645
        module.resource.clientCallDelayCounter().record(
1 ✔
646
            delayNanos * SECONDS_PER_NANO,
647
            callLevelBaseAttributes.toBuilder().put("grpc.delay_type", delayType).build(),
1 ✔
648
            otelContext);
649
      }
650
    }
1 ✔
651
  }
652

653
  private static final class ServerTracer extends ServerStreamTracer
654
      implements ServerCallMethodListener {
655
    @Nullable private static final AtomicIntegerFieldUpdater<ServerTracer> streamClosedUpdater;
656
    @Nullable private static final AtomicLongFieldUpdater<ServerTracer> outboundWireSizeUpdater;
657
    @Nullable private static final AtomicLongFieldUpdater<ServerTracer> inboundWireSizeUpdater;
658

659
    /*
660
     * When using Atomic*FieldUpdater, some Samsung Android 5.0.x devices encounter a bug in their
661
     * JDK reflection API that triggers a NoSuchFieldException. When this occurs, we fall back to
662
     * (potentially racy) direct updates of the volatile variables.
663
     */
664
    static {
665
      AtomicIntegerFieldUpdater<ServerTracer> tmpStreamClosedUpdater;
666
      AtomicLongFieldUpdater<ServerTracer> tmpOutboundWireSizeUpdater;
667
      AtomicLongFieldUpdater<ServerTracer> tmpInboundWireSizeUpdater;
668
      try {
669
        tmpStreamClosedUpdater =
1 ✔
670
            AtomicIntegerFieldUpdater.newUpdater(ServerTracer.class, "streamClosed");
1 ✔
671
        tmpOutboundWireSizeUpdater =
1 ✔
672
            AtomicLongFieldUpdater.newUpdater(ServerTracer.class, "outboundWireSize");
1 ✔
673
        tmpInboundWireSizeUpdater =
1 ✔
674
            AtomicLongFieldUpdater.newUpdater(ServerTracer.class, "inboundWireSize");
1 ✔
675
      } catch (Throwable t) {
×
676
        logger.log(Level.SEVERE, "Creating atomic field updaters failed", t);
×
677
        tmpStreamClosedUpdater = null;
×
678
        tmpOutboundWireSizeUpdater = null;
×
679
        tmpInboundWireSizeUpdater = null;
×
680
      }
1 ✔
681
      streamClosedUpdater = tmpStreamClosedUpdater;
1 ✔
682
      outboundWireSizeUpdater = tmpOutboundWireSizeUpdater;
1 ✔
683
      inboundWireSizeUpdater = tmpInboundWireSizeUpdater;
1 ✔
684
    }
1 ✔
685

686
    private final OpenTelemetryMetricsModule module;
687
    private final String fullMethodName;
688
    private final List<OpenTelemetryPlugin.ServerStreamPlugin> streamPlugins;
689
    private Context otelContext = Context.root();
1 ✔
690
    private volatile boolean isGeneratedMethod;
691
    private volatile int streamClosed;
692
    private final Stopwatch stopwatch;
693
    private volatile long outboundWireSize;
694
    private volatile long inboundWireSize;
695

696
    ServerTracer(OpenTelemetryMetricsModule module, String fullMethodName,
697
        List<OpenTelemetryPlugin.ServerStreamPlugin> streamPlugins) {
1 ✔
698
      this.module = checkNotNull(module, "module");
1 ✔
699
      this.fullMethodName = fullMethodName;
1 ✔
700
      this.streamPlugins = checkNotNull(streamPlugins, "streamPlugins");
1 ✔
701
      this.stopwatch = module.stopwatchSupplier.get().start();
1 ✔
702
    }
1 ✔
703

704
    @Override
705
    public io.grpc.Context filterContext(io.grpc.Context context) {
706
      Baggage baggage = BAGGAGE_KEY.get(context);
1 ✔
707
      if (baggage != null) {
1 ✔
708
        otelContext = Context.current().with(baggage);
1 ✔
709
      } else {
710
        otelContext = Context.current();
1 ✔
711
      }
712
      return context;
1 ✔
713
    }
714

715
    @Override
716
    public void serverCallMethodResolved(MethodDescriptor<?, ?> method) {
717
      isGeneratedMethod = method.isSampledToLocalTracing();
1 ✔
718
    }
1 ✔
719

720
    @Override
721
    public void serverCallStarted(ServerCallInfo<?, ?> callInfo) {
722
      // Only record method name as an attribute if isSampledToLocalTracing is set to true,
723
      // which is true for all generated methods. Otherwise, programmatically
724
      // created methods result in high cardinality metrics.
725
      boolean isSampledToLocalTracing = callInfo.getMethodDescriptor().isSampledToLocalTracing();
1 ✔
726
      isGeneratedMethod = isSampledToLocalTracing;
1 ✔
727

728
      io.opentelemetry.api.common.Attributes attribute =
1 ✔
729
          io.opentelemetry.api.common.Attributes.of(
1 ✔
730
              METHOD_KEY, recordMethodName(fullMethodName, isSampledToLocalTracing));
1 ✔
731

732
      if (module.resource.serverCallCountCounter() != null) {
1 ✔
733
        module.resource.serverCallCountCounter().add(1, attribute, otelContext);
1 ✔
734
      }
735
    }
1 ✔
736

737
    @Override
738
    @SuppressWarnings("NonAtomicVolatileUpdate")
739
    public void outboundWireSize(long bytes) {
740
      if (outboundWireSizeUpdater != null) {
1 ✔
741
        outboundWireSizeUpdater.getAndAdd(this, bytes);
1 ✔
742
      } else {
743
        outboundWireSize += bytes;
×
744
      }
745
    }
1 ✔
746

747
    @Override
748
    @SuppressWarnings("NonAtomicVolatileUpdate")
749
    public void inboundWireSize(long bytes) {
750
      if (inboundWireSizeUpdater != null) {
1 ✔
751
        inboundWireSizeUpdater.getAndAdd(this, bytes);
1 ✔
752
      } else {
753
        inboundWireSize += bytes;
×
754
      }
755
    }
1 ✔
756

757
    /**
758
     * Record a finished stream and mark the current time as the end time.
759
     *
760
     * <p>Can be called from any thread without synchronization.  Calling it the second time or more
761
     * is a no-op.
762
     */
763
    @Override
764
    public void streamClosed(Status status) {
765
      if (streamClosedUpdater != null) {
1 ✔
766
        if (streamClosedUpdater.getAndSet(this, 1) != 0) {
1 ✔
767
          return;
×
768
        }
769
      } else {
770
        if (streamClosed != 0) {
×
771
          return;
×
772
        }
773
        streamClosed = 1;
×
774
      }
775
      stopwatch.stop();
1 ✔
776
      long elapsedTimeNanos = stopwatch.elapsed(TimeUnit.NANOSECONDS);
1 ✔
777
      recordClosedStream(
1 ✔
778
          status,
779
          elapsedTimeNanos,
780
          outboundWireSize,
781
          inboundWireSize,
782
          isGeneratedMethod);
783
    }
1 ✔
784

785
    private void recordClosedStream(
786
        Status status,
787
        long elapsedTimeNanos,
788
        long closedOutboundWireSize,
789
        long closedInboundWireSize,
790
        boolean generatedMethod) {
791
      AttributesBuilder builder =
792
          io.opentelemetry.api.common.Attributes.builder()
1 ✔
793
              .put(METHOD_KEY, recordMethodName(fullMethodName, generatedMethod))
1 ✔
794
              .put(STATUS_KEY, status.getCode().toString());
1 ✔
795
      for (OpenTelemetryPlugin.ServerStreamPlugin plugin : streamPlugins) {
1 ✔
796
        plugin.addLabels(builder);
1 ✔
797
      }
1 ✔
798
      io.opentelemetry.api.common.Attributes attributes = builder.build();
1 ✔
799

800
      if (module.resource.serverCallDurationCounter() != null) {
1 ✔
801
        module.resource.serverCallDurationCounter()
1 ✔
802
            .record(elapsedTimeNanos * SECONDS_PER_NANO, attributes, otelContext);
1 ✔
803
      }
804
      if (module.resource.serverTotalSentCompressedMessageSizeCounter() != null) {
1 ✔
805
        module.resource.serverTotalSentCompressedMessageSizeCounter()
1 ✔
806
            .record(closedOutboundWireSize, attributes, otelContext);
1 ✔
807
      }
808
      if (module.resource.serverTotalReceivedCompressedMessageSizeCounter() != null) {
1 ✔
809
        module.resource.serverTotalReceivedCompressedMessageSizeCounter()
1 ✔
810
            .record(closedInboundWireSize, attributes, otelContext);
1 ✔
811
      }
812
    }
1 ✔
813
  }
814

815
  @VisibleForTesting
816
  final class ServerTracerFactory extends ServerStreamTracer.Factory {
1 ✔
817
    @Override
818
    public ServerStreamTracer newServerStreamTracer(String fullMethodName, Metadata headers) {
819
      final List<OpenTelemetryPlugin.ServerStreamPlugin> streamPlugins;
820
      if (plugins.isEmpty()) {
1 ✔
821
        streamPlugins = Collections.emptyList();
1 ✔
822
      } else {
823
        List<OpenTelemetryPlugin.ServerStreamPlugin> streamPluginsMutable =
1 ✔
824
            new ArrayList<>(plugins.size());
1 ✔
825
        for (OpenTelemetryPlugin plugin : plugins) {
1 ✔
826
          streamPluginsMutable.add(plugin.newServerStreamPlugin(headers));
1 ✔
827
        }
1 ✔
828
        streamPlugins = Collections.unmodifiableList(streamPluginsMutable);
1 ✔
829
      }
830
      return new ServerTracer(OpenTelemetryMetricsModule.this, fullMethodName,
1 ✔
831
          streamPlugins);
832
    }
833
  }
834

835
  @VisibleForTesting
836
  final class MetricsClientInterceptor implements ClientInterceptor {
837
    private final String target;
838
    private final ImmutableList<OpenTelemetryPlugin> plugins;
839

840
    MetricsClientInterceptor(String target, ImmutableList<OpenTelemetryPlugin> plugins) {
1 ✔
841
      this.target = checkNotNull(target, "target");
1 ✔
842
      this.plugins = checkNotNull(plugins, "plugins");
1 ✔
843
    }
1 ✔
844

845
    @Override
846
    public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(
847
        MethodDescriptor<ReqT, RespT> method, CallOptions callOptions, Channel next) {
848
      final List<OpenTelemetryPlugin.ClientCallPlugin> callPlugins;
849
      if (plugins.isEmpty()) {
1 ✔
850
        callPlugins = Collections.emptyList();
1 ✔
851
      } else {
852
        List<OpenTelemetryPlugin.ClientCallPlugin> callPluginsMutable =
1 ✔
853
            new ArrayList<>(plugins.size());
1 ✔
854
        for (OpenTelemetryPlugin plugin : plugins) {
1 ✔
855
          callPluginsMutable.add(plugin.newClientCallPlugin());
1 ✔
856
        }
1 ✔
857
        callPlugins = Collections.unmodifiableList(callPluginsMutable);
1 ✔
858
        for (OpenTelemetryPlugin.ClientCallPlugin plugin : callPlugins) {
1 ✔
859
          callOptions = plugin.filterCallOptions(callOptions);
1 ✔
860
        }
1 ✔
861
      }
862
      final CallOptions finalCallOptions = callOptions;
1 ✔
863
      // Only record method name as an attribute if isSampledToLocalTracing is set to true,
864
      // which is true for all generated methods. Otherwise, programatically
865
      // created methods result in high cardinality metrics.
866
      final CallAttemptsTracerFactory tracerFactory = new CallAttemptsTracerFactory(
1 ✔
867
          OpenTelemetryMetricsModule.this, target, callOptions,
868
          recordMethodName(method.getFullMethodName(), method.isSampledToLocalTracing()),
1 ✔
869
          callPlugins, Context.current());
1 ✔
870
      ClientCall<ReqT, RespT> call =
1 ✔
871
          next.newCall(method, callOptions.withStreamTracerFactory(tracerFactory));
1 ✔
872
      return new SimpleForwardingClientCall<ReqT, RespT>(call) {
1 ✔
873
        @Override
874
        public void start(Listener<RespT> responseListener, Metadata headers) {
875
          for (OpenTelemetryPlugin.ClientCallPlugin plugin : callPlugins) {
1 ✔
876
            plugin.addMetadata(headers);
1 ✔
877
          }
1 ✔
878
          delegate().start(
1 ✔
879
              new SimpleForwardingClientCallListener<RespT>(responseListener) {
1 ✔
880
                @Override
881
                public void onClose(Status status, Metadata trailers) {
882
                  tracerFactory.callEnded(status, finalCallOptions);
1 ✔
883
                  super.onClose(status, trailers);
1 ✔
884
                }
1 ✔
885
              },
886
              headers);
887
        }
1 ✔
888
      };
889
    }
890
  }
891
}
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