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

grpc / grpc-java / #20494

30 Sep 2026 08:29AM UTC coverage: 89.351% (+0.05%) from 89.306%
#20494

push

github

web-flow
core: Implement [A121](https://github.com/grpc/proposal/pull/556) (#12893)

39251 of 43929 relevant lines covered (89.35%)

0.89 hits per line

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

93.55
/../core/src/main/java/io/grpc/internal/ManagedChannelImpl.java
1
/*
2
 * Copyright 2016 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.internal;
18

19
import static com.google.common.base.Preconditions.checkArgument;
20
import static com.google.common.base.Preconditions.checkNotNull;
21
import static com.google.common.base.Preconditions.checkState;
22
import static io.grpc.ClientStreamTracer.NAME_RESOLUTION_DELAYED;
23
import static io.grpc.ConnectivityState.CONNECTING;
24
import static io.grpc.ConnectivityState.IDLE;
25
import static io.grpc.ConnectivityState.SHUTDOWN;
26
import static io.grpc.ConnectivityState.TRANSIENT_FAILURE;
27
import static io.grpc.EquivalentAddressGroup.ATTR_AUTHORITY_OVERRIDE;
28

29
import com.google.common.annotations.VisibleForTesting;
30
import com.google.common.base.MoreObjects;
31
import com.google.common.base.Stopwatch;
32
import com.google.common.base.Supplier;
33
import com.google.common.util.concurrent.ListenableFuture;
34
import com.google.common.util.concurrent.SettableFuture;
35
import com.google.errorprone.annotations.concurrent.GuardedBy;
36
import io.grpc.Attributes;
37
import io.grpc.CallCredentials;
38
import io.grpc.CallOptions;
39
import io.grpc.Channel;
40
import io.grpc.ChannelConfigurator;
41
import io.grpc.ChannelCredentials;
42
import io.grpc.ChannelLogger;
43
import io.grpc.ChannelLogger.ChannelLogLevel;
44
import io.grpc.ClientCall;
45
import io.grpc.ClientInterceptor;
46
import io.grpc.ClientInterceptors;
47
import io.grpc.ClientStreamTracer;
48
import io.grpc.ClientTransportFilter;
49
import io.grpc.CompressorRegistry;
50
import io.grpc.ConnectivityState;
51
import io.grpc.ConnectivityStateInfo;
52
import io.grpc.Context;
53
import io.grpc.Deadline;
54
import io.grpc.DecompressorRegistry;
55
import io.grpc.EquivalentAddressGroup;
56
import io.grpc.ForwardingChannelBuilder2;
57
import io.grpc.ForwardingClientCall;
58
import io.grpc.Grpc;
59
import io.grpc.InternalChannelz;
60
import io.grpc.InternalChannelz.ChannelStats;
61
import io.grpc.InternalChannelz.ChannelTrace;
62
import io.grpc.InternalConfigSelector;
63
import io.grpc.InternalInstrumented;
64
import io.grpc.InternalLogId;
65
import io.grpc.InternalWithLogId;
66
import io.grpc.LoadBalancer;
67
import io.grpc.LoadBalancer.CreateSubchannelArgs;
68
import io.grpc.LoadBalancer.PickResult;
69
import io.grpc.LoadBalancer.PickSubchannelArgs;
70
import io.grpc.LoadBalancer.ResolvedAddresses;
71
import io.grpc.LoadBalancer.SubchannelPicker;
72
import io.grpc.LoadBalancer.SubchannelStateListener;
73
import io.grpc.LoadBalancerProvider;
74
import io.grpc.ManagedChannel;
75
import io.grpc.ManagedChannelBuilder;
76
import io.grpc.Metadata;
77
import io.grpc.MethodDescriptor;
78
import io.grpc.MetricInstrumentRegistry;
79
import io.grpc.MetricRecorder;
80
import io.grpc.NameResolver;
81
import io.grpc.NameResolver.ConfigOrError;
82
import io.grpc.NameResolver.ResolutionResult;
83
import io.grpc.NameResolverProvider;
84
import io.grpc.NameResolverRegistry;
85
import io.grpc.ProxyDetector;
86
import io.grpc.Status;
87
import io.grpc.StatusOr;
88
import io.grpc.SynchronizationContext;
89
import io.grpc.SynchronizationContext.ScheduledHandle;
90
import io.grpc.internal.ClientCallImpl.ClientStreamProvider;
91
import io.grpc.internal.ClientTransportFactory.SwapChannelCredentialsResult;
92
import io.grpc.internal.ManagedChannelImplBuilder.ClientTransportFactoryBuilder;
93
import io.grpc.internal.ManagedChannelImplBuilder.FixedPortProvider;
94
import io.grpc.internal.ManagedChannelServiceConfig.MethodInfo;
95
import io.grpc.internal.ManagedChannelServiceConfig.ServiceConfigConvertedSelector;
96
import io.grpc.internal.RetriableStream.ChannelBufferMeter;
97
import io.grpc.internal.RetriableStream.Throttle;
98
import java.net.URI;
99
import java.net.URISyntaxException;
100
import java.util.ArrayList;
101
import java.util.Collection;
102
import java.util.Collections;
103
import java.util.HashSet;
104
import java.util.LinkedHashSet;
105
import java.util.List;
106
import java.util.Map;
107
import java.util.Set;
108
import java.util.concurrent.Callable;
109
import java.util.concurrent.CountDownLatch;
110
import java.util.concurrent.ExecutionException;
111
import java.util.concurrent.Executor;
112
import java.util.concurrent.Future;
113
import java.util.concurrent.ScheduledExecutorService;
114
import java.util.concurrent.ScheduledFuture;
115
import java.util.concurrent.TimeUnit;
116
import java.util.concurrent.TimeoutException;
117
import java.util.concurrent.atomic.AtomicBoolean;
118
import java.util.concurrent.atomic.AtomicReference;
119
import java.util.logging.Level;
120
import java.util.logging.Logger;
121
import javax.annotation.Nullable;
122
import javax.annotation.concurrent.ThreadSafe;
123

124
/** A communication channel for making outgoing RPCs. */
125
@ThreadSafe
126
final class ManagedChannelImpl extends ManagedChannel implements
127
    InternalInstrumented<ChannelStats> {
128
  @VisibleForTesting
129
  static final Logger logger = Logger.getLogger(ManagedChannelImpl.class.getName());
1 ✔
130

131
  static final long IDLE_TIMEOUT_MILLIS_DISABLE = -1;
132

133
  static final long SUBCHANNEL_SHUTDOWN_DELAY_SECONDS = 5;
134

135
  @VisibleForTesting
136
  static final Status SHUTDOWN_NOW_STATUS =
1 ✔
137
      Status.UNAVAILABLE.withDescription("Channel shutdownNow invoked");
1 ✔
138

139
  @VisibleForTesting
140
  static final Status SHUTDOWN_STATUS =
1 ✔
141
      Status.UNAVAILABLE.withDescription("Channel shutdown invoked");
1 ✔
142

143
  @VisibleForTesting
144
  static final Status SUBCHANNEL_SHUTDOWN_STATUS =
1 ✔
145
      Status.UNAVAILABLE.withDescription("Subchannel shutdown invoked");
1 ✔
146

147
  private static final ManagedChannelServiceConfig EMPTY_SERVICE_CONFIG =
148
      ManagedChannelServiceConfig.empty();
1 ✔
149
  private static final InternalConfigSelector INITIAL_PENDING_SELECTOR =
1 ✔
150
      new InternalConfigSelector() {
1 ✔
151
        @Override
152
        public Result selectConfig(PickSubchannelArgs args) {
153
          throw new IllegalStateException("Resolution is pending");
×
154
        }
155
      };
156
  private static final LoadBalancer.PickDetailsConsumer NOOP_PICK_DETAILS_CONSUMER =
1 ✔
157
      new LoadBalancer.PickDetailsConsumer() {};
1 ✔
158

159
  /**
160
   * Retrieves the user-provided configuration function for internal child channels.
161
   *
162
   * <p>This is intended for use by gRPC internal components
163
   * that are responsible for creating auxiliary {@code ManagedChannel} instances.
164
   */
165
  private final ChannelConfigurator channelConfigurator;
166

167
  private final InternalLogId logId;
168
  private final String target;
169
  @Nullable
170
  private final String authorityOverride;
171
  private final NameResolverRegistry nameResolverRegistry;
172
  private final UriWrapper targetUri;
173
  private final NameResolverProvider nameResolverProvider;
174
  private final NameResolver.Args nameResolverArgs;
175
  private final LoadBalancerProvider loadBalancerFactory;
176
  private final RefCountedClientTransportFactory originalTransportFactory;
177
  @Nullable
178
  private final ChannelCredentials originalChannelCreds;
179
  private final ClientTransportFactory transportFactory;
180
  private final RestrictedScheduledExecutor scheduledExecutor;
181
  private final Executor executor;
182
  private final ObjectPool<? extends Executor> executorPool;
183
  private final ExecutorHolder balancerRpcExecutorHolder;
184
  private final ExecutorHolder offloadExecutorHolder;
185
  private final TimeProvider timeProvider;
186
  private final int maxTraceEvents;
187

188
  @VisibleForTesting
1 ✔
189
  final SynchronizationContext syncContext = new SynchronizationContext(
190
      new Thread.UncaughtExceptionHandler() {
1 ✔
191
        @Override
192
        public void uncaughtException(Thread t, Throwable e) {
193
          logger.log(
1 ✔
194
              Level.SEVERE,
195
              "[" + getLogId() + "] Uncaught exception in the SynchronizationContext. Panic!",
1 ✔
196
              e);
197
          try {
198
            panic(e);
1 ✔
199
          } catch (Throwable anotherT) {
×
200
            logger.log(
×
201
                Level.SEVERE, "[" + getLogId() + "] Uncaught exception while panicking", anotherT);
×
202
          }
1 ✔
203
        }
1 ✔
204
      });
205

206
  private boolean fullStreamDecompression;
207

208
  private final DecompressorRegistry decompressorRegistry;
209
  private final CompressorRegistry compressorRegistry;
210

211
  private final Supplier<Stopwatch> stopwatchSupplier;
212
  /** The timeout before entering idle mode. */
213
  private final long idleTimeoutMillis;
214

215
  private final ConnectivityStateManager channelStateManager = new ConnectivityStateManager();
1 ✔
216
  private final BackoffPolicy.Provider backoffPolicyProvider;
217

218
  /**
219
   * We delegate to this channel, so that we can have interceptors as necessary. If there aren't
220
   * any interceptors and the {@link io.grpc.BinaryLog} is {@code null} then this will just be a
221
   * {@link RealChannel}.
222
   */
223
  private final Channel interceptorChannel;
224

225
  private final List<ClientTransportFilter> transportFilters;
226
  @Nullable private final String userAgent;
227

228
  // Only null after channel is terminated. Must be assigned from the syncContext.
229
  private NameResolver nameResolver;
230

231
  // Must be accessed from the syncContext.
232
  private boolean nameResolverStarted;
233

234
  // null when channel is in idle mode.  Must be assigned from syncContext.
235
  @Nullable
236
  private LbHelperImpl lbHelper;
237

238
  // Must be accessed from the syncContext
239
  private boolean panicMode;
240

241
  // Must be mutated from syncContext
242
  // If any monitoring hook to be added later needs to get a snapshot of this Set, we could
243
  // switch to a ConcurrentHashMap.
244
  private final Set<InternalSubchannel> subchannels = new HashSet<>(16, .75f);
1 ✔
245

246
  // Must be accessed from syncContext
247
  @Nullable
248
  private Collection<RealChannel.PendingCall<?, ?>> pendingCalls;
249
  private final Object pendingCallsInUseObject = new Object();
1 ✔
250

251
  // reprocess() must be run from syncContext
252
  private final DelayedClientTransport delayedTransport;
253
  private final UncommittedRetriableStreamsRegistry uncommittedRetriableStreamsRegistry
1 ✔
254
      = new UncommittedRetriableStreamsRegistry();
255

256
  // Shutdown states.
257
  //
258
  // Channel's shutdown process:
259
  // 1. shutdown(): stop accepting new calls from applications
260
  //   1a shutdown <- true
261
  //   1b delayedTransport.shutdown()
262
  // 2. delayedTransport terminated: stop stream-creation functionality
263
  //   2a terminating <- true
264
  //   2b loadBalancer.shutdown()
265
  //     * LoadBalancer will shutdown subchannels and OOB channels
266
  //   2c loadBalancer <- null
267
  //   2d nameResolver.shutdown()
268
  //   2e nameResolver <- null
269
  // 3. All subchannels and OOB channels terminated: Channel considered terminated
270

271
  private final AtomicBoolean shutdown = new AtomicBoolean(false);
1 ✔
272
  // Must only be mutated and read from syncContext
273
  private boolean shutdownNowed;
274
  // Must only be mutated from syncContext
275
  private boolean terminating;
276
  // Must be mutated from syncContext
277
  private volatile boolean terminated;
278
  private final CountDownLatch terminatedLatch = new CountDownLatch(1);
1 ✔
279

280
  private final CallTracer.Factory callTracerFactory;
281
  private final CallTracer channelCallTracer;
282
  private final ChannelTracer channelTracer;
283
  private final ChannelLogger channelLogger;
284
  private final InternalChannelz channelz;
285
  private final RealChannel realChannel;
286
  // Must be mutated and read from syncContext
287
  // a flag for doing channel tracing when flipped
288
  private ResolutionState lastResolutionState = ResolutionState.NO_RESOLUTION;
1 ✔
289
  // Must be mutated and read from constructor or syncContext
290
  // used for channel tracing when value changed
291
  private ManagedChannelServiceConfig lastServiceConfig = EMPTY_SERVICE_CONFIG;
1 ✔
292

293
  @Nullable
294
  private final ManagedChannelServiceConfig defaultServiceConfig;
295
  // Must be mutated and read from constructor or syncContext
296
  private boolean serviceConfigUpdated = false;
1 ✔
297
  private final boolean lookUpServiceConfig;
298

299
  // One instance per channel.
300
  private final ChannelBufferMeter channelBufferUsed = new ChannelBufferMeter();
1 ✔
301

302
  private final long perRpcBufferLimit;
303
  private final long channelBufferLimit;
304

305
  // Temporary false flag that can skip the retry code path.
306
  private final boolean retryEnabled;
307

308
  private final Deadline.Ticker ticker = Deadline.getSystemTicker();
1 ✔
309

310
  // Called from syncContext
311
  private final ManagedClientTransport.Listener delayedTransportListener =
1 ✔
312
      new DelayedTransportListener();
313

314
  // Must be called from syncContext
315
  private void maybeShutdownNowSubchannels() {
316
    if (shutdownNowed) {
1 ✔
317
      for (InternalSubchannel subchannel : subchannels) {
1 ✔
318
        subchannel.shutdownNow(SHUTDOWN_NOW_STATUS);
1 ✔
319
      }
1 ✔
320
    }
321
  }
1 ✔
322

323
  // Must be accessed from syncContext
324
  @VisibleForTesting
1 ✔
325
  final InUseStateAggregator<Object> inUseStateAggregator = new IdleModeStateAggregator();
326

327
  @Override
328
  public ListenableFuture<ChannelStats> getStats() {
329
    final SettableFuture<ChannelStats> ret = SettableFuture.create();
1 ✔
330
    final class StatsFetcher implements Runnable {
1 ✔
331
      @Override
332
      public void run() {
333
        ChannelStats.Builder builder = new InternalChannelz.ChannelStats.Builder();
1 ✔
334
        channelCallTracer.updateBuilder(builder);
1 ✔
335
        channelTracer.updateBuilder(builder);
1 ✔
336
        builder.setTarget(target).setState(channelStateManager.getState());
1 ✔
337
        List<InternalWithLogId> children = new ArrayList<>();
1 ✔
338
        children.addAll(subchannels);
1 ✔
339
        builder.setSubchannels(children);
1 ✔
340
        ret.set(builder.build());
1 ✔
341
      }
1 ✔
342
    }
343

344
    // subchannels and oobchannels can only be accessed from syncContext
345
    syncContext.execute(new StatsFetcher());
1 ✔
346
    return ret;
1 ✔
347
  }
348

349
  @Override
350
  public InternalLogId getLogId() {
351
    return logId;
1 ✔
352
  }
353

354
  // Run from syncContext
355
  private class IdleModeTimer implements Runnable {
1 ✔
356

357
    @Override
358
    public void run() {
359
      // Workaround timer scheduled while in idle mode. This can happen from handleNotInUse() after
360
      // an explicit enterIdleMode() by the user. Protecting here as other locations are a bit too
361
      // subtle to change rapidly to resolve the channel panic. See #8714
362
      if (lbHelper == null) {
1 ✔
363
        return;
×
364
      }
365
      enterIdleMode();
1 ✔
366
    }
1 ✔
367
  }
368

369
  // Must be called from syncContext
370
  private void shutdownNameResolverAndLoadBalancer(boolean channelIsActive) {
371
    syncContext.throwIfNotInThisSynchronizationContext();
1 ✔
372
    if (channelIsActive) {
1 ✔
373
      checkState(nameResolverStarted, "nameResolver is not started");
1 ✔
374
      checkState(lbHelper != null, "lbHelper is null");
1 ✔
375
    }
376
    if (nameResolver != null) {
1 ✔
377
      nameResolver.shutdown();
1 ✔
378
      nameResolverStarted = false;
1 ✔
379
      if (channelIsActive) {
1 ✔
380
        nameResolver = getNameResolver(
1 ✔
381
            targetUri, authorityOverride, nameResolverProvider, nameResolverArgs);
382
      } else {
383
        nameResolver = null;
1 ✔
384
      }
385
    }
386
    if (lbHelper != null) {
1 ✔
387
      lbHelper.lb.shutdown();
1 ✔
388
      lbHelper = null;
1 ✔
389
    }
390
  }
1 ✔
391

392
  /**
393
   * Make the channel exit idle mode, if it's in it.
394
   *
395
   * <p>Must be called from syncContext
396
   */
397
  @VisibleForTesting
398
  void exitIdleMode() {
399
    syncContext.throwIfNotInThisSynchronizationContext();
1 ✔
400
    if (shutdown.get() || panicMode) {
1 ✔
401
      return;
1 ✔
402
    }
403
    if (inUseStateAggregator.isInUse()) {
1 ✔
404
      // Cancel the timer now, so that a racing due timer will not put Channel on idleness
405
      // when the caller of exitIdleMode() is about to use the returned loadBalancer.
406
      cancelIdleTimer(false);
1 ✔
407
    } else {
408
      // exitIdleMode() may be called outside of inUseStateAggregator.handleNotInUse() while
409
      // isInUse() == false, in which case we still need to schedule the timer.
410
      rescheduleIdleTimer();
1 ✔
411
    }
412
    if (lbHelper != null) {
1 ✔
413
      return;
1 ✔
414
    }
415
    channelLogger.log(ChannelLogLevel.INFO, "Exiting idle mode");
1 ✔
416
    LbHelperImpl lbHelper = new LbHelperImpl();
1 ✔
417
    lbHelper.lb = loadBalancerFactory.newLoadBalancer(lbHelper);
1 ✔
418
    // Delay setting lbHelper until fully initialized, since loadBalancerFactory is user code and
419
    // may throw. We don't want to confuse our state, even if we enter panic mode.
420
    this.lbHelper = lbHelper;
1 ✔
421

422
    channelStateManager.gotoState(CONNECTING);
1 ✔
423
    NameResolverListener listener = new NameResolverListener(lbHelper, nameResolver);
1 ✔
424
    nameResolver.start(listener);
1 ✔
425
    nameResolverStarted = true;
1 ✔
426
  }
1 ✔
427

428
  // Must be run from syncContext
429
  private void enterIdleMode() {
430
    // nameResolver and loadBalancer are guaranteed to be non-null.  If any of them were null,
431
    // either the idleModeTimer ran twice without exiting the idle mode, or the task in shutdown()
432
    // did not cancel idleModeTimer, or enterIdle() ran while shutdown or in idle, all of
433
    // which are bugs.
434
    shutdownNameResolverAndLoadBalancer(true);
1 ✔
435
    delayedTransport.reprocess(null);
1 ✔
436
    realChannel.updateConfigSelector(INITIAL_PENDING_SELECTOR);
1 ✔
437
    channelLogger.log(ChannelLogLevel.INFO, "Entering IDLE state");
1 ✔
438
    channelStateManager.gotoState(IDLE);
1 ✔
439
    // If the inUseStateAggregator still considers pending calls to be queued up or the delayed
440
    // transport to be holding some we need to exit idle mode to give these calls a chance to
441
    // be processed.
442
    if (inUseStateAggregator.anyObjectInUse(pendingCallsInUseObject, delayedTransport)) {
1 ✔
443
      exitIdleMode();
1 ✔
444
    }
445
  }
1 ✔
446

447
  // Must be run from syncContext
448
  private void cancelIdleTimer(boolean permanent) {
449
    idleTimer.cancel(permanent);
1 ✔
450
  }
1 ✔
451

452
  // Always run from syncContext
453
  private void rescheduleIdleTimer() {
454
    if (idleTimeoutMillis == IDLE_TIMEOUT_MILLIS_DISABLE) {
1 ✔
455
      return;
1 ✔
456
    }
457
    idleTimer.reschedule(idleTimeoutMillis, TimeUnit.MILLISECONDS);
1 ✔
458
  }
1 ✔
459

460
  /**
461
   * Force name resolution refresh to happen immediately. Must be run
462
   * from syncContext.
463
   */
464
  private void refreshNameResolution() {
465
    syncContext.throwIfNotInThisSynchronizationContext();
1 ✔
466
    if (nameResolverStarted) {
1 ✔
467
      nameResolver.refresh();
1 ✔
468
    }
469
  }
1 ✔
470

471
  private final class ChannelStreamProvider implements ClientStreamProvider {
1 ✔
472
    volatile Throttle throttle;
473

474
    @Override
475
    public ClientStream newStream(
476
        final MethodDescriptor<?, ?> method,
477
        final CallOptions callOptions,
478
        final Metadata headers,
479
        final Context context) {
480
      // There is no need to reschedule the idle timer here. If the channel isn't shut down, either
481
      // the delayed transport or a real transport will go in-use and cancel the idle timer.
482
      if (!retryEnabled) {
1 ✔
483
        ClientStreamTracer[] tracers = GrpcUtil.getClientStreamTracers(
1 ✔
484
            callOptions, headers, 0, /* isTransparentRetry= */ false,
485
            /* isHedging= */false);
486
        Context origContext = context.attach();
1 ✔
487
        try {
488
          return delayedTransport.newStream(method, headers, callOptions, tracers);
1 ✔
489
        } finally {
490
          context.detach(origContext);
1 ✔
491
        }
492
      } else {
493
        MethodInfo methodInfo = callOptions.getOption(MethodInfo.KEY);
1 ✔
494
        final RetryPolicy retryPolicy = methodInfo == null ? null : methodInfo.retryPolicy;
1 ✔
495
        final HedgingPolicy hedgingPolicy = methodInfo == null ? null : methodInfo.hedgingPolicy;
1 ✔
496
        final class RetryStream<ReqT> extends RetriableStream<ReqT> {
497
          @SuppressWarnings("unchecked")
498
          RetryStream() {
1 ✔
499
            super(
1 ✔
500
                (MethodDescriptor<ReqT, ?>) method,
501
                headers,
502
                channelBufferUsed,
1 ✔
503
                perRpcBufferLimit,
1 ✔
504
                channelBufferLimit,
1 ✔
505
                getCallExecutor(callOptions),
1 ✔
506
                transportFactory.getScheduledExecutorService(),
1 ✔
507
                retryPolicy,
508
                hedgingPolicy,
509
                throttle);
510
          }
1 ✔
511

512
          @Override
513
          Status prestart() {
514
            return uncommittedRetriableStreamsRegistry.add(this);
1 ✔
515
          }
516

517
          @Override
518
          void postCommit() {
519
            uncommittedRetriableStreamsRegistry.remove(this);
1 ✔
520
          }
1 ✔
521

522
          @Override
523
          ClientStream newSubstream(
524
              Metadata newHeaders, ClientStreamTracer.Factory factory, int previousAttempts,
525
              boolean isTransparentRetry, boolean isHedgedStream) {
526
            CallOptions newOptions = callOptions.withStreamTracerFactory(factory);
1 ✔
527
            ClientStreamTracer[] tracers = GrpcUtil.getClientStreamTracers(
1 ✔
528
                newOptions, newHeaders, previousAttempts, isTransparentRetry, isHedgedStream);
529
            Context origContext = context.attach();
1 ✔
530
            try {
531
              return delayedTransport.newStream(method, newHeaders, newOptions, tracers);
1 ✔
532
            } finally {
533
              context.detach(origContext);
1 ✔
534
            }
535
          }
536
        }
537

538
        return new RetryStream<>();
1 ✔
539
      }
540
    }
541
  }
542

543
  private final ChannelStreamProvider transportProvider = new ChannelStreamProvider();
1 ✔
544

545
  private final Rescheduler idleTimer;
546
  private final MetricRecorder metricRecorder;
547

548
  ManagedChannelImpl(
549
      ManagedChannelImplBuilder builder,
550
      ClientTransportFactory clientTransportFactory,
551
      UriWrapper targetUri,
552
      NameResolverProvider nameResolverProvider,
553
      BackoffPolicy.Provider backoffPolicyProvider,
554
      ObjectPool<? extends Executor> balancerRpcExecutorPool,
555
      Supplier<Stopwatch> stopwatchSupplier,
556
      List<ClientInterceptor> interceptors,
557
      final TimeProvider timeProvider) {
1 ✔
558
    this.channelConfigurator = checkNotNull(builder.channelConfigurator,
1 ✔
559
            "channelConfigurator");
560
    this.target = checkNotNull(builder.target, "target");
1 ✔
561
    this.logId = InternalLogId.allocate("Channel", target);
1 ✔
562
    this.timeProvider = checkNotNull(timeProvider, "timeProvider");
1 ✔
563
    this.executorPool = checkNotNull(builder.executorPool, "executorPool");
1 ✔
564
    this.executor = checkNotNull(executorPool.getObject(), "executor");
1 ✔
565
    this.originalChannelCreds = builder.channelCredentials;
1 ✔
566
    if (clientTransportFactory instanceof RefCountedClientTransportFactory) {
1 ✔
567
      this.originalTransportFactory = (RefCountedClientTransportFactory) clientTransportFactory;
1 ✔
568
    } else {
569
      this.originalTransportFactory = new RefCountedClientTransportFactory(clientTransportFactory);
1 ✔
570
    }
571
    this.offloadExecutorHolder =
1 ✔
572
        new ExecutorHolder(checkNotNull(builder.offloadExecutorPool, "offloadExecutorPool"));
1 ✔
573
    this.transportFactory = new CallCredentialsApplyingTransportFactory(
1 ✔
574
        originalTransportFactory, builder.callCredentials, this.offloadExecutorHolder);
575
    this.scheduledExecutor =
1 ✔
576
        new RestrictedScheduledExecutor(transportFactory.getScheduledExecutorService());
1 ✔
577
    maxTraceEvents = builder.maxTraceEvents;
1 ✔
578
    channelTracer = new ChannelTracer(
1 ✔
579
        logId, builder.maxTraceEvents, timeProvider.currentTimeNanos(),
1 ✔
580
        "Channel for '" + target + "'");
581
    channelLogger = new ChannelLoggerImpl(channelTracer, timeProvider);
1 ✔
582
    ProxyDetector proxyDetector =
583
        builder.proxyDetector != null ? builder.proxyDetector : GrpcUtil.DEFAULT_PROXY_DETECTOR;
1 ✔
584
    this.retryEnabled = builder.retryEnabled;
1 ✔
585
    this.loadBalancerFactory = new AutoConfiguredLoadBalancerFactory(builder.defaultLbPolicy);
1 ✔
586
    this.nameResolverRegistry = builder.nameResolverRegistry;
1 ✔
587
    this.targetUri = checkNotNull(targetUri, "targetUri");
1 ✔
588
    this.nameResolverProvider = checkNotNull(nameResolverProvider, "nameResolverProvider");
1 ✔
589
    ScParser serviceConfigParser =
1 ✔
590
        new ScParser(
591
            retryEnabled,
592
            builder.maxRetryAttempts,
593
            builder.maxHedgedAttempts,
594
            loadBalancerFactory);
595
    this.authorityOverride = builder.authorityOverride;
1 ✔
596
    this.metricRecorder = new MetricRecorderImpl(builder.metricSinks,
1 ✔
597
        MetricInstrumentRegistry.getDefaultRegistry());
1 ✔
598
    NameResolver.Args.Builder nameResolverArgsBuilder = NameResolver.Args.newBuilder()
1 ✔
599
            .setDefaultPort(builder.getDefaultPort())
1 ✔
600
            .setProxyDetector(proxyDetector)
1 ✔
601
            .setSynchronizationContext(syncContext)
1 ✔
602
            .setScheduledExecutorService(scheduledExecutor)
1 ✔
603
            .setServiceConfigParser(serviceConfigParser)
1 ✔
604
            .setChannelLogger(channelLogger)
1 ✔
605
            .setOffloadExecutor(this.offloadExecutorHolder)
1 ✔
606
            .setOverrideAuthority(this.authorityOverride)
1 ✔
607
            .setMetricRecorder(this.metricRecorder)
1 ✔
608
            .setNameResolverRegistry(builder.nameResolverRegistry)
1 ✔
609
            .setChildChannelConfigurator(this.channelConfigurator);
1 ✔
610
    builder.copyAllNameResolverCustomArgsTo(nameResolverArgsBuilder);
1 ✔
611
    this.nameResolverArgs = nameResolverArgsBuilder.build();
1 ✔
612
    this.nameResolver = getNameResolver(
1 ✔
613
        targetUri, authorityOverride, nameResolverProvider, nameResolverArgs);
614
    this.balancerRpcExecutorHolder = new ExecutorHolder(
1 ✔
615
        checkNotNull(balancerRpcExecutorPool, "balancerRpcExecutorPool"));
1 ✔
616
    this.delayedTransport = new DelayedClientTransport(this.executor, this.syncContext);
1 ✔
617
    this.delayedTransport.start(delayedTransportListener);
1 ✔
618
    this.backoffPolicyProvider = backoffPolicyProvider;
1 ✔
619

620
    if (builder.defaultServiceConfig != null) {
1 ✔
621
      ConfigOrError parsedDefaultServiceConfig =
1 ✔
622
          serviceConfigParser.parseServiceConfig(builder.defaultServiceConfig);
1 ✔
623
      checkState(
1 ✔
624
          parsedDefaultServiceConfig.getError() == null,
1 ✔
625
          "Default config is invalid: %s",
626
          parsedDefaultServiceConfig.getError());
1 ✔
627
      this.defaultServiceConfig =
1 ✔
628
          (ManagedChannelServiceConfig) parsedDefaultServiceConfig.getConfig();
1 ✔
629
      this.transportProvider.throttle = this.defaultServiceConfig.getRetryThrottling();
1 ✔
630
    } else {
1 ✔
631
      this.defaultServiceConfig = null;
1 ✔
632
    }
633
    this.lookUpServiceConfig = builder.lookUpServiceConfig;
1 ✔
634
    realChannel = new RealChannel(nameResolver.getServiceAuthority());
1 ✔
635
    Channel channel = realChannel;
1 ✔
636
    if (builder.binlog != null) {
1 ✔
637
      channel = builder.binlog.wrapChannel(channel);
1 ✔
638
    }
639
    this.interceptorChannel = ClientInterceptors.intercept(channel, interceptors);
1 ✔
640
    this.transportFilters = new ArrayList<>(builder.transportFilters);
1 ✔
641
    this.stopwatchSupplier = checkNotNull(stopwatchSupplier, "stopwatchSupplier");
1 ✔
642
    if (builder.idleTimeoutMillis == IDLE_TIMEOUT_MILLIS_DISABLE) {
1 ✔
643
      this.idleTimeoutMillis = builder.idleTimeoutMillis;
1 ✔
644
    } else {
645
      checkArgument(
1 ✔
646
          builder.idleTimeoutMillis
647
              >= ManagedChannelImplBuilder.IDLE_MODE_MIN_TIMEOUT_MILLIS,
648
          "invalid idleTimeoutMillis %s", builder.idleTimeoutMillis);
649
      this.idleTimeoutMillis = builder.idleTimeoutMillis;
1 ✔
650
    }
651

652
    idleTimer = new Rescheduler(
1 ✔
653
        new IdleModeTimer(),
654
        syncContext,
655
        transportFactory.getScheduledExecutorService(),
1 ✔
656
        stopwatchSupplier.get());
1 ✔
657
    this.fullStreamDecompression = builder.fullStreamDecompression;
1 ✔
658
    this.decompressorRegistry = checkNotNull(builder.decompressorRegistry, "decompressorRegistry");
1 ✔
659
    this.compressorRegistry = checkNotNull(builder.compressorRegistry, "compressorRegistry");
1 ✔
660
    this.userAgent = builder.userAgent;
1 ✔
661

662
    this.channelBufferLimit = builder.retryBufferSize;
1 ✔
663
    this.perRpcBufferLimit = builder.perRpcBufferLimit;
1 ✔
664
    final class ChannelCallTracerFactory implements CallTracer.Factory {
1 ✔
665
      @Override
666
      public CallTracer create() {
667
        return new CallTracer(timeProvider);
1 ✔
668
      }
669
    }
670

671
    this.callTracerFactory = new ChannelCallTracerFactory();
1 ✔
672
    channelCallTracer = callTracerFactory.create();
1 ✔
673
    this.channelz = checkNotNull(builder.channelz);
1 ✔
674
    channelz.addRootChannel(this);
1 ✔
675

676
    if (!lookUpServiceConfig) {
1 ✔
677
      if (defaultServiceConfig != null) {
1 ✔
678
        channelLogger.log(
1 ✔
679
            ChannelLogLevel.INFO, "Service config look-up disabled, using default service config");
680
      }
681
      serviceConfigUpdated = true;
1 ✔
682
    }
683
  }
1 ✔
684

685
  @VisibleForTesting
686
  static NameResolver getNameResolver(
687
      UriWrapper targetUri, @Nullable final String overrideAuthority,
688
      NameResolverProvider provider, NameResolver.Args nameResolverArgs) {
689
    NameResolver resolver = targetUri.newNameResolver(provider, nameResolverArgs);
1 ✔
690
    if (resolver == null) {
1 ✔
691
      throw new IllegalArgumentException("cannot create a NameResolver for " + targetUri);
1 ✔
692
    }
693

694
    // We wrap the name resolver in a RetryingNameResolver to give it the ability to retry failures.
695
    // TODO: After a transition period, all NameResolver implementations that need retry should use
696
    //       RetryingNameResolver directly and this step can be removed.
697
    NameResolver usedNameResolver = RetryingNameResolver.wrap(resolver, nameResolverArgs);
1 ✔
698

699
    if (overrideAuthority == null) {
1 ✔
700
      return usedNameResolver;
1 ✔
701
    }
702

703
    return new ForwardingNameResolver(usedNameResolver) {
1 ✔
704
      @Override
705
      public String getServiceAuthority() {
706
        return overrideAuthority;
1 ✔
707
      }
708
    };
709
  }
710

711
  @VisibleForTesting
712
  InternalConfigSelector getConfigSelector() {
713
    return realChannel.configSelector.get();
1 ✔
714
  }
715
  
716
  @VisibleForTesting
717
  boolean hasThrottle() {
718
    return this.transportProvider.throttle != null;
1 ✔
719
  }
720

721
  /**
722
   * Initiates an orderly shutdown in which preexisting calls continue but new calls are immediately
723
   * cancelled.
724
   */
725
  @Override
726
  public ManagedChannelImpl shutdown() {
727
    channelLogger.log(ChannelLogLevel.DEBUG, "shutdown() called");
1 ✔
728
    if (!shutdown.compareAndSet(false, true)) {
1 ✔
729
      return this;
1 ✔
730
    }
731
    final class Shutdown implements Runnable {
1 ✔
732
      @Override
733
      public void run() {
734
        channelLogger.log(ChannelLogLevel.INFO, "Entering SHUTDOWN state");
1 ✔
735
        channelStateManager.gotoState(SHUTDOWN);
1 ✔
736
      }
1 ✔
737
    }
738

739
    syncContext.execute(new Shutdown());
1 ✔
740
    realChannel.shutdown();
1 ✔
741
    final class CancelIdleTimer implements Runnable {
1 ✔
742
      @Override
743
      public void run() {
744
        cancelIdleTimer(/* permanent= */ true);
1 ✔
745
      }
1 ✔
746
    }
747

748
    syncContext.execute(new CancelIdleTimer());
1 ✔
749
    return this;
1 ✔
750
  }
751

752
  /**
753
   * Initiates a forceful shutdown in which preexisting and new calls are cancelled. Although
754
   * forceful, the shutdown process is still not instantaneous; {@link #isTerminated()} will likely
755
   * return {@code false} immediately after this method returns.
756
   */
757
  @Override
758
  public ManagedChannelImpl shutdownNow() {
759
    channelLogger.log(ChannelLogLevel.DEBUG, "shutdownNow() called");
1 ✔
760
    shutdown();
1 ✔
761
    realChannel.shutdownNow();
1 ✔
762
    final class ShutdownNow implements Runnable {
1 ✔
763
      @Override
764
      public void run() {
765
        if (shutdownNowed) {
1 ✔
766
          return;
1 ✔
767
        }
768
        shutdownNowed = true;
1 ✔
769
        maybeShutdownNowSubchannels();
1 ✔
770
      }
1 ✔
771
    }
772

773
    syncContext.execute(new ShutdownNow());
1 ✔
774
    return this;
1 ✔
775
  }
776

777
  // Called from syncContext
778
  @VisibleForTesting
779
  void panic(final Throwable t) {
780
    if (panicMode) {
1 ✔
781
      // Preserve the first panic information
782
      return;
×
783
    }
784
    panicMode = true;
1 ✔
785
    try {
786
      cancelIdleTimer(/* permanent= */ true);
1 ✔
787
      shutdownNameResolverAndLoadBalancer(false);
1 ✔
788
    } finally {
789
      updateSubchannelPicker(new LoadBalancer.FixedResultPicker(PickResult.withDrop(
1 ✔
790
          Status.INTERNAL.withDescription("Panic! This is a bug!").withCause(t))));
1 ✔
791
      realChannel.updateConfigSelector(null);
1 ✔
792
      channelLogger.log(ChannelLogLevel.ERROR, "PANIC! Entering TRANSIENT_FAILURE");
1 ✔
793
      channelStateManager.gotoState(TRANSIENT_FAILURE);
1 ✔
794
    }
795
  }
1 ✔
796

797
  @VisibleForTesting
798
  boolean isInPanicMode() {
799
    return panicMode;
1 ✔
800
  }
801

802
  // Called from syncContext
803
  private void updateSubchannelPicker(SubchannelPicker newPicker) {
804
    delayedTransport.reprocess(newPicker);
1 ✔
805
  }
1 ✔
806

807
  @Override
808
  public boolean isShutdown() {
809
    return shutdown.get();
1 ✔
810
  }
811

812
  @Override
813
  public boolean awaitTermination(long timeout, TimeUnit unit) throws InterruptedException {
814
    return terminatedLatch.await(timeout, unit);
1 ✔
815
  }
816

817
  @Override
818
  public boolean isTerminated() {
819
    return terminated;
1 ✔
820
  }
821

822
  /*
823
   * Creates a new outgoing call on the channel.
824
   */
825
  @Override
826
  public <ReqT, RespT> ClientCall<ReqT, RespT> newCall(MethodDescriptor<ReqT, RespT> method,
827
      CallOptions callOptions) {
828
    return interceptorChannel.newCall(method, callOptions);
1 ✔
829
  }
830

831
  @Override
832
  public String authority() {
833
    return interceptorChannel.authority();
1 ✔
834
  }
835

836
  private Executor getCallExecutor(CallOptions callOptions) {
837
    Executor executor = callOptions.getExecutor();
1 ✔
838
    if (executor == null) {
1 ✔
839
      executor = this.executor;
1 ✔
840
    }
841
    return executor;
1 ✔
842
  }
843

844
  private class RealChannel extends Channel {
845
    // Reference to null if no config selector is available from resolution result
846
    // Reference must be set() from syncContext
847
    private final AtomicReference<InternalConfigSelector> configSelector =
1 ✔
848
        new AtomicReference<>(INITIAL_PENDING_SELECTOR);
1 ✔
849
    // Set when the NameResolver is initially created. When we create a new NameResolver for the
850
    // same target, the new instance must have the same value.
851
    private final String authority;
852

853
    private final Channel clientCallImplChannel = new Channel() {
1 ✔
854
      @Override
855
      public <RequestT, ResponseT> ClientCall<RequestT, ResponseT> newCall(
856
          MethodDescriptor<RequestT, ResponseT> method, CallOptions callOptions) {
857
        return new ClientCallImpl<>(
1 ✔
858
            method,
859
            getCallExecutor(callOptions),
1 ✔
860
            callOptions,
861
            transportProvider,
1 ✔
862
            terminated ? null : transportFactory.getScheduledExecutorService(),
1 ✔
863
            channelCallTracer,
1 ✔
864
            null)
865
            .setFullStreamDecompression(fullStreamDecompression)
1 ✔
866
            .setDecompressorRegistry(decompressorRegistry)
1 ✔
867
            .setCompressorRegistry(compressorRegistry);
1 ✔
868
      }
869

870
      @Override
871
      public String authority() {
872
        return authority;
×
873
      }
874
    };
875

876
    private RealChannel(String authority) {
1 ✔
877
      this.authority =  checkNotNull(authority, "authority");
1 ✔
878
    }
1 ✔
879

880
    @Override
881
    public <ReqT, RespT> ClientCall<ReqT, RespT> newCall(
882
        MethodDescriptor<ReqT, RespT> method, CallOptions callOptions) {
883
      if (configSelector.get() != INITIAL_PENDING_SELECTOR) {
1 ✔
884
        return newClientCall(method, callOptions);
1 ✔
885
      }
886
      syncContext.execute(new Runnable() {
1 ✔
887
        @Override
888
        public void run() {
889
          exitIdleMode();
1 ✔
890
        }
1 ✔
891
      });
892
      if (configSelector.get() != INITIAL_PENDING_SELECTOR) {
1 ✔
893
        // This is an optimization for the case (typically with InProcessTransport) when name
894
        // resolution result is immediately available at this point. Otherwise, some users'
895
        // tests might observe slight behavior difference from earlier grpc versions.
896
        return newClientCall(method, callOptions);
1 ✔
897
      }
898
      if (shutdown.get()) {
1 ✔
899
        // Return a failing ClientCall.
900
        return new ClientCall<ReqT, RespT>() {
×
901
          @Override
902
          public void start(Listener<RespT> responseListener, Metadata headers) {
903
            responseListener.onClose(SHUTDOWN_STATUS, new Metadata());
×
904
          }
×
905

906
          @Override public void request(int numMessages) {}
×
907

908
          @Override public void cancel(@Nullable String message, @Nullable Throwable cause) {}
×
909

910
          @Override public void halfClose() {}
×
911

912
          @Override public void sendMessage(ReqT message) {}
×
913
        };
914
      }
915
      Context context = Context.current();
1 ✔
916
      final PendingCall<ReqT, RespT> pendingCall = new PendingCall<>(context, method, callOptions);
1 ✔
917
      syncContext.execute(new Runnable() {
1 ✔
918
        @Override
919
        public void run() {
920
          if (configSelector.get() == INITIAL_PENDING_SELECTOR) {
1 ✔
921
            if (pendingCalls == null) {
1 ✔
922
              pendingCalls = new LinkedHashSet<>();
1 ✔
923
              inUseStateAggregator.updateObjectInUse(pendingCallsInUseObject, true);
1 ✔
924
            }
925
            pendingCalls.add(pendingCall);
1 ✔
926
            pendingCall.notifyQueuedForNameResolution();
1 ✔
927
          } else {
928
            pendingCall.reprocess();
1 ✔
929
          }
930
        }
1 ✔
931
      });
932
      return pendingCall;
1 ✔
933
    }
934

935
    // Must run in SynchronizationContext.
936
    void updateConfigSelector(@Nullable InternalConfigSelector config) {
937
      InternalConfigSelector prevConfig = configSelector.get();
1 ✔
938
      configSelector.set(config);
1 ✔
939
      if (prevConfig == INITIAL_PENDING_SELECTOR
1 ✔
940
          && config != INITIAL_PENDING_SELECTOR && pendingCalls != null) {
1 ✔
941
        for (RealChannel.PendingCall<?, ?> pendingCall : pendingCalls) {
1 ✔
942
          pendingCall.reprocess();
1 ✔
943
        }
1 ✔
944
      }
945
    }
1 ✔
946

947
    // Must run in SynchronizationContext.
948
    void onConfigError() {
949
      if (configSelector.get() == INITIAL_PENDING_SELECTOR) {
1 ✔
950
        // Apply Default Service Config if initial name resolution fails.
951
        if (defaultServiceConfig != null) {
1 ✔
952
          updateConfigSelector(defaultServiceConfig.getDefaultConfigSelector());
1 ✔
953
          lastServiceConfig = defaultServiceConfig;
1 ✔
954
          channelLogger.log(ChannelLogLevel.ERROR,
1 ✔
955
              "Initial Name Resolution error, using default service config");
956
        } else {
957
          updateConfigSelector(null);
1 ✔
958
        }
959
      }
960
    }
1 ✔
961

962
    void shutdown() {
963
      final class RealChannelShutdown implements Runnable {
1 ✔
964
        @Override
965
        public void run() {
966
          if (pendingCalls == null) {
1 ✔
967
            if (configSelector.get() == INITIAL_PENDING_SELECTOR) {
1 ✔
968
              configSelector.set(null);
1 ✔
969
            }
970
            uncommittedRetriableStreamsRegistry.onShutdown(SHUTDOWN_STATUS);
1 ✔
971
          }
972
        }
1 ✔
973
      }
974

975
      syncContext.execute(new RealChannelShutdown());
1 ✔
976
    }
1 ✔
977

978
    void shutdownNow() {
979
      final class RealChannelShutdownNow implements Runnable {
1 ✔
980
        @Override
981
        public void run() {
982
          if (configSelector.get() == INITIAL_PENDING_SELECTOR) {
1 ✔
983
            configSelector.set(null);
1 ✔
984
          }
985
          if (pendingCalls != null) {
1 ✔
986
            for (RealChannel.PendingCall<?, ?> pendingCall : pendingCalls) {
1 ✔
987
              pendingCall.cancel("Channel is forcefully shutdown", null);
1 ✔
988
            }
1 ✔
989
          }
990
          uncommittedRetriableStreamsRegistry.onShutdownNow(SHUTDOWN_NOW_STATUS);
1 ✔
991
        }
1 ✔
992
      }
993

994
      syncContext.execute(new RealChannelShutdownNow());
1 ✔
995
    }
1 ✔
996

997
    @Override
998
    public String authority() {
999
      return authority;
1 ✔
1000
    }
1001

1002
    private final class PendingCall<ReqT, RespT> extends DelayedClientCall<ReqT, RespT> {
1003
      final Context context;
1004
      final MethodDescriptor<ReqT, RespT> method;
1005
      final CallOptions callOptions;
1006
      private final long callCreationTime;
1007
      @GuardedBy("this")
1008
      private boolean queuedForResolution;
1009
      @GuardedBy("this")
1010
      private boolean delayEnded;
1011

1012
      PendingCall(Context context, MethodDescriptor<ReqT, RespT> method, CallOptions callOptions) {
1 ✔
1013
        super(
1 ✔
1014
            "name_resolver",
1015
            getCallExecutor(callOptions),
1 ✔
1016
            scheduledExecutor,
1 ✔
1017
            callOptions.getDeadline());
1 ✔
1018
        this.context = context;
1 ✔
1019
        this.method = method;
1 ✔
1020
        this.callOptions = callOptions;
1 ✔
1021
        this.callCreationTime = ticker.nanoTime();
1 ✔
1022
      }
1 ✔
1023

1024
      private synchronized void notifyQueuedForNameResolution() {
1025
        if (delayEnded || queuedForResolution) {
1 ✔
1026
          return;
1 ✔
1027
        }
1028
        queuedForResolution = true;
1 ✔
1029
        for (ClientStreamTracer.Factory factory : callOptions.getStreamTracerFactories()) {
1 ✔
1030
          if (delayEnded) {
1 ✔
1031
            break;
1 ✔
1032
          }
1033
          factory.recordDelayStart(
1 ✔
1034
              "resolving", "waiting for name resolution or service config");
1035
        }
1 ✔
1036
      }
1 ✔
1037

1038
      private synchronized void endDelayIfNeeded() {
1039
        if (delayEnded) {
1 ✔
1040
          return;
1 ✔
1041
        }
1042
        delayEnded = true;
1 ✔
1043
        if (queuedForResolution) {
1 ✔
1044
          for (ClientStreamTracer.Factory factory : callOptions.getStreamTracerFactories()) {
1 ✔
1045
            factory.recordDelayEnd("resolving");
1 ✔
1046
          }
1 ✔
1047
        }
1048
      }
1 ✔
1049

1050
      /** Called when it's ready to create a real call and reprocess the pending call. */
1051
      void reprocess() {
1052
        endDelayIfNeeded();
1 ✔
1053
        ClientCall<ReqT, RespT> realCall;
1054
        Context previous = context.attach();
1 ✔
1055
        try {
1056
          CallOptions delayResolutionOption = callOptions.withOption(NAME_RESOLUTION_DELAYED,
1 ✔
1057
              ticker.nanoTime() - callCreationTime);
1 ✔
1058
          realCall = newClientCall(method, delayResolutionOption);
1 ✔
1059
        } finally {
1060
          context.detach(previous);
1 ✔
1061
        }
1062
        Runnable toRun = setCall(realCall);
1 ✔
1063
        if (toRun == null) {
1 ✔
1064
          syncContext.execute(new PendingCallRemoval());
1 ✔
1065
        } else {
1066
          getCallExecutor(callOptions).execute(new Runnable() {
1 ✔
1067
            @Override
1068
            public void run() {
1069
              toRun.run();
1 ✔
1070
              syncContext.execute(new PendingCallRemoval());
1 ✔
1071
            }
1 ✔
1072
          });
1073
        }
1074
      }
1 ✔
1075

1076
      @Override
1077
      protected void callCancelled() {
1078
        endDelayIfNeeded();
1 ✔
1079
        super.callCancelled();
1 ✔
1080
        syncContext.execute(new PendingCallRemoval());
1 ✔
1081
      }
1 ✔
1082

1083
      final class PendingCallRemoval implements Runnable {
1 ✔
1084
        @Override
1085
        public void run() {
1086
          if (pendingCalls != null) {
1 ✔
1087
            pendingCalls.remove(PendingCall.this);
1 ✔
1088
            if (pendingCalls.isEmpty()) {
1 ✔
1089
              inUseStateAggregator.updateObjectInUse(pendingCallsInUseObject, false);
1 ✔
1090
              pendingCalls = null;
1 ✔
1091
              if (shutdown.get()) {
1 ✔
1092
                uncommittedRetriableStreamsRegistry.onShutdown(SHUTDOWN_STATUS);
1 ✔
1093
              }
1094
            }
1095
          }
1096
        }
1 ✔
1097
      }
1098
    }
1099

1100
    private <ReqT, RespT> ClientCall<ReqT, RespT> newClientCall(
1101
        MethodDescriptor<ReqT, RespT> method, CallOptions callOptions) {
1102
      InternalConfigSelector selector = configSelector.get();
1 ✔
1103
      if (selector == null) {
1 ✔
1104
        return clientCallImplChannel.newCall(method, callOptions);
1 ✔
1105
      }
1106
      if (selector instanceof ServiceConfigConvertedSelector) {
1 ✔
1107
        MethodInfo methodInfo =
1 ✔
1108
            ((ServiceConfigConvertedSelector) selector).config.getMethodConfig(method);
1 ✔
1109
        if (methodInfo != null) {
1 ✔
1110
          callOptions = callOptions.withOption(MethodInfo.KEY, methodInfo);
1 ✔
1111
        }
1112
        return clientCallImplChannel.newCall(method, callOptions);
1 ✔
1113
      }
1114
      return new ConfigSelectingClientCall<>(
1 ✔
1115
          selector, clientCallImplChannel, executor, method, callOptions);
1 ✔
1116
    }
1117
  }
1118

1119
  /**
1120
   * A client call for a given channel that applies a given config selector when it starts.
1121
   */
1122
  static final class ConfigSelectingClientCall<ReqT, RespT>
1123
      extends ForwardingClientCall<ReqT, RespT> {
1124

1125
    private final InternalConfigSelector configSelector;
1126
    private final Channel channel;
1127
    private final Executor callExecutor;
1128
    private final MethodDescriptor<ReqT, RespT> method;
1129
    private final Context context;
1130
    private CallOptions callOptions;
1131

1132
    private ClientCall<ReqT, RespT> delegate;
1133

1134
    ConfigSelectingClientCall(
1135
        InternalConfigSelector configSelector, Channel channel, Executor channelExecutor,
1136
        MethodDescriptor<ReqT, RespT> method,
1137
        CallOptions callOptions) {
1 ✔
1138
      this.configSelector = configSelector;
1 ✔
1139
      this.channel = channel;
1 ✔
1140
      this.method = method;
1 ✔
1141
      this.callExecutor =
1 ✔
1142
          callOptions.getExecutor() == null ? channelExecutor : callOptions.getExecutor();
1 ✔
1143
      this.callOptions = callOptions.withExecutor(callExecutor);
1 ✔
1144
      this.context = Context.current();
1 ✔
1145
    }
1 ✔
1146

1147
    @Override
1148
    protected ClientCall<ReqT, RespT> delegate() {
1149
      return delegate;
1 ✔
1150
    }
1151

1152
    @SuppressWarnings("unchecked")
1153
    @Override
1154
    public void start(Listener<RespT> observer, Metadata headers) {
1155
      PickSubchannelArgs args =
1 ✔
1156
          new PickSubchannelArgsImpl(method, headers, callOptions, NOOP_PICK_DETAILS_CONSUMER);
1 ✔
1157
      InternalConfigSelector.Result result = configSelector.selectConfig(args);
1 ✔
1158
      Status status = result.getStatus();
1 ✔
1159
      if (!status.isOk()) {
1 ✔
1160
        executeCloseObserverInContext(observer,
1 ✔
1161
            GrpcUtil.replaceInappropriateControlPlaneStatus(status));
1 ✔
1162
        delegate = (ClientCall<ReqT, RespT>) NOOP_CALL;
1 ✔
1163
        return;
1 ✔
1164
      }
1165
      ClientInterceptor interceptor = result.getInterceptor();
1 ✔
1166
      ManagedChannelServiceConfig config = (ManagedChannelServiceConfig) result.getConfig();
1 ✔
1167
      MethodInfo methodInfo = config.getMethodConfig(method);
1 ✔
1168
      if (methodInfo != null) {
1 ✔
1169
        callOptions = callOptions.withOption(MethodInfo.KEY, methodInfo);
1 ✔
1170
      }
1171
      if (interceptor != null) {
1 ✔
1172
        delegate = interceptor.interceptCall(method, callOptions, channel);
1 ✔
1173
      } else {
1174
        delegate = channel.newCall(method, callOptions);
1 ✔
1175
      }
1176
      delegate.start(observer, headers);
1 ✔
1177
    }
1 ✔
1178

1179
    private void executeCloseObserverInContext(
1180
        final Listener<RespT> observer, final Status status) {
1181
      class CloseInContext extends ContextRunnable {
1182
        CloseInContext() {
1 ✔
1183
          super(context);
1 ✔
1184
        }
1 ✔
1185

1186
        @Override
1187
        public void runInContext() {
1188
          observer.onClose(status, new Metadata());
1 ✔
1189
        }
1 ✔
1190
      }
1191

1192
      callExecutor.execute(new CloseInContext());
1 ✔
1193
    }
1 ✔
1194

1195
    @Override
1196
    public void cancel(@Nullable String message, @Nullable Throwable cause) {
1197
      if (delegate != null) {
×
1198
        delegate.cancel(message, cause);
×
1199
      }
1200
    }
×
1201
  }
1202

1203
  private static final ClientCall<Object, Object> NOOP_CALL = new ClientCall<Object, Object>() {
1 ✔
1204
    @Override
1205
    public void start(Listener<Object> responseListener, Metadata headers) {}
×
1206

1207
    @Override
1208
    public void request(int numMessages) {}
1 ✔
1209

1210
    @Override
1211
    public void cancel(String message, Throwable cause) {}
×
1212

1213
    @Override
1214
    public void halfClose() {}
×
1215

1216
    @Override
1217
    public void sendMessage(Object message) {}
×
1218

1219
    // Always returns {@code false}, since this is only used when the startup of the call fails.
1220
    @Override
1221
    public boolean isReady() {
1222
      return false;
×
1223
    }
1224
  };
1225

1226
  /**
1227
   * Terminate the channel if termination conditions are met.
1228
   */
1229
  // Must be run from syncContext
1230
  private void maybeTerminateChannel() {
1231
    if (terminated) {
1 ✔
1232
      return;
×
1233
    }
1234
    if (shutdown.get() && subchannels.isEmpty()) {
1 ✔
1235
      channelLogger.log(ChannelLogLevel.INFO, "Terminated");
1 ✔
1236
      channelz.removeRootChannel(this);
1 ✔
1237
      executorPool.returnObject(executor);
1 ✔
1238
      balancerRpcExecutorHolder.release();
1 ✔
1239
      offloadExecutorHolder.release();
1 ✔
1240
      // Release the transport factory so that it can deallocate any resources.
1241
      transportFactory.close();
1 ✔
1242

1243
      terminated = true;
1 ✔
1244
      terminatedLatch.countDown();
1 ✔
1245
    }
1246
  }
1 ✔
1247

1248
  @Override
1249
  public ConnectivityState getState(boolean requestConnection) {
1250
    ConnectivityState savedChannelState = channelStateManager.getState();
1 ✔
1251
    if (requestConnection && savedChannelState == IDLE) {
1 ✔
1252
      final class RequestConnection implements Runnable {
1 ✔
1253
        @Override
1254
        public void run() {
1255
          exitIdleMode();
1 ✔
1256
          if (lbHelper != null) {
1 ✔
1257
            lbHelper.lb.requestConnection();
1 ✔
1258
          }
1259
        }
1 ✔
1260
      }
1261

1262
      syncContext.execute(new RequestConnection());
1 ✔
1263
    }
1264
    return savedChannelState;
1 ✔
1265
  }
1266

1267
  @Override
1268
  public void notifyWhenStateChanged(final ConnectivityState source, final Runnable callback) {
1269
    final class NotifyStateChanged implements Runnable {
1 ✔
1270
      @Override
1271
      public void run() {
1272
        channelStateManager.notifyWhenStateChanged(callback, executor, source);
1 ✔
1273
      }
1 ✔
1274
    }
1275

1276
    syncContext.execute(new NotifyStateChanged());
1 ✔
1277
  }
1 ✔
1278

1279
  @Override
1280
  public void resetConnectBackoff() {
1281
    final class ResetConnectBackoff implements Runnable {
1 ✔
1282
      @Override
1283
      public void run() {
1284
        if (shutdown.get()) {
1 ✔
1285
          return;
1 ✔
1286
        }
1287
        if (nameResolverStarted) {
1 ✔
1288
          refreshNameResolution();
1 ✔
1289
        }
1290
        for (InternalSubchannel subchannel : subchannels) {
1 ✔
1291
          subchannel.resetConnectBackoff();
1 ✔
1292
        }
1 ✔
1293
      }
1 ✔
1294
    }
1295

1296
    syncContext.execute(new ResetConnectBackoff());
1 ✔
1297
  }
1 ✔
1298

1299
  @Override
1300
  public void enterIdle() {
1301
    final class PrepareToLoseNetworkRunnable implements Runnable {
1 ✔
1302
      @Override
1303
      public void run() {
1304
        if (shutdown.get() || lbHelper == null) {
1 ✔
1305
          return;
1 ✔
1306
        }
1307
        cancelIdleTimer(/* permanent= */ false);
1 ✔
1308
        enterIdleMode();
1 ✔
1309
      }
1 ✔
1310
    }
1311

1312
    syncContext.execute(new PrepareToLoseNetworkRunnable());
1 ✔
1313
  }
1 ✔
1314

1315
  /**
1316
   * A registry that prevents channel shutdown from killing existing retry attempts that are in
1317
   * backoff.
1318
   */
1319
  private final class UncommittedRetriableStreamsRegistry {
1 ✔
1320
    // TODO(zdapeng): This means we would acquire a lock for each new retry-able stream,
1321
    // it's worthwhile to look for a lock-free approach.
1322
    final Object lock = new Object();
1 ✔
1323

1324
    @GuardedBy("lock")
1 ✔
1325
    Collection<ClientStream> uncommittedRetriableStreams = new HashSet<>();
1326

1327
    @GuardedBy("lock")
1328
    Status shutdownStatus;
1329

1330
    void onShutdown(Status reason) {
1331
      boolean shouldShutdownDelayedTransport = false;
1 ✔
1332
      synchronized (lock) {
1 ✔
1333
        if (shutdownStatus != null) {
1 ✔
1334
          return;
1 ✔
1335
        }
1336
        shutdownStatus = reason;
1 ✔
1337
        // Keep the delayedTransport open until there is no more uncommitted streams, b/c those
1338
        // retriable streams, which may be in backoff and not using any transport, are already
1339
        // started RPCs.
1340
        if (uncommittedRetriableStreams.isEmpty()) {
1 ✔
1341
          shouldShutdownDelayedTransport = true;
1 ✔
1342
        }
1343
      }
1 ✔
1344

1345
      if (shouldShutdownDelayedTransport) {
1 ✔
1346
        delayedTransport.shutdown(reason);
1 ✔
1347
      }
1348
    }
1 ✔
1349

1350
    void onShutdownNow(Status reason) {
1351
      onShutdown(reason);
1 ✔
1352
      Collection<ClientStream> streams;
1353

1354
      synchronized (lock) {
1 ✔
1355
        streams = new ArrayList<>(uncommittedRetriableStreams);
1 ✔
1356
      }
1 ✔
1357

1358
      for (ClientStream stream : streams) {
1 ✔
1359
        stream.cancel(reason);
1 ✔
1360
      }
1 ✔
1361
      delayedTransport.shutdownNow(reason);
1 ✔
1362
    }
1 ✔
1363

1364
    /**
1365
     * Registers a RetriableStream and return null if not shutdown, otherwise just returns the
1366
     * shutdown Status.
1367
     */
1368
    @Nullable
1369
    Status add(RetriableStream<?> retriableStream) {
1370
      synchronized (lock) {
1 ✔
1371
        if (shutdownStatus != null) {
1 ✔
1372
          return shutdownStatus;
1 ✔
1373
        }
1374
        uncommittedRetriableStreams.add(retriableStream);
1 ✔
1375
        return null;
1 ✔
1376
      }
1377
    }
1378

1379
    void remove(RetriableStream<?> retriableStream) {
1380
      Status shutdownStatusCopy = null;
1 ✔
1381

1382
      synchronized (lock) {
1 ✔
1383
        uncommittedRetriableStreams.remove(retriableStream);
1 ✔
1384
        if (uncommittedRetriableStreams.isEmpty()) {
1 ✔
1385
          shutdownStatusCopy = shutdownStatus;
1 ✔
1386
          // Because retriable transport is long-lived, we take this opportunity to down-size the
1387
          // hashmap.
1388
          uncommittedRetriableStreams = new HashSet<>();
1 ✔
1389
        }
1390
      }
1 ✔
1391

1392
      if (shutdownStatusCopy != null) {
1 ✔
1393
        delayedTransport.shutdown(shutdownStatusCopy);
1 ✔
1394
      }
1395
    }
1 ✔
1396
  }
1397

1398
  private final class LbHelperImpl extends LoadBalancer.Helper {
1 ✔
1399
    LoadBalancer lb;
1400

1401
    @Override
1402
    public AbstractSubchannel createSubchannel(CreateSubchannelArgs args) {
1403
      syncContext.throwIfNotInThisSynchronizationContext();
1 ✔
1404
      // No new subchannel should be created after load balancer has been shutdown.
1405
      checkState(!terminating, "Channel is being terminated");
1 ✔
1406
      return new SubchannelImpl(args);
1 ✔
1407
    }
1408

1409
    @Override
1410
    public void updateBalancingState(
1411
        final ConnectivityState newState, final SubchannelPicker newPicker) {
1412
      syncContext.throwIfNotInThisSynchronizationContext();
1 ✔
1413
      checkNotNull(newState, "newState");
1 ✔
1414
      checkNotNull(newPicker, "newPicker");
1 ✔
1415

1416
      if (LbHelperImpl.this != lbHelper || panicMode) {
1 ✔
1417
        return;
1 ✔
1418
      }
1419
      updateSubchannelPicker(newPicker);
1 ✔
1420
      // It's not appropriate to report SHUTDOWN state from lb.
1421
      // Ignore the case of newState == SHUTDOWN for now.
1422
      if (newState != SHUTDOWN) {
1 ✔
1423
        channelLogger.log(
1 ✔
1424
            ChannelLogLevel.INFO, "Entering {0} state with picker: {1}", newState, newPicker);
1425
        channelStateManager.gotoState(newState);
1 ✔
1426
      }
1427
    }
1 ✔
1428

1429
    @Override
1430
    public void refreshNameResolution() {
1431
      syncContext.throwIfNotInThisSynchronizationContext();
1 ✔
1432
      final class LoadBalancerRefreshNameResolution implements Runnable {
1 ✔
1433
        @Override
1434
        public void run() {
1435
          ManagedChannelImpl.this.refreshNameResolution();
1 ✔
1436
        }
1 ✔
1437
      }
1438

1439
      syncContext.execute(new LoadBalancerRefreshNameResolution());
1 ✔
1440
    }
1 ✔
1441

1442
    @Override
1443
    public ManagedChannel createOobChannel(EquivalentAddressGroup addressGroup, String authority) {
1444
      return createOobChannel(Collections.singletonList(addressGroup), authority);
×
1445
    }
1446

1447
    @Override
1448
    public ManagedChannel createOobChannel(List<EquivalentAddressGroup> addressGroup,
1449
        String authority) {
1450
      NameResolverRegistry nameResolverRegistry = new NameResolverRegistry();
1 ✔
1451
      OobNameResolverProvider resolverProvider =
1 ✔
1452
          new OobNameResolverProvider(authority, addressGroup, syncContext);
1453
      nameResolverRegistry.register(resolverProvider);
1 ✔
1454
      // We could use a hard-coded target, as the name resolver won't actually use this string.
1455
      // However, that would make debugging less clear, as we use the target to identify the
1456
      // channel.
1457
      String target;
1458
      try {
1459
        target = new URI("oob", "", "/" + authority, null, null).toString();
1 ✔
1460
      } catch (URISyntaxException ex) {
×
1461
        // Any special characters in the path will be percent encoded. So this should be impossible.
1462
        throw new AssertionError(ex);
×
1463
      }
1 ✔
1464
      ManagedChannel delegate = createResolvingOobChannelBuilder(
1 ✔
1465
          target, new DefaultChannelCreds(), nameResolverRegistry)
1466
          // TODO(zdapeng): executors should not outlive the parent channel.
1467
          .executor(balancerRpcExecutorHolder.getExecutor())
1 ✔
1468
          .idleTimeout(Integer.MAX_VALUE, TimeUnit.SECONDS)
1 ✔
1469
          .disableRetry()
1 ✔
1470
          .build();
1 ✔
1471
      return new OobChannel(delegate, resolverProvider);
1 ✔
1472
    }
1473

1474
    @Deprecated
1475
    @Override
1476
    public ManagedChannelBuilder<?> createResolvingOobChannelBuilder(String target) {
1477
      return createResolvingOobChannelBuilder(target, new DefaultChannelCreds())
1 ✔
1478
          // Override authority to keep the old behavior.
1479
          // createResolvingOobChannelBuilder(String target) will be deleted soon.
1480
          .overrideAuthority(getAuthority());
1 ✔
1481
    }
1482

1483
    @Override
1484
    public ManagedChannelBuilder<?> createResolvingOobChannelBuilder(
1485
        final String target, final ChannelCredentials channelCreds) {
1486
      return createResolvingOobChannelBuilder(target, channelCreds, nameResolverRegistry);
1 ✔
1487
    }
1488

1489
    // TODO(creamsoup) prevent main channel to shutdown if oob channel is not terminated
1490
    // TODO(zdapeng) register the channel as a subchannel of the parent channel in channelz.
1491
    private ManagedChannelBuilder<?> createResolvingOobChannelBuilder(
1492
        final String target, final ChannelCredentials channelCreds,
1493
        NameResolverRegistry nameResolverRegistry) {
1494
      checkNotNull(channelCreds, "channelCreds");
1 ✔
1495

1496
      final class ResolvingOobChannelBuilder
1497
          extends ForwardingChannelBuilder2<ResolvingOobChannelBuilder> {
1498
        final ManagedChannelBuilder<?> delegate;
1499

1500
        ResolvingOobChannelBuilder() {
1 ✔
1501
          final ClientTransportFactory transportFactory;
1502
          CallCredentials callCredentials;
1503
          if (channelCreds instanceof DefaultChannelCreds) {
1 ✔
1504
            // TODO(kannanjgithub) We should eventually refactor ManagedChannelImplBuilder so
1505
            // callCredentials can be resolved lazily at build() time, allowing transport factory
1506
            // retention to happen strictly inside buildClientTransportFactory().
1507
            transportFactory = originalTransportFactory.retain();
1 ✔
1508
            callCredentials = null;
1 ✔
1509
          } else {
1510
            SwapChannelCredentialsResult swapResult =
1 ✔
1511
                originalTransportFactory.swapChannelCredentials(channelCreds);
1 ✔
1512
            if (swapResult == null) {
1 ✔
1513
              delegate = Grpc.newChannelBuilder(target, channelCreds);
×
1514
              return;
×
1515
            } else {
1516
              transportFactory = swapResult.transportFactory;
1 ✔
1517
              callCredentials = swapResult.callCredentials;
1 ✔
1518
            }
1519
          }
1520
          ClientTransportFactoryBuilder transportFactoryBuilder =
1 ✔
1521
              new ClientTransportFactoryBuilder() {
1 ✔
1522
                @Override
1523
                public ClientTransportFactory buildClientTransportFactory() {
1524
                  return transportFactory;
1 ✔
1525
                }
1526
              };
1527
          delegate = new ManagedChannelImplBuilder(
1 ✔
1528
              target,
1529
              channelCreds,
1530
              callCredentials,
1531
              transportFactoryBuilder,
1532
              new FixedPortProvider(nameResolverArgs.getDefaultPort()))
1 ✔
1533
              .nameResolverRegistry(nameResolverRegistry);
1 ✔
1534
        }
1 ✔
1535

1536
        @Override
1537
        protected ManagedChannelBuilder<?> delegate() {
1538
          return delegate;
1 ✔
1539
        }
1540
      }
1541

1542
      checkState(!terminated, "Channel is terminated");
1 ✔
1543

1544
      ResolvingOobChannelBuilder builder = new ResolvingOobChannelBuilder();
1 ✔
1545

1546
      // Note that we follow the global configurator pattern and try to fuse the configurations as
1547
      // soon as the builder gets created
1548
      channelConfigurator.configureChannelBuilder(builder);
1 ✔
1549
      builder.childChannelConfigurator(channelConfigurator);
1 ✔
1550

1551
      return builder
1 ✔
1552
          // TODO(zdapeng): executors should not outlive the parent channel.
1553
          .executor(executor)
1 ✔
1554
          .offloadExecutor(offloadExecutorHolder.getExecutor())
1 ✔
1555
          .maxTraceEvents(maxTraceEvents)
1 ✔
1556
          .proxyDetector(nameResolverArgs.getProxyDetector())
1 ✔
1557
          .userAgent(userAgent);
1 ✔
1558
    }
1559

1560
    @Override
1561
    public ChannelCredentials getUnsafeChannelCredentials() {
1562
      if (originalChannelCreds == null) {
1 ✔
1563
        return new DefaultChannelCreds();
1 ✔
1564
      }
1565
      return originalChannelCreds;
×
1566
    }
1567

1568
    @Override
1569
    public void updateOobChannelAddresses(ManagedChannel channel, EquivalentAddressGroup eag) {
1570
      updateOobChannelAddresses(channel, Collections.singletonList(eag));
×
1571
    }
×
1572

1573
    @Override
1574
    public void updateOobChannelAddresses(ManagedChannel channel,
1575
        List<EquivalentAddressGroup> eag) {
1576
      checkArgument(channel instanceof OobChannel,
1 ✔
1577
          "channel must have been returned from createOobChannel");
1578
      ((OobChannel) channel).updateAddresses(eag);
1 ✔
1579
    }
1 ✔
1580

1581
    @Override
1582
    public String getAuthority() {
1583
      return ManagedChannelImpl.this.authority();
1 ✔
1584
    }
1585

1586
    @Override
1587
    public String getChannelTarget() {
1588
      return targetUri.toString();
1 ✔
1589
    }
1590

1591
    @Override
1592
    public SynchronizationContext getSynchronizationContext() {
1593
      return syncContext;
1 ✔
1594
    }
1595

1596
    @Override
1597
    public ScheduledExecutorService getScheduledExecutorService() {
1598
      return scheduledExecutor;
1 ✔
1599
    }
1600

1601
    @Override
1602
    public ChannelLogger getChannelLogger() {
1603
      return channelLogger;
1 ✔
1604
    }
1605

1606
    @Override
1607
    public NameResolver.Args getNameResolverArgs() {
1608
      return nameResolverArgs;
1 ✔
1609
    }
1610

1611
    @Override
1612
    public NameResolverRegistry getNameResolverRegistry() {
1613
      return nameResolverRegistry;
1 ✔
1614
    }
1615

1616
    @Override
1617
    public MetricRecorder getMetricRecorder() {
1618
      return metricRecorder;
1 ✔
1619
    }
1620

1621
    /**
1622
     * A placeholder for channel creds if user did not specify channel creds for the channel.
1623
     */
1624
    // TODO(zdapeng): get rid of this class and let all ChannelBuilders always provide a non-null
1625
    //     channel creds.
1626
    final class DefaultChannelCreds extends ChannelCredentials {
1 ✔
1627
      @Override
1628
      public ChannelCredentials withoutBearerTokens() {
1629
        return this;
×
1630
      }
1631
    }
1632
  }
1633

1634
  static final class OobChannel extends ForwardingManagedChannel {
1635
    private final OobNameResolverProvider resolverProvider;
1636

1637
    public OobChannel(ManagedChannel delegate, OobNameResolverProvider resolverProvider) {
1638
      super(delegate);
1 ✔
1639
      this.resolverProvider = checkNotNull(resolverProvider, "resolverProvider");
1 ✔
1640
    }
1 ✔
1641

1642
    public void updateAddresses(List<EquivalentAddressGroup> eags) {
1643
      resolverProvider.updateAddresses(eags);
1 ✔
1644
    }
1 ✔
1645
  }
1646

1647
  final class NameResolverListener extends NameResolver.Listener2 {
1648
    final LbHelperImpl helper;
1649
    final NameResolver resolver;
1650

1651
    NameResolverListener(LbHelperImpl helperImpl, NameResolver resolver) {
1 ✔
1652
      this.helper = checkNotNull(helperImpl, "helperImpl");
1 ✔
1653
      this.resolver = checkNotNull(resolver, "resolver");
1 ✔
1654
    }
1 ✔
1655

1656
    @Override
1657
    public void onResult(final ResolutionResult resolutionResult) {
1658
      syncContext.execute(() -> onResult2(resolutionResult));
×
1659
    }
×
1660

1661
    @SuppressWarnings("ReferenceEquality")
1662
    @Override
1663
    public Status onResult2(final ResolutionResult resolutionResult) {
1664
      syncContext.throwIfNotInThisSynchronizationContext();
1 ✔
1665
      if (ManagedChannelImpl.this.nameResolver != resolver) {
1 ✔
1666
        return Status.OK;
1 ✔
1667
      }
1668

1669
      StatusOr<List<EquivalentAddressGroup>> serversOrError =
1 ✔
1670
          resolutionResult.getAddressesOrError();
1 ✔
1671
      if (!serversOrError.hasValue()) {
1 ✔
1672
        handleErrorInSyncContext(serversOrError.getStatus());
1 ✔
1673
        return serversOrError.getStatus();
1 ✔
1674
      }
1675
      List<EquivalentAddressGroup> servers = serversOrError.getValue();
1 ✔
1676
      channelLogger.log(
1 ✔
1677
          ChannelLogLevel.DEBUG,
1678
          "Resolved address: {0}, config={1}",
1679
          servers,
1680
          resolutionResult.getAttributes());
1 ✔
1681

1682
      if (lastResolutionState != ResolutionState.SUCCESS) {
1 ✔
1683
        channelLogger.log(ChannelLogLevel.INFO, "Address resolved: {0}",
1 ✔
1684
            servers);
1685
        lastResolutionState = ResolutionState.SUCCESS;
1 ✔
1686
      }
1687
      ConfigOrError configOrError = resolutionResult.getServiceConfig();
1 ✔
1688
      InternalConfigSelector resolvedConfigSelector =
1 ✔
1689
          resolutionResult.getAttributes().get(InternalConfigSelector.KEY);
1 ✔
1690
      ManagedChannelServiceConfig validServiceConfig =
1691
          configOrError != null && configOrError.getConfig() != null
1 ✔
1692
              ? (ManagedChannelServiceConfig) configOrError.getConfig()
1 ✔
1693
              : null;
1 ✔
1694
      Status serviceConfigError = configOrError != null ? configOrError.getError() : null;
1 ✔
1695

1696
      ManagedChannelServiceConfig effectiveServiceConfig;
1697
      if (!lookUpServiceConfig) {
1 ✔
1698
        if (validServiceConfig != null) {
1 ✔
1699
          channelLogger.log(
1 ✔
1700
              ChannelLogLevel.INFO,
1701
              "Service config from name resolver discarded by channel settings");
1702
        }
1703
        effectiveServiceConfig =
1704
            defaultServiceConfig == null ? EMPTY_SERVICE_CONFIG : defaultServiceConfig;
1 ✔
1705
        if (resolvedConfigSelector != null) {
1 ✔
1706
          channelLogger.log(
1 ✔
1707
              ChannelLogLevel.INFO,
1708
              "Config selector from name resolver discarded by channel settings");
1709
        }
1710
        realChannel.updateConfigSelector(effectiveServiceConfig.getDefaultConfigSelector());
1 ✔
1711
      } else {
1712
        // Try to use config if returned from name resolver
1713
        // Otherwise, try to use the default config if available
1714
        if (validServiceConfig != null) {
1 ✔
1715
          effectiveServiceConfig = validServiceConfig;
1 ✔
1716
          if (resolvedConfigSelector != null) {
1 ✔
1717
            realChannel.updateConfigSelector(resolvedConfigSelector);
1 ✔
1718
            if (effectiveServiceConfig.getDefaultConfigSelector() != null) {
1 ✔
1719
              channelLogger.log(
×
1720
                  ChannelLogLevel.DEBUG,
1721
                  "Method configs in service config will be discarded due to presence of"
1722
                      + "config-selector");
1723
            }
1724
          } else {
1725
            realChannel.updateConfigSelector(effectiveServiceConfig.getDefaultConfigSelector());
1 ✔
1726
          }
1727
        } else if (defaultServiceConfig != null) {
1 ✔
1728
          effectiveServiceConfig = defaultServiceConfig;
1 ✔
1729
          realChannel.updateConfigSelector(effectiveServiceConfig.getDefaultConfigSelector());
1 ✔
1730
          channelLogger.log(
1 ✔
1731
              ChannelLogLevel.INFO,
1732
              "Received no service config, using default service config");
1733
        } else if (serviceConfigError != null) {
1 ✔
1734
          if (!serviceConfigUpdated) {
1 ✔
1735
            // First DNS lookup has invalid service config, and cannot fall back to default
1736
            channelLogger.log(
1 ✔
1737
                ChannelLogLevel.INFO,
1738
                "Fallback to error due to invalid first service config without default config");
1739
            // This error could be an "inappropriate" control plane error that should not bleed
1740
            // through to client code using gRPC. We let them flow through here to the LB as
1741
            // we later check for these error codes when investigating pick results in
1742
            // GrpcUtil.getTransportFromPickResult().
1743
            onError(configOrError.getError());
1 ✔
1744
            return configOrError.getError();
1 ✔
1745
          } else {
1746
            effectiveServiceConfig = lastServiceConfig;
1 ✔
1747
          }
1748
        } else {
1749
          effectiveServiceConfig = EMPTY_SERVICE_CONFIG;
1 ✔
1750
          realChannel.updateConfigSelector(null);
1 ✔
1751
        }
1752
        if (!effectiveServiceConfig.equals(lastServiceConfig)) {
1 ✔
1753
          channelLogger.log(
1 ✔
1754
              ChannelLogLevel.INFO,
1755
              "Service config changed{0}",
1756
              effectiveServiceConfig == EMPTY_SERVICE_CONFIG ? " to empty" : "");
1 ✔
1757
          lastServiceConfig = effectiveServiceConfig;
1 ✔
1758
          transportProvider.throttle = effectiveServiceConfig.getRetryThrottling();
1 ✔
1759
        }
1760

1761
        try {
1762
          // TODO(creamsoup): when `serversOrError` is empty and lastResolutionStateCopy == SUCCESS
1763
          //  and lbNeedAddress, it shouldn't call the handleServiceConfigUpdate. But,
1764
          //  lbNeedAddress is not deterministic
1765
          serviceConfigUpdated = true;
1 ✔
1766
        } catch (RuntimeException re) {
×
1767
          logger.log(
×
1768
              Level.WARNING,
1769
              "[" + getLogId() + "] Unexpected exception from parsing service config",
×
1770
              re);
1771
        }
1 ✔
1772
      }
1773

1774
      Attributes effectiveAttrs = resolutionResult.getAttributes();
1 ✔
1775
      // Call LB only if it's not shutdown.  If LB is shutdown, lbHelper won't match.
1776
      if (NameResolverListener.this.helper == ManagedChannelImpl.this.lbHelper) {
1 ✔
1777
        Attributes.Builder attrBuilder =
1 ✔
1778
            effectiveAttrs.toBuilder().discard(InternalConfigSelector.KEY);
1 ✔
1779
        Map<String, ?> healthCheckingConfig =
1 ✔
1780
            effectiveServiceConfig.getHealthCheckingConfig();
1 ✔
1781
        if (healthCheckingConfig != null) {
1 ✔
1782
          attrBuilder
1 ✔
1783
              .set(LoadBalancer.ATTR_HEALTH_CHECKING_CONFIG, healthCheckingConfig)
1 ✔
1784
              .build();
1 ✔
1785
        }
1786
        Attributes attributes = attrBuilder.build();
1 ✔
1787

1788
        ResolvedAddresses.Builder resolvedAddresses = ResolvedAddresses.newBuilder()
1 ✔
1789
            .setAddresses(serversOrError.getValue())
1 ✔
1790
            .setAttributes(attributes)
1 ✔
1791
            .setLoadBalancingPolicyConfig(effectiveServiceConfig.getLoadBalancingConfig());
1 ✔
1792
        Status addressAcceptanceStatus = helper.lb.acceptResolvedAddresses(
1 ✔
1793
            resolvedAddresses.build());
1 ✔
1794
        return addressAcceptanceStatus;
1 ✔
1795
      }
1796
      return Status.OK;
×
1797
    }
1798

1799
    @Override
1800
    public void onError(final Status error) {
1801
      checkArgument(!error.isOk(), "the error status must not be OK");
1 ✔
1802
      final class NameResolverErrorHandler implements Runnable {
1 ✔
1803
        @Override
1804
        public void run() {
1805
          handleErrorInSyncContext(error);
1 ✔
1806
        }
1 ✔
1807
      }
1808

1809
      syncContext.execute(new NameResolverErrorHandler());
1 ✔
1810
    }
1 ✔
1811

1812
    private void handleErrorInSyncContext(Status error) {
1813
      logger.log(Level.WARNING, "[{0}] Failed to resolve name. status={1}",
1 ✔
1814
          new Object[] {getLogId(), error});
1 ✔
1815
      realChannel.onConfigError();
1 ✔
1816
      if (lastResolutionState != ResolutionState.ERROR) {
1 ✔
1817
        channelLogger.log(ChannelLogLevel.WARNING, "Failed to resolve name: {0}", error);
1 ✔
1818
        lastResolutionState = ResolutionState.ERROR;
1 ✔
1819
      }
1820
      // Call LB only if it's not shutdown.  If LB is shutdown, lbHelper won't match.
1821
      if (NameResolverListener.this.helper != ManagedChannelImpl.this.lbHelper) {
1 ✔
1822
        return;
1 ✔
1823
      }
1824

1825
      helper.lb.handleNameResolutionError(error);
1 ✔
1826
    }
1 ✔
1827
  }
1828

1829
  private final class SubchannelImpl extends AbstractSubchannel {
1830
    final CreateSubchannelArgs args;
1831
    final InternalLogId subchannelLogId;
1832
    final ChannelLoggerImpl subchannelLogger;
1833
    final ChannelTracer subchannelTracer;
1834
    List<EquivalentAddressGroup> addressGroups;
1835
    InternalSubchannel subchannel;
1836
    boolean started;
1837
    boolean shutdown;
1838
    ScheduledHandle delayedShutdownTask;
1839

1840
    SubchannelImpl(CreateSubchannelArgs args) {
1 ✔
1841
      checkNotNull(args, "args");
1 ✔
1842
      addressGroups = args.getAddresses();
1 ✔
1843
      if (authorityOverride != null) {
1 ✔
1844
        List<EquivalentAddressGroup> eagsWithoutOverrideAttr =
1 ✔
1845
            stripOverrideAuthorityAttributes(args.getAddresses());
1 ✔
1846
        args = args.toBuilder().setAddresses(eagsWithoutOverrideAttr).build();
1 ✔
1847
      }
1848
      this.args = args;
1 ✔
1849
      subchannelLogId = InternalLogId.allocate("Subchannel", /*details=*/ authority());
1 ✔
1850
      subchannelTracer = new ChannelTracer(
1 ✔
1851
          subchannelLogId, maxTraceEvents, timeProvider.currentTimeNanos(),
1 ✔
1852
          "Subchannel for " + args.getAddresses());
1 ✔
1853
      subchannelLogger = new ChannelLoggerImpl(subchannelTracer, timeProvider);
1 ✔
1854
    }
1 ✔
1855

1856
    @Override
1857
    public void start(final SubchannelStateListener listener) {
1858
      syncContext.throwIfNotInThisSynchronizationContext();
1 ✔
1859
      checkState(!started, "already started");
1 ✔
1860
      checkState(!shutdown, "already shutdown");
1 ✔
1861
      checkState(!terminating, "Channel is being terminated");
1 ✔
1862
      started = true;
1 ✔
1863
      final class ManagedInternalSubchannelCallback extends InternalSubchannel.Callback {
1 ✔
1864
        // All callbacks are run in syncContext
1865
        @Override
1866
        void onTerminated(InternalSubchannel is) {
1867
          subchannels.remove(is);
1 ✔
1868
          channelz.removeSubchannel(is);
1 ✔
1869
          maybeTerminateChannel();
1 ✔
1870
        }
1 ✔
1871

1872
        @Override
1873
        void onStateChange(InternalSubchannel is, ConnectivityStateInfo newState) {
1874
          checkState(listener != null, "listener is null");
1 ✔
1875
          listener.onSubchannelState(newState);
1 ✔
1876
        }
1 ✔
1877

1878
        @Override
1879
        void onInUse(InternalSubchannel is) {
1880
          inUseStateAggregator.updateObjectInUse(is, true);
1 ✔
1881
        }
1 ✔
1882

1883
        @Override
1884
        void onNotInUse(InternalSubchannel is) {
1885
          inUseStateAggregator.updateObjectInUse(is, false);
1 ✔
1886
        }
1 ✔
1887
      }
1888

1889
      final InternalSubchannel internalSubchannel = new InternalSubchannel(
1 ✔
1890
          args,
1891
          authority(),
1 ✔
1892
          userAgent,
1 ✔
1893
          backoffPolicyProvider,
1 ✔
1894
          transportFactory,
1 ✔
1895
          transportFactory.getScheduledExecutorService(),
1 ✔
1896
          stopwatchSupplier,
1 ✔
1897
          syncContext,
1898
          new ManagedInternalSubchannelCallback(),
1899
          channelz,
1 ✔
1900
          callTracerFactory.create(),
1 ✔
1901
          subchannelTracer,
1902
          subchannelLogId,
1903
          subchannelLogger,
1904
          transportFilters, target,
1 ✔
1905
          lbHelper.getMetricRecorder());
1 ✔
1906

1907
      channelTracer.reportEvent(new ChannelTrace.Event.Builder()
1 ✔
1908
          .setDescription("Child Subchannel started")
1 ✔
1909
          .setSeverity(ChannelTrace.Event.Severity.CT_INFO)
1 ✔
1910
          .setTimestampNanos(timeProvider.currentTimeNanos())
1 ✔
1911
          .setSubchannelRef(internalSubchannel)
1 ✔
1912
          .build());
1 ✔
1913

1914
      this.subchannel = internalSubchannel;
1 ✔
1915
      channelz.addSubchannel(internalSubchannel);
1 ✔
1916
      subchannels.add(internalSubchannel);
1 ✔
1917
    }
1 ✔
1918

1919
    @Override
1920
    InternalInstrumented<ChannelStats> getInstrumentedInternalSubchannel() {
1921
      checkState(started, "not started");
1 ✔
1922
      return subchannel;
1 ✔
1923
    }
1924

1925
    @Override
1926
    public void shutdown() {
1927
      syncContext.throwIfNotInThisSynchronizationContext();
1 ✔
1928
      if (subchannel == null) {
1 ✔
1929
        // start() was not successful
1930
        shutdown = true;
×
1931
        return;
×
1932
      }
1933
      if (shutdown) {
1 ✔
1934
        if (terminating && delayedShutdownTask != null) {
1 ✔
1935
          // shutdown() was previously called when terminating == false, thus a delayed shutdown()
1936
          // was scheduled.  Now since terminating == true, We should expedite the shutdown.
1937
          delayedShutdownTask.cancel();
×
1938
          delayedShutdownTask = null;
×
1939
          // Will fall through to the subchannel.shutdown() at the end.
1940
        } else {
1941
          return;
1 ✔
1942
        }
1943
      } else {
1944
        shutdown = true;
1 ✔
1945
      }
1946
      // Add a delay to shutdown to deal with the race between 1) a transport being picked and
1947
      // newStream() being called on it, and 2) its Subchannel is shut down by LoadBalancer (e.g.,
1948
      // because of address change, or because LoadBalancer is shutdown by Channel entering idle
1949
      // mode). If (2) wins, the app will see a spurious error. We work around this by delaying
1950
      // shutdown of Subchannel for a few seconds here.
1951
      //
1952
      // TODO(zhangkun83): consider a better approach
1953
      // (https://github.com/grpc/grpc-java/issues/2562).
1954
      if (!terminating) {
1 ✔
1955
        final class ShutdownSubchannel implements Runnable {
1 ✔
1956
          @Override
1957
          public void run() {
1958
            subchannel.shutdown(SUBCHANNEL_SHUTDOWN_STATUS);
1 ✔
1959
          }
1 ✔
1960
        }
1961

1962
        delayedShutdownTask = syncContext.schedule(
1 ✔
1963
            new LogExceptionRunnable(new ShutdownSubchannel()),
1964
            SUBCHANNEL_SHUTDOWN_DELAY_SECONDS, TimeUnit.SECONDS,
1965
            transportFactory.getScheduledExecutorService());
1 ✔
1966
        return;
1 ✔
1967
      }
1968
      // When terminating == true, no more real streams will be created. It's safe and also
1969
      // desirable to shutdown timely.
1970
      subchannel.shutdown(SHUTDOWN_STATUS);
1 ✔
1971
    }
1 ✔
1972

1973
    @Override
1974
    public void requestConnection() {
1975
      syncContext.throwIfNotInThisSynchronizationContext();
1 ✔
1976
      checkState(started, "not started");
1 ✔
1977
      if (shutdown) {
1 ✔
1978
        return;
1 ✔
1979
      }
1980
      subchannel.obtainActiveTransport();
1 ✔
1981
    }
1 ✔
1982

1983
    @Override
1984
    public List<EquivalentAddressGroup> getAllAddresses() {
1985
      syncContext.throwIfNotInThisSynchronizationContext();
1 ✔
1986
      checkState(started, "not started");
1 ✔
1987
      return addressGroups;
1 ✔
1988
    }
1989

1990
    @Override
1991
    public Attributes getAttributes() {
1992
      return args.getAttributes();
1 ✔
1993
    }
1994

1995
    @Override
1996
    public String toString() {
1997
      return subchannelLogId.toString();
1 ✔
1998
    }
1999

2000
    @Override
2001
    public Channel asChannel() {
2002
      checkState(started, "not started");
1 ✔
2003
      return new SubchannelChannel(
1 ✔
2004
          subchannel, balancerRpcExecutorHolder.getExecutor(),
1 ✔
2005
          transportFactory.getScheduledExecutorService(),
1 ✔
2006
          callTracerFactory.create(),
1 ✔
2007
          new AtomicReference<InternalConfigSelector>(null));
2008
    }
2009

2010
    @Override
2011
    public Object getInternalSubchannel() {
2012
      checkState(started, "Subchannel is not started");
1 ✔
2013
      return subchannel;
1 ✔
2014
    }
2015

2016
    @Override
2017
    public ChannelLogger getChannelLogger() {
2018
      return subchannelLogger;
1 ✔
2019
    }
2020

2021
    @Override
2022
    public void updateAddresses(List<EquivalentAddressGroup> addrs) {
2023
      syncContext.throwIfNotInThisSynchronizationContext();
1 ✔
2024
      addressGroups = addrs;
1 ✔
2025
      if (authorityOverride != null) {
1 ✔
2026
        addrs = stripOverrideAuthorityAttributes(addrs);
1 ✔
2027
      }
2028
      subchannel.updateAddresses(addrs);
1 ✔
2029
    }
1 ✔
2030

2031
    @Override
2032
    public Attributes getConnectedAddressAttributes() {
2033
      return subchannel.getConnectedAddressAttributes();
1 ✔
2034
    }
2035

2036
    private List<EquivalentAddressGroup> stripOverrideAuthorityAttributes(
2037
        List<EquivalentAddressGroup> eags) {
2038
      List<EquivalentAddressGroup> eagsWithoutOverrideAttr = new ArrayList<>();
1 ✔
2039
      for (EquivalentAddressGroup eag : eags) {
1 ✔
2040
        EquivalentAddressGroup eagWithoutOverrideAttr = new EquivalentAddressGroup(
1 ✔
2041
            eag.getAddresses(),
1 ✔
2042
            eag.getAttributes().toBuilder().discard(ATTR_AUTHORITY_OVERRIDE).build());
1 ✔
2043
        eagsWithoutOverrideAttr.add(eagWithoutOverrideAttr);
1 ✔
2044
      }
1 ✔
2045
      return Collections.unmodifiableList(eagsWithoutOverrideAttr);
1 ✔
2046
    }
2047
  }
2048

2049
  @Override
2050
  public String toString() {
2051
    return MoreObjects.toStringHelper(this)
1 ✔
2052
        .add("logId", logId.getId())
1 ✔
2053
        .add("target", target)
1 ✔
2054
        .toString();
1 ✔
2055
  }
2056

2057
  /**
2058
   * Called from syncContext.
2059
   */
2060
  private final class DelayedTransportListener implements ManagedClientTransport.Listener {
1 ✔
2061
    @Override
2062
    public void transportShutdown(Status s, DisconnectError e) {
2063
      checkState(shutdown.get(), "Channel must have been shut down");
1 ✔
2064
    }
1 ✔
2065

2066
    @Override
2067
    public void transportReady() {
2068
      // Don't care
2069
    }
×
2070

2071
    @Override
2072
    public Attributes filterTransport(Attributes attributes) {
2073
      return attributes;
×
2074
    }
2075

2076
    @Override
2077
    public void transportInUse(final boolean inUse) {
2078
      inUseStateAggregator.updateObjectInUse(delayedTransport, inUse);
1 ✔
2079
      if (inUse) {
1 ✔
2080
        // It's possible to be in idle mode while inUseStateAggregator is in-use, if one of the
2081
        // subchannels is in use. But we should never be in idle mode when delayed transport is in
2082
        // use.
2083
        exitIdleMode();
1 ✔
2084
      }
2085
    }
1 ✔
2086

2087
    @Override
2088
    public void transportTerminated() {
2089
      checkState(shutdown.get(), "Channel must have been shut down");
1 ✔
2090
      terminating = true;
1 ✔
2091
      shutdownNameResolverAndLoadBalancer(false);
1 ✔
2092
      // No need to call channelStateManager since we are already in SHUTDOWN state.
2093
      // Until LoadBalancer is shutdown, it may still create new subchannels.  We catch them
2094
      // here.
2095
      maybeShutdownNowSubchannels();
1 ✔
2096
      maybeTerminateChannel();
1 ✔
2097
    }
1 ✔
2098
  }
2099

2100
  /**
2101
   * Must be accessed from syncContext.
2102
   */
2103
  private final class IdleModeStateAggregator extends InUseStateAggregator<Object> {
1 ✔
2104
    @Override
2105
    protected void handleInUse() {
2106
      exitIdleMode();
1 ✔
2107
    }
1 ✔
2108

2109
    @Override
2110
    protected void handleNotInUse() {
2111
      if (shutdown.get()) {
1 ✔
2112
        return;
1 ✔
2113
      }
2114
      rescheduleIdleTimer();
1 ✔
2115
    }
1 ✔
2116
  }
2117

2118
  /**
2119
   * Lazily request for Executor from an executor pool.
2120
   * Also act as an Executor directly to simply run a cmd
2121
   */
2122
  @VisibleForTesting
2123
  static final class ExecutorHolder implements Executor {
2124
    private final ObjectPool<? extends Executor> pool;
2125
    private Executor executor;
2126

2127
    ExecutorHolder(ObjectPool<? extends Executor> executorPool) {
1 ✔
2128
      this.pool = checkNotNull(executorPool, "executorPool");
1 ✔
2129
    }
1 ✔
2130

2131
    synchronized Executor getExecutor() {
2132
      if (executor == null) {
1 ✔
2133
        executor = checkNotNull(pool.getObject(), "%s.getObject()", executor);
1 ✔
2134
      }
2135
      return executor;
1 ✔
2136
    }
2137

2138
    synchronized void release() {
2139
      if (executor != null) {
1 ✔
2140
        executor = pool.returnObject(executor);
1 ✔
2141
      }
2142
    }
1 ✔
2143

2144
    @Override
2145
    public void execute(Runnable command) {
2146
      getExecutor().execute(command);
1 ✔
2147
    }
1 ✔
2148
  }
2149

2150
  private static final class RestrictedScheduledExecutor implements ScheduledExecutorService {
2151
    final ScheduledExecutorService delegate;
2152

2153
    private RestrictedScheduledExecutor(ScheduledExecutorService delegate) {
1 ✔
2154
      this.delegate = checkNotNull(delegate, "delegate");
1 ✔
2155
    }
1 ✔
2156

2157
    @Override
2158
    public <V> ScheduledFuture<V> schedule(Callable<V> callable, long delay, TimeUnit unit) {
2159
      return delegate.schedule(callable, delay, unit);
×
2160
    }
2161

2162
    @Override
2163
    public ScheduledFuture<?> schedule(Runnable cmd, long delay, TimeUnit unit) {
2164
      return delegate.schedule(cmd, delay, unit);
1 ✔
2165
    }
2166

2167
    @Override
2168
    public ScheduledFuture<?> scheduleAtFixedRate(
2169
        Runnable command, long initialDelay, long period, TimeUnit unit) {
2170
      return delegate.scheduleAtFixedRate(command, initialDelay, period, unit);
1 ✔
2171
    }
2172

2173
    @Override
2174
    public ScheduledFuture<?> scheduleWithFixedDelay(
2175
        Runnable command, long initialDelay, long delay, TimeUnit unit) {
2176
      return delegate.scheduleWithFixedDelay(command, initialDelay, delay, unit);
×
2177
    }
2178

2179
    @Override
2180
    public boolean awaitTermination(long timeout, TimeUnit unit)
2181
        throws InterruptedException {
2182
      return delegate.awaitTermination(timeout, unit);
×
2183
    }
2184

2185
    @Override
2186
    public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks)
2187
        throws InterruptedException {
2188
      return delegate.invokeAll(tasks);
×
2189
    }
2190

2191
    @Override
2192
    public <T> List<Future<T>> invokeAll(
2193
        Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit)
2194
        throws InterruptedException {
2195
      return delegate.invokeAll(tasks, timeout, unit);
×
2196
    }
2197

2198
    @Override
2199
    public <T> T invokeAny(Collection<? extends Callable<T>> tasks)
2200
        throws InterruptedException, ExecutionException {
2201
      return delegate.invokeAny(tasks);
×
2202
    }
2203

2204
    @Override
2205
    public <T> T invokeAny(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit)
2206
        throws InterruptedException, ExecutionException, TimeoutException {
2207
      return delegate.invokeAny(tasks, timeout, unit);
×
2208
    }
2209

2210
    @Override
2211
    public boolean isShutdown() {
2212
      return delegate.isShutdown();
×
2213
    }
2214

2215
    @Override
2216
    public boolean isTerminated() {
2217
      return delegate.isTerminated();
×
2218
    }
2219

2220
    @Override
2221
    public void shutdown() {
2222
      throw new UnsupportedOperationException("Restricted: shutdown() is not allowed");
1 ✔
2223
    }
2224

2225
    @Override
2226
    public List<Runnable> shutdownNow() {
2227
      throw new UnsupportedOperationException("Restricted: shutdownNow() is not allowed");
1 ✔
2228
    }
2229

2230
    @Override
2231
    public <T> Future<T> submit(Callable<T> task) {
2232
      return delegate.submit(task);
×
2233
    }
2234

2235
    @Override
2236
    public Future<?> submit(Runnable task) {
2237
      return delegate.submit(task);
×
2238
    }
2239

2240
    @Override
2241
    public <T> Future<T> submit(Runnable task, T result) {
2242
      return delegate.submit(task, result);
×
2243
    }
2244

2245
    @Override
2246
    public void execute(Runnable command) {
2247
      delegate.execute(command);
×
2248
    }
×
2249
  }
2250

2251
  /**
2252
   * A ResolutionState indicates the status of last name resolution.
2253
   */
2254
  enum ResolutionState {
1 ✔
2255
    NO_RESOLUTION,
1 ✔
2256
    SUCCESS,
1 ✔
2257
    ERROR
1 ✔
2258
  }
2259
}
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