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

grpc / grpc-java / #20414

20 Aug 2026 06:02AM UTC coverage: 89.192% (-0.008%) from 89.2%
#20414

push

github

web-flow
netty: Support never-indexed metadata keys (#12976)

Add NettyChannelBuilder.neverIndexMetadataKey() and
neverIndexMetadataKeys() so callers can mark selected outbound metadata
keys for HPACK's never-indexed literal representation.

High-cardinality metadata values provide little compression benefit and
can churn the server's dynamic HPACK table. Keeping them out of the
table avoids unnecessary insertion and eviction work while preserving
dynamic indexing for other headers.

Propagate an immutable set of normalized metadata names through the
client transport and use it in Netty's HPACK sensitivity detector. Add
unit and interoperability coverage.

Generated with AI using OpenAI Codex (GPT-5).

Co-authored-by: Codex <noreply@openai.com>

38580 of 43255 relevant lines covered (89.19%)

0.89 hits per line

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

94.87
/../netty/src/main/java/io/grpc/netty/NettyClientTransport.java
1
/*
2
 * Copyright 2014 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.netty;
18

19
import static io.grpc.internal.GrpcUtil.KEEPALIVE_TIME_NANOS_DISABLED;
20
import static io.netty.channel.ChannelOption.ALLOCATOR;
21
import static io.netty.channel.ChannelOption.SO_KEEPALIVE;
22

23
import com.google.common.annotations.VisibleForTesting;
24
import com.google.common.base.MoreObjects;
25
import com.google.common.base.Preconditions;
26
import com.google.common.base.Ticker;
27
import com.google.common.util.concurrent.ListenableFuture;
28
import com.google.common.util.concurrent.SettableFuture;
29
import io.grpc.Attributes;
30
import io.grpc.CallOptions;
31
import io.grpc.ChannelLogger;
32
import io.grpc.ClientStreamTracer;
33
import io.grpc.InternalChannelz.SocketStats;
34
import io.grpc.InternalLogId;
35
import io.grpc.Metadata;
36
import io.grpc.MethodDescriptor;
37
import io.grpc.MetricRecorder;
38
import io.grpc.Status;
39
import io.grpc.internal.ClientStream;
40
import io.grpc.internal.ConnectionClientTransport;
41
import io.grpc.internal.DisconnectError;
42
import io.grpc.internal.FailingClientStream;
43
import io.grpc.internal.GrpcUtil;
44
import io.grpc.internal.Http2Ping;
45
import io.grpc.internal.KeepAliveManager;
46
import io.grpc.internal.KeepAliveManager.ClientKeepAlivePinger;
47
import io.grpc.internal.SimpleDisconnectError;
48
import io.grpc.internal.StatsTraceContext;
49
import io.grpc.internal.TransportTracer;
50
import io.grpc.netty.NettyChannelBuilder.LocalSocketPicker;
51
import io.netty.bootstrap.Bootstrap;
52
import io.netty.channel.Channel;
53
import io.netty.channel.ChannelFactory;
54
import io.netty.channel.ChannelFuture;
55
import io.netty.channel.ChannelFutureListener;
56
import io.netty.channel.ChannelHandler;
57
import io.netty.channel.ChannelOption;
58
import io.netty.channel.EventLoop;
59
import io.netty.channel.EventLoopGroup;
60
import io.netty.handler.codec.http2.StreamBufferingEncoder.Http2ChannelClosedException;
61
import io.netty.util.AsciiString;
62
import io.netty.util.concurrent.Future;
63
import io.netty.util.concurrent.GenericFutureListener;
64
import java.net.SocketAddress;
65
import java.nio.channels.ClosedChannelException;
66
import java.util.Map;
67
import java.util.Set;
68
import java.util.concurrent.Executor;
69
import java.util.concurrent.TimeUnit;
70
import javax.annotation.Nullable;
71

72
/**
73
 * A Netty-based {@link ConnectionClientTransport} implementation.
74
 */
75
class NettyClientTransport implements ConnectionClientTransport,
1 ✔
76
    ClientKeepAlivePinger.TransportWithDisconnectReason {
77

78
  private final InternalLogId logId;
79
  private final Map<ChannelOption<?>, ?> channelOptions;
80
  private final SocketAddress remoteAddress;
81
  private final ChannelFactory<? extends Channel> channelFactory;
82
  private final EventLoopGroup group;
83
  private final ProtocolNegotiator negotiator;
84
  private final String authorityString;
85
  private final AsciiString authority;
86
  private final AsciiString userAgent;
87
  private final boolean autoFlowControl;
88
  private final int flowControlWindow;
89
  private final Set<AsciiString> neverIndexedMetadataKeys;
90
  private final int maxMessageSize;
91
  private final int maxHeaderListSize;
92
  private final int softLimitHeaderListSize;
93
  private KeepAliveManager keepAliveManager;
94
  private final long keepAliveTimeNanos;
95
  private final long keepAliveTimeoutNanos;
96
  private final boolean keepAliveWithoutCalls;
97
  private final AsciiString negotiationScheme;
98
  private final Runnable tooManyPingsRunnable;
99
  private NettyClientHandler handler;
100
  // We should not send on the channel until negotiation completes. This is a hard requirement
101
  // by SslHandler but is appropriate for HTTP/1.1 Upgrade as well.
102
  private Channel channel;
103
  /** If {@link #start} has been called, non-{@code null} if channel is {@code null}. */
104
  private Status statusExplainingWhyTheChannelIsNull;
105
  /** Since not thread-safe, may only be used from event loop. */
106
  private ClientTransportLifecycleManager lifecycleManager;
107
  /** Since not thread-safe, may only be used from event loop. */
108
  private final TransportTracer transportTracer;
109
  private final Attributes eagAttributes;
110
  private final LocalSocketPicker localSocketPicker;
111
  private final ChannelLogger channelLogger;
112
  private final boolean useGetForSafeMethods;
113
  private final Ticker ticker;
114
  private final MetricRecorder metricRecorder;
115

116

117
  NettyClientTransport(
118
      SocketAddress address,
119
      ChannelFactory<? extends Channel> channelFactory,
120
      Map<ChannelOption<?>, ?> channelOptions,
121
      EventLoopGroup group,
122
      ProtocolNegotiator negotiator,
123
      boolean autoFlowControl,
124
      int flowControlWindow,
125
      Set<AsciiString> neverIndexedMetadataKeys,
126
      int maxMessageSize,
127
      int maxHeaderListSize,
128
      int softLimitHeaderListSize,
129
      long keepAliveTimeNanos,
130
      long keepAliveTimeoutNanos,
131
      boolean keepAliveWithoutCalls,
132
      String authority,
133
      @Nullable String userAgent,
134
      Runnable tooManyPingsRunnable,
135
      TransportTracer transportTracer,
136
      Attributes eagAttributes,
137
      LocalSocketPicker localSocketPicker,
138
      ChannelLogger channelLogger,
139
      boolean useGetForSafeMethods,
140
      MetricRecorder metricRecorder,
141
      Ticker ticker) {
1 ✔
142

143
    this.negotiator = Preconditions.checkNotNull(negotiator, "negotiator");
1 ✔
144
    this.negotiationScheme = this.negotiator.scheme();
1 ✔
145
    this.remoteAddress = Preconditions.checkNotNull(address, "address");
1 ✔
146
    this.group = Preconditions.checkNotNull(group, "group");
1 ✔
147
    this.channelFactory = channelFactory;
1 ✔
148
    this.channelOptions = Preconditions.checkNotNull(channelOptions, "channelOptions");
1 ✔
149
    this.autoFlowControl = autoFlowControl;
1 ✔
150
    this.flowControlWindow = flowControlWindow;
1 ✔
151
    this.neverIndexedMetadataKeys =
1 ✔
152
        Preconditions.checkNotNull(neverIndexedMetadataKeys, "neverIndexedMetadataKeys");
1 ✔
153
    this.maxMessageSize = maxMessageSize;
1 ✔
154
    this.maxHeaderListSize = maxHeaderListSize;
1 ✔
155
    this.softLimitHeaderListSize = softLimitHeaderListSize;
1 ✔
156
    this.keepAliveTimeNanos = keepAliveTimeNanos;
1 ✔
157
    this.keepAliveTimeoutNanos = keepAliveTimeoutNanos;
1 ✔
158
    this.keepAliveWithoutCalls = keepAliveWithoutCalls;
1 ✔
159
    this.authorityString = authority;
1 ✔
160
    this.authority = new AsciiString(authority);
1 ✔
161
    this.userAgent = new AsciiString(GrpcUtil.getGrpcUserAgent("netty", userAgent));
1 ✔
162
    this.tooManyPingsRunnable =
1 ✔
163
        Preconditions.checkNotNull(tooManyPingsRunnable, "tooManyPingsRunnable");
1 ✔
164
    this.transportTracer = Preconditions.checkNotNull(transportTracer, "transportTracer");
1 ✔
165
    this.eagAttributes = Preconditions.checkNotNull(eagAttributes, "eagAttributes");
1 ✔
166
    this.localSocketPicker = Preconditions.checkNotNull(localSocketPicker, "localSocketPicker");
1 ✔
167
    this.logId = InternalLogId.allocate(getClass(), remoteAddress.toString());
1 ✔
168
    this.channelLogger = Preconditions.checkNotNull(channelLogger, "channelLogger");
1 ✔
169
    this.useGetForSafeMethods = useGetForSafeMethods;
1 ✔
170
    this.metricRecorder = metricRecorder;
1 ✔
171
    this.ticker = Preconditions.checkNotNull(ticker, "ticker");
1 ✔
172
  }
1 ✔
173

174
  @Override
175
  public void ping(final PingCallback callback, final Executor executor) {
176
    if (channel == null) {
1 ✔
177
      executor.execute(new Runnable() {
1 ✔
178
        @Override
179
        public void run() {
180
          callback.onFailure(statusExplainingWhyTheChannelIsNull);
1 ✔
181
        }
1 ✔
182
      });
183
      return;
1 ✔
184
    }
185
    // The promise and listener always succeed in NettyClientHandler. So this listener handles the
186
    // error case, when the channel is closed and the NettyClientHandler no longer in the pipeline.
187
    ChannelFutureListener failureListener = new ChannelFutureListener() {
1 ✔
188
      @Override
189
      public void operationComplete(ChannelFuture future) throws Exception {
190
        if (!future.isSuccess()) {
1 ✔
191
          Status s = statusFromFailedFuture(future);
1 ✔
192
          Http2Ping.notifyFailed(callback, executor, s);
1 ✔
193
        }
194
      }
1 ✔
195
    };
196
    // Write the command requesting the ping
197
    handler.getWriteQueue().enqueue(new SendPingCommand(callback, executor), true)
1 ✔
198
        .addListener(failureListener);
1 ✔
199
  }
1 ✔
200

201
  @Override
202
  public ClientStream newStream(
203
      MethodDescriptor<?, ?> method, Metadata headers, CallOptions callOptions,
204
      ClientStreamTracer[] tracers) {
205
    Preconditions.checkNotNull(method, "method");
1 ✔
206
    Preconditions.checkNotNull(headers, "headers");
1 ✔
207
    if (channel == null) {
1 ✔
208
      return new FailingClientStream(statusExplainingWhyTheChannelIsNull, tracers);
1 ✔
209
    }
210
    StatsTraceContext statsTraceCtx =
1 ✔
211
        StatsTraceContext.newClientContext(tracers, getAttributes(), headers);
1 ✔
212
    return new NettyClientStream(
1 ✔
213
        new NettyClientStream.TransportState(
214
            handler,
215
            channel.eventLoop(),
1 ✔
216
            maxMessageSize,
217
            statsTraceCtx,
218
            transportTracer,
219
            method.getFullMethodName(),
1 ✔
220
            callOptions) {
1 ✔
221
          @Override
222
          protected Status statusFromFailedFuture(ChannelFuture f) {
223
            return NettyClientTransport.this.statusFromFailedFuture(f);
1 ✔
224
          }
225
        },
226
        method,
227
        headers,
228
        channel,
229
        authority,
230
        negotiationScheme,
231
        userAgent,
232
        statsTraceCtx,
233
        transportTracer,
234
        callOptions,
235
        useGetForSafeMethods);
236
  }
237

238
  @SuppressWarnings("unchecked")
239
  @Override
240
  public Runnable start(Listener transportListener) {
241
    lifecycleManager = new ClientTransportLifecycleManager(
1 ✔
242
        Preconditions.checkNotNull(transportListener, "listener"));
1 ✔
243
    EventLoop eventLoop = group.next();
1 ✔
244
    if (keepAliveTimeNanos != KEEPALIVE_TIME_NANOS_DISABLED) {
1 ✔
245
      keepAliveManager = new KeepAliveManager(
1 ✔
246
          new ClientKeepAlivePinger(this), eventLoop, keepAliveTimeNanos,
247
          keepAliveTimeoutNanos, keepAliveWithoutCalls);
248
    }
249

250
    handler = NettyClientHandler.newHandler(
1 ✔
251
            lifecycleManager,
252
            keepAliveManager,
253
            autoFlowControl,
254
            flowControlWindow,
255
            neverIndexedMetadataKeys,
256
            maxHeaderListSize,
257
            softLimitHeaderListSize,
258
            GrpcUtil.STOPWATCH_SUPPLIER,
259
            tooManyPingsRunnable,
260
            transportTracer,
261
            eagAttributes,
262
            authorityString,
263
            channelLogger,
264
            ticker,
265
            metricRecorder);
266

267
    ChannelHandler negotiationHandler = negotiator.newHandler(handler);
1 ✔
268

269
    Bootstrap b = new Bootstrap();
1 ✔
270
    b.option(ALLOCATOR, Utils.getByteBufAllocator(false));
1 ✔
271
    b.group(eventLoop);
1 ✔
272
    b.channelFactory(channelFactory);
1 ✔
273
    // For non-socket based channel, the option will be ignored.
274
    b.option(SO_KEEPALIVE, true);
1 ✔
275
    for (Map.Entry<ChannelOption<?>, ?> entry : channelOptions.entrySet()) {
1 ✔
276
      // Every entry in the map is obtained from
277
      // NettyChannelBuilder#withOption(ChannelOption<T> option, T value)
278
      // so it is safe to pass the key-value pair to b.option().
279
      b.option((ChannelOption<Object>) entry.getKey(), entry.getValue());
1 ✔
280
    }
1 ✔
281

282
    ChannelHandler bufferingHandler = new WriteBufferingAndExceptionHandler(negotiationHandler);
1 ✔
283

284
    /*
285
     * We don't use a ChannelInitializer in the client bootstrap because its "initChannel" method
286
     * is executed in the event loop and we need this handler to be in the pipeline immediately so
287
     * that it may begin buffering writes.
288
     */
289
    b.handler(bufferingHandler);
1 ✔
290
    ChannelFuture regFuture = b.register();
1 ✔
291
    if (regFuture.isDone() && !regFuture.isSuccess()) {
1 ✔
292
      channel = null;
1 ✔
293
      // Initialization has failed badly. All new streams should be made to fail.
294
      Throwable t = regFuture.cause();
1 ✔
295
      if (t == null) {
1 ✔
296
        t = new IllegalStateException("Channel is null, but future doesn't have a cause");
×
297
      }
298
      statusExplainingWhyTheChannelIsNull = Utils.statusFromThrowable(t);
1 ✔
299
      // Use a Runnable since lifecycleManager calls transportListener
300
      return new Runnable() {
1 ✔
301
        @Override
302
        public void run() {
303
          // NOTICE: we not are calling lifecycleManager from the event loop. But there isn't really
304
          // an event loop in this case, so nothing should be accessing the lifecycleManager. We
305
          // could use GlobalEventExecutor (which is what regFuture would use for notifying
306
          // listeners in this case), but avoiding on-demand thread creation in an error case seems
307
          // a good idea and is probably clearer threading.
308
          lifecycleManager.notifyTerminated(statusExplainingWhyTheChannelIsNull,
1 ✔
309
              SimpleDisconnectError.UNKNOWN);
310
        }
1 ✔
311
      };
312
    }
313
    channel = regFuture.channel();
1 ✔
314
    // For non-epoll based channel, the option will be ignored.
315
    try {
316
      if (keepAliveTimeNanos != KEEPALIVE_TIME_NANOS_DISABLED
1 ✔
317
              && Class.forName("io.netty.channel.epoll.AbstractEpollChannel").isInstance(channel)) {
1 ✔
318
        ChannelOption<Integer> tcpUserTimeout = Utils.maybeGetTcpUserTimeoutOption();
1 ✔
319
        if (tcpUserTimeout != null) {
1 ✔
320
          int tcpUserTimeoutMs = (int) TimeUnit.NANOSECONDS.toMillis(keepAliveTimeoutNanos);
1 ✔
321
          channel.config().setOption(tcpUserTimeout, tcpUserTimeoutMs);
1 ✔
322
        }
323
      }
324
    } catch (ClassNotFoundException ignored) {
×
325
      // JVM did not load AbstractEpollChannel, so the current channel will not be of epoll type,
326
      // so there is no need to set TCP_USER_TIMEOUT
327
    }
1 ✔
328
    // Start the write queue as soon as the channel is constructed
329
    handler.startWriteQueue(channel);
1 ✔
330
    // This write will have no effect, yet it will only complete once the negotiationHandler
331
    // flushes any pending writes. We need it to be staged *before* the `connect` so that
332
    // the channel can't have been closed yet, removing all handlers. This write will sit in the
333
    // AbstractBufferingHandler's buffer, and will either be flushed on a successful connection,
334
    // or failed if the connection fails.
335
    channel.writeAndFlush(NettyClientHandler.NOOP_MESSAGE).addListener(new ChannelFutureListener() {
1 ✔
336
      @Override
337
      public void operationComplete(ChannelFuture future) throws Exception {
338
        if (!future.isSuccess()) {
1 ✔
339
          // Need to notify of this failure, because NettyClientHandler may not have been added to
340
          // the pipeline before the error occurred.
341
          lifecycleManager.notifyTerminated(Utils.statusFromThrowable(future.cause()),
1 ✔
342
              SimpleDisconnectError.UNKNOWN);
343
        }
344
      }
1 ✔
345
    });
346
    // Start the connection operation to the server.
347
    SocketAddress localAddress =
1 ✔
348
        localSocketPicker.createSocketAddress(remoteAddress, eagAttributes);
1 ✔
349
    if (localAddress != null) {
1 ✔
350
      channel.connect(remoteAddress, localAddress);
×
351
    } else {
352
      channel.connect(remoteAddress);
1 ✔
353
    }
354

355
    if (keepAliveManager != null) {
1 ✔
356
      keepAliveManager.onTransportStarted();
1 ✔
357
    }
358

359
    return null;
1 ✔
360
  }
361

362
  @Override
363
  public void shutdown(Status reason) {
364
    // start() could have failed
365
    if (channel == null) {
1 ✔
366
      return;
1 ✔
367
    }
368
    // Notifying of termination is automatically done when the channel closes.
369
    if (channel.isOpen()) {
1 ✔
370
      handler.getWriteQueue().enqueue(new GracefulCloseCommand(reason), true);
1 ✔
371
    }
372
  }
1 ✔
373

374
  @Override
375
  public void shutdownNow(final Status reason) {
376
    shutdownNow(reason, SimpleDisconnectError.SUBCHANNEL_SHUTDOWN);
1 ✔
377
  }
1 ✔
378

379
  @Override
380
  public void shutdownNow(final Status reason, DisconnectError disconnectError) {
381
    // Notifying of termination is automatically done when the channel closes.
382
    if (channel != null && channel.isOpen()) {
1 ✔
383
      handler.getWriteQueue().enqueue(new Runnable() {
1 ✔
384
        @Override
385
        public void run() {
386
          lifecycleManager.notifyShutdown(reason, disconnectError);
1 ✔
387
          channel.write(new ForcefulCloseCommand(reason));
1 ✔
388
        }
1 ✔
389
      }, true);
390
    }
391
  }
1 ✔
392

393
  @Override
394
  public String toString() {
395
    return MoreObjects.toStringHelper(this)
1 ✔
396
        .add("logId", logId.getId())
1 ✔
397
        .add("remoteAddress", remoteAddress)
1 ✔
398
        .add("channel", channel)
1 ✔
399
        .toString();
1 ✔
400
  }
401

402
  @Override
403
  public InternalLogId getLogId() {
404
    return logId;
1 ✔
405
  }
406

407
  @Override
408
  public Attributes getAttributes() {
409
    return handler.getAttributes();
1 ✔
410
  }
411

412
  @Override
413
  public ListenableFuture<SocketStats> getStats() {
414
    final SettableFuture<SocketStats> result = SettableFuture.create();
1 ✔
415
    if (channel.eventLoop().inEventLoop()) {
1 ✔
416
      // This is necessary, otherwise we will block forever if we get the future from inside
417
      // the event loop.
418
      result.set(getStatsHelper(channel));
×
419
      return result;
×
420
    }
421
    channel.eventLoop().submit(
1 ✔
422
        new Runnable() {
1 ✔
423
          @Override
424
          public void run() {
425
            result.set(getStatsHelper(channel));
1 ✔
426
          }
1 ✔
427
        })
428
        .addListener(
1 ✔
429
            new GenericFutureListener<Future<Object>>() {
1 ✔
430
              @Override
431
              public void operationComplete(Future<Object> future) throws Exception {
432
                if (!future.isSuccess()) {
1 ✔
433
                  result.setException(future.cause());
×
434
                }
435
              }
1 ✔
436
            });
437
    return result;
1 ✔
438
  }
439

440
  private SocketStats getStatsHelper(Channel ch) {
441
    assert ch.eventLoop().inEventLoop();
1 ✔
442
    return new SocketStats(
1 ✔
443
        transportTracer.getStats(),
1 ✔
444
        channel.localAddress(),
1 ✔
445
        channel.remoteAddress(),
1 ✔
446
        Utils.getSocketOptions(ch),
1 ✔
447
        handler == null ? null : handler.getSecurityInfo());
1 ✔
448
  }
449

450
  @VisibleForTesting
451
  Channel channel() {
452
    return channel;
1 ✔
453
  }
454

455
  @VisibleForTesting
456
  KeepAliveManager keepAliveManager() {
457
    return keepAliveManager;
1 ✔
458
  }
459

460
  /**
461
   * Convert ChannelFuture.cause() to a Status, taking into account that all handlers are removed
462
   * from the pipeline when the channel is closed. Since handlers are removed, you may get an
463
   * unhelpful exception like ClosedChannelException.
464
   *
465
   * <p>This method must only be called on the event loop.
466
   */
467
  private Status statusFromFailedFuture(ChannelFuture f) {
468
    Throwable t = f.cause();
1 ✔
469
    if (t instanceof ClosedChannelException
1 ✔
470
        // Exception thrown by the StreamBufferingEncoder if the channel is closed while there
471
        // are still streams buffered. This exception is not helpful. Replace it by the real
472
        // cause of the shutdown (if available).
473
        || t instanceof Http2ChannelClosedException) {
474
      Status shutdownStatus = lifecycleManager.getShutdownStatus();
1 ✔
475
      if (shutdownStatus == null) {
1 ✔
476
        return Status.UNKNOWN.withDescription("Channel closed but for unknown reason")
×
477
            .withCause(new ClosedChannelException().initCause(t));
×
478
      }
479
      return shutdownStatus;
1 ✔
480
    }
481
    return Utils.statusFromThrowable(t);
1 ✔
482
  }
483
}
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