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

grpc / grpc-java / #20499

05 Oct 2026 10:38AM UTC coverage: 89.356% (+0.03%) from 89.331%
#20499

push

github

web-flow
okhttp: Set socket read timeout during TLS handshake (#13081)

In `OkHttpClientTransport`, direct TLS connection setup calls
`OkHttpTlsUpgrader.upgrade(...)` without configuring a socket read
timeout (`setSoTimeout`). When a middlebox or unresponsive server
completes the TCP three-way handshake and acknowledges the TLS
`ClientHello` without sending a `ServerHello` or `RST`,
`sslSocket.startHandshake()` blocks indefinitely in socket read. Because
`this.socket` and `asyncSink` are only assigned after
`OkHttpTlsUpgrader.upgrade(...)` returns, channel shutdown cannot close
the underlying socket either, leaving the subchannel permanently stuck
in `CONNECTING`. Configure `sock.setSoTimeout(proxySocketTimeout)`
before `OkHttpTlsUpgrader.upgrade(...)` and reset it to `0` once the TLS
upgrade completes, matching the timeout handling in
`createHttpProxySocket`. Also close `sock` via
`GrpcUtil.closeQuietly(sock)` in `catch (Exception e)` so the socket is
not leaked when handshake setup fails before
`asyncSink.becomeConnected(...)` takes ownership.

39256 of 43932 relevant lines covered (89.36%)

0.89 hits per line

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

93.24
/../okhttp/src/main/java/io/grpc/okhttp/OkHttpClientTransport.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.okhttp;
18

19
import static com.google.common.base.Preconditions.checkState;
20
import static io.grpc.okhttp.Utils.DEFAULT_WINDOW_SIZE;
21
import static io.grpc.okhttp.Utils.DEFAULT_WINDOW_UPDATE_RATIO;
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.Stopwatch;
27
import com.google.common.base.Supplier;
28
import com.google.common.util.concurrent.ListenableFuture;
29
import com.google.common.util.concurrent.SettableFuture;
30
import com.google.errorprone.annotations.concurrent.GuardedBy;
31
import io.grpc.Attributes;
32
import io.grpc.CallOptions;
33
import io.grpc.ChannelCredentials;
34
import io.grpc.ClientStreamTracer;
35
import io.grpc.Grpc;
36
import io.grpc.HttpConnectProxiedSocketAddress;
37
import io.grpc.InternalChannelz;
38
import io.grpc.InternalChannelz.SocketStats;
39
import io.grpc.InternalLogId;
40
import io.grpc.Metadata;
41
import io.grpc.MethodDescriptor;
42
import io.grpc.MethodDescriptor.MethodType;
43
import io.grpc.SecurityLevel;
44
import io.grpc.Status;
45
import io.grpc.Status.Code;
46
import io.grpc.StatusException;
47
import io.grpc.TlsChannelCredentials;
48
import io.grpc.internal.CertificateUtils;
49
import io.grpc.internal.ClientStreamListener.RpcProgress;
50
import io.grpc.internal.ConnectionClientTransport;
51
import io.grpc.internal.DisconnectError;
52
import io.grpc.internal.GoAwayDisconnectError;
53
import io.grpc.internal.GrpcAttributes;
54
import io.grpc.internal.GrpcUtil;
55
import io.grpc.internal.Http2Ping;
56
import io.grpc.internal.InUseStateAggregator;
57
import io.grpc.internal.KeepAliveManager;
58
import io.grpc.internal.KeepAliveManager.ClientKeepAlivePinger;
59
import io.grpc.internal.NoopSslSession;
60
import io.grpc.internal.SerializingExecutor;
61
import io.grpc.internal.SimpleDisconnectError;
62
import io.grpc.internal.StatsTraceContext;
63
import io.grpc.internal.TransportTracer;
64
import io.grpc.okhttp.ExceptionHandlingFrameWriter.TransportExceptionHandler;
65
import io.grpc.okhttp.OkHttpChannelBuilder.OkHttpTransportFactory;
66
import io.grpc.okhttp.internal.ConnectionSpec;
67
import io.grpc.okhttp.internal.Credentials;
68
import io.grpc.okhttp.internal.OkHostnameVerifier;
69
import io.grpc.okhttp.internal.StatusLine;
70
import io.grpc.okhttp.internal.framed.ErrorCode;
71
import io.grpc.okhttp.internal.framed.FrameReader;
72
import io.grpc.okhttp.internal.framed.FrameWriter;
73
import io.grpc.okhttp.internal.framed.Header;
74
import io.grpc.okhttp.internal.framed.HeadersMode;
75
import io.grpc.okhttp.internal.framed.Http2;
76
import io.grpc.okhttp.internal.framed.Settings;
77
import io.grpc.okhttp.internal.framed.Variant;
78
import io.grpc.okhttp.internal.proxy.HttpUrl;
79
import io.grpc.okhttp.internal.proxy.Request;
80
import io.perfmark.PerfMark;
81
import java.io.EOFException;
82
import java.io.IOException;
83
import java.lang.reflect.InvocationTargetException;
84
import java.lang.reflect.Method;
85
import java.net.InetSocketAddress;
86
import java.net.Socket;
87
import java.net.URI;
88
import java.security.GeneralSecurityException;
89
import java.security.KeyStore;
90
import java.security.cert.Certificate;
91
import java.security.cert.X509Certificate;
92
import java.util.Collections;
93
import java.util.Deque;
94
import java.util.EnumMap;
95
import java.util.HashMap;
96
import java.util.Iterator;
97
import java.util.LinkedHashMap;
98
import java.util.LinkedList;
99
import java.util.List;
100
import java.util.Locale;
101
import java.util.Map;
102
import java.util.Random;
103
import java.util.concurrent.BrokenBarrierException;
104
import java.util.concurrent.CountDownLatch;
105
import java.util.concurrent.CyclicBarrier;
106
import java.util.concurrent.Executor;
107
import java.util.concurrent.ScheduledExecutorService;
108
import java.util.concurrent.TimeUnit;
109
import java.util.concurrent.TimeoutException;
110
import java.util.logging.Level;
111
import java.util.logging.Logger;
112
import javax.annotation.Nullable;
113
import javax.net.SocketFactory;
114
import javax.net.ssl.HostnameVerifier;
115
import javax.net.ssl.SSLParameters;
116
import javax.net.ssl.SSLPeerUnverifiedException;
117
import javax.net.ssl.SSLSession;
118
import javax.net.ssl.SSLSocket;
119
import javax.net.ssl.SSLSocketFactory;
120
import javax.net.ssl.TrustManager;
121
import javax.net.ssl.TrustManagerFactory;
122
import javax.net.ssl.X509TrustManager;
123
import okio.Buffer;
124
import okio.BufferedSink;
125
import okio.BufferedSource;
126
import okio.ByteString;
127
import okio.Okio;
128
import okio.Source;
129
import okio.Timeout;
130

131
/**
132
 * A okhttp-based {@link ConnectionClientTransport} implementation.
133
 */
134
class OkHttpClientTransport implements ConnectionClientTransport, TransportExceptionHandler,
135
      OutboundFlowController.Transport, ClientKeepAlivePinger.TransportWithDisconnectReason {
136
  private static final Map<ErrorCode, Status> ERROR_CODE_TO_STATUS = buildErrorCodeToStatusMap();
1 ✔
137
  private static final Logger log = Logger.getLogger(OkHttpClientTransport.class.getName());
1 ✔
138
  private static final String GRPC_ENABLE_PER_RPC_AUTHORITY_CHECK =
139
          "GRPC_ENABLE_PER_RPC_AUTHORITY_CHECK";
140
  static boolean enablePerRpcAuthorityCheck =
1 ✔
141
          GrpcUtil.getFlag(GRPC_ENABLE_PER_RPC_AUTHORITY_CHECK, false);
1 ✔
142
  private Socket sock;
143
  private SSLSession sslSession;
144

145
  private static Map<ErrorCode, Status> buildErrorCodeToStatusMap() {
146
    Map<ErrorCode, Status> errorToStatus = new EnumMap<>(ErrorCode.class);
1 ✔
147
    errorToStatus.put(ErrorCode.NO_ERROR,
1 ✔
148
        Status.INTERNAL.withDescription("No error: A GRPC status of OK should have been sent"));
1 ✔
149
    errorToStatus.put(ErrorCode.PROTOCOL_ERROR,
1 ✔
150
        Status.INTERNAL.withDescription("Protocol error"));
1 ✔
151
    errorToStatus.put(ErrorCode.INTERNAL_ERROR,
1 ✔
152
        Status.INTERNAL.withDescription("Internal error"));
1 ✔
153
    errorToStatus.put(ErrorCode.FLOW_CONTROL_ERROR,
1 ✔
154
        Status.INTERNAL.withDescription("Flow control error"));
1 ✔
155
    errorToStatus.put(ErrorCode.STREAM_CLOSED,
1 ✔
156
        Status.INTERNAL.withDescription("Stream closed"));
1 ✔
157
    errorToStatus.put(ErrorCode.FRAME_TOO_LARGE,
1 ✔
158
        Status.INTERNAL.withDescription("Frame too large"));
1 ✔
159
    errorToStatus.put(ErrorCode.REFUSED_STREAM,
1 ✔
160
        Status.UNAVAILABLE.withDescription("Refused stream"));
1 ✔
161
    errorToStatus.put(ErrorCode.CANCEL,
1 ✔
162
        Status.CANCELLED.withDescription("Cancelled"));
1 ✔
163
    errorToStatus.put(ErrorCode.COMPRESSION_ERROR,
1 ✔
164
        Status.INTERNAL.withDescription("Compression error"));
1 ✔
165
    errorToStatus.put(ErrorCode.CONNECT_ERROR,
1 ✔
166
        Status.INTERNAL.withDescription("Connect error"));
1 ✔
167
    errorToStatus.put(ErrorCode.ENHANCE_YOUR_CALM,
1 ✔
168
        Status.RESOURCE_EXHAUSTED.withDescription("Enhance your calm"));
1 ✔
169
    errorToStatus.put(ErrorCode.INADEQUATE_SECURITY,
1 ✔
170
        Status.PERMISSION_DENIED.withDescription("Inadequate security"));
1 ✔
171
    return Collections.unmodifiableMap(errorToStatus);
1 ✔
172
  }
173

174
  private static final Class<?> x509ExtendedTrustManagerClass;
175
  private static final Method checkServerTrustedMethod;
176

177
  static {
178
    Class<?> x509ExtendedTrustManagerClass1 = null;
1 ✔
179
    Method checkServerTrustedMethod1 = null;
1 ✔
180
    try {
181
      x509ExtendedTrustManagerClass1 = Class.forName("javax.net.ssl.X509ExtendedTrustManager");
1 ✔
182
      checkServerTrustedMethod1 = x509ExtendedTrustManagerClass1.getMethod("checkServerTrusted",
1 ✔
183
              X509Certificate[].class, String.class, Socket.class);
184
    } catch (ClassNotFoundException e) {
×
185
      // Per-rpc authority override via call options will be disallowed.
186
    } catch (NoSuchMethodException e) {
×
187
      // Should never happen since X509ExtendedTrustManager was introduced in Android API level 24
188
      // along with checkServerTrusted.
189
    }
1 ✔
190
    x509ExtendedTrustManagerClass = x509ExtendedTrustManagerClass1;
1 ✔
191
    checkServerTrustedMethod = checkServerTrustedMethod1;
1 ✔
192
  }
1 ✔
193

194
  private final InetSocketAddress address;
195
  private final String defaultAuthority;
196
  private final String userAgent;
197
  private final Random random = new Random();
1 ✔
198
  // Returns new unstarted stopwatches
199
  private final Supplier<Stopwatch> stopwatchFactory;
200
  private final int initialWindowSize;
201
  private final Variant variant;
202
  private Listener listener;
203
  @GuardedBy("lock")
204
  private ExceptionHandlingFrameWriter frameWriter;
205
  private OutboundFlowController outboundFlow;
206
  private final Object lock = new Object();
1 ✔
207
  private final InternalLogId logId;
208
  @GuardedBy("lock")
209
  private int nextStreamId;
210
  @GuardedBy("lock")
1 ✔
211
  private final Map<Integer, OkHttpClientStream> streams = new HashMap<>();
212
  private final Executor executor;
213
  // Wrap on executor, to guarantee some operations be executed serially.
214
  private final SerializingExecutor serializingExecutor;
215
  private final ScheduledExecutorService scheduler;
216
  private final int maxMessageSize;
217
  private int connectionUnacknowledgedBytesRead;
218
  private ClientFrameHandler clientFrameHandler;
219
  // Caution: Not synchronized, new value can only be safely read after the connection is complete.
220
  private Attributes attributes;
221
  /**
222
   * Indicates the transport is in go-away state: no new streams will be processed, but existing
223
   * streams may continue.
224
   */
225
  @GuardedBy("lock")
226
  private Status goAwayStatus;
227
  @GuardedBy("lock")
228
  private boolean goAwaySent;
229
  @GuardedBy("lock")
230
  private Http2Ping ping;
231
  @GuardedBy("lock")
232
  private boolean stopped;
233
  @GuardedBy("lock")
234
  private boolean hasStream;
235
  private final SocketFactory socketFactory;
236
  private SSLSocketFactory sslSocketFactory;
237
  private HostnameVerifier hostnameVerifier;
238
  private Socket socket;
239
  @GuardedBy("lock")
1 ✔
240
  private int maxConcurrentStreams = 0;
241
  @SuppressWarnings("JdkObsolete") // Usage is bursty; want low memory usage when empty
1 ✔
242
  @GuardedBy("lock")
243
  private final Deque<OkHttpClientStream> pendingStreams = new LinkedList<>();
244
  private final ConnectionSpec connectionSpec;
245
  private KeepAliveManager keepAliveManager;
246
  private boolean enableKeepAlive;
247
  private long keepAliveTimeNanos;
248
  private long keepAliveTimeoutNanos;
249
  private boolean keepAliveWithoutCalls;
250
  private final Runnable tooManyPingsRunnable;
251
  private final int maxInboundMetadataSize;
252
  private final boolean useGetForSafeMethods;
253
  @GuardedBy("lock")
254
  private final TransportTracer transportTracer;
255
  private final TrustManager x509TrustManager;
256

257
  @SuppressWarnings("serial")
258
  private static class LruCache<T> extends LinkedHashMap<String, T> {
259
    @Override
260
    protected boolean removeEldestEntry(Map.Entry<String, T> eldest) {
261
      return size() > 100;
1 ✔
262
    }
263
  }
264

265
  @GuardedBy("lock")
1 ✔
266
  private final Map<String, Status> authorityVerificationResults = new LruCache<>();
267

268
  @GuardedBy("lock")
1 ✔
269
  private final InUseStateAggregator<OkHttpClientStream> inUseState =
270
      new InUseStateAggregator<OkHttpClientStream>() {
1 ✔
271
        @Override
272
        protected void handleInUse() {
273
          listener.transportInUse(true);
1 ✔
274
        }
1 ✔
275

276
        @Override
277
        protected void handleNotInUse() {
278
          listener.transportInUse(false);
1 ✔
279
        }
1 ✔
280
      };
281
  @GuardedBy("lock")
282
  private InternalChannelz.Security securityInfo;
283

284
  @VisibleForTesting
285
  @Nullable
286
  final HttpConnectProxiedSocketAddress proxiedAddr;
287

288
  @VisibleForTesting
1 ✔
289
  int proxySocketTimeout = 30000;
290

291
  // The following fields should only be used for test.
292
  Runnable connectingCallback;
293
  SettableFuture<Void> connectedFuture;
294

295
  public OkHttpClientTransport(
296
          OkHttpTransportFactory transportFactory,
297
          InetSocketAddress address,
298
          String authority,
299
          @Nullable String userAgent,
300
          Attributes eagAttrs,
301
          @Nullable HttpConnectProxiedSocketAddress proxiedAddr,
302
          Runnable tooManyPingsRunnable,
303
          ChannelCredentials channelCredentials) {
304
    this(
1 ✔
305
        transportFactory,
306
        address,
307
        authority,
308
        userAgent,
309
        eagAttrs,
310
        GrpcUtil.STOPWATCH_SUPPLIER,
311
        new Http2(),
312
        proxiedAddr,
313
        tooManyPingsRunnable,
314
        channelCredentials);
315
  }
1 ✔
316

317
  private OkHttpClientTransport(
318
          OkHttpTransportFactory transportFactory,
319
          InetSocketAddress address,
320
          String authority,
321
          @Nullable String userAgent,
322
          Attributes eagAttrs,
323
          Supplier<Stopwatch> stopwatchFactory,
324
          Variant variant,
325
          @Nullable HttpConnectProxiedSocketAddress proxiedAddr,
326
          Runnable tooManyPingsRunnable,
327
          ChannelCredentials channelCredentials) {
1 ✔
328
    this.address = Preconditions.checkNotNull(address, "address");
1 ✔
329
    this.defaultAuthority = authority;
1 ✔
330
    this.maxMessageSize = transportFactory.maxMessageSize;
1 ✔
331
    this.initialWindowSize = transportFactory.flowControlWindow;
1 ✔
332
    this.executor = Preconditions.checkNotNull(transportFactory.executor, "executor");
1 ✔
333
    serializingExecutor = new SerializingExecutor(transportFactory.executor);
1 ✔
334
    this.scheduler = Preconditions.checkNotNull(
1 ✔
335
        transportFactory.scheduledExecutorService, "scheduledExecutorService");
336
    // Client initiated streams are odd, server initiated ones are even. Server should not need to
337
    // use it. We start clients at 3 to avoid conflicting with HTTP negotiation.
338
    nextStreamId = 3;
1 ✔
339
    this.socketFactory = transportFactory.socketFactory == null
1 ✔
340
        ? SocketFactory.getDefault() : transportFactory.socketFactory;
1 ✔
341
    this.sslSocketFactory = transportFactory.sslSocketFactory;
1 ✔
342
    this.hostnameVerifier = transportFactory.hostnameVerifier != null
1 ✔
343
        ? transportFactory.hostnameVerifier : OkHostnameVerifier.INSTANCE;
1 ✔
344
    this.connectionSpec = Preconditions.checkNotNull(
1 ✔
345
        transportFactory.connectionSpec, "connectionSpec");
346
    this.stopwatchFactory = Preconditions.checkNotNull(stopwatchFactory, "stopwatchFactory");
1 ✔
347
    this.variant = Preconditions.checkNotNull(variant, "variant");
1 ✔
348
    this.userAgent = GrpcUtil.getGrpcUserAgent("okhttp", userAgent);
1 ✔
349
    this.proxiedAddr = proxiedAddr;
1 ✔
350
    this.tooManyPingsRunnable =
1 ✔
351
        Preconditions.checkNotNull(tooManyPingsRunnable, "tooManyPingsRunnable");
1 ✔
352
    this.maxInboundMetadataSize = transportFactory.maxInboundMetadataSize;
1 ✔
353
    this.transportTracer = transportFactory.transportTracerFactory.create();
1 ✔
354
    this.logId = InternalLogId.allocate(getClass(), address.toString());
1 ✔
355
    this.attributes = Attributes.newBuilder()
1 ✔
356
        .set(GrpcAttributes.ATTR_CLIENT_EAG_ATTRS, eagAttrs).build();
1 ✔
357
    this.useGetForSafeMethods = transportFactory.useGetForSafeMethods;
1 ✔
358
    initTransportTracer();
1 ✔
359
    TrustManager tempX509TrustManager;
360
    if (channelCredentials instanceof TlsChannelCredentials
1 ✔
361
        && x509ExtendedTrustManagerClass != null) {
362
      try {
363
        tempX509TrustManager = getTrustManager(
1 ✔
364
                (TlsChannelCredentials) channelCredentials);
365
      } catch (GeneralSecurityException e) {
×
366
        tempX509TrustManager = null;
×
367
        log.log(Level.WARNING, "Obtaining X509ExtendedTrustManager for the transport failed."
×
368
            + "Per-rpc authority overrides will be disallowed.", e);
369
      }
1 ✔
370
    } else {
371
      tempX509TrustManager = null;
1 ✔
372
    }
373
    x509TrustManager = tempX509TrustManager;
1 ✔
374
  }
1 ✔
375

376
  /**
377
   * Create a transport connected to a fake peer for test.
378
   */
379
  @SuppressWarnings("AddressSelection") // An IP address always returns one address
380
  @VisibleForTesting
381
  OkHttpClientTransport(
382
      OkHttpTransportFactory transportFactory,
383
      String userAgent,
384
      Supplier<Stopwatch> stopwatchFactory,
385
      Variant variant,
386
      @Nullable Runnable connectingCallback,
387
      SettableFuture<Void> connectedFuture,
388
      Runnable tooManyPingsRunnable) {
389
    this(
1 ✔
390
        transportFactory,
391
        new InetSocketAddress("127.0.0.1", 80),
392
        "notarealauthority:80",
393
        userAgent,
394
        Attributes.EMPTY,
395
        stopwatchFactory,
396
        variant,
397
        null,
398
        tooManyPingsRunnable,
399
        null);
400
    this.connectingCallback = connectingCallback;
1 ✔
401
    this.connectedFuture = Preconditions.checkNotNull(connectedFuture, "connectedFuture");
1 ✔
402
  }
1 ✔
403

404
  // sslSocketFactory is set to null when use plaintext.
405
  boolean isUsingPlaintext() {
406
    return sslSocketFactory == null;
1 ✔
407
  }
408

409
  private void initTransportTracer() {
410
    synchronized (lock) { // to make @GuardedBy linter happy
1 ✔
411
      transportTracer.setFlowControlWindowReader(new TransportTracer.FlowControlReader() {
1 ✔
412
        @Override
413
        public TransportTracer.FlowControlWindows read() {
414
          synchronized (lock) {
1 ✔
415
            long local = outboundFlow == null ? -1 : outboundFlow.windowUpdate(null, 0);
1 ✔
416
            // connectionUnacknowledgedBytesRead is only readable by ClientFrameHandler, so we
417
            // provide a lower bound.
418
            long remote = (long) (initialWindowSize * DEFAULT_WINDOW_UPDATE_RATIO);
1 ✔
419
            return new TransportTracer.FlowControlWindows(local, remote);
1 ✔
420
          }
421
        }
422
      });
423
    }
1 ✔
424
  }
1 ✔
425

426
  /**
427
   * Enable keepalive with custom delay and timeout.
428
   */
429
  void enableKeepAlive(boolean enable, long keepAliveTimeNanos,
430
      long keepAliveTimeoutNanos, boolean keepAliveWithoutCalls) {
431
    enableKeepAlive = enable;
×
432
    this.keepAliveTimeNanos = keepAliveTimeNanos;
×
433
    this.keepAliveTimeoutNanos = keepAliveTimeoutNanos;
×
434
    this.keepAliveWithoutCalls = keepAliveWithoutCalls;
×
435
  }
×
436

437
  @Override
438
  public void ping(final PingCallback callback, Executor executor) {
439
    long data = 0;
1 ✔
440
    Http2Ping p;
441
    boolean writePing;
442
    synchronized (lock) {
1 ✔
443
      checkState(frameWriter != null);
1 ✔
444
      if (stopped) {
1 ✔
445
        Http2Ping.notifyFailed(callback, executor, getPingFailure());
1 ✔
446
        return;
1 ✔
447
      }
448
      if (ping != null) {
1 ✔
449
        // we only allow one outstanding ping at a time, so just add the callback to
450
        // any outstanding operation
451
        p = ping;
1 ✔
452
        writePing = false;
1 ✔
453
      } else {
454
        // set outstanding operation and then write the ping after releasing lock
455
        data = random.nextLong();
1 ✔
456
        Stopwatch stopwatch = stopwatchFactory.get();
1 ✔
457
        stopwatch.start();
1 ✔
458
        p = ping = new Http2Ping(data, stopwatch);
1 ✔
459
        writePing = true;
1 ✔
460
        transportTracer.reportKeepAliveSent();
1 ✔
461
      }
462
      if (writePing) {
1 ✔
463
        frameWriter.ping(false, (int) (data >>> 32), (int) data);
1 ✔
464
      }
465
    }
1 ✔
466
    // If transport concurrently failed/stopped since we released the lock above, this could
467
    // immediately invoke callback (which we shouldn't do while holding a lock)
468
    p.addCallback(callback, executor);
1 ✔
469
  }
1 ✔
470

471
  @Override
472
  public OkHttpClientStream newStream(
473
      MethodDescriptor<?, ?> method, Metadata headers, CallOptions callOptions,
474
      ClientStreamTracer[] tracers) {
475
    Preconditions.checkNotNull(method, "method");
1 ✔
476
    Preconditions.checkNotNull(headers, "headers");
1 ✔
477
    StatsTraceContext statsTraceContext =
1 ✔
478
        StatsTraceContext.newClientContext(tracers, getAttributes(), headers);
1 ✔
479

480
    // FIXME: it is likely wrong to pass the transportTracer here as it'll exit the lock's scope
481
    synchronized (lock) { // to make @GuardedBy linter happy
1 ✔
482
      return new OkHttpClientStream(
1 ✔
483
          method,
484
          headers,
485
          frameWriter,
486
          OkHttpClientTransport.this,
487
          outboundFlow,
488
          lock,
489
          maxMessageSize,
490
          initialWindowSize,
491
          defaultAuthority,
492
          userAgent,
493
          statsTraceContext,
494
          transportTracer,
495
          callOptions,
496
          useGetForSafeMethods);
497
    }
498
  }
499

500
  private TrustManager getTrustManager(TlsChannelCredentials tlsCreds)
501
      throws GeneralSecurityException {
502
    TrustManager[] tm;
503
    // Using the same way of creating TrustManager from OkHttpChannelBuilder.sslSocketFactoryFrom()
504
    if (tlsCreds.getTrustManagers() != null) {
1 ✔
505
      tm = tlsCreds.getTrustManagers().toArray(new TrustManager[0]);
1 ✔
506
    } else if (tlsCreds.getRootCertificates() != null) {
1 ✔
507
      tm = CertificateUtils.createTrustManager(tlsCreds.getRootCertificates());
1 ✔
508
    } else { // else use system default
509
      TrustManagerFactory tmf = TrustManagerFactory.getInstance(
1 ✔
510
          TrustManagerFactory.getDefaultAlgorithm());
1 ✔
511
      tmf.init((KeyStore) null);
1 ✔
512
      tm = tmf.getTrustManagers();
1 ✔
513
    }
514
    for (TrustManager trustManager: tm) {
1 ✔
515
      if (trustManager instanceof X509TrustManager) {
1 ✔
516
        return trustManager;
1 ✔
517
      }
518
    }
519
    return null;
×
520
  }
521

522
  @GuardedBy("lock")
523
  void streamReadyToStart(OkHttpClientStream clientStream, String authority) {
524
    if (goAwayStatus != null) {
1 ✔
525
      clientStream.transportState().transportReportStatus(
1 ✔
526
          goAwayStatus, RpcProgress.MISCARRIED, true, new Metadata());
527
    } else {
528
      if (socket instanceof SSLSocket && !authority.equals(defaultAuthority)) {
1 ✔
529
        Status authorityVerificationResult;
530
        if (authorityVerificationResults.containsKey(authority)) {
1 ✔
531
          authorityVerificationResult = authorityVerificationResults.get(authority);
×
532
        } else {
533
          authorityVerificationResult = verifyAuthority(authority);
1 ✔
534
          authorityVerificationResults.put(authority, authorityVerificationResult);
1 ✔
535
        }
536
        if (!authorityVerificationResult.isOk()) {
1 ✔
537
          if (enablePerRpcAuthorityCheck) {
1 ✔
538
            clientStream.transportState().transportReportStatus(
1 ✔
539
                    authorityVerificationResult, RpcProgress.PROCESSED, true, new Metadata());
540
            return;
1 ✔
541
          }
542
        }
543
      }
544
      if (streams.size() >= maxConcurrentStreams) {
1 ✔
545
        pendingStreams.add(clientStream);
1 ✔
546
        setInUse(clientStream);
1 ✔
547
      } else {
548
        startStream(clientStream);
1 ✔
549
      }
550
    }
551
  }
1 ✔
552

553
  private Status verifyAuthority(String authority) {
554
    Status authorityVerificationResult;
555
    if (hostnameVerifier.verify(authority, ((SSLSocket) socket).getSession())) {
1 ✔
556
      authorityVerificationResult = Status.OK;
1 ✔
557
    } else {
558
      authorityVerificationResult = Status.UNAVAILABLE.withDescription(String.format(
1 ✔
559
              "HostNameVerifier verification failed for authority '%s'",
560
              authority));
561
    }
562
    if (!authorityVerificationResult.isOk() && !enablePerRpcAuthorityCheck) {
1 ✔
563
      log.log(Level.WARNING, String.format("HostNameVerifier verification failed for "
1 ✔
564
                      + "authority '%s'. This will be an error in the future.",
565
              authority));
566
    }
567
    if (authorityVerificationResult.isOk()) {
1 ✔
568
      // The status is trivially assigned in this case, but we are still making use of the
569
      // cache to keep track that a warning log had been logged for the authority when
570
      // enablePerRpcAuthorityCheck is false. When we permanently enable the feature, the
571
      // status won't need to be cached for case when x509TrustManager is null.
572
      if (x509TrustManager == null) {
1 ✔
573
        authorityVerificationResult = Status.UNAVAILABLE.withDescription(
1 ✔
574
                String.format("Could not verify authority '%s' for the rpc with no "
1 ✔
575
                                + "X509TrustManager available",
576
                        authority));
577
      } else if (x509ExtendedTrustManagerClass.isInstance(x509TrustManager)) {
1 ✔
578
        try {
579
          Certificate[] peerCertificates = sslSession.getPeerCertificates();
1 ✔
580
          X509Certificate[] x509PeerCertificates =
1 ✔
581
                  new X509Certificate[peerCertificates.length];
582
          for (int i = 0; i < peerCertificates.length; i++) {
1 ✔
583
            x509PeerCertificates[i] = (X509Certificate) peerCertificates[i];
1 ✔
584
          }
585
          checkServerTrustedMethod.invoke(x509TrustManager, x509PeerCertificates,
1 ✔
586
                  "RSA", new SslSocketWrapper((SSLSocket) socket, authority));
587
          authorityVerificationResult = Status.OK;
1 ✔
588
        } catch (SSLPeerUnverifiedException | InvocationTargetException
1 ✔
589
                 | IllegalAccessException e) {
590
          authorityVerificationResult = Status.UNAVAILABLE.withCause(e).withDescription(
1 ✔
591
                  "Peer verification failed");
592
        }
1 ✔
593
        if (authorityVerificationResult.getCause() != null) {
1 ✔
594
          log.log(Level.WARNING, authorityVerificationResult.getDescription()
1 ✔
595
                          + ". This will be an error in the future.",
596
                  authorityVerificationResult.getCause());
1 ✔
597
        } else {
598
          log.log(Level.WARNING, authorityVerificationResult.getDescription()
1 ✔
599
                  + ". This will be an error in the future.");
600
        }
601
      }
602
    }
603
    return authorityVerificationResult;
1 ✔
604
  }
605

606
  @SuppressWarnings("GuardedBy")
607
  @GuardedBy("lock")
608
  private void startStream(OkHttpClientStream stream) {
609
    checkState(
1 ✔
610
        stream.transportState().id() == OkHttpClientStream.ABSENT_ID, "StreamId already assigned");
1 ✔
611
    streams.put(nextStreamId, stream);
1 ✔
612
    setInUse(stream);
1 ✔
613
    // TODO(b/145386688): This access should be guarded by 'stream.transportState().lock'; instead
614
    // found: 'this.lock'
615
    stream.transportState().start(nextStreamId);
1 ✔
616
    // For unary and server streaming, there will be a data frame soon, no need to flush the header.
617
    if ((stream.getType() != MethodType.UNARY && stream.getType() != MethodType.SERVER_STREAMING)
1 ✔
618
        || stream.useGet()) {
1 ✔
619
      frameWriter.flush();
1 ✔
620
    }
621
    if (nextStreamId >= Integer.MAX_VALUE - 2) {
1 ✔
622
      // Make sure nextStreamId greater than all used id, so that mayHaveCreatedStream() performs
623
      // correctly.
624
      nextStreamId = Integer.MAX_VALUE;
1 ✔
625
      startGoAway(Integer.MAX_VALUE, ErrorCode.NO_ERROR,
1 ✔
626
          Status.UNAVAILABLE.withDescription("Stream ids exhausted"));
1 ✔
627
    } else {
628
      nextStreamId += 2;
1 ✔
629
    }
630
  }
1 ✔
631

632
  /**
633
   * Starts pending streams, returns true if at least one pending stream is started.
634
   */
635
  @GuardedBy("lock")
636
  private boolean startPendingStreams() {
637
    boolean hasStreamStarted = false;
1 ✔
638
    while (!pendingStreams.isEmpty() && streams.size() < maxConcurrentStreams) {
1 ✔
639
      OkHttpClientStream stream = pendingStreams.poll();
1 ✔
640
      startStream(stream);
1 ✔
641
      hasStreamStarted = true;
1 ✔
642
    }
1 ✔
643
    return hasStreamStarted;
1 ✔
644
  }
645

646
  /**
647
   * Removes given pending stream, used when a pending stream is cancelled.
648
   */
649
  @GuardedBy("lock")
650
  void removePendingStream(OkHttpClientStream pendingStream) {
651
    pendingStreams.remove(pendingStream);
1 ✔
652
    maybeClearInUse(pendingStream);
1 ✔
653
  }
1 ✔
654

655
  @Override
656
  public Runnable start(Listener listener) {
657
    this.listener = Preconditions.checkNotNull(listener, "listener");
1 ✔
658

659
    if (enableKeepAlive) {
1 ✔
660
      keepAliveManager = new KeepAliveManager(
×
661
          new ClientKeepAlivePinger(this), scheduler, keepAliveTimeNanos, keepAliveTimeoutNanos,
662
          keepAliveWithoutCalls);
663
      keepAliveManager.onTransportStarted();
×
664
    }
665

666
    int maxQueuedControlFrames = 10000;
1 ✔
667
    final AsyncSink asyncSink = AsyncSink.sink(serializingExecutor, this, maxQueuedControlFrames);
1 ✔
668
    FrameWriter rawFrameWriter = asyncSink.limitControlFramesWriter(
1 ✔
669
        variant.newWriter(Okio.buffer(asyncSink), true));
1 ✔
670

671
    synchronized (lock) {
1 ✔
672
      // Handle FrameWriter exceptions centrally, since there are many callers. Note that errors
673
      // coming from rawFrameWriter are generally broken invariants/bugs, as AsyncSink does not
674
      // propagate syscall errors through the FrameWriter. But we handle the AsyncSink failures with
675
      // the same TransportExceptionHandler instance so it is all mixed back together.
676
      frameWriter = new ExceptionHandlingFrameWriter(this, rawFrameWriter);
1 ✔
677
      outboundFlow = new OutboundFlowController(this, frameWriter);
1 ✔
678
    }
1 ✔
679
    final CountDownLatch latch = new CountDownLatch(1);
1 ✔
680
    final CountDownLatch latchForExtraThread = new CountDownLatch(1);
1 ✔
681
    // The transport needs up to two threads to function once started,
682
    // but only needs one during handshaking. Start another thread during handshaking
683
    // to make sure there's still a free thread available. If the number of threads is exhausted,
684
    // it is better to kill the transport than for all the transports to hang unable to send.
685
    CyclicBarrier barrier = new CyclicBarrier(2);
1 ✔
686
    // Connecting in the serializingExecutor, so that some stream operations like synStream
687
    // will be executed after connected.
688

689
    serializingExecutor.execute(new Runnable() {
1 ✔
690
      @Override
691
      public void run() {
692
        // Use closed source on failure so that the reader immediately shuts down.
693
        BufferedSource source = Okio.buffer(new Source() {
1 ✔
694
          @Override
695
          public long read(Buffer sink, long byteCount) {
696
            return -1;
1 ✔
697
          }
698

699
          @Override
700
          public Timeout timeout() {
701
            return Timeout.NONE;
×
702
          }
703

704
          @Override
705
          public void close() {
706
          }
1 ✔
707
        });
708
        try {
709
          // This is a hack to make sure the connection preface and initial settings to be sent out
710
          // without blocking the start. By doing this essentially prevents potential deadlock when
711
          // network is not available during startup while another thread holding lock to send the
712
          // initial preface.
713
          try {
714
            latch.await();
1 ✔
715
            barrier.await(1000, TimeUnit.MILLISECONDS);
1 ✔
716
          } catch (InterruptedException e) {
×
717
            Thread.currentThread().interrupt();
×
718
          } catch (TimeoutException | BrokenBarrierException e) {
1 ✔
719
            startGoAway(0, ErrorCode.INTERNAL_ERROR, Status.UNAVAILABLE
1 ✔
720
                .withDescription("Timed out waiting for second handshake thread. "
1 ✔
721
                    + "The transport executor pool may have run out of threads"));
722
            return;
1 ✔
723
          }
1 ✔
724

725
          if (proxiedAddr == null) {
1 ✔
726
            sock = socketFactory.createSocket(address.getAddress(), address.getPort());
1 ✔
727
          } else {
728
            if (proxiedAddr.getProxyAddress() instanceof InetSocketAddress) {
1 ✔
729
              sock = createHttpProxySocket(
1 ✔
730
                  proxiedAddr.getTargetAddress(),
1 ✔
731
                  (InetSocketAddress) proxiedAddr.getProxyAddress(),
1 ✔
732
                  proxiedAddr.getUsername(),
1 ✔
733
                  proxiedAddr.getPassword()
1 ✔
734
              );
735
            } else {
736
              throw Status.INTERNAL.withDescription(
×
737
                  "Unsupported SocketAddress implementation "
738
                  + proxiedAddr.getProxyAddress().getClass()).asException();
×
739
            }
740
          }
741
          if (sslSocketFactory != null) {
1 ✔
742
            sock.setSoTimeout(proxySocketTimeout);
1 ✔
743
            SSLSocket sslSocket = OkHttpTlsUpgrader.upgrade(
1 ✔
744
                sslSocketFactory, hostnameVerifier, sock, getOverridenHost(), getOverridenPort(),
1 ✔
745
                connectionSpec);
1 ✔
746
            // As the socket will be used for RPCs from here on, we want the socket 
747
            // timeout back to zero.
748
            sock.setSoTimeout(0);
1 ✔
749
            sslSession = sslSocket.getSession();
1 ✔
750
            sock = sslSocket;
1 ✔
751
          }
752
          sock.setTcpNoDelay(true);
1 ✔
753
          source = Okio.buffer(Okio.source(sock));
1 ✔
754
          asyncSink.becomeConnected(Okio.sink(sock), sock);
1 ✔
755

756
          // The return value of OkHttpTlsUpgrader.upgrade is an SSLSocket that has this info
757
          attributes = attributes.toBuilder()
1 ✔
758
              .set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, sock.getRemoteSocketAddress())
1 ✔
759
              .set(Grpc.TRANSPORT_ATTR_LOCAL_ADDR, sock.getLocalSocketAddress())
1 ✔
760
              .set(Grpc.TRANSPORT_ATTR_SSL_SESSION, sslSession)
1 ✔
761
              .set(GrpcAttributes.ATTR_SECURITY_LEVEL,
1 ✔
762
                  sslSession == null ? SecurityLevel.NONE : SecurityLevel.PRIVACY_AND_INTEGRITY)
1 ✔
763
              .build();
1 ✔
764
        } catch (StatusException e) {
1 ✔
765
          startGoAway(0, ErrorCode.INTERNAL_ERROR, e.getStatus());
1 ✔
766
          return;
1 ✔
767
        } catch (Exception e) {
1 ✔
768
          GrpcUtil.closeQuietly(sock);
1 ✔
769
          onException(e);
1 ✔
770
          return;
1 ✔
771
        } finally {
772
          clientFrameHandler = new ClientFrameHandler(variant.newReader(source, true));
1 ✔
773
          latchForExtraThread.countDown();
1 ✔
774
        }
775
        synchronized (lock) {
1 ✔
776
          socket = Preconditions.checkNotNull(sock, "socket");
1 ✔
777
          if (sslSession != null) {
1 ✔
778
            securityInfo = new InternalChannelz.Security(new InternalChannelz.Tls(sslSession));
1 ✔
779
          }
780
        }
1 ✔
781
      }
1 ✔
782
    });
783

784
    executor.execute(new Runnable() {
1 ✔
785
      @Override
786
      public void run() {
787
        try {
788
          barrier.await(1000, TimeUnit.MILLISECONDS);
1 ✔
789
          latchForExtraThread.await();
1 ✔
790
        } catch (BrokenBarrierException | TimeoutException e) {
1 ✔
791
          // Something bad happened, maybe too few threads available!
792
          // This will be handled in the handshake thread.
793
        } catch (InterruptedException e) {
1 ✔
794
          Thread.currentThread().interrupt();
1 ✔
795
        }
1 ✔
796
      }
1 ✔
797
    });
798
    // Schedule to send connection preface & settings before any other write.
799
    try {
800
      sendConnectionPrefaceAndSettings();
1 ✔
801
    } finally {
802
      latch.countDown();
1 ✔
803
    }
804

805
    serializingExecutor.execute(new Runnable() {
1 ✔
806
      @Override
807
      public void run() {
808
        if (connectingCallback != null) {
1 ✔
809
          connectingCallback.run();
1 ✔
810
        }
811
        synchronized (lock) {
1 ✔
812
          maxConcurrentStreams = Integer.MAX_VALUE;
1 ✔
813
          checkState(pendingStreams.isEmpty(),
1 ✔
814
              "Pending streams detected during transport start."
815
                  + " RPCs should not be started before transport is ready.");
816
        }
1 ✔
817
        // ClientFrameHandler need to be started after connectionPreface / settings, otherwise it
818
        // may send goAway immediately.
819
        executor.execute(clientFrameHandler);
1 ✔
820
        if (connectedFuture != null) {
1 ✔
821
          connectedFuture.set(null);
1 ✔
822
        }
823
      }
1 ✔
824
    });
825
    return null;
1 ✔
826
  }
827

828
  /**
829
   * Should only be called once when the transport is first established.
830
   */
831
  private void sendConnectionPrefaceAndSettings() {
832
    synchronized (lock) {
1 ✔
833
      frameWriter.connectionPreface();
1 ✔
834
      Settings settings = new Settings();
1 ✔
835
      OkHttpSettingsUtil.set(settings, OkHttpSettingsUtil.INITIAL_WINDOW_SIZE, initialWindowSize);
1 ✔
836
      frameWriter.settings(settings);
1 ✔
837
      if (initialWindowSize > DEFAULT_WINDOW_SIZE) {
1 ✔
838
        frameWriter.windowUpdate(
1 ✔
839
                Utils.CONNECTION_STREAM_ID, initialWindowSize - DEFAULT_WINDOW_SIZE);
840
      }
841
    }
1 ✔
842
  }
1 ✔
843

844
  private Socket createHttpProxySocket(InetSocketAddress address, InetSocketAddress proxyAddress,
845
      String proxyUsername, String proxyPassword) throws StatusException {
846
    Socket sock = null;
1 ✔
847
    try {
848
      // The proxy address may not be resolved
849
      if (proxyAddress.getAddress() != null) {
1 ✔
850
        sock = socketFactory.createSocket(proxyAddress.getAddress(), proxyAddress.getPort());
1 ✔
851
      } else {
852
        sock =
×
853
            socketFactory.createSocket(proxyAddress.getHostName(), proxyAddress.getPort());
×
854
      }
855
      sock.setTcpNoDelay(true);
1 ✔
856
      // A socket timeout is needed because lost network connectivity while reading from the proxy,
857
      // can cause reading from the socket to hang.
858
      sock.setSoTimeout(proxySocketTimeout);
1 ✔
859

860
      Source source = Okio.source(sock);
1 ✔
861
      BufferedSink sink = Okio.buffer(Okio.sink(sock));
1 ✔
862

863
      // Prepare headers and request method line
864
      Request proxyRequest = createHttpProxyRequest(address, proxyUsername, proxyPassword);
1 ✔
865
      HttpUrl url = proxyRequest.httpUrl();
1 ✔
866
      String requestLine =
1 ✔
867
          String.format(Locale.US, "CONNECT %s:%d HTTP/1.1", url.host(), url.port());
1 ✔
868

869
      // Write request to socket
870
      sink.writeUtf8(requestLine).writeUtf8("\r\n");
1 ✔
871
      for (int i = 0, size = proxyRequest.headers().size(); i < size; i++) {
1 ✔
872
        sink.writeUtf8(proxyRequest.headers().name(i))
1 ✔
873
            .writeUtf8(": ")
1 ✔
874
            .writeUtf8(proxyRequest.headers().value(i))
1 ✔
875
            .writeUtf8("\r\n");
1 ✔
876
      }
877
      sink.writeUtf8("\r\n");
1 ✔
878
      // Flush buffer (flushes socket and sends request)
879
      sink.flush();
1 ✔
880

881
      // Read status line, check if 2xx was returned
882
      StatusLine statusLine = StatusLine.parse(readUtf8LineStrictUnbuffered(source));
1 ✔
883
      // Drain rest of headers
884
      while (!readUtf8LineStrictUnbuffered(source).equals("")) {}
1 ✔
885
      if (statusLine.code < 200 || statusLine.code >= 300) {
1 ✔
886
        Buffer body = new Buffer();
1 ✔
887
        try {
888
          sock.shutdownOutput();
1 ✔
889
          source.read(body, 1024);
1 ✔
890
        } catch (IOException ex) {
×
891
          body.writeUtf8("Unable to read body: " + ex.toString());
×
892
        }
1 ✔
893
        try {
894
          sock.close();
1 ✔
895
        } catch (IOException ignored) {
×
896
          // ignored
897
        }
1 ✔
898
        String message = String.format(
1 ✔
899
            Locale.US,
900
            "Response returned from proxy was not successful (expected 2xx, got %d %s). "
901
              + "Response body:\n%s",
902
            statusLine.code, statusLine.message, body.readUtf8());
1 ✔
903
        throw Status.UNAVAILABLE.withDescription(message).asException();
1 ✔
904
      }
905
      // As the socket will be used for RPCs from here on, we want the socket timeout back to zero.
906
      sock.setSoTimeout(0);
1 ✔
907
      return sock;
1 ✔
908
    } catch (IOException e) {
1 ✔
909
      if (sock != null) {
1 ✔
910
        GrpcUtil.closeQuietly(sock);
1 ✔
911
      }
912
      throw Status.UNAVAILABLE.withDescription("Failed trying to connect with proxy").withCause(e)
1 ✔
913
          .asException();
1 ✔
914
    }
915
  }
916

917
  private Request createHttpProxyRequest(InetSocketAddress address, String proxyUsername,
918
                                         String proxyPassword) {
919
    HttpUrl tunnelUrl = new HttpUrl.Builder()
1 ✔
920
        .scheme("https")
1 ✔
921
        .host(address.getHostName())
1 ✔
922
        .port(address.getPort())
1 ✔
923
        .build();
1 ✔
924

925
    Request.Builder request = new Request.Builder()
1 ✔
926
        .url(tunnelUrl)
1 ✔
927
        .header("Host", tunnelUrl.host() + ":" + tunnelUrl.port())
1 ✔
928
        .header("User-Agent", userAgent);
1 ✔
929

930
    // If we have proxy credentials, set them right away
931
    if (proxyUsername != null && proxyPassword != null) {
1 ✔
932
      request.header("Proxy-Authorization", Credentials.basic(proxyUsername, proxyPassword));
×
933
    }
934
    return request.build();
1 ✔
935
  }
936

937
  private static String readUtf8LineStrictUnbuffered(Source source) throws IOException {
938
    Buffer buffer = new Buffer();
1 ✔
939
    while (true) {
940
      if (source.read(buffer, 1) == -1) {
1 ✔
941
        throw new EOFException("\\n not found: " + buffer.readByteString().hex());
×
942
      }
943
      if (buffer.getByte(buffer.size() - 1) == '\n') {
1 ✔
944
        return buffer.readUtf8LineStrict();
1 ✔
945
      }
946
    }
947
  }
948

949
  @Override
950
  public String toString() {
951
    return MoreObjects.toStringHelper(this)
1 ✔
952
        .add("logId", logId.getId())
1 ✔
953
        .add("address", address)
1 ✔
954
        .toString();
1 ✔
955
  }
956

957
  @Override
958
  public InternalLogId getLogId() {
959
    return logId;
1 ✔
960
  }
961

962
  /**
963
   * Gets the overridden authority hostname.  If the authority is overridden to be an invalid
964
   * authority, uri.getHost() will (rightly) return null, since the authority is no longer
965
   * an actual service.  This method overrides the behavior for practical reasons.  For example,
966
   * if an authority is in the form "invalid_authority" (note the "_"), rather than return null,
967
   * we return the input.  This is because the return value, in conjunction with getOverridenPort,
968
   * are used by the SSL library to reconstruct the actual authority.  It /already/ has a
969
   * connection to the port, independent of this function.
970
   *
971
   * <p>Note: if the defaultAuthority has a port number in it and is also bad, this code will do
972
   * the wrong thing.  An example wrong behavior would be "invalid_host:443".   Registry based
973
   * authorities do not have ports, so this is even more wrong than before.  Sorry.
974
   */
975
  @VisibleForTesting
976
  String getOverridenHost() {
977
    URI uri = GrpcUtil.authorityToUri(defaultAuthority);
1 ✔
978
    if (uri.getHost() != null) {
1 ✔
979
      return uri.getHost();
1 ✔
980
    }
981

982
    return defaultAuthority;
1 ✔
983
  }
984

985
  @VisibleForTesting
986
  int getOverridenPort() {
987
    URI uri = GrpcUtil.authorityToUri(defaultAuthority);
1 ✔
988
    if (uri.getPort() != -1) {
1 ✔
989
      return uri.getPort();
×
990
    }
991

992
    return address.getPort();
1 ✔
993
  }
994

995
  @Override
996
  public void shutdown(Status reason) {
997
    synchronized (lock) {
1 ✔
998
      if (goAwayStatus != null) {
1 ✔
999
        return;
1 ✔
1000
      }
1001

1002
      goAwayStatus = reason;
1 ✔
1003
      listener.transportShutdown(goAwayStatus, SimpleDisconnectError.SUBCHANNEL_SHUTDOWN);
1 ✔
1004
      stopIfNecessary();
1 ✔
1005
    }
1 ✔
1006
  }
1 ✔
1007

1008
  @Override
1009
  public void shutdownNow(Status reason) {
1010
    shutdownNow(reason, SimpleDisconnectError.SUBCHANNEL_SHUTDOWN);
1 ✔
1011
  }
1 ✔
1012

1013
  @Override
1014
  public void shutdownNow(Status reason, DisconnectError disconnectError) {
1015
    shutdown(reason);
1 ✔
1016
    synchronized (lock) {
1 ✔
1017
      Iterator<Map.Entry<Integer, OkHttpClientStream>> it = streams.entrySet().iterator();
1 ✔
1018
      while (it.hasNext()) {
1 ✔
1019
        Map.Entry<Integer, OkHttpClientStream> entry = it.next();
1 ✔
1020
        it.remove();
1 ✔
1021
        entry.getValue().transportState().transportReportStatus(reason, false, new Metadata());
1 ✔
1022
        maybeClearInUse(entry.getValue());
1 ✔
1023
      }
1 ✔
1024

1025
      for (OkHttpClientStream stream : pendingStreams) {
1 ✔
1026
        // in cases such as the connection fails to ACK keep-alive, pending streams should have a
1027
        // chance to retry and be routed to another connection.
1028
        stream.transportState().transportReportStatus(
1 ✔
1029
            reason, RpcProgress.MISCARRIED, true, new Metadata());
1030
        maybeClearInUse(stream);
1 ✔
1031
      }
1 ✔
1032
      pendingStreams.clear();
1 ✔
1033

1034
      stopIfNecessary();
1 ✔
1035
    }
1 ✔
1036
  }
1 ✔
1037

1038
  @Override
1039
  public Attributes getAttributes() {
1040
    return attributes;
1 ✔
1041
  }
1042

1043
  /**
1044
   * Gets all active streams as an array.
1045
   */
1046
  @Override
1047
  public OutboundFlowController.StreamState[] getActiveStreams() {
1048
    synchronized (lock) {
1 ✔
1049
      OutboundFlowController.StreamState[] flowStreams =
1 ✔
1050
          new OutboundFlowController.StreamState[streams.size()];
1 ✔
1051
      int i = 0;
1 ✔
1052
      for (OkHttpClientStream stream : streams.values()) {
1 ✔
1053
        flowStreams[i++] = stream.transportState().getOutboundFlowState();
1 ✔
1054
      }
1 ✔
1055
      return flowStreams;
1 ✔
1056
    }
1057
  }
1058

1059
  @VisibleForTesting
1060
  ClientFrameHandler getHandler() {
1061
    return clientFrameHandler;
1 ✔
1062
  }
1063

1064
  @VisibleForTesting
1065
  SocketFactory getSocketFactory() {
1066
    return socketFactory;
1 ✔
1067
  }
1068

1069
  @VisibleForTesting
1070
  int getPendingStreamSize() {
1071
    synchronized (lock) {
1 ✔
1072
      return pendingStreams.size();
1 ✔
1073
    }
1074
  }
1075

1076
  @VisibleForTesting
1077
  void setNextStreamId(int nextStreamId) {
1078
    synchronized (lock) {
1 ✔
1079
      this.nextStreamId = nextStreamId;
1 ✔
1080
    }
1 ✔
1081
  }
1 ✔
1082

1083
  /**
1084
   * Finish all active streams due to an IOException, then close the transport.
1085
   */
1086
  @Override
1087
  public void onException(Throwable failureCause) {
1088
    Preconditions.checkNotNull(failureCause, "failureCause");
1 ✔
1089
    Status status = Status.UNAVAILABLE.withCause(failureCause);
1 ✔
1090
    startGoAway(0, ErrorCode.INTERNAL_ERROR, status);
1 ✔
1091
  }
1 ✔
1092

1093
  /**
1094
   * Send GOAWAY to the server, then finish all active streams and close the transport.
1095
   */
1096
  private void onError(ErrorCode errorCode, String moreDetail) {
1097
    startGoAway(0, errorCode, toGrpcStatus(errorCode).augmentDescription(moreDetail));
1 ✔
1098
  }
1 ✔
1099

1100
  private void startGoAway(int lastKnownStreamId, ErrorCode errorCode, Status status) {
1101
    synchronized (lock) {
1 ✔
1102
      if (goAwayStatus == null) {
1 ✔
1103
        goAwayStatus = status;
1 ✔
1104
        GrpcUtil.Http2Error http2Error;
1105
        if (errorCode == null) {
1 ✔
1106
          http2Error = GrpcUtil.Http2Error.NO_ERROR;
1 ✔
1107
        } else {
1108
          http2Error = GrpcUtil.Http2Error.forCode(errorCode.httpCode);
1 ✔
1109
        }
1110
        listener.transportShutdown(status, new GoAwayDisconnectError(http2Error));
1 ✔
1111
      }
1112
      if (errorCode != null && !goAwaySent) {
1 ✔
1113
        // Send GOAWAY with lastGoodStreamId of 0, since we don't expect any server-initiated
1114
        // streams. The GOAWAY is part of graceful shutdown.
1115
        goAwaySent = true;
1 ✔
1116
        frameWriter.goAway(0, errorCode, new byte[0]);
1 ✔
1117
      }
1118

1119
      Iterator<Map.Entry<Integer, OkHttpClientStream>> it = streams.entrySet().iterator();
1 ✔
1120
      while (it.hasNext()) {
1 ✔
1121
        Map.Entry<Integer, OkHttpClientStream> entry = it.next();
1 ✔
1122
        if (entry.getKey() > lastKnownStreamId) {
1 ✔
1123
          it.remove();
1 ✔
1124
          entry.getValue().transportState().transportReportStatus(
1 ✔
1125
              status, RpcProgress.REFUSED, false, new Metadata());
1126
          maybeClearInUse(entry.getValue());
1 ✔
1127
        }
1128
      }
1 ✔
1129

1130
      for (OkHttpClientStream stream : pendingStreams) {
1 ✔
1131
        stream.transportState().transportReportStatus(
1 ✔
1132
            status, RpcProgress.MISCARRIED, true, new Metadata());
1133
        maybeClearInUse(stream);
1 ✔
1134
      }
1 ✔
1135
      pendingStreams.clear();
1 ✔
1136

1137
      stopIfNecessary();
1 ✔
1138
    }
1 ✔
1139
  }
1 ✔
1140

1141
  /**
1142
   * Called when a stream is closed. We do things like:
1143
   * <ul>
1144
   * <li>Removing the stream from the map.
1145
   * <li>Optionally reporting the status.
1146
   * <li>Starting pending streams if we can.
1147
   * <li>Stopping the transport if this is the last live stream under a go-away status.
1148
   * </ul>
1149
   *
1150
   * @param streamId the Id of the stream.
1151
   * @param status the final status of this stream, null means no need to report.
1152
   * @param stopDelivery interrupt queued messages in the deframer
1153
   * @param errorCode reset the stream with this ErrorCode if not null.
1154
   * @param trailers the trailers received if not null
1155
   */
1156
  void finishStream(
1157
      int streamId,
1158
      @Nullable Status status,
1159
      RpcProgress rpcProgress,
1160
      boolean stopDelivery,
1161
      @Nullable ErrorCode errorCode,
1162
      @Nullable Metadata trailers) {
1163
    synchronized (lock) {
1 ✔
1164
      OkHttpClientStream stream = streams.remove(streamId);
1 ✔
1165
      if (stream != null) {
1 ✔
1166
        if (errorCode != null) {
1 ✔
1167
          frameWriter.rstStream(streamId, ErrorCode.CANCEL);
1 ✔
1168
        }
1169
        if (status != null) {
1 ✔
1170
          stream
1 ✔
1171
              .transportState()
1 ✔
1172
              .transportReportStatus(
1 ✔
1173
                  status,
1174
                  rpcProgress,
1175
                  stopDelivery,
1176
                  trailers != null ? trailers : new Metadata());
1 ✔
1177
        }
1178
        if (!startPendingStreams()) {
1 ✔
1179
          stopIfNecessary();
1 ✔
1180
        }
1181
        maybeClearInUse(stream);
1 ✔
1182
      }
1183
    }
1 ✔
1184
  }
1 ✔
1185

1186
  /**
1187
   * When the transport is in goAway state, we should stop it once all active streams finish.
1188
   */
1189
  @GuardedBy("lock")
1190
  private void stopIfNecessary() {
1191
    if (!(goAwayStatus != null && streams.isEmpty() && pendingStreams.isEmpty())) {
1 ✔
1192
      return;
1 ✔
1193
    }
1194
    if (stopped) {
1 ✔
1195
      return;
1 ✔
1196
    }
1197
    stopped = true;
1 ✔
1198

1199
    if (keepAliveManager != null) {
1 ✔
1200
      keepAliveManager.onTransportTermination();
×
1201
    }
1202

1203
    if (ping != null) {
1 ✔
1204
      ping.failed(getPingFailure());
1 ✔
1205
      ping = null;
1 ✔
1206
    }
1207

1208
    if (!goAwaySent) {
1 ✔
1209
      // Send GOAWAY with lastGoodStreamId of 0, since we don't expect any server-initiated
1210
      // streams. The GOAWAY is part of graceful shutdown.
1211
      goAwaySent = true;
1 ✔
1212
      frameWriter.goAway(0, ErrorCode.NO_ERROR, new byte[0]);
1 ✔
1213
    }
1214

1215
    // We will close the underlying socket in the writing thread to break out the reader
1216
    // thread, which will close the frameReader and notify the listener.
1217
    frameWriter.close();
1 ✔
1218
  }
1 ✔
1219

1220
  @GuardedBy("lock")
1221
  private void maybeClearInUse(OkHttpClientStream stream) {
1222
    if (hasStream) {
1 ✔
1223
      if (pendingStreams.isEmpty() && streams.isEmpty()) {
1 ✔
1224
        hasStream = false;
1 ✔
1225
        if (keepAliveManager != null) {
1 ✔
1226
          // We don't have any active streams. No need to do keepalives any more.
1227
          // Again, we have to call this inside the lock to avoid the race between onTransportIdle
1228
          // and onTransportActive.
1229
          keepAliveManager.onTransportIdle();
×
1230
        }
1231
      }
1232
    }
1233
    if (stream.shouldBeCountedForInUse()) {
1 ✔
1234
      inUseState.updateObjectInUse(stream, false);
1 ✔
1235
    }
1236
  }
1 ✔
1237

1238
  @GuardedBy("lock")
1239
  private void setInUse(OkHttpClientStream stream) {
1240
    if (!hasStream) {
1 ✔
1241
      hasStream = true;
1 ✔
1242
      if (keepAliveManager != null) {
1 ✔
1243
        // We have a new stream. We might need to do keepalives now.
1244
        // Note that we have to do this inside the lock to avoid calling
1245
        // KeepAliveManager.onTransportActive and KeepAliveManager.onTransportIdle in the wrong
1246
        // order.
1247
        keepAliveManager.onTransportActive();
×
1248
      }
1249
    }
1250
    if (stream.shouldBeCountedForInUse()) {
1 ✔
1251
      inUseState.updateObjectInUse(stream, true);
1 ✔
1252
    }
1253
  }
1 ✔
1254

1255
  private Status getPingFailure() {
1256
    synchronized (lock) {
1 ✔
1257
      if (goAwayStatus != null) {
1 ✔
1258
        return goAwayStatus;
1 ✔
1259
      } else {
1260
        return Status.UNAVAILABLE.withDescription("Connection closed");
×
1261
      }
1262
    }
1263
  }
1264

1265
  boolean mayHaveCreatedStream(int streamId) {
1266
    synchronized (lock) {
1 ✔
1267
      return streamId < nextStreamId && (streamId & 1) == 1;
1 ✔
1268
    }
1269
  }
1270

1271
  OkHttpClientStream getStream(int streamId) {
1272
    synchronized (lock) {
1 ✔
1273
      return streams.get(streamId);
1 ✔
1274
    }
1275
  }
1276

1277
  /**
1278
   * Returns a Grpc status corresponding to the given ErrorCode.
1279
   */
1280
  @VisibleForTesting
1281
  static Status toGrpcStatus(ErrorCode code) {
1282
    Status status = ERROR_CODE_TO_STATUS.get(code);
1 ✔
1283
    return status != null ? status : Status.UNKNOWN.withDescription(
1 ✔
1284
        "Unknown http2 error code: " + code.httpCode);
1285
  }
1286

1287
  @Override
1288
  public ListenableFuture<SocketStats> getStats() {
1289
    SettableFuture<SocketStats> ret = SettableFuture.create();
1 ✔
1290
    synchronized (lock) {
1 ✔
1291
      if (socket == null) {
1 ✔
1292
        ret.set(new SocketStats(
×
1293
            transportTracer.getStats(),
×
1294
            /*local=*/ null,
1295
            /*remote=*/ null,
1296
            new InternalChannelz.SocketOptions.Builder().build(),
×
1297
            /*security=*/ null));
1298
      } else {
1299
        ret.set(new SocketStats(
1 ✔
1300
            transportTracer.getStats(),
1 ✔
1301
            socket.getLocalSocketAddress(),
1 ✔
1302
            socket.getRemoteSocketAddress(),
1 ✔
1303
            Utils.getSocketOptions(socket),
1 ✔
1304
            securityInfo));
1305
      }
1306
      return ret;
1 ✔
1307
    }
1308
  }
1309

1310
  /**
1311
   * Runnable which reads frames and dispatches them to in flight calls.
1312
   */
1313
  class ClientFrameHandler implements FrameReader.Handler, Runnable {
1314

1315
    private final OkHttpFrameLogger logger =
1 ✔
1316
        new OkHttpFrameLogger(Level.FINE, OkHttpClientTransport.class);
1317
    FrameReader frameReader;
1318
    boolean firstSettings = true;
1 ✔
1319

1320
    ClientFrameHandler(FrameReader frameReader) {
1 ✔
1321
      this.frameReader = frameReader;
1 ✔
1322
    }
1 ✔
1323

1324
    @Override
1325
    @SuppressWarnings("Finally")
1326
    public void run() {
1327
      String threadName = Thread.currentThread().getName();
1 ✔
1328
      Thread.currentThread().setName("OkHttpClientTransport");
1 ✔
1329
      try {
1330
        // Read until the underlying socket closes.
1331
        while (frameReader.nextFrame(this)) {
1 ✔
1332
          if (keepAliveManager != null) {
1 ✔
1333
            keepAliveManager.onDataReceived();
×
1334
          }
1335
        }
1336
        // frameReader.nextFrame() returns false when the underlying read encounters an IOException,
1337
        // it may be triggered by the socket closing, in such case, the startGoAway() will do
1338
        // nothing, otherwise, we finish all streams since it's a real IO issue.
1339
        Status status;
1340
        synchronized (lock) {
1 ✔
1341
          status = goAwayStatus;
1 ✔
1342
        }
1 ✔
1343
        if (status == null) {
1 ✔
1344
          status = Status.UNAVAILABLE.withDescription("End of stream or IOException");
1 ✔
1345
        }
1346
        startGoAway(0, ErrorCode.INTERNAL_ERROR, status);
1 ✔
1347
      } catch (Throwable t) {
1 ✔
1348
        // TODO(madongfly): Send the exception message to the server.
1349
        startGoAway(
1 ✔
1350
            0,
1351
            ErrorCode.PROTOCOL_ERROR,
1352
            Status.INTERNAL.withDescription("error in frame handler").withCause(t));
1 ✔
1353
      } finally {
1354
        try {
1355
          frameReader.close();
1 ✔
1356
        } catch (IOException ex) {
×
1357
          log.log(Level.INFO, "Exception closing frame reader", ex);
×
1358
        } catch (RuntimeException e) {
×
1359
          // This same check is done in okhttp proper:
1360
          // https://github.com/square/okhttp/blob/3cc0f4917cbda03cb31617f8ead1e0aeb19de2fb/okhttp/src/main/kotlin/okhttp3/internal/-UtilJvm.kt#L270
1361

1362
          // Conscrypt in Android 10 and 11 may throw closing an SSLSocket. This is safe to ignore.
1363
          // https://issuetracker.google.com/issues/177450597
1364
          if (!"bio == null".equals(e.getMessage())) {
×
1365
            throw e;
×
1366
          }
1367
        }
1 ✔
1368
        listener.transportTerminated();
1 ✔
1369
        Thread.currentThread().setName(threadName);
1 ✔
1370
      }
1371
    }
1 ✔
1372

1373
    /**
1374
     * Handle an HTTP2 DATA frame.
1375
     */
1376
    @SuppressWarnings("GuardedBy")
1377
    @Override
1378
    public void data(boolean inFinished, int streamId, BufferedSource in, int length,
1379
                     int paddedLength)
1380
        throws IOException {
1381
      logger.logData(OkHttpFrameLogger.Direction.INBOUND,
1 ✔
1382
          streamId, in.getBuffer(), length, inFinished);
1 ✔
1383
      OkHttpClientStream stream = getStream(streamId);
1 ✔
1384
      if (stream == null) {
1 ✔
1385
        if (mayHaveCreatedStream(streamId)) {
1 ✔
1386
          synchronized (lock) {
1 ✔
1387
            frameWriter.rstStream(streamId, ErrorCode.STREAM_CLOSED);
1 ✔
1388
          }
1 ✔
1389
          in.skip(length);
1 ✔
1390
        } else {
1391
          onError(ErrorCode.PROTOCOL_ERROR, "Received data for unknown stream: " + streamId);
1 ✔
1392
          return;
1 ✔
1393
        }
1394
      } else {
1395
        // Wait until the frame is complete.
1396
        in.require(length);
1 ✔
1397

1398
        Buffer buf = new Buffer();
1 ✔
1399
        buf.write(in.getBuffer(), length);
1 ✔
1400
        PerfMark.event("OkHttpClientTransport$ClientFrameHandler.data",
1 ✔
1401
            stream.transportState().tag());
1 ✔
1402
        synchronized (lock) {
1 ✔
1403
          // TODO(b/145386688): This access should be guarded by 'stream.transportState().lock';
1404
          // instead found: 'OkHttpClientTransport.this.lock'
1405
          stream.transportState().transportDataReceived(buf, inFinished, paddedLength - length);
1 ✔
1406
        }
1 ✔
1407
      }
1408

1409
      // connection window update
1410
      connectionUnacknowledgedBytesRead += paddedLength;
1 ✔
1411
      if (connectionUnacknowledgedBytesRead >= initialWindowSize * DEFAULT_WINDOW_UPDATE_RATIO) {
1 ✔
1412
        synchronized (lock) {
1 ✔
1413
          frameWriter.windowUpdate(0, connectionUnacknowledgedBytesRead);
1 ✔
1414
        }
1 ✔
1415
        connectionUnacknowledgedBytesRead = 0;
1 ✔
1416
      }
1417
    }
1 ✔
1418

1419
    /**
1420
     * Handle HTTP2 HEADER and CONTINUATION frames.
1421
     */
1422
    @SuppressWarnings("GuardedBy")
1423
    @Override
1424
    public void headers(boolean outFinished,
1425
        boolean inFinished,
1426
        int streamId,
1427
        int associatedStreamId,
1428
        List<Header> headerBlock,
1429
        HeadersMode headersMode) {
1430
      logger.logHeaders(OkHttpFrameLogger.Direction.INBOUND, streamId, headerBlock, inFinished);
1 ✔
1431
      boolean unknownStream = false;
1 ✔
1432
      Status failedStatus = null;
1 ✔
1433
      if (maxInboundMetadataSize != Integer.MAX_VALUE) {
1 ✔
1434
        int metadataSize = headerBlockSize(headerBlock);
1 ✔
1435
        if (metadataSize > maxInboundMetadataSize) {
1 ✔
1436
          failedStatus = Status.RESOURCE_EXHAUSTED.withDescription(
1 ✔
1437
              String.format(
1 ✔
1438
                  Locale.US,
1439
                  "Response %s metadata larger than %d: %d",
1440
                  inFinished ? "trailer" : "header",
1 ✔
1441
                  maxInboundMetadataSize,
1 ✔
1442
                  metadataSize));
1 ✔
1443
        }
1444
      }
1445
      synchronized (lock) {
1 ✔
1446
        OkHttpClientStream stream = streams.get(streamId);
1 ✔
1447
        if (stream == null) {
1 ✔
1448
          if (mayHaveCreatedStream(streamId)) {
1 ✔
1449
            frameWriter.rstStream(streamId, ErrorCode.STREAM_CLOSED);
1 ✔
1450
          } else {
1451
            unknownStream = true;
1 ✔
1452
          }
1453
        } else {
1454
          if (failedStatus == null) {
1 ✔
1455
            PerfMark.event("OkHttpClientTransport$ClientFrameHandler.headers",
1 ✔
1456
                stream.transportState().tag());
1 ✔
1457
            // TODO(b/145386688): This access should be guarded by 'stream.transportState().lock';
1458
            // instead found: 'OkHttpClientTransport.this.lock'
1459
            stream.transportState().transportHeadersReceived(headerBlock, inFinished);
1 ✔
1460
          } else {
1461
            if (!inFinished) {
1 ✔
1462
              frameWriter.rstStream(streamId, ErrorCode.CANCEL);
1 ✔
1463
            }
1464
            stream.transportState().transportReportStatus(failedStatus, false, new Metadata());
1 ✔
1465
          }
1466
        }
1467
      }
1 ✔
1468
      if (unknownStream) {
1 ✔
1469
        // We don't expect any server-initiated streams.
1470
        onError(ErrorCode.PROTOCOL_ERROR, "Received header for unknown stream: " + streamId);
1 ✔
1471
      }
1472
    }
1 ✔
1473

1474
    private int headerBlockSize(List<Header> headerBlock) {
1475
      // Calculate as defined for SETTINGS_MAX_HEADER_LIST_SIZE in RFC 7540 ยง6.5.2.
1476
      long size = 0;
1 ✔
1477
      for (int i = 0; i < headerBlock.size(); i++) {
1 ✔
1478
        Header header = headerBlock.get(i);
1 ✔
1479
        size += 32 + header.name.size() + header.value.size();
1 ✔
1480
      }
1481
      size = Math.min(size, Integer.MAX_VALUE);
1 ✔
1482
      return (int) size;
1 ✔
1483
    }
1484

1485
    @Override
1486
    public void rstStream(int streamId, ErrorCode errorCode) {
1487
      logger.logRstStream(OkHttpFrameLogger.Direction.INBOUND, streamId, errorCode);
1 ✔
1488
      Status status = toGrpcStatus(errorCode).augmentDescription("Rst Stream");
1 ✔
1489
      boolean stopDelivery =
1 ✔
1490
          (status.getCode() == Code.CANCELLED || status.getCode() == Code.DEADLINE_EXCEEDED);
1 ✔
1491
      synchronized (lock) {
1 ✔
1492
        OkHttpClientStream stream = streams.get(streamId);
1 ✔
1493
        if (stream != null) {
1 ✔
1494
          PerfMark.event("OkHttpClientTransport$ClientFrameHandler.rstStream",
1 ✔
1495
              stream.transportState().tag());
1 ✔
1496
          finishStream(
1 ✔
1497
              streamId, status,
1498
              errorCode == ErrorCode.REFUSED_STREAM ? RpcProgress.REFUSED : RpcProgress.PROCESSED,
1 ✔
1499
              stopDelivery, null, null);
1500
        }
1501
      }
1 ✔
1502
    }
1 ✔
1503

1504
    @Override
1505
    public void settings(boolean clearPrevious, Settings settings) {
1506
      logger.logSettings(OkHttpFrameLogger.Direction.INBOUND, settings);
1 ✔
1507
      boolean outboundWindowSizeIncreased = false;
1 ✔
1508
      synchronized (lock) {
1 ✔
1509
        if (OkHttpSettingsUtil.isSet(settings, OkHttpSettingsUtil.MAX_CONCURRENT_STREAMS)) {
1 ✔
1510
          int receivedMaxConcurrentStreams = OkHttpSettingsUtil.get(
1 ✔
1511
              settings, OkHttpSettingsUtil.MAX_CONCURRENT_STREAMS);
1512
          maxConcurrentStreams = receivedMaxConcurrentStreams;
1 ✔
1513
        }
1514

1515
        if (OkHttpSettingsUtil.isSet(settings, OkHttpSettingsUtil.INITIAL_WINDOW_SIZE)) {
1 ✔
1516
          int initialWindowSize = OkHttpSettingsUtil.get(
1 ✔
1517
              settings, OkHttpSettingsUtil.INITIAL_WINDOW_SIZE);
1518
          outboundWindowSizeIncreased = outboundFlow.initialOutboundWindowSize(initialWindowSize);
1 ✔
1519
        }
1520
        if (firstSettings) {
1 ✔
1521
          attributes = listener.filterTransport(attributes);
1 ✔
1522
          listener.transportReady();
1 ✔
1523
          firstSettings = false;
1 ✔
1524
        }
1525

1526
        // The changed settings are not finalized until SETTINGS acknowledgment frame is sent. Any
1527
        // writes due to update in settings must be sent after SETTINGS acknowledgment frame,
1528
        // otherwise it will cause a stream error (RST_STREAM).
1529
        frameWriter.ackSettings(settings);
1 ✔
1530

1531
        // send any pending bytes / streams
1532
        if (outboundWindowSizeIncreased) {
1 ✔
1533
          outboundFlow.writeStreams();
1 ✔
1534
        }
1535
        startPendingStreams();
1 ✔
1536
      }
1 ✔
1537
    }
1 ✔
1538

1539
    @Override
1540
    public void ping(boolean ack, int payload1, int payload2) {
1541
      long ackPayload = (((long) payload1) << 32) | (payload2 & 0xffffffffL);
1 ✔
1542
      logger.logPing(OkHttpFrameLogger.Direction.INBOUND, ackPayload);
1 ✔
1543
      if (!ack) {
1 ✔
1544
        synchronized (lock) {
1 ✔
1545
          frameWriter.ping(true, payload1, payload2);
1 ✔
1546
        }
1 ✔
1547
      } else {
1548
        Http2Ping p = null;
1 ✔
1549
        synchronized (lock) {
1 ✔
1550
          if (ping != null) {
1 ✔
1551
            if (ping.payload() == ackPayload) {
1 ✔
1552
              p = ping;
1 ✔
1553
              ping = null;
1 ✔
1554
            } else {
1555
              log.log(Level.WARNING, String.format(
1 ✔
1556
                  Locale.US, "Received unexpected ping ack. Expecting %d, got %d",
1557
                  ping.payload(), ackPayload));
1 ✔
1558
            }
1559
          } else {
1560
            log.warning("Received unexpected ping ack. No ping outstanding");
×
1561
          }
1562
        }
1 ✔
1563
        // don't complete it while holding lock since callbacks could run immediately
1564
        if (p != null) {
1 ✔
1565
          p.complete();
1 ✔
1566
        }
1567
      }
1568
    }
1 ✔
1569

1570
    @Override
1571
    public void ackSettings() {
1572
      // Do nothing currently.
1573
    }
1 ✔
1574

1575
    @Override
1576
    public void goAway(int lastGoodStreamId, ErrorCode errorCode, ByteString debugData) {
1577
      logger.logGoAway(OkHttpFrameLogger.Direction.INBOUND, lastGoodStreamId, errorCode, debugData);
1 ✔
1578
      if (errorCode == ErrorCode.ENHANCE_YOUR_CALM) {
1 ✔
1579
        String data = debugData.utf8();
1 ✔
1580
        log.log(Level.WARNING, String.format(
1 ✔
1581
            "%s: Received GOAWAY with ENHANCE_YOUR_CALM. Debug data: %s", this, data));
1582
        if ("too_many_pings".equals(data)) {
1 ✔
1583
          tooManyPingsRunnable.run();
1 ✔
1584
        }
1585
      }
1586
      Status status = GrpcUtil.Http2Error.statusForCode(errorCode.httpCode)
1 ✔
1587
          .augmentDescription("Received Goaway");
1 ✔
1588
      if (debugData.size() > 0) {
1 ✔
1589
        // If a debug message was provided, use it.
1590
        status = status.augmentDescription(debugData.utf8());
1 ✔
1591
      }
1592
      startGoAway(lastGoodStreamId, null, status);
1 ✔
1593
    }
1 ✔
1594

1595
    @Override
1596
    public void pushPromise(int streamId, int promisedStreamId, List<Header> requestHeaders)
1597
        throws IOException {
1598
      logger.logPushPromise(OkHttpFrameLogger.Direction.INBOUND,
1 ✔
1599
          streamId, promisedStreamId, requestHeaders);
1600
      // We don't accept server initiated stream.
1601
      synchronized (lock) {
1 ✔
1602
        frameWriter.rstStream(streamId, ErrorCode.PROTOCOL_ERROR);
1 ✔
1603
      }
1 ✔
1604
    }
1 ✔
1605

1606
    @Override
1607
    public void windowUpdate(int streamId, long delta) {
1608
      logger.logWindowsUpdate(OkHttpFrameLogger.Direction.INBOUND, streamId, delta);
1 ✔
1609
      if (delta == 0) {
1 ✔
1610
        String errorMsg = "Received 0 flow control window increment.";
×
1611
        if (streamId == 0) {
×
1612
          onError(ErrorCode.PROTOCOL_ERROR, errorMsg);
×
1613
        } else {
1614
          finishStream(
×
1615
              streamId, Status.INTERNAL.withDescription(errorMsg), RpcProgress.PROCESSED, false,
×
1616
              ErrorCode.PROTOCOL_ERROR, null);
1617
        }
1618
        return;
×
1619
      }
1620

1621
      boolean unknownStream = false;
1 ✔
1622
      synchronized (lock) {
1 ✔
1623
        if (streamId == Utils.CONNECTION_STREAM_ID) {
1 ✔
1624
          outboundFlow.windowUpdate(null, (int) delta);
1 ✔
1625
          return;
1 ✔
1626
        }
1627

1628
        OkHttpClientStream stream = streams.get(streamId);
1 ✔
1629
        if (stream != null) {
1 ✔
1630
          outboundFlow.windowUpdate(stream.transportState().getOutboundFlowState(), (int) delta);
1 ✔
1631
        } else if (!mayHaveCreatedStream(streamId)) {
1 ✔
1632
          unknownStream = true;
1 ✔
1633
        }
1634
      }
1 ✔
1635
      if (unknownStream) {
1 ✔
1636
        onError(ErrorCode.PROTOCOL_ERROR,
1 ✔
1637
            "Received window_update for unknown stream: " + streamId);
1638
      }
1639
    }
1 ✔
1640

1641
    @Override
1642
    public void priority(int streamId, int streamDependency, int weight, boolean exclusive) {
1643
      // Ignore priority change.
1644
      // TODO(madongfly): log
1645
    }
×
1646

1647
    @Override
1648
    public void alternateService(int streamId, String origin, ByteString protocol, String host,
1649
        int port, long maxAge) {
1650
      // TODO(madongfly): Deal with alternateService propagation
1651
    }
×
1652
  }
1653

1654
  /**
1655
   * SSLSocket wrapper that provides a fake SSLSession for handshake session.
1656
   */
1657
  static final class SslSocketWrapper extends NoopSslSocket {
1658

1659
    private final SSLSession sslSession;
1660
    private final SSLSocket sslSocket;
1661

1662
    SslSocketWrapper(SSLSocket sslSocket, String peerHost) {
1 ✔
1663
      this.sslSocket = sslSocket;
1 ✔
1664
      this.sslSession = new FakeSslSession(peerHost);
1 ✔
1665
    }
1 ✔
1666

1667
    @Override
1668
    public SSLSession getHandshakeSession() {
1669
      return this.sslSession;
1 ✔
1670
    }
1671

1672
    @Override
1673
    public boolean isConnected() {
1674
      return sslSocket.isConnected();
1 ✔
1675
    }
1676

1677
    @Override
1678
    public SSLParameters getSSLParameters() {
1679
      return sslSocket.getSSLParameters();
1 ✔
1680
    }
1681
  }
1682

1683
  /**
1684
   * Fake SSLSession instance that provides the peer host name to verify for per-rpc check.
1685
   */
1686
  static class FakeSslSession extends NoopSslSession {
1687

1688
    private final String peerHost;
1689

1690
    FakeSslSession(String peerHost) {
1 ✔
1691
      this.peerHost = peerHost;
1 ✔
1692
    }
1 ✔
1693

1694
    @Override
1695
    public String getPeerHost() {
1696
      return peerHost;
1 ✔
1697
    }
1698
  }
1699
}
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