• 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

93.51
/../netty/src/main/java/io/grpc/netty/NettyClientHandler.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.netty.handler.codec.http2.DefaultHttp2LocalFlowController.DEFAULT_WINDOW_UPDATE_RATIO;
20
import static io.netty.util.CharsetUtil.UTF_8;
21
import static io.netty.util.internal.ObjectUtil.checkNotNull;
22

23
import com.google.common.annotations.VisibleForTesting;
24
import com.google.common.base.Preconditions;
25
import com.google.common.base.Stopwatch;
26
import com.google.common.base.Supplier;
27
import com.google.common.base.Ticker;
28
import io.grpc.Attributes;
29
import io.grpc.ChannelLogger;
30
import io.grpc.InternalChannelz;
31
import io.grpc.InternalStatus;
32
import io.grpc.Metadata;
33
import io.grpc.MetricRecorder;
34
import io.grpc.Status;
35
import io.grpc.StatusException;
36
import io.grpc.internal.ClientStreamListener.RpcProgress;
37
import io.grpc.internal.ClientTransport.PingCallback;
38
import io.grpc.internal.DisconnectError;
39
import io.grpc.internal.GoAwayDisconnectError;
40
import io.grpc.internal.GrpcAttributes;
41
import io.grpc.internal.GrpcUtil;
42
import io.grpc.internal.Http2Ping;
43
import io.grpc.internal.InUseStateAggregator;
44
import io.grpc.internal.KeepAliveManager;
45
import io.grpc.internal.SimpleDisconnectError;
46
import io.grpc.internal.TransportTracer;
47
import io.grpc.netty.GrpcHttp2HeadersUtils.GrpcHttp2ClientHeadersDecoder;
48
import io.netty.buffer.ByteBuf;
49
import io.netty.buffer.ByteBufUtil;
50
import io.netty.buffer.Unpooled;
51
import io.netty.channel.Channel;
52
import io.netty.channel.ChannelFuture;
53
import io.netty.channel.ChannelFutureListener;
54
import io.netty.channel.ChannelHandlerContext;
55
import io.netty.channel.ChannelPromise;
56
import io.netty.handler.codec.http2.DecoratingHttp2FrameWriter;
57
import io.netty.handler.codec.http2.DefaultHttp2Connection;
58
import io.netty.handler.codec.http2.DefaultHttp2ConnectionDecoder;
59
import io.netty.handler.codec.http2.DefaultHttp2ConnectionEncoder;
60
import io.netty.handler.codec.http2.DefaultHttp2FrameReader;
61
import io.netty.handler.codec.http2.DefaultHttp2FrameWriter;
62
import io.netty.handler.codec.http2.DefaultHttp2HeadersEncoder;
63
import io.netty.handler.codec.http2.DefaultHttp2LocalFlowController;
64
import io.netty.handler.codec.http2.DefaultHttp2RemoteFlowController;
65
import io.netty.handler.codec.http2.Http2CodecUtil;
66
import io.netty.handler.codec.http2.Http2Connection;
67
import io.netty.handler.codec.http2.Http2ConnectionAdapter;
68
import io.netty.handler.codec.http2.Http2ConnectionDecoder;
69
import io.netty.handler.codec.http2.Http2ConnectionEncoder;
70
import io.netty.handler.codec.http2.Http2Error;
71
import io.netty.handler.codec.http2.Http2Exception;
72
import io.netty.handler.codec.http2.Http2FrameAdapter;
73
import io.netty.handler.codec.http2.Http2FrameLogger;
74
import io.netty.handler.codec.http2.Http2FrameReader;
75
import io.netty.handler.codec.http2.Http2FrameWriter;
76
import io.netty.handler.codec.http2.Http2Headers;
77
import io.netty.handler.codec.http2.Http2HeadersDecoder;
78
import io.netty.handler.codec.http2.Http2HeadersEncoder;
79
import io.netty.handler.codec.http2.Http2HeadersEncoder.SensitivityDetector;
80
import io.netty.handler.codec.http2.Http2InboundFrameLogger;
81
import io.netty.handler.codec.http2.Http2OutboundFrameLogger;
82
import io.netty.handler.codec.http2.Http2Settings;
83
import io.netty.handler.codec.http2.Http2Stream;
84
import io.netty.handler.codec.http2.Http2StreamVisitor;
85
import io.netty.handler.codec.http2.StreamBufferingEncoder;
86
import io.netty.handler.codec.http2.UniformStreamByteDistributor;
87
import io.netty.handler.logging.LogLevel;
88
import io.netty.util.AsciiString;
89
import io.perfmark.PerfMark;
90
import io.perfmark.Tag;
91
import io.perfmark.TaskCloseable;
92
import java.nio.channels.ClosedChannelException;
93
import java.util.LinkedHashMap;
94
import java.util.Map;
95
import java.util.Set;
96
import java.util.concurrent.Executor;
97
import java.util.logging.Level;
98
import java.util.logging.Logger;
99
import javax.annotation.Nullable;
100

101
/**
102
 * Client-side Netty handler for GRPC processing. All event handlers are executed entirely within
103
 * the context of the Netty Channel thread.
104
 */
105
class NettyClientHandler extends AbstractNettyHandler {
106
  private static final Logger logger = Logger.getLogger(NettyClientHandler.class.getName());
1 ✔
107
  static boolean enablePerRpcAuthorityCheck =
1 ✔
108
      GrpcUtil.getFlag("GRPC_ENABLE_PER_RPC_AUTHORITY_CHECK", false);
1 ✔
109

110
  /**
111
   * A message that simply passes through the channel without any real processing. It is useful to
112
   * check if buffers have been drained and test the health of the channel in a single operation.
113
   */
114
  static final Object NOOP_MESSAGE = new Object();
1 ✔
115

116
  /**
117
   * Status used when the transport has exhausted the number of streams.
118
   */
119
  private static final Status EXHAUSTED_STREAMS_STATUS =
1 ✔
120
          Status.UNAVAILABLE.withDescription("Stream IDs have been exhausted");
1 ✔
121
  private static final long USER_PING_PAYLOAD = 1111;
122

123
  private final Http2Connection.PropertyKey streamKey;
124
  private final ClientTransportLifecycleManager lifecycleManager;
125
  private final KeepAliveManager keepAliveManager;
126
  // Returns new unstarted stopwatches
127
  private final Supplier<Stopwatch> stopwatchFactory;
128
  private final TransportTracer transportTracer;
129
  private final Attributes eagAttributes;
130
  private final TcpMetrics tcpMetrics;
131
  private final String authority;
132
  private final InUseStateAggregator<Http2Stream> inUseState =
1 ✔
133
      new InUseStateAggregator<Http2Stream>() {
1 ✔
134
        @Override
135
        protected void handleInUse() {
136
          lifecycleManager.notifyInUse(true);
1 ✔
137
        }
1 ✔
138

139
        @Override
140
        protected void handleNotInUse() {
141
          lifecycleManager.notifyInUse(false);
1 ✔
142
        }
1 ✔
143
      };
144
  private final Map<String, Status> peerVerificationResults =
1 ✔
145
      new LinkedHashMap<String, Status>() {
1 ✔
146
        @Override
147
        protected boolean removeEldestEntry(Map.Entry<String, Status> eldest) {
148
          return size() > 100;
1 ✔
149
        }
150
      };
151

152
  private WriteQueue clientWriteQueue;
153
  private Http2Ping ping;
154
  private Attributes attributes;
155
  private InternalChannelz.Security securityInfo;
156
  private Status abruptGoAwayStatus;
157
  private Status channelInactiveReason;
158

159
  static NettyClientHandler newHandler(
160
      ClientTransportLifecycleManager lifecycleManager,
161
      @Nullable KeepAliveManager keepAliveManager,
162
      boolean autoFlowControl,
163
      int flowControlWindow,
164
      Set<AsciiString> neverIndexedMetadataKeys,
165
      int maxHeaderListSize,
166
      int softLimitHeaderListSize,
167
      Supplier<Stopwatch> stopwatchFactory,
168
      Runnable tooManyPingsRunnable,
169
      TransportTracer transportTracer,
170
      Attributes eagAttributes,
171
      String authority,
172
      ChannelLogger negotiationLogger,
173
      Ticker ticker,
174
      MetricRecorder metricRecorder) {
175
    Preconditions.checkArgument(maxHeaderListSize > 0, "maxHeaderListSize must be positive");
1 ✔
176
    Http2HeadersDecoder headersDecoder = new GrpcHttp2ClientHeadersDecoder(maxHeaderListSize);
1 ✔
177
    Http2FrameReader frameReader = new DefaultHttp2FrameReader(headersDecoder);
1 ✔
178
    Http2HeadersEncoder encoder = new DefaultHttp2HeadersEncoder(
1 ✔
179
        sensitivityDetector(neverIndexedMetadataKeys), false, 16, Integer.MAX_VALUE);
1 ✔
180
    Http2FrameWriter frameWriter = new DefaultHttp2FrameWriter(encoder);
1 ✔
181
    Http2Connection connection = new DefaultHttp2Connection(false);
1 ✔
182
    UniformStreamByteDistributor dist = new UniformStreamByteDistributor(connection);
1 ✔
183
    dist.minAllocationChunk(MIN_ALLOCATED_CHUNK); // Increased for benchmarks performance.
1 ✔
184
    DefaultHttp2RemoteFlowController controller =
1 ✔
185
        new DefaultHttp2RemoteFlowController(connection, dist);
186
    connection.remote().flowController(controller);
1 ✔
187

188
    return newHandler(
1 ✔
189
        connection,
190
        frameReader,
191
        frameWriter,
192
        lifecycleManager,
193
        keepAliveManager,
194
        autoFlowControl,
195
        flowControlWindow,
196
        maxHeaderListSize,
197
        softLimitHeaderListSize,
198
        stopwatchFactory,
199
        tooManyPingsRunnable,
200
        transportTracer,
201
        eagAttributes,
202
        authority,
203
        negotiationLogger,
204
        ticker,
205
        metricRecorder);
206
  }
207

208
  @VisibleForTesting
209
  static NettyClientHandler newHandler(
210
      final Http2Connection connection,
211
      Http2FrameReader frameReader,
212
      Http2FrameWriter frameWriter,
213
      ClientTransportLifecycleManager lifecycleManager,
214
      KeepAliveManager keepAliveManager,
215
      boolean autoFlowControl,
216
      int flowControlWindow,
217
      int maxHeaderListSize,
218
      int softLimitHeaderListSize,
219
      Supplier<Stopwatch> stopwatchFactory,
220
      Runnable tooManyPingsRunnable,
221
      TransportTracer transportTracer,
222
      Attributes eagAttributes,
223
      String authority,
224
      ChannelLogger negotiationLogger,
225
      Ticker ticker,
226
      MetricRecorder metricRecorder) {
227
    Preconditions.checkNotNull(connection, "connection");
1 ✔
228
    Preconditions.checkNotNull(frameReader, "frameReader");
1 ✔
229
    Preconditions.checkNotNull(lifecycleManager, "lifecycleManager");
1 ✔
230
    Preconditions.checkArgument(flowControlWindow > 0, "flowControlWindow must be positive");
1 ✔
231
    Preconditions.checkArgument(maxHeaderListSize > 0, "maxHeaderListSize must be positive");
1 ✔
232
    Preconditions.checkArgument(softLimitHeaderListSize > 0,
1 ✔
233
        "softLimitHeaderListSize must be positive");
234
    Preconditions.checkNotNull(stopwatchFactory, "stopwatchFactory");
1 ✔
235
    Preconditions.checkNotNull(tooManyPingsRunnable, "tooManyPingsRunnable");
1 ✔
236
    Preconditions.checkNotNull(eagAttributes, "eagAttributes");
1 ✔
237
    Preconditions.checkNotNull(authority, "authority");
1 ✔
238

239
    Http2FrameLogger frameLogger = new Http2FrameLogger(LogLevel.DEBUG, NettyClientHandler.class);
1 ✔
240
    frameReader = new Http2InboundFrameLogger(frameReader, frameLogger);
1 ✔
241
    frameWriter = new Http2OutboundFrameLogger(frameWriter, frameLogger);
1 ✔
242

243
    PingCountingFrameWriter pingCounter;
244
    frameWriter = pingCounter = new PingCountingFrameWriter(frameWriter);
1 ✔
245

246
    StreamBufferingEncoder encoder =
1 ✔
247
        new StreamBufferingEncoder(
248
            new DefaultHttp2ConnectionEncoder(connection, frameWriter));
249

250
    // Create the local flow controller configured to auto-refill the connection window.
251
    connection.local().flowController(
1 ✔
252
        new DefaultHttp2LocalFlowController(connection, DEFAULT_WINDOW_UPDATE_RATIO, true));
253

254
    Http2ConnectionDecoder decoder = new DefaultHttp2ConnectionDecoder(connection, encoder,
1 ✔
255
        frameReader);
256

257
    transportTracer.setFlowControlWindowReader(new Utils.FlowControlReader(connection));
1 ✔
258

259
    Http2Settings settings = new Http2Settings();
1 ✔
260
    settings.pushEnabled(false);
1 ✔
261
    settings.initialWindowSize(flowControlWindow);
1 ✔
262
    settings.maxConcurrentStreams(0);
1 ✔
263
    settings.maxHeaderListSize(maxHeaderListSize);
1 ✔
264

265
    return new NettyClientHandler(
1 ✔
266
        decoder,
267
        encoder,
268
        settings,
269
        negotiationLogger,
270
        lifecycleManager,
271
        keepAliveManager,
272
        stopwatchFactory,
273
        tooManyPingsRunnable,
274
        transportTracer,
275
        eagAttributes,
276
        authority,
277
        autoFlowControl,
278
        pingCounter,
279
        ticker,
280
        maxHeaderListSize,
281
        softLimitHeaderListSize,
282
        metricRecorder);
283
  }
284

285
  @VisibleForTesting
286
  static SensitivityDetector sensitivityDetector(
287
      final Set<AsciiString> neverIndexedMetadataKeys) {
288
    if (neverIndexedMetadataKeys.isEmpty()) {
1 ✔
289
      return Http2HeadersEncoder.NEVER_SENSITIVE;
1 ✔
290
    }
291
    return new SensitivityDetector() {
1 ✔
292
      @Override
293
      public boolean isSensitive(CharSequence name, CharSequence value) {
294
        return neverIndexedMetadataKeys.contains(AsciiString.of(name));
1 ✔
295
      }
296
    };
297
  }
298

299
  private NettyClientHandler(
300
      Http2ConnectionDecoder decoder,
301
      Http2ConnectionEncoder encoder,
302
      Http2Settings settings,
303
      ChannelLogger negotiationLogger,
304
      ClientTransportLifecycleManager lifecycleManager,
305
      KeepAliveManager keepAliveManager,
306
      Supplier<Stopwatch> stopwatchFactory,
307
      final Runnable tooManyPingsRunnable,
308
      TransportTracer transportTracer,
309
      Attributes eagAttributes,
310
      String authority,
311
      boolean autoFlowControl,
312
      PingLimiter pingLimiter,
313
      Ticker ticker,
314
      int maxHeaderListSize,
315
      int softLimitHeaderListSize,
316
      MetricRecorder metricRecorder) {
317
    super(
1 ✔
318
        /* channelUnused= */ null,
319
        decoder,
320
        encoder,
321
        settings,
322
        negotiationLogger,
323
        autoFlowControl,
324
        pingLimiter,
325
        ticker,
326
        maxHeaderListSize,
327
        softLimitHeaderListSize);
328
    this.lifecycleManager = lifecycleManager;
1 ✔
329
    this.keepAliveManager = keepAliveManager;
1 ✔
330
    this.stopwatchFactory = stopwatchFactory;
1 ✔
331
    this.transportTracer = Preconditions.checkNotNull(transportTracer);
1 ✔
332
    this.eagAttributes = eagAttributes;
1 ✔
333
    this.authority = authority;
1 ✔
334
    this.attributes = Attributes.newBuilder()
1 ✔
335
        .set(GrpcAttributes.ATTR_CLIENT_EAG_ATTRS, eagAttributes).build();
1 ✔
336
    this.tcpMetrics = new TcpMetrics(metricRecorder);
1 ✔
337

338
    // Set the frame listener on the decoder.
339
    decoder().frameListener(new FrameListener());
1 ✔
340

341
    Http2Connection connection = encoder.connection();
1 ✔
342
    streamKey = connection.newKey();
1 ✔
343

344
    connection.addListener(new Http2ConnectionAdapter() {
1 ✔
345
      @Override
346
      public void onGoAwayReceived(int lastStreamId, long errorCode, ByteBuf debugData) {
347
        byte[] debugDataBytes = ByteBufUtil.getBytes(debugData);
1 ✔
348
        goingAway(errorCode, debugDataBytes);
1 ✔
349
        if (errorCode == Http2Error.ENHANCE_YOUR_CALM.code()) {
1 ✔
350
          String data = new String(debugDataBytes, UTF_8);
1 ✔
351
          logger.log(
1 ✔
352
              Level.WARNING, "Received GOAWAY with ENHANCE_YOUR_CALM. Debug data: {0}", data);
353
          if ("too_many_pings".equals(data)) {
1 ✔
354
            tooManyPingsRunnable.run();
1 ✔
355
          }
356
        }
357
      }
1 ✔
358

359
      @Override
360
      public void onStreamActive(Http2Stream stream) {
361
        if (connection().numActiveStreams() == 1
1 ✔
362
            && NettyClientHandler.this.keepAliveManager != null) {
1 ✔
363
          NettyClientHandler.this.keepAliveManager.onTransportActive();
1 ✔
364
        }
365
      }
1 ✔
366

367
      @Override
368
      public void onStreamClosed(Http2Stream stream) {
369
        // Although streams with CALL_OPTIONS_RPC_OWNED_BY_BALANCER are not marked as "in-use" in
370
        // the first place, we don't propagate that option here, and it's safe to reset the in-use
371
        // state for them, which will be a cheap no-op.
372
        inUseState.updateObjectInUse(stream, false);
1 ✔
373
        if (connection().numActiveStreams() == 0
1 ✔
374
            && NettyClientHandler.this.keepAliveManager != null) {
1 ✔
375
          NettyClientHandler.this.keepAliveManager.onTransportIdle();
1 ✔
376
        }
377
      }
1 ✔
378
    });
379
  }
1 ✔
380

381
  /**
382
   * The protocol negotiation attributes, available once the protocol negotiation completes;
383
   * otherwise returns {@code Attributes.EMPTY}.
384
   */
385
  Attributes getAttributes() {
386
    return attributes;
1 ✔
387
  }
388

389
  /**
390
   * Handler for commands sent from the stream.
391
   */
392
  @Override
393
  public void write(ChannelHandlerContext ctx, Object msg, ChannelPromise promise)
394
          throws Exception {
395
    if (msg instanceof CreateStreamCommand) {
1 ✔
396
      createStream((CreateStreamCommand) msg, promise);
1 ✔
397
    } else if (msg instanceof SendGrpcFrameCommand) {
1 ✔
398
      sendGrpcFrame(ctx, (SendGrpcFrameCommand) msg, promise);
1 ✔
399
    } else if (msg instanceof CancelClientStreamCommand) {
1 ✔
400
      cancelStream(ctx, (CancelClientStreamCommand) msg, promise);
1 ✔
401
    } else if (msg instanceof SendPingCommand) {
1 ✔
402
      sendPingFrame(ctx, (SendPingCommand) msg, promise);
1 ✔
403
    } else if (msg instanceof GracefulCloseCommand) {
1 ✔
404
      gracefulClose(ctx, (GracefulCloseCommand) msg, promise);
1 ✔
405
    } else if (msg instanceof ForcefulCloseCommand) {
1 ✔
406
      forcefulClose(ctx, (ForcefulCloseCommand) msg, promise);
1 ✔
407
    } else if (msg == NOOP_MESSAGE) {
1 ✔
408
      ctx.write(Unpooled.EMPTY_BUFFER, promise);
1 ✔
409
    } else {
410
      throw new AssertionError("Write called for unexpected type: " + msg.getClass().getName());
×
411
    }
412
  }
1 ✔
413

414
  void startWriteQueue(Channel channel) {
415
    clientWriteQueue = new WriteQueue(channel);
1 ✔
416
  }
1 ✔
417

418
  WriteQueue getWriteQueue() {
419
    return clientWriteQueue;
1 ✔
420
  }
421

422
  ClientTransportLifecycleManager getLifecycleManager() {
423
    return lifecycleManager;
1 ✔
424
  }
425

426
  /**
427
   * Returns the given processed bytes back to inbound flow control.
428
   */
429
  void returnProcessedBytes(Http2Stream stream, int bytes) {
430
    try {
431
      decoder().flowController().consumeBytes(stream, bytes);
1 ✔
432
    } catch (Http2Exception e) {
×
433
      throw new RuntimeException(e);
×
434
    }
1 ✔
435
  }
1 ✔
436

437
  private void onHeadersRead(int streamId, Http2Headers headers, boolean endStream) {
438
    // Stream 1 is reserved for the Upgrade response, so we should ignore its headers here:
439
    if (streamId != Http2CodecUtil.HTTP_UPGRADE_STREAM_ID) {
1 ✔
440
      NettyClientStream.TransportState stream = clientStream(requireHttp2Stream(streamId));
1 ✔
441
      PerfMark.event("NettyClientHandler.onHeadersRead", stream.tag());
1 ✔
442
      // check metadata size vs soft limit
443
      int h2HeadersSize = Utils.getH2HeadersSize(headers);
1 ✔
444
      boolean shouldFail =
1 ✔
445
          Utils.shouldRejectOnMetadataSizeSoftLimitExceeded(
1 ✔
446
              h2HeadersSize, softLimitHeaderListSize, maxHeaderListSize);
447
      if (shouldFail && endStream) {
1 ✔
448
        stream.transportReportStatus(Status.RESOURCE_EXHAUSTED
×
449
            .withDescription(
×
450
                String.format(
×
451
                    "Server Status + Trailers of size %d exceeded Metadata size soft limit: %d",
452
                    h2HeadersSize,
×
453
                    softLimitHeaderListSize)), true, new Metadata());
×
454
        return;
×
455
      } else if (shouldFail) {
1 ✔
456
        stream.transportReportStatus(Status.RESOURCE_EXHAUSTED
1 ✔
457
            .withDescription(
1 ✔
458
                String.format(
1 ✔
459
                    "Server Headers of size %d exceeded Metadata size soft limit: %d",
460
                    h2HeadersSize,
1 ✔
461
                    softLimitHeaderListSize)), true, new Metadata());
1 ✔
462
        return;
1 ✔
463
      }
464
      stream.transportHeadersReceived(headers, endStream);
1 ✔
465
    }
466

467
    if (keepAliveManager != null) {
1 ✔
468
      keepAliveManager.onDataReceived();
1 ✔
469
    }
470
  }
1 ✔
471

472
  /**
473
   * Handler for an inbound HTTP/2 DATA frame.
474
   */
475
  private void onDataRead(int streamId, ByteBuf data, int padding, boolean endOfStream) {
476
    flowControlPing().onDataRead(data.readableBytes(), padding);
1 ✔
477
    NettyClientStream.TransportState stream = clientStream(requireHttp2Stream(streamId));
1 ✔
478
    PerfMark.event("NettyClientHandler.onDataRead", stream.tag());
1 ✔
479
    stream.transportDataReceived(data, endOfStream);
1 ✔
480
    if (keepAliveManager != null) {
1 ✔
481
      keepAliveManager.onDataReceived();
1 ✔
482
    }
483
  }
1 ✔
484

485
  /**
486
   * Handler for an inbound HTTP/2 RST_STREAM frame, terminating a stream.
487
   */
488
  private void onRstStreamRead(int streamId, long errorCode) {
489
    NettyClientStream.TransportState stream = clientStream(connection().stream(streamId));
1 ✔
490
    if (stream != null) {
1 ✔
491
      PerfMark.event("NettyClientHandler.onRstStreamRead", stream.tag());
1 ✔
492
      Status status = statusFromH2Error(null, "RST_STREAM closed stream", errorCode, null);
1 ✔
493
      stream.transportReportStatus(
1 ✔
494
          status,
495
          errorCode == Http2Error.REFUSED_STREAM.code()
1 ✔
496
              ? RpcProgress.REFUSED : RpcProgress.PROCESSED,
1 ✔
497
          false /*stop delivery*/,
498
          new Metadata());
499
      if (keepAliveManager != null) {
1 ✔
500
        keepAliveManager.onDataReceived();
×
501
      }
502
    }
503
  }
1 ✔
504

505
  @Override
506
  public void close(ChannelHandlerContext ctx, ChannelPromise promise) throws Exception {
507
    tcpMetrics.recordTcpInfo(ctx.channel());
1 ✔
508
    logger.fine("Network channel being closed by the application.");
1 ✔
509
    if (ctx.channel().isActive()) { // Ignore notification that the socket was closed
1 ✔
510
      lifecycleManager.notifyShutdown(
1 ✔
511
          Status.UNAVAILABLE.withDescription("Transport closed for unknown reason"),
1 ✔
512
          SimpleDisconnectError.UNKNOWN);
513
    }
514
    super.close(ctx, promise);
1 ✔
515
  }
1 ✔
516

517
  /**
518
   * Handler for the Channel shutting down.
519
   */
520
  @Override
521
  public void channelActive(ChannelHandlerContext ctx) throws Exception {
522
    tcpMetrics.channelActive(ctx.channel());
1 ✔
523
    super.channelActive(ctx);
1 ✔
524
  }
1 ✔
525

526
  @Override
527
  public void channelInactive(ChannelHandlerContext ctx) throws Exception {
528
    try {
529
      logger.fine("Network channel is closed");
1 ✔
530
      tcpMetrics.channelInactive(ctx.channel());
1 ✔
531
      Status status = Status.UNAVAILABLE.withDescription("Network closed for unknown reason");
1 ✔
532
      lifecycleManager.notifyShutdown(status, SimpleDisconnectError.UNKNOWN);
1 ✔
533
      final Status streamStatus;
534
      if (channelInactiveReason != null) {
1 ✔
535
        streamStatus = channelInactiveReason;
1 ✔
536
      } else {
537
        streamStatus = lifecycleManager.getShutdownStatus();
1 ✔
538
      }
539
      try {
540
        cancelPing(lifecycleManager.getShutdownStatus());
1 ✔
541
        // Report status to the application layer for any open streams
542
        connection().forEachActiveStream(new Http2StreamVisitor() {
1 ✔
543
          @Override
544
          public boolean visit(Http2Stream stream) throws Http2Exception {
545
            NettyClientStream.TransportState clientStream = clientStream(stream);
1 ✔
546
            if (clientStream != null) {
1 ✔
547
              clientStream.transportReportStatus(streamStatus, false, new Metadata());
1 ✔
548
            }
549
            return true;
1 ✔
550
          }
551
        });
552
      } finally {
553
        lifecycleManager.notifyTerminated(status, SimpleDisconnectError.UNKNOWN);
1 ✔
554
      }
555
    } finally {
556
      // Close any open streams
557
      super.channelInactive(ctx);
1 ✔
558
      if (keepAliveManager != null) {
1 ✔
559
        keepAliveManager.onTransportTermination();
1 ✔
560
      }
561
    }
562
  }
1 ✔
563

564
  @Override
565
  public void handleProtocolNegotiationCompleted(
566
      Attributes attributes, InternalChannelz.Security securityInfo) {
567
    this.attributes = this.attributes.toBuilder().setAll(attributes).build();
1 ✔
568
    this.securityInfo = securityInfo;
1 ✔
569
    super.handleProtocolNegotiationCompleted(attributes, securityInfo);
1 ✔
570
    writeBufferingAndRemove(ctx().channel());
1 ✔
571
  }
1 ✔
572

573
  static void writeBufferingAndRemove(Channel channel) {
574
    checkNotNull(channel, "channel");
1 ✔
575
    ChannelHandlerContext handlerCtx =
1 ✔
576
        channel.pipeline().context(WriteBufferingAndExceptionHandler.class);
1 ✔
577
    if (handlerCtx == null) {
1 ✔
578
      return;
1 ✔
579
    }
580
    ((WriteBufferingAndExceptionHandler) handlerCtx.handler()).writeBufferedAndRemove(handlerCtx);
1 ✔
581
  }
1 ✔
582

583
  @Override
584
  public Attributes getEagAttributes() {
585
    return eagAttributes;
1 ✔
586
  }
587

588
  @Override
589
  public String getAuthority() {
590
    return authority;
1 ✔
591
  }
592

593
  InternalChannelz.Security getSecurityInfo() {
594
    return securityInfo;
1 ✔
595
  }
596

597
  @Override
598
  protected void onConnectionError(ChannelHandlerContext ctx,  boolean outbound, Throwable cause,
599
      Http2Exception http2Ex) {
600
    logger.log(Level.FINE, "Caught a connection error", cause);
1 ✔
601
    lifecycleManager.notifyShutdown(Utils.statusFromThrowable(cause),
1 ✔
602
        SimpleDisconnectError.SOCKET_ERROR);
603
    // Parent class will shut down the Channel
604
    super.onConnectionError(ctx, outbound, cause, http2Ex);
1 ✔
605
  }
1 ✔
606

607
  @Override
608
  protected void onStreamError(ChannelHandlerContext ctx, boolean outbound, Throwable cause,
609
      Http2Exception.StreamException http2Ex) {
610
    // Close the stream with a status that contains the cause.
611
    NettyClientStream.TransportState stream = clientStream(connection().stream(http2Ex.streamId()));
1 ✔
612
    if (stream != null) {
1 ✔
613
      stream.transportReportStatus(Utils.statusFromThrowable(cause), false, new Metadata());
1 ✔
614
    } else {
615
      logger.log(Level.FINE, "Stream error for unknown stream " + http2Ex.streamId(), cause);
1 ✔
616
    }
617

618
    // Delegate to the base class to send a RST_STREAM.
619
    super.onStreamError(ctx, outbound, cause, http2Ex);
1 ✔
620
  }
1 ✔
621

622
  @Override
623
  protected boolean isGracefulShutdownComplete() {
624
    // Only allow graceful shutdown to complete after all pending streams have completed.
625
    return super.isGracefulShutdownComplete()
1 ✔
626
        && ((StreamBufferingEncoder) encoder()).numBufferedStreams() == 0;
1 ✔
627
  }
628

629
  /**
630
   * Attempts to create a new stream from the given command. If there are too many active streams,
631
   * the creation request is queued.
632
   */
633
  private void createStream(CreateStreamCommand command, ChannelPromise promise)
634
          throws Exception {
635
    if (lifecycleManager.getShutdownStatus() != null) {
1 ✔
636
      command.stream().setNonExistent();
1 ✔
637
      // The connection is going away (it is really the GOAWAY case),
638
      // just terminate the stream now.
639
      command.stream().transportReportStatus(
1 ✔
640
          lifecycleManager.getShutdownStatus(), RpcProgress.MISCARRIED, true, new Metadata());
1 ✔
641
      promise.setFailure(InternalStatus.asRuntimeExceptionWithoutStacktrace(
1 ✔
642
              lifecycleManager.getShutdownStatus(), null));
1 ✔
643
      return;
1 ✔
644
    }
645

646
    CharSequence authorityHeader = command.headers().authority();
1 ✔
647
    if (authorityHeader == null) {
1 ✔
648
      Status authorityVerificationStatus = Status.UNAVAILABLE.withDescription(
1 ✔
649
              "Missing authority header");
650
      command.stream().setNonExistent();
1 ✔
651
      command.stream().transportReportStatus(
1 ✔
652
              Status.UNAVAILABLE, RpcProgress.PROCESSED, true, new Metadata());
653
      promise.setFailure(InternalStatus.asRuntimeExceptionWithoutStacktrace(
1 ✔
654
              authorityVerificationStatus, null));
655
      return;
1 ✔
656
    }
657
    // No need to verify authority for the rpc outgoing header if it is same as the authority
658
    // for the transport
659
    if (!authority.contentEquals(authorityHeader)) {
1 ✔
660
      Status authorityVerificationStatus = peerVerificationResults.get(
1 ✔
661
              authorityHeader.toString());
1 ✔
662
      if (authorityVerificationStatus == null) {
1 ✔
663
        if (attributes.get(GrpcAttributes.ATTR_AUTHORITY_VERIFIER) == null) {
1 ✔
664
          authorityVerificationStatus = Status.UNAVAILABLE.withDescription(
1 ✔
665
                  "Authority verifier not found to verify authority");
666
          command.stream().setNonExistent();
1 ✔
667
          command.stream().transportReportStatus(
1 ✔
668
                  authorityVerificationStatus, RpcProgress.PROCESSED, true, new Metadata());
669
          promise.setFailure(InternalStatus.asRuntimeExceptionWithoutStacktrace(
1 ✔
670
                  authorityVerificationStatus, null));
671
          return;
1 ✔
672
        }
673
        authorityVerificationStatus = attributes.get(GrpcAttributes.ATTR_AUTHORITY_VERIFIER)
1 ✔
674
                .verifyAuthority(authorityHeader.toString());
1 ✔
675
        peerVerificationResults.put(authorityHeader.toString(), authorityVerificationStatus);
1 ✔
676
        if (!authorityVerificationStatus.isOk() && !enablePerRpcAuthorityCheck) {
1 ✔
677
          logger.log(Level.WARNING, String.format("%s.%s",
1 ✔
678
                          authorityVerificationStatus.getDescription(),
1 ✔
679
                          enablePerRpcAuthorityCheck
1 ✔
680
                                  ? "" : " This will be an error in the future."),
1 ✔
681
                  InternalStatus.asRuntimeExceptionWithoutStacktrace(
1 ✔
682
                          authorityVerificationStatus, null));
683
        }
684
      }
685
      if (!authorityVerificationStatus.isOk()) {
1 ✔
686
        if (enablePerRpcAuthorityCheck) {
1 ✔
687
          command.stream().setNonExistent();
1 ✔
688
          command.stream().transportReportStatus(
1 ✔
689
                  authorityVerificationStatus, RpcProgress.PROCESSED, true, new Metadata());
690
          promise.setFailure(InternalStatus.asRuntimeExceptionWithoutStacktrace(
1 ✔
691
                  authorityVerificationStatus, null));
692
          return;
1 ✔
693
        }
694
      }
695
    }
696
    // Get the stream ID for the new stream.
697
    int streamId;
698
    try {
699
      streamId = incrementAndGetNextStreamId();
1 ✔
700
    } catch (StatusException e) {
1 ✔
701
      command.stream().setNonExistent();
1 ✔
702
      // Stream IDs have been exhausted for this connection. Fail the promise immediately.
703
      promise.setFailure(e);
1 ✔
704

705
      // Initiate a graceful shutdown if we haven't already.
706
      if (!connection().goAwaySent()) {
1 ✔
707
        logger.fine("Stream IDs have been exhausted for this connection. "
1 ✔
708
                + "Initiating graceful shutdown of the connection.");
709
        lifecycleManager.notifyShutdown(e.getStatus(), SimpleDisconnectError.UNKNOWN);
1 ✔
710
        close(ctx(), ctx().newPromise());
1 ✔
711
      }
712
      return;
1 ✔
713
    }
1 ✔
714
    if (connection().goAwayReceived()) {
1 ✔
715
      Status s = abruptGoAwayStatus;
1 ✔
716
      int maxActiveStreams = connection().local().maxActiveStreams();
1 ✔
717
      int lastStreamId = connection().local().lastStreamKnownByPeer();
1 ✔
718
      if (s == null) {
1 ✔
719
        // Should be impossible, but handle pseudo-gracefully
720
        s = Status.INTERNAL.withDescription(
×
721
            "Failed due to abrupt GOAWAY, but can't find GOAWAY details");
722
      } else if (streamId > lastStreamId) {
1 ✔
723
        s = s.augmentDescription(
1 ✔
724
            "stream id: " + streamId + ", GOAWAY Last-Stream-ID:" + lastStreamId);
725
      } else if (connection().local().numActiveStreams() == maxActiveStreams) {
1 ✔
726
        s = s.augmentDescription("At MAX_CONCURRENT_STREAMS limit. limit: " + maxActiveStreams);
1 ✔
727
      }
728
      if (streamId > lastStreamId || connection().local().numActiveStreams() == maxActiveStreams) {
1 ✔
729
        // This should only be reachable during onGoAwayReceived, as otherwise
730
        // getShutdownThrowable() != null
731
        command.stream().setNonExistent();
1 ✔
732
        command.stream().transportReportStatus(s, RpcProgress.MISCARRIED, true, new Metadata());
1 ✔
733
        promise.setFailure(s.asRuntimeException());
1 ✔
734
        return;
1 ✔
735
      }
736
    }
737

738
    NettyClientStream.TransportState stream = command.stream();
1 ✔
739
    Http2Headers headers = command.headers();
1 ✔
740
    stream.setId(streamId);
1 ✔
741

742
    try (TaskCloseable ignore = PerfMark.traceTask("NettyClientHandler.createStream")) {
1 ✔
743
      PerfMark.linkIn(command.getLink());
1 ✔
744
      PerfMark.attachTag(stream.tag());
1 ✔
745
      createStreamTraced(
1 ✔
746
          streamId, stream, headers, command.isGet(), command.shouldBeCountedForInUse(), promise);
1 ✔
747
    }
748
  }
1 ✔
749

750
  private void createStreamTraced(
751
      final int streamId,
752
      final NettyClientStream.TransportState stream,
753
      final Http2Headers headers,
754
      boolean isGet,
755
      final boolean shouldBeCountedForInUse,
756
      final ChannelPromise promise) {
757
    // Create an intermediate promise so that we can intercept the failure reported back to the
758
    // application.
759
    ChannelPromise tempPromise = ctx().newPromise();
1 ✔
760
    encoder().writeHeaders(ctx(), streamId, headers, 0, isGet, tempPromise)
1 ✔
761
        .addListener(new ChannelFutureListener() {
1 ✔
762
          @Override
763
          public void operationComplete(ChannelFuture future) throws Exception {
764
            if (future.isSuccess()) {
1 ✔
765
              // The http2Stream will be null in case a stream buffered in the encoder
766
              // was canceled via RST_STREAM.
767
              Http2Stream http2Stream = connection().stream(streamId);
1 ✔
768
              if (http2Stream != null) {
1 ✔
769
                stream.getStatsTraceContext().clientOutboundHeaders();
1 ✔
770
                http2Stream.setProperty(streamKey, stream);
1 ✔
771

772
                // This delays the in-use state until the I/O completes, which technically may
773
                // be later than we would like.
774
                if (shouldBeCountedForInUse) {
1 ✔
775
                  inUseState.updateObjectInUse(http2Stream, true);
1 ✔
776
                }
777

778
                // Attach the client stream to the HTTP/2 stream object as user data.
779
                stream.setHttp2Stream(http2Stream);
1 ✔
780
                promise.setSuccess();
1 ✔
781
              } else {
782
                // Otherwise, the stream has been cancelled and Netty is sending a
783
                // RST_STREAM frame which causes it to purge pending writes from the
784
                // flow-controller and delete the http2Stream. The stream listener has already
785
                // been notified of cancellation so there is nothing to do.
786
                //
787
                // This process has been observed to fail in some circumstances, leaving listeners
788
                // unanswered. Ensure that some exception has been delivered consistent with the
789
                // implied RST_STREAM result above.
790
                Status status = Status.INTERNAL.withDescription("unknown stream for connection");
1 ✔
791
                promise.setFailure(status.asRuntimeException());
1 ✔
792
              }
793
            } else {
1 ✔
794
              Throwable cause = future.cause();
1 ✔
795
              if (cause instanceof StreamBufferingEncoder.Http2GoAwayException) {
1 ✔
796
                StreamBufferingEncoder.Http2GoAwayException e =
1 ✔
797
                    (StreamBufferingEncoder.Http2GoAwayException) cause;
798
                Status status = statusFromH2Error(
1 ✔
799
                    Status.Code.UNAVAILABLE, "GOAWAY closed buffered stream",
800
                    e.errorCode(), e.debugData());
1 ✔
801
                cause = status.asRuntimeException();
1 ✔
802
                stream.transportReportStatus(status, RpcProgress.MISCARRIED, true, new Metadata());
1 ✔
803
              } else if (cause instanceof StreamBufferingEncoder.Http2ChannelClosedException) {
1 ✔
804
                Status status = lifecycleManager.getShutdownStatus();
1 ✔
805
                if (status == null) {
1 ✔
806
                  status = Status.UNAVAILABLE.withCause(cause)
×
807
                      .withDescription("Connection closed while stream is buffered");
×
808
                }
809
                stream.transportReportStatus(status, RpcProgress.MISCARRIED, true, new Metadata());
1 ✔
810
              }
811
              promise.setFailure(cause);
1 ✔
812
            }
813
          }
1 ✔
814
        });
815
    // When the HEADERS are not buffered because of MAX_CONCURRENT_STREAMS in
816
    // StreamBufferingEncoder, the stream is created immediately even if the bytes of the HEADERS
817
    // are delayed because the OS may have too much buffered and isn't accepting the write. The
818
    // write promise is also delayed until flush(). However, we need to associate the netty stream
819
    // with the transport state so that goingAway() and forcefulClose() and able to notify the
820
    // stream of failures.
821
    //
822
    // This leaves a hole when MAX_CONCURRENT_STREAMS is reached, as http2Stream will be null, but
823
    // it is better than nothing.
824
    Http2Stream http2Stream = connection().stream(streamId);
1 ✔
825
    if (http2Stream != null) {
1 ✔
826
      http2Stream.setProperty(streamKey, stream);
1 ✔
827
    }
828
  }
1 ✔
829

830
  /**
831
   * Cancels this stream.
832
   */
833
  private void cancelStream(ChannelHandlerContext ctx, CancelClientStreamCommand cmd,
834
      ChannelPromise promise) {
835
    NettyClientStream.TransportState stream = cmd.stream();
1 ✔
836
    try (TaskCloseable ignore = PerfMark.traceTask("NettyClientHandler.cancelStream")) {
1 ✔
837
      PerfMark.attachTag(stream.tag());
1 ✔
838
      PerfMark.linkIn(cmd.getLink());
1 ✔
839
      Status reason = cmd.reason();
1 ✔
840
      if (reason != null) {
1 ✔
841
        stream.transportReportStatus(reason, true, new Metadata());
1 ✔
842
      }
843
      if (!cmd.stream().isNonExistent()) {
1 ✔
844
        encoder().writeRstStream(ctx, stream.id(), Http2Error.CANCEL.code(), promise);
1 ✔
845
      } else {
846
        promise.setSuccess();
1 ✔
847
      }
848
    }
849
  }
1 ✔
850

851
  /**
852
   * Sends the given GRPC frame for the stream.
853
   */
854
  private void sendGrpcFrame(ChannelHandlerContext ctx, SendGrpcFrameCommand cmd,
855
      ChannelPromise promise) {
856
    try (TaskCloseable ignore = PerfMark.traceTask("NettyClientHandler.sendGrpcFrame")) {
1 ✔
857
      PerfMark.attachTag(cmd.stream().tag());
1 ✔
858
      PerfMark.linkIn(cmd.getLink());
1 ✔
859
      // Call the base class to write the HTTP/2 DATA frame.
860
      // Note: no need to flush since this is handled by the outbound flow controller.
861
      encoder().writeData(ctx, cmd.stream().id(), cmd.content(), 0, cmd.endStream(), promise);
1 ✔
862
    }
863
  }
1 ✔
864

865
  private void sendPingFrame(ChannelHandlerContext ctx, SendPingCommand msg,
866
      ChannelPromise promise) {
867
    try (TaskCloseable ignore = PerfMark.traceTask("NettyClientHandler.sendPingFrame")) {
1 ✔
868
      PerfMark.linkIn(msg.getLink());
1 ✔
869
      sendPingFrameTraced(ctx, msg, promise);
1 ✔
870
    }
871
  }
1 ✔
872

873
  /**
874
   * Sends a PING frame. If a ping operation is already outstanding, the callback in the message is
875
   * registered to be called when the existing operation completes, and no new frame is sent.
876
   */
877
  private void sendPingFrameTraced(ChannelHandlerContext ctx, SendPingCommand msg,
878
      ChannelPromise promise) {
879
    // Don't check lifecycleManager.getShutdownStatus() since we want to allow pings after shutdown
880
    // but before termination. After termination, messages will no longer arrive because the
881
    // pipeline clears all handlers on channel close.
882

883
    PingCallback callback = msg.callback();
1 ✔
884
    Executor executor = msg.executor();
1 ✔
885
    // we only allow one outstanding ping at a time, so just add the callback to
886
    // any outstanding operation
887
    if (ping != null) {
1 ✔
888
      promise.setSuccess();
1 ✔
889
      ping.addCallback(callback, executor);
1 ✔
890
      return;
1 ✔
891
    }
892

893
    // Use a new promise to prevent calling the callback twice on write failure: here and in
894
    // NettyClientTransport.ping(). It may appear strange, but it will behave the same as if
895
    // ping != null above.
896
    promise.setSuccess();
1 ✔
897
    promise = ctx().newPromise();
1 ✔
898
    // set outstanding operation
899
    long data = USER_PING_PAYLOAD;
1 ✔
900
    Stopwatch stopwatch = stopwatchFactory.get();
1 ✔
901
    stopwatch.start();
1 ✔
902
    ping = new Http2Ping(data, stopwatch);
1 ✔
903
    ping.addCallback(callback, executor);
1 ✔
904
    // and then write the ping
905
    encoder().writePing(ctx, false, USER_PING_PAYLOAD, promise);
1 ✔
906
    ctx.flush();
1 ✔
907
    final Http2Ping finalPing = ping;
1 ✔
908
    promise.addListener(new ChannelFutureListener() {
1 ✔
909
      @Override
910
      public void operationComplete(ChannelFuture future) throws Exception {
911
        if (future.isSuccess()) {
1 ✔
912
          transportTracer.reportKeepAliveSent();
1 ✔
913
          return;
1 ✔
914
        }
915
        Throwable cause = future.cause();
×
916
        Status status = lifecycleManager.getShutdownStatus();
×
917
        if (cause instanceof ClosedChannelException) {
×
918
          if (status == null) {
×
919
            status = Status.UNKNOWN.withDescription("Ping failed but for unknown reason.")
×
920
                    .withCause(future.cause());
×
921
          }
922
        } else {
923
          status = Utils.statusFromThrowable(cause);
×
924
        }
925
        finalPing.failed(status);
×
926
        if (ping == finalPing) {
×
927
          ping = null;
×
928
        }
929
      }
×
930
    });
931
  }
1 ✔
932

933
  private void gracefulClose(ChannelHandlerContext ctx, GracefulCloseCommand msg,
934
      ChannelPromise promise) throws Exception {
935
    lifecycleManager.notifyShutdown(msg.getStatus(), SimpleDisconnectError.SUBCHANNEL_SHUTDOWN);
1 ✔
936
    // Explicitly flush to create any buffered streams before sending GOAWAY.
937
    // TODO(ejona): determine if the need to flush is a bug in Netty
938
    flush(ctx);
1 ✔
939
    close(ctx, promise);
1 ✔
940
  }
1 ✔
941

942
  private void forcefulClose(final ChannelHandlerContext ctx, final ForcefulCloseCommand msg,
943
      ChannelPromise promise) throws Exception {
944
    connection().forEachActiveStream(new Http2StreamVisitor() {
1 ✔
945
      @Override
946
      public boolean visit(Http2Stream stream) throws Http2Exception {
947
        NettyClientStream.TransportState clientStream = clientStream(stream);
1 ✔
948
        Tag tag = clientStream != null ? clientStream.tag() : PerfMark.createTag();
1 ✔
949
        try (TaskCloseable ignore = PerfMark.traceTask("NettyClientHandler.forcefulClose")) {
1 ✔
950
          PerfMark.linkIn(msg.getLink());
1 ✔
951
          PerfMark.attachTag(tag);
1 ✔
952
          if (clientStream != null) {
1 ✔
953
            clientStream.transportReportStatus(msg.getStatus(), true, new Metadata());
1 ✔
954
            resetStream(ctx, stream.id(), Http2Error.CANCEL.code(), ctx.newPromise());
1 ✔
955
          }
956
          stream.close();
1 ✔
957
          return true;
1 ✔
958
        }
959
      }
960
    });
961
    close(ctx, promise);
1 ✔
962
  }
1 ✔
963

964
  /**
965
   * Handler for a GOAWAY being received. Fails any streams created after the
966
   * last known stream. May only be called during a read.
967
   */
968
  private void goingAway(long errorCode, byte[] debugData) {
969
    Status finalStatus = statusFromH2Error(
1 ✔
970
        Status.Code.UNAVAILABLE, "GOAWAY shut down transport", errorCode, debugData);
971
    DisconnectError disconnectError = new GoAwayDisconnectError(
1 ✔
972
        GrpcUtil.Http2Error.forCode(errorCode));
1 ✔
973
    lifecycleManager.notifyGracefulShutdown(finalStatus, disconnectError);
1 ✔
974
    abruptGoAwayStatus = statusFromH2Error(
1 ✔
975
        Status.Code.UNAVAILABLE, "Abrupt GOAWAY closed unsent stream", errorCode, debugData);
976
    // While this _should_ be UNAVAILABLE, Netty uses the wrong stream id in the GOAWAY when it
977
    // fails streams due to HPACK failures (e.g., header list too large). To be more conservative,
978
    // we assume any sent streams may be related to the GOAWAY. This should rarely impact users
979
    // since the main time servers should use abrupt GOAWAYs if there is a protocol error, and if
980
    // there wasn't a protocol error the error code was probably NO_ERROR which is mapped to
981
    // UNAVAILABLE. https://github.com/netty/netty/issues/10670
982
    final Status abruptGoAwayStatusConservative = statusFromH2Error(
1 ✔
983
        null, "Abrupt GOAWAY closed sent stream", errorCode, debugData);
984
    final boolean mayBeHittingNettyBug = errorCode != Http2Error.NO_ERROR.code();
1 ✔
985
    // Try to allocate as many in-flight streams as possible, to reduce race window of
986
    // https://github.com/grpc/grpc-java/issues/2562 . To be of any help, the server has to
987
    // gracefully shut down the connection with two GOAWAYs. gRPC servers generally send a PING
988
    // after the first GOAWAY, so they can very precisely detect when the GOAWAY has been
989
    // processed and thus this processing must be in-line before processing additional reads.
990

991
    // This can cause reentrancy, but should be minor since it is normal to handle writes in
992
    // response to a read. Also, the call stack is rather shallow at this point
993
    clientWriteQueue.drainNow();
1 ✔
994
    if (lifecycleManager.notifyShutdown(finalStatus, disconnectError)) {
1 ✔
995
      // This is for the only RPCs that are actually covered by the GOAWAY error code. All other
996
      // RPCs were not observed by the remote and so should be UNAVAILABLE.
997
      channelInactiveReason = statusFromH2Error(
1 ✔
998
          null, "Connection closed after GOAWAY", errorCode, debugData);
999
    }
1000

1001
    final int lastKnownStream = connection().local().lastStreamKnownByPeer();
1 ✔
1002
    try {
1003
      connection().forEachActiveStream(new Http2StreamVisitor() {
1 ✔
1004
        @Override
1005
        public boolean visit(Http2Stream stream) throws Http2Exception {
1006
          if (stream.id() > lastKnownStream) {
1 ✔
1007
            NettyClientStream.TransportState clientStream = clientStream(stream);
1 ✔
1008
            if (clientStream != null) {
1 ✔
1009
              // RpcProgress _should_ be REFUSED, but are being conservative. See comment for
1010
              // abruptGoAwayStatusConservative. This does reduce our ability to perform transparent
1011
              // retries, but only if something else caused a connection failure.
1012
              RpcProgress progress = mayBeHittingNettyBug
1 ✔
1013
                  ? RpcProgress.PROCESSED
1 ✔
1014
                  : RpcProgress.REFUSED;
1 ✔
1015
              clientStream.transportReportStatus(
1 ✔
1016
                  abruptGoAwayStatusConservative, progress, false, new Metadata());
1017
            }
1018
            stream.close();
1 ✔
1019
          }
1020
          return true;
1 ✔
1021
        }
1022
      });
1023
    } catch (Http2Exception e) {
×
1024
      throw new RuntimeException(e);
×
1025
    }
1 ✔
1026
  }
1 ✔
1027

1028
  private void cancelPing(Status s) {
1029
    if (ping != null) {
1 ✔
1030
      ping.failed(s);
1 ✔
1031
      ping = null;
1 ✔
1032
    }
1033
  }
1 ✔
1034

1035
  /** If {@code statusCode} is non-null, it will be used instead of the http2 error code mapping. */
1036
  private Status statusFromH2Error(
1037
      Status.Code statusCode, String context, long errorCode, byte[] debugData) {
1038
    Status status = GrpcUtil.Http2Error.statusForCode(errorCode);
1 ✔
1039
    if (statusCode == null) {
1 ✔
1040
      statusCode = status.getCode();
1 ✔
1041
    }
1042
    String debugString = "";
1 ✔
1043
    if (debugData != null && debugData.length > 0) {
1 ✔
1044
      // If a debug message was provided, use it.
1045
      debugString = ", debug data: " + new String(debugData, UTF_8);
1 ✔
1046
    }
1047
    return statusCode.toStatus()
1 ✔
1048
        .withDescription(context + ". " + status.getDescription() + debugString);
1 ✔
1049
  }
1050

1051
  /**
1052
   * Gets the client stream associated to the given HTTP/2 stream object.
1053
   */
1054
  private NettyClientStream.TransportState clientStream(Http2Stream stream) {
1055
    return stream == null ? null : (NettyClientStream.TransportState) stream.getProperty(streamKey);
1 ✔
1056
  }
1057

1058
  private int incrementAndGetNextStreamId() throws StatusException {
1059
    int nextStreamId = connection().local().incrementAndGetNextStreamId();
1 ✔
1060
    if (nextStreamId < 0) {
1 ✔
1061
      logger.fine("Stream IDs have been exhausted for this connection. "
1 ✔
1062
              + "Initiating graceful shutdown of the connection.");
1063
      throw EXHAUSTED_STREAMS_STATUS.asException();
1 ✔
1064
    }
1065
    return nextStreamId;
1 ✔
1066
  }
1067

1068
  private Http2Stream requireHttp2Stream(int streamId) {
1069
    Http2Stream stream = connection().stream(streamId);
1 ✔
1070
    if (stream == null) {
1 ✔
1071
      // This should never happen.
1072
      throw new AssertionError("Stream does not exist: " + streamId);
×
1073
    }
1074
    return stream;
1 ✔
1075
  }
1076

1077
  private class FrameListener extends Http2FrameAdapter {
1 ✔
1078
    private boolean firstSettings = true;
1 ✔
1079

1080
    @Override
1081
    public void onSettingsRead(ChannelHandlerContext ctx, Http2Settings settings) {
1082
      if (firstSettings) {
1 ✔
1083
        firstSettings = false;
1 ✔
1084
        attributes = lifecycleManager.filterAttributes(attributes);
1 ✔
1085
        lifecycleManager.notifyReady();
1 ✔
1086
      }
1087
    }
1 ✔
1088

1089
    @Override
1090
    public int onDataRead(ChannelHandlerContext ctx, int streamId, ByteBuf data, int padding,
1091
        boolean endOfStream) throws Http2Exception {
1092
      NettyClientHandler.this.onDataRead(streamId, data, padding, endOfStream);
1 ✔
1093
      return padding;
1 ✔
1094
    }
1095

1096
    @Override
1097
    public void onHeadersRead(ChannelHandlerContext ctx,
1098
        int streamId,
1099
        Http2Headers headers,
1100
        int streamDependency,
1101
        short weight,
1102
        boolean exclusive,
1103
        int padding,
1104
        boolean endStream) throws Http2Exception {
1105
      NettyClientHandler.this.onHeadersRead(streamId, headers, endStream);
1 ✔
1106
    }
1 ✔
1107

1108
    @Override
1109
    public void onRstStreamRead(ChannelHandlerContext ctx, int streamId, long errorCode)
1110
        throws Http2Exception {
1111
      NettyClientHandler.this.onRstStreamRead(streamId, errorCode);
1 ✔
1112
    }
1 ✔
1113

1114
    @Override
1115
    public void onPingAckRead(ChannelHandlerContext ctx, long ackPayload) throws Http2Exception {
1116
      Http2Ping p = ping;
1 ✔
1117
      if (ackPayload == flowControlPing().payload()) {
1 ✔
1118
        flowControlPing().updateWindow();
1 ✔
1119
        logger.log(Level.FINE, "Window: {0}",
1 ✔
1120
            decoder().flowController().initialWindowSize(connection().connectionStream()));
1 ✔
1121
      } else if (p != null) {
1 ✔
1122
        if (p.payload() == ackPayload) {
1 ✔
1123
          p.complete();
1 ✔
1124
          ping = null;
1 ✔
1125
        } else {
1126
          logger.log(Level.WARNING,
1 ✔
1127
              "Received unexpected ping ack. Expecting {0}, got {1}",
1128
              new Object[] {p.payload(), ackPayload});
1 ✔
1129
        }
1130
      } else {
1131
        logger.warning("Received unexpected ping ack. No ping outstanding");
×
1132
      }
1133
      if (keepAliveManager != null) {
1 ✔
1134
        keepAliveManager.onDataReceived();
1 ✔
1135
      }
1136
    }
1 ✔
1137

1138
    @Override
1139
    public void onPingRead(ChannelHandlerContext ctx, long data) throws Http2Exception {
1140
      if (keepAliveManager != null) {
1 ✔
1141
        keepAliveManager.onDataReceived();
×
1142
      }
1143
    }
1 ✔
1144
  }
1145

1146
  private static class PingCountingFrameWriter extends DecoratingHttp2FrameWriter
1147
      implements AbstractNettyHandler.PingLimiter {
1148
    private int pingCount;
1149

1150
    public PingCountingFrameWriter(Http2FrameWriter delegate) {
1151
      super(delegate);
1 ✔
1152
    }
1 ✔
1153

1154
    @Override
1155
    public boolean isPingAllowed() {
1156
      // "3 strikes" may cause the server to complain, so we limit ourselves to 2 or below.
1157
      return pingCount < 2;
1 ✔
1158
    }
1159

1160
    @Override
1161
    public ChannelFuture writeHeaders(
1162
        ChannelHandlerContext ctx, int streamId, Http2Headers headers,
1163
        int padding, boolean endStream, ChannelPromise promise) {
1164
      pingCount = 0;
1 ✔
1165
      return super.writeHeaders(ctx, streamId, headers, padding, endStream, promise);
1 ✔
1166
    }
1167

1168
    @Override
1169
    public ChannelFuture writeHeaders(
1170
        ChannelHandlerContext ctx, int streamId, Http2Headers headers,
1171
        int streamDependency, short weight, boolean exclusive,
1172
        int padding, boolean endStream, ChannelPromise promise) {
1173
      pingCount = 0;
×
1174
      return super.writeHeaders(ctx, streamId, headers, streamDependency, weight, exclusive,
×
1175
          padding, endStream, promise);
1176
    }
1177

1178
    @Override
1179
    public ChannelFuture writeWindowUpdate(
1180
        ChannelHandlerContext ctx, int streamId, int windowSizeIncrement, ChannelPromise promise) {
1181
      pingCount = 0;
1 ✔
1182
      return super.writeWindowUpdate(ctx, streamId, windowSizeIncrement, promise);
1 ✔
1183
    }
1184

1185
    @Override
1186
    public ChannelFuture writePing(
1187
        ChannelHandlerContext ctx, boolean ack, long data, ChannelPromise promise) {
1188
      if (!ack) {
1 ✔
1189
        pingCount++;
1 ✔
1190
      }
1191
      return super.writePing(ctx, ack, data, promise);
1 ✔
1192
    }
1193

1194
    @Override
1195
    public ChannelFuture writeData(
1196
        ChannelHandlerContext ctx, int streamId, ByteBuf data, int padding, boolean endStream,
1197
        ChannelPromise promise) {
1198
      if (data.isReadable()) {
1 ✔
1199
        pingCount = 0;
1 ✔
1200
      }
1201
      return super.writeData(ctx, streamId, data, padding, endStream, promise);
1 ✔
1202
    }
1203
  }
1204
}
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