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

grpc / grpc-java / #20365

27 Jul 2026 06:04AM UTC coverage: 89.199% (+0.07%) from 89.125%
#20365

push

github

web-flow
core: Implement LB Delay Observability (Proposal A121) (#12807)

This PR implements **Attempt-Level RPC Delay Observability** across the core channel transport, built-in load balancers, xDS policies, and the OpenTelemetry telemetry plugin, aligned with [gRPC Proposal A121](https://github.com/grpc/proposal/pull/556).

38276 of 42911 relevant lines covered (89.2%)

0.89 hits per line

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

95.26
/../opentelemetry/src/main/java/io/grpc/opentelemetry/GrpcOpenTelemetry.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.internal.GrpcUtil.IMPLEMENTATION_VERSION;
21
import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.HEDGE_BUCKETS;
22
import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.LATENCY_BUCKETS;
23
import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.RETRY_BUCKETS;
24
import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.SIZE_BUCKETS;
25
import static io.grpc.opentelemetry.internal.OpenTelemetryConstants.TRANSPARENT_RETRY_BUCKETS;
26

27
import com.google.common.annotations.VisibleForTesting;
28
import com.google.common.base.Stopwatch;
29
import com.google.common.base.Supplier;
30
import com.google.common.collect.ImmutableList;
31
import com.google.common.collect.ImmutableMap;
32
import io.grpc.ExperimentalApi;
33
import io.grpc.InternalConfigurator;
34
import io.grpc.InternalConfiguratorRegistry;
35
import io.grpc.InternalManagedChannelBuilder;
36
import io.grpc.ManagedChannelBuilder;
37
import io.grpc.MetricSink;
38
import io.grpc.ServerBuilder;
39
import io.grpc.internal.GrpcUtil;
40
import io.grpc.opentelemetry.internal.OpenTelemetryConstants;
41
import io.opentelemetry.api.OpenTelemetry;
42
import io.opentelemetry.api.metrics.Meter;
43
import io.opentelemetry.api.metrics.MeterProvider;
44
import io.opentelemetry.api.trace.Tracer;
45
import java.util.ArrayList;
46
import java.util.Collection;
47
import java.util.Collections;
48
import java.util.HashMap;
49
import java.util.List;
50
import java.util.Map;
51
import java.util.function.Predicate;
52
import javax.annotation.Nullable;
53
import org.codehaus.mojo.animal_sniffer.IgnoreJRERequirement;
54

55
/**
56
 *  The entrypoint for OpenTelemetry metrics functionality in gRPC.
57
 *
58
 *  <p>GrpcOpenTelemetry uses {@link io.opentelemetry.api.OpenTelemetry} APIs for instrumentation.
59
 *  When no SDK is explicitly added no telemetry data will be collected. See
60
 *  {@code io.opentelemetry.sdk.OpenTelemetrySdk} for information on how to construct the SDK.
61
 *
62
 */
63
public final class GrpcOpenTelemetry {
64

65
  private static final Supplier<Stopwatch> STOPWATCH_SUPPLIER = new Supplier<Stopwatch>() {
1✔
66
    @Override
67
    public Stopwatch get() {
68
      return Stopwatch.createUnstarted();
1✔
69
    }
70
  };
71

72
  @VisibleForTesting
73
  static boolean ENABLE_OTEL_TRACING =
1✔
74
      GrpcUtil.getFlag("GRPC_EXPERIMENTAL_ENABLE_OTEL_TRACING", false);
1✔
75

76
  private final OpenTelemetry openTelemetrySdk;
77
  private final MeterProvider meterProvider;
78
  private final Meter meter;
79
  private final Map<String, Boolean> enableMetrics;
80
  private final boolean disableDefault;
81
  private final OpenTelemetryMetricsResource resource;
82
  private final OpenTelemetryMetricsModule openTelemetryMetricsModule;
83
  private final OpenTelemetryTracingModule openTelemetryTracingModule;
84
  private final List<String> optionalLabels;
85
  private final MetricSink sink;
86

87
  public static Builder newBuilder() {
88
    return new Builder();
1✔
89
  }
90

91
  private GrpcOpenTelemetry(Builder builder) {
1✔
92
    this.openTelemetrySdk = checkNotNull(builder.openTelemetrySdk, "openTelemetrySdk");
1✔
93
    this.meterProvider = checkNotNull(openTelemetrySdk.getMeterProvider(), "meterProvider");
1✔
94
    this.meter = this.meterProvider
1✔
95
        .meterBuilder(OpenTelemetryConstants.INSTRUMENTATION_SCOPE)
1✔
96
        .setInstrumentationVersion(IMPLEMENTATION_VERSION)
1✔
97
        .build();
1✔
98
    this.enableMetrics = ImmutableMap.copyOf(builder.enableMetrics);
1✔
99
    this.disableDefault = builder.disableAll;
1✔
100
    this.resource = createMetricInstruments(meter, enableMetrics, disableDefault);
1✔
101
    this.optionalLabels = ImmutableList.copyOf(builder.optionalLabels);
1✔
102
    this.openTelemetryMetricsModule = new OpenTelemetryMetricsModule(
1✔
103
        STOPWATCH_SUPPLIER, resource, optionalLabels, builder.plugins,
1✔
104
        builder.targetFilter);
1✔
105
    this.openTelemetryTracingModule = new OpenTelemetryTracingModule(openTelemetrySdk);
1✔
106
    this.sink = new OpenTelemetryMetricSink(meter, enableMetrics, disableDefault, optionalLabels);
1✔
107
  }
1✔
108

109
  @VisibleForTesting
110
  OpenTelemetry getOpenTelemetryInstance() {
111
    return this.openTelemetrySdk;
1✔
112
  }
113

114
  @VisibleForTesting
115
  MeterProvider getMeterProvider() {
116
    return this.meterProvider;
1✔
117
  }
118

119
  @VisibleForTesting
120
  Meter getMeter() {
121
    return this.meter;
1✔
122
  }
123

124
  @VisibleForTesting
125
  OpenTelemetryMetricsResource getResource() {
126
    return this.resource;
×
127
  }
128

129
  @VisibleForTesting
130
  Map<String, Boolean> getEnableMetrics() {
131
    return this.enableMetrics;
1✔
132
  }
133

134
  @VisibleForTesting
135
  List<String> getOptionalLabels() {
136
    return optionalLabels;
1✔
137
  }
138

139
  MetricSink getSink() {
140
    return sink;
1✔
141
  }
142

143
  @VisibleForTesting
144
  Tracer getTracer() {
145
    return this.openTelemetryTracingModule.getTracer();
1✔
146
  }
147

148
  @VisibleForTesting
149
  TargetFilter getTargetAttributeFilter() {
150
    return this.openTelemetryMetricsModule.getTargetAttributeFilter();
1✔
151
  }
152

153
  /**
154
   * Registers GrpcOpenTelemetry globally, applying its configuration to all subsequently created
155
   * gRPC channels and servers.
156
   */
157
  @ExperimentalApi("https://github.com/grpc/grpc-java/issues/10591")
158
  public void registerGlobal() {
159
    InternalConfiguratorRegistry.setConfigurators(Collections.singletonList(
×
160
        new InternalConfigurator() {
×
161
          @Override
162
          public void configureChannelBuilder(ManagedChannelBuilder<?> channelBuilder) {
163
            GrpcOpenTelemetry.this.configureChannelBuilder(channelBuilder);
×
164
          }
×
165

166
          @Override
167
          public void configureServerBuilder(ServerBuilder<?> serverBuilder) {
168
            GrpcOpenTelemetry.this.configureServerBuilder(serverBuilder);
×
169
          }
×
170
        }));
171
  }
×
172

173
  /**
174
   * Configures the given {@link ManagedChannelBuilder} with OpenTelemetry metrics instrumentation.
175
   */
176
  public void configureChannelBuilder(ManagedChannelBuilder<?> builder) {
177
    InternalManagedChannelBuilder.addMetricSink(builder, sink);
1✔
178
    InternalManagedChannelBuilder.interceptWithTarget(
1✔
179
        builder, openTelemetryMetricsModule::getClientInterceptor);
180
    if (ENABLE_OTEL_TRACING) {
1✔
181
      builder.intercept(openTelemetryTracingModule.getClientInterceptor());
1✔
182
    }
183
  }
1✔
184

185
  /**
186
   * Configures the given {@link ServerBuilder} with OpenTelemetry metrics instrumentation.
187
   *
188
   * @param serverBuilder the server builder to configure
189
   */
190
  public void configureServerBuilder(ServerBuilder<?> serverBuilder) {
191
    /* To ensure baggage propagation to metrics, we need the tracing
192
    tracers to be initialised before metrics */
193
    if (ENABLE_OTEL_TRACING) {
1✔
194
      serverBuilder.addStreamTracerFactory(
1✔
195
          openTelemetryTracingModule.getServerTracerFactory());
1✔
196
      serverBuilder.intercept(openTelemetryTracingModule.getServerSpanPropagationInterceptor());
1✔
197
    }
198
    serverBuilder.addStreamTracerFactory(openTelemetryMetricsModule.getServerTracerFactory());
1✔
199
    serverBuilder.addMetricSink(sink);
1✔
200
  }
1✔
201

202
  @VisibleForTesting
203
  static OpenTelemetryMetricsResource createMetricInstruments(Meter meter,
204
      Map<String, Boolean> enableMetrics, boolean disableDefault) {
205
    OpenTelemetryMetricsResource.Builder builder = OpenTelemetryMetricsResource.builder();
1✔
206

207
    if (isMetricEnabled("grpc.client.call.duration", enableMetrics, disableDefault)) {
1✔
208
      builder.clientCallDurationCounter(
1✔
209
          meter.histogramBuilder("grpc.client.call.duration")
1✔
210
              .setUnit("s")
1✔
211
              .setDescription(
1✔
212
                  "Time taken by gRPC to complete an RPC from application's perspective")
213
              .setExplicitBucketBoundariesAdvice(LATENCY_BUCKETS)
1✔
214
              .build());
1✔
215
    }
216

217
    if (isMetricEnabled("grpc.client.attempt.started", enableMetrics, disableDefault)) {
1✔
218
      builder.clientAttemptCountCounter(
1✔
219
          meter.counterBuilder("grpc.client.attempt.started")
1✔
220
              .setUnit("{attempt}")
1✔
221
              .setDescription("Number of client call attempts started")
1✔
222
              .build());
1✔
223
    }
224

225
    if (isMetricEnabled("grpc.client.attempt.duration", enableMetrics, disableDefault)) {
1✔
226
      builder.clientAttemptDurationCounter(
1✔
227
          meter.histogramBuilder(
1✔
228
                  "grpc.client.attempt.duration")
229
              .setUnit("s")
1✔
230
              .setDescription("Time taken to complete a client call attempt")
1✔
231
              .setExplicitBucketBoundariesAdvice(LATENCY_BUCKETS)
1✔
232
              .build());
1✔
233
    }
234

235
    if (isDelayObservabilityEnabled()
1✔
236
        && isMetricEnabled("grpc.client.attempt.delay.duration", enableMetrics, disableDefault)) {
1✔
237
      builder.clientAttemptDelayCounter(
1✔
238
          meter.histogramBuilder(
1✔
239
                  "grpc.client.attempt.delay.duration")
240
              .setUnit("s")
1✔
241
              .setDescription("Time taken before a client call attempt starts")
1✔
242
              .setExplicitBucketBoundariesAdvice(LATENCY_BUCKETS)
1✔
243
              .build());
1✔
244
    }
245

246
    if (isMetricEnabled("grpc.client.attempt.sent_total_compressed_message_size", enableMetrics,
1✔
247
        disableDefault)) {
248
      builder.clientTotalSentCompressedMessageSizeCounter(
1✔
249
          meter.histogramBuilder(
1✔
250
                  "grpc.client.attempt.sent_total_compressed_message_size")
251
              .setUnit("By")
1✔
252
              .setDescription("Compressed message bytes sent per client call attempt")
1✔
253
              .ofLongs()
1✔
254
              .setExplicitBucketBoundariesAdvice(SIZE_BUCKETS)
1✔
255
              .build());
1✔
256
    }
257

258
    if (isMetricEnabled("grpc.client.attempt.rcvd_total_compressed_message_size", enableMetrics,
1✔
259
        disableDefault)) {
260
      builder.clientTotalReceivedCompressedMessageSizeCounter(
1✔
261
          meter.histogramBuilder(
1✔
262
                  "grpc.client.attempt.rcvd_total_compressed_message_size")
263
              .setUnit("By")
1✔
264
              .setDescription("Compressed message bytes received per call attempt")
1✔
265
              .ofLongs()
1✔
266
              .setExplicitBucketBoundariesAdvice(SIZE_BUCKETS)
1✔
267
              .build());
1✔
268
    }
269

270
    if (isMetricEnabled("grpc.client.call.retries", enableMetrics, disableDefault)) {
1✔
271
      builder.clientCallRetriesCounter(
1✔
272
          meter.histogramBuilder(
1✔
273
                  "grpc.client.call.retries")
274
              .setUnit("{retry}")
1✔
275
              .setDescription("Number of retries during the client call. "
1✔
276
                  + "If there were no retries, 0 is not reported.")
277
              .ofLongs()
1✔
278
              .setExplicitBucketBoundariesAdvice(RETRY_BUCKETS)
1✔
279
              .build());
1✔
280
    }
281

282
    if (isMetricEnabled("grpc.client.call.transparent_retries", enableMetrics,
1✔
283
        disableDefault)) {
284
      builder.clientCallTransparentRetriesCounter(
1✔
285
          meter.histogramBuilder(
1✔
286
                  "grpc.client.call.transparent_retries")
287
              .setUnit("{transparent_retry}")
1✔
288
              .setDescription("Number of transparent retries during the client call. "
1✔
289
                  + "If there were no transparent retries, 0 is not reported.")
290
              .ofLongs()
1✔
291
              .setExplicitBucketBoundariesAdvice(TRANSPARENT_RETRY_BUCKETS)
1✔
292
              .build());
1✔
293
    }
294

295
    if (isMetricEnabled("grpc.client.call.hedges", enableMetrics, disableDefault)) {
1✔
296
      builder.clientCallHedgesCounter(
1✔
297
          meter.histogramBuilder(
1✔
298
                  "grpc.client.call.hedges")
299
              .setUnit("{hedge}")
1✔
300
              .setDescription("Number of hedges during the client call. "
1✔
301
                  + "If there were no hedges, 0 is not reported.")
302
              .ofLongs()
1✔
303
              .setExplicitBucketBoundariesAdvice(HEDGE_BUCKETS)
1✔
304
              .build());
1✔
305
    }
306

307
    if (isMetricEnabled("grpc.client.call.retry_delay", enableMetrics, disableDefault)) {
1✔
308
      builder.clientCallRetryDelayCounter(
1✔
309
          meter.histogramBuilder(
1✔
310
                  "grpc.client.call.retry_delay")
311
              .setUnit("s")
1✔
312
              .setDescription("Total time of delay while there is no active attempt during the "
1✔
313
                  + "client call")
314
              .setExplicitBucketBoundariesAdvice(LATENCY_BUCKETS)
1✔
315
              .build());
1✔
316
    }
317

318
    if (isMetricEnabled("grpc.server.call.started", enableMetrics, disableDefault)) {
1✔
319
      builder.serverCallCountCounter(
1✔
320
          meter.counterBuilder("grpc.server.call.started")
1✔
321
              .setUnit("{call}")
1✔
322
              .setDescription("Number of server calls started")
1✔
323
              .build());
1✔
324
    }
325

326
    if (isMetricEnabled("grpc.server.call.duration", enableMetrics, disableDefault)) {
1✔
327
      builder.serverCallDurationCounter(
1✔
328
          meter.histogramBuilder("grpc.server.call.duration")
1✔
329
              .setUnit("s")
1✔
330
              .setDescription(
1✔
331
                  "Time taken to complete a call from server transport's perspective")
332
              .setExplicitBucketBoundariesAdvice(LATENCY_BUCKETS)
1✔
333
              .build());
1✔
334
    }
335

336
    if (isMetricEnabled("grpc.server.call.sent_total_compressed_message_size",
1✔
337
        enableMetrics, disableDefault)) {
338
      builder.serverTotalSentCompressedMessageSizeCounter(
1✔
339
          meter.histogramBuilder(
1✔
340
                  "grpc.server.call.sent_total_compressed_message_size")
341
              .setUnit("By")
1✔
342
              .setDescription("Compressed message bytes sent per server call")
1✔
343
              .ofLongs()
1✔
344
              .setExplicitBucketBoundariesAdvice(SIZE_BUCKETS)
1✔
345
              .build());
1✔
346
    }
347

348
    if (isMetricEnabled("grpc.server.call.rcvd_total_compressed_message_size",
1✔
349
        enableMetrics, disableDefault)) {
350
      builder.serverTotalReceivedCompressedMessageSizeCounter(
1✔
351
          meter.histogramBuilder(
1✔
352
                  "grpc.server.call.rcvd_total_compressed_message_size")
353
              .setUnit("By")
1✔
354
              .setDescription("Compressed message bytes received per server call")
1✔
355
              .ofLongs()
1✔
356
              .setExplicitBucketBoundariesAdvice(SIZE_BUCKETS)
1✔
357
              .build());
1✔
358
    }
359

360
    return builder.build();
1✔
361
  }
362

363
  /**
364
   * Checks whether experimental client attempt and call delay observability is globally enabled.
365
   *
366
   * <p>Guarded strictly by the {@code GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY} environment
367
   * variable (defaults to {@code false}). When disabled, delay spans and
368
   * duration histograms are suppressed to avoid runtime overhead.
369
   */
370
  static boolean isDelayObservabilityEnabled() {
371
    return GrpcUtil.getFlag("GRPC_EXPERIMENTAL_ENABLE_DELAY_OBSERVABILITY", false);
1✔
372
  }
373

374
  static boolean isMetricEnabled(String metricName, Map<String, Boolean> enableMetrics,
375
      boolean disableDefault) {
376
    Boolean explicitlyEnabled = enableMetrics.get(metricName);
1✔
377
    if (explicitlyEnabled != null) {
1✔
378
      return explicitlyEnabled;
1✔
379
    }
380
    return OpenTelemetryMetricsModule.DEFAULT_PER_CALL_METRICS_SET.contains(metricName)
1✔
381
        && !disableDefault;
382
  }
383

384
  /**
385
   * Internal interface to avoid storing a {@link java.util.function.Predicate} directly, ensuring
386
   * compatibility with Android devices (API level < 24) that do not use library desugaring.
387
   */
388
  interface TargetFilter {
389
    boolean test(String target);
390
  }
391

392
  /**
393
   * Builder for configuring {@link GrpcOpenTelemetry}.
394
   */
395
  public static class Builder {
396
    private OpenTelemetry openTelemetrySdk = OpenTelemetry.noop();
1✔
397
    private final List<OpenTelemetryPlugin> plugins = new ArrayList<>();
1✔
398
    private final Collection<String> optionalLabels = new ArrayList<>();
1✔
399
    private final Map<String, Boolean> enableMetrics = new HashMap<>();
1✔
400
    private boolean disableAll;
401
    @Nullable
402
    private TargetFilter targetFilter;
403

404
    private Builder() {}
1✔
405

406
    /**
407
     * Sets the {@link io.opentelemetry.api.OpenTelemetry} entrypoint to use. This can be used to
408
     * configure OpenTelemetry by returning the instance created by a
409
     * {@code io.opentelemetry.sdk.OpenTelemetrySdkBuilder}.
410
     */
411
    public Builder sdk(OpenTelemetry sdk) {
412
      this.openTelemetrySdk = sdk;
1✔
413
      return this;
1✔
414
    }
415

416
    Builder plugin(OpenTelemetryPlugin plugin) {
417
      plugins.add(checkNotNull(plugin, "plugin"));
1✔
418
      return this;
1✔
419
    }
420

421
    /**
422
     * Adds optionalLabelKey to all the metrics that can provide value for the
423
     * optionalLabelKey.
424
     */
425
    public Builder addOptionalLabel(String optionalLabelKey) {
426
      this.optionalLabels.add(optionalLabelKey);
1✔
427
      return this;
1✔
428
    }
429

430
    /**
431
     * Enables the specified metrics for collection and export. By default, only a subset of
432
     * metrics are enabled.
433
     */
434
    public Builder enableMetrics(Collection<String> enableMetrics) {
435
      for (String metric : enableMetrics) {
1✔
436
        this.enableMetrics.put(metric, true);
1✔
437
      }
1✔
438
      return this;
1✔
439
    }
440

441
    /**
442
     * Disables the specified metrics from being collected and exported.
443
     */
444
    public Builder disableMetrics(Collection<String> disableMetrics) {
445
      for (String metric : disableMetrics) {
1✔
446
        this.enableMetrics.put(metric, false);
1✔
447
      }
1✔
448
      return this;
1✔
449
    }
450

451
    /**
452
     * Disable all metrics. If set to true all metrics must be explicitly enabled.
453
     */
454
    public Builder disableAllMetrics() {
455
      this.enableMetrics.clear();
1✔
456
      this.disableAll = true;
1✔
457
      return this;
1✔
458
    }
459

460
    Builder enableTracing(boolean enable) {
461
      ENABLE_OTEL_TRACING = enable;
1✔
462
      return this;
1✔
463
    }
464

465
    /**
466
     * Sets an optional filter to control recording of the {@code grpc.target} metric
467
     * attribute.
468
     *
469
     * <p>If the predicate returns {@code true}, the original target is recorded. Otherwise,
470
     * the target is recorded as {@code "other"} to limit metric cardinality.
471
     *
472
     * <p>If unset, all targets are recorded as-is.
473
     */
474
    @ExperimentalApi("https://github.com/grpc/grpc-java/issues/12595")
475
    @IgnoreJRERequirement
476
    public Builder targetAttributeFilter(@Nullable Predicate<String> filter) {
477
      if (filter == null) {
1✔
478
        this.targetFilter = null;
×
479
      } else {
480
        this.targetFilter = filter::test;
1✔
481
      }
482
      return this;
1✔
483
    }
484

485
    /**
486
     * Returns a new {@link GrpcOpenTelemetry} built with the configuration of this {@link
487
     * Builder}.
488
     */
489
    public GrpcOpenTelemetry build() {
490
      return new GrpcOpenTelemetry(this);
1✔
491
    }
492
  }
493
}
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