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

grpc / grpc-java / #20427

01 Sep 2026 08:34AM UTC coverage: 89.141% (-0.04%) from 89.185%
#20427

push

github

web-flow
Revert "Implement gRFC A97: xDS JWT Call Credentials (#12951)" (v1.84.x backport) (#13019)

Backport of #13017 to v1.84.x.
---
This reverts commit 1053f4be9. We need to
address the concerns raised in #13006.

38335 of 43005 relevant lines covered (89.14%)

0.89 hits per line

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

93.28
/../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
    channelLogger.log(ChannelLogLevel.INFO, "Entering IDLE state");
1✔
437
    channelStateManager.gotoState(IDLE);
1✔
438
    // If the inUseStateAggregator still considers pending calls to be queued up or the delayed
439
    // transport to be holding some we need to exit idle mode to give these calls a chance to
440
    // be processed.
441
    if (inUseStateAggregator.anyObjectInUse(pendingCallsInUseObject, delayedTransport)) {
1✔
442
      exitIdleMode();
1✔
443
    }
444
  }
1✔
445

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

972
      syncContext.execute(new RealChannelShutdown());
1✔
973
    }
1✔
974

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

991
      syncContext.execute(new RealChannelShutdownNow());
1✔
992
    }
1✔
993

994
    @Override
995
    public String authority() {
996
      return authority;
1✔
997
    }
998

999
    private final class PendingCall<ReqT, RespT> extends DelayedClientCall<ReqT, RespT> {
1000
      final Context context;
1001
      final MethodDescriptor<ReqT, RespT> method;
1002
      final CallOptions callOptions;
1003
      private final long callCreationTime;
1004

1005
      PendingCall(Context context, MethodDescriptor<ReqT, RespT> method, CallOptions callOptions) {
1✔
1006
        super(
1✔
1007
            "name_resolver",
1008
            getCallExecutor(callOptions),
1✔
1009
            scheduledExecutor,
1✔
1010
            callOptions.getDeadline());
1✔
1011
        this.context = context;
1✔
1012
        this.method = method;
1✔
1013
        this.callOptions = callOptions;
1✔
1014
        this.callCreationTime = ticker.nanoTime();
1✔
1015
      }
1✔
1016

1017
      /** Called when it's ready to create a real call and reprocess the pending call. */
1018
      void reprocess() {
1019
        ClientCall<ReqT, RespT> realCall;
1020
        Context previous = context.attach();
1✔
1021
        try {
1022
          CallOptions delayResolutionOption = callOptions.withOption(NAME_RESOLUTION_DELAYED,
1✔
1023
              ticker.nanoTime() - callCreationTime);
1✔
1024
          realCall = newClientCall(method, delayResolutionOption);
1✔
1025
        } finally {
1026
          context.detach(previous);
1✔
1027
        }
1028
        Runnable toRun = setCall(realCall);
1✔
1029
        if (toRun == null) {
1✔
1030
          syncContext.execute(new PendingCallRemoval());
1✔
1031
        } else {
1032
          getCallExecutor(callOptions).execute(new Runnable() {
1✔
1033
            @Override
1034
            public void run() {
1035
              toRun.run();
1✔
1036
              syncContext.execute(new PendingCallRemoval());
1✔
1037
            }
1✔
1038
          });
1039
        }
1040
      }
1✔
1041

1042
      @Override
1043
      protected void callCancelled() {
1044
        super.callCancelled();
1✔
1045
        syncContext.execute(new PendingCallRemoval());
1✔
1046
      }
1✔
1047

1048
      final class PendingCallRemoval implements Runnable {
1✔
1049
        @Override
1050
        public void run() {
1051
          if (pendingCalls != null) {
1✔
1052
            pendingCalls.remove(PendingCall.this);
1✔
1053
            if (pendingCalls.isEmpty()) {
1✔
1054
              inUseStateAggregator.updateObjectInUse(pendingCallsInUseObject, false);
1✔
1055
              pendingCalls = null;
1✔
1056
              if (shutdown.get()) {
1✔
1057
                uncommittedRetriableStreamsRegistry.onShutdown(SHUTDOWN_STATUS);
1✔
1058
              }
1059
            }
1060
          }
1061
        }
1✔
1062
      }
1063
    }
1064

1065
    private <ReqT, RespT> ClientCall<ReqT, RespT> newClientCall(
1066
        MethodDescriptor<ReqT, RespT> method, CallOptions callOptions) {
1067
      InternalConfigSelector selector = configSelector.get();
1✔
1068
      if (selector == null) {
1✔
1069
        return clientCallImplChannel.newCall(method, callOptions);
1✔
1070
      }
1071
      if (selector instanceof ServiceConfigConvertedSelector) {
1✔
1072
        MethodInfo methodInfo =
1✔
1073
            ((ServiceConfigConvertedSelector) selector).config.getMethodConfig(method);
1✔
1074
        if (methodInfo != null) {
1✔
1075
          callOptions = callOptions.withOption(MethodInfo.KEY, methodInfo);
1✔
1076
        }
1077
        return clientCallImplChannel.newCall(method, callOptions);
1✔
1078
      }
1079
      return new ConfigSelectingClientCall<>(
1✔
1080
          selector, clientCallImplChannel, executor, method, callOptions);
1✔
1081
    }
1082
  }
1083

1084
  /**
1085
   * A client call for a given channel that applies a given config selector when it starts.
1086
   */
1087
  static final class ConfigSelectingClientCall<ReqT, RespT>
1088
      extends ForwardingClientCall<ReqT, RespT> {
1089

1090
    private final InternalConfigSelector configSelector;
1091
    private final Channel channel;
1092
    private final Executor callExecutor;
1093
    private final MethodDescriptor<ReqT, RespT> method;
1094
    private final Context context;
1095
    private CallOptions callOptions;
1096

1097
    private ClientCall<ReqT, RespT> delegate;
1098

1099
    ConfigSelectingClientCall(
1100
        InternalConfigSelector configSelector, Channel channel, Executor channelExecutor,
1101
        MethodDescriptor<ReqT, RespT> method,
1102
        CallOptions callOptions) {
1✔
1103
      this.configSelector = configSelector;
1✔
1104
      this.channel = channel;
1✔
1105
      this.method = method;
1✔
1106
      this.callExecutor =
1✔
1107
          callOptions.getExecutor() == null ? channelExecutor : callOptions.getExecutor();
1✔
1108
      this.callOptions = callOptions.withExecutor(callExecutor);
1✔
1109
      this.context = Context.current();
1✔
1110
    }
1✔
1111

1112
    @Override
1113
    protected ClientCall<ReqT, RespT> delegate() {
1114
      return delegate;
1✔
1115
    }
1116

1117
    @SuppressWarnings("unchecked")
1118
    @Override
1119
    public void start(Listener<RespT> observer, Metadata headers) {
1120
      PickSubchannelArgs args =
1✔
1121
          new PickSubchannelArgsImpl(method, headers, callOptions, NOOP_PICK_DETAILS_CONSUMER);
1✔
1122
      InternalConfigSelector.Result result = configSelector.selectConfig(args);
1✔
1123
      Status status = result.getStatus();
1✔
1124
      if (!status.isOk()) {
1✔
1125
        executeCloseObserverInContext(observer,
1✔
1126
            GrpcUtil.replaceInappropriateControlPlaneStatus(status));
1✔
1127
        delegate = (ClientCall<ReqT, RespT>) NOOP_CALL;
1✔
1128
        return;
1✔
1129
      }
1130
      ClientInterceptor interceptor = result.getInterceptor();
1✔
1131
      ManagedChannelServiceConfig config = (ManagedChannelServiceConfig) result.getConfig();
1✔
1132
      MethodInfo methodInfo = config.getMethodConfig(method);
1✔
1133
      if (methodInfo != null) {
1✔
1134
        callOptions = callOptions.withOption(MethodInfo.KEY, methodInfo);
1✔
1135
      }
1136
      if (interceptor != null) {
1✔
1137
        delegate = interceptor.interceptCall(method, callOptions, channel);
1✔
1138
      } else {
1139
        delegate = channel.newCall(method, callOptions);
×
1140
      }
1141
      delegate.start(observer, headers);
1✔
1142
    }
1✔
1143

1144
    private void executeCloseObserverInContext(
1145
        final Listener<RespT> observer, final Status status) {
1146
      class CloseInContext extends ContextRunnable {
1147
        CloseInContext() {
1✔
1148
          super(context);
1✔
1149
        }
1✔
1150

1151
        @Override
1152
        public void runInContext() {
1153
          observer.onClose(status, new Metadata());
1✔
1154
        }
1✔
1155
      }
1156

1157
      callExecutor.execute(new CloseInContext());
1✔
1158
    }
1✔
1159

1160
    @Override
1161
    public void cancel(@Nullable String message, @Nullable Throwable cause) {
1162
      if (delegate != null) {
×
1163
        delegate.cancel(message, cause);
×
1164
      }
1165
    }
×
1166
  }
1167

1168
  private static final ClientCall<Object, Object> NOOP_CALL = new ClientCall<Object, Object>() {
1✔
1169
    @Override
1170
    public void start(Listener<Object> responseListener, Metadata headers) {}
×
1171

1172
    @Override
1173
    public void request(int numMessages) {}
1✔
1174

1175
    @Override
1176
    public void cancel(String message, Throwable cause) {}
×
1177

1178
    @Override
1179
    public void halfClose() {}
×
1180

1181
    @Override
1182
    public void sendMessage(Object message) {}
×
1183

1184
    // Always returns {@code false}, since this is only used when the startup of the call fails.
1185
    @Override
1186
    public boolean isReady() {
1187
      return false;
×
1188
    }
1189
  };
1190

1191
  /**
1192
   * Terminate the channel if termination conditions are met.
1193
   */
1194
  // Must be run from syncContext
1195
  private void maybeTerminateChannel() {
1196
    if (terminated) {
1✔
1197
      return;
×
1198
    }
1199
    if (shutdown.get() && subchannels.isEmpty()) {
1✔
1200
      channelLogger.log(ChannelLogLevel.INFO, "Terminated");
1✔
1201
      channelz.removeRootChannel(this);
1✔
1202
      executorPool.returnObject(executor);
1✔
1203
      balancerRpcExecutorHolder.release();
1✔
1204
      offloadExecutorHolder.release();
1✔
1205
      // Release the transport factory so that it can deallocate any resources.
1206
      transportFactory.close();
1✔
1207

1208
      terminated = true;
1✔
1209
      terminatedLatch.countDown();
1✔
1210
    }
1211
  }
1✔
1212

1213
  @Override
1214
  public ConnectivityState getState(boolean requestConnection) {
1215
    ConnectivityState savedChannelState = channelStateManager.getState();
1✔
1216
    if (requestConnection && savedChannelState == IDLE) {
1✔
1217
      final class RequestConnection implements Runnable {
1✔
1218
        @Override
1219
        public void run() {
1220
          exitIdleMode();
1✔
1221
          if (lbHelper != null) {
1✔
1222
            lbHelper.lb.requestConnection();
1✔
1223
          }
1224
        }
1✔
1225
      }
1226

1227
      syncContext.execute(new RequestConnection());
1✔
1228
    }
1229
    return savedChannelState;
1✔
1230
  }
1231

1232
  @Override
1233
  public void notifyWhenStateChanged(final ConnectivityState source, final Runnable callback) {
1234
    final class NotifyStateChanged implements Runnable {
1✔
1235
      @Override
1236
      public void run() {
1237
        channelStateManager.notifyWhenStateChanged(callback, executor, source);
1✔
1238
      }
1✔
1239
    }
1240

1241
    syncContext.execute(new NotifyStateChanged());
1✔
1242
  }
1✔
1243

1244
  @Override
1245
  public void resetConnectBackoff() {
1246
    final class ResetConnectBackoff implements Runnable {
1✔
1247
      @Override
1248
      public void run() {
1249
        if (shutdown.get()) {
1✔
1250
          return;
1✔
1251
        }
1252
        if (nameResolverStarted) {
1✔
1253
          refreshNameResolution();
1✔
1254
        }
1255
        for (InternalSubchannel subchannel : subchannels) {
1✔
1256
          subchannel.resetConnectBackoff();
1✔
1257
        }
1✔
1258
      }
1✔
1259
    }
1260

1261
    syncContext.execute(new ResetConnectBackoff());
1✔
1262
  }
1✔
1263

1264
  @Override
1265
  public void enterIdle() {
1266
    final class PrepareToLoseNetworkRunnable implements Runnable {
1✔
1267
      @Override
1268
      public void run() {
1269
        if (shutdown.get() || lbHelper == null) {
1✔
1270
          return;
1✔
1271
        }
1272
        cancelIdleTimer(/* permanent= */ false);
1✔
1273
        enterIdleMode();
1✔
1274
      }
1✔
1275
    }
1276

1277
    syncContext.execute(new PrepareToLoseNetworkRunnable());
1✔
1278
  }
1✔
1279

1280
  /**
1281
   * A registry that prevents channel shutdown from killing existing retry attempts that are in
1282
   * backoff.
1283
   */
1284
  private final class UncommittedRetriableStreamsRegistry {
1✔
1285
    // TODO(zdapeng): This means we would acquire a lock for each new retry-able stream,
1286
    // it's worthwhile to look for a lock-free approach.
1287
    final Object lock = new Object();
1✔
1288

1289
    @GuardedBy("lock")
1✔
1290
    Collection<ClientStream> uncommittedRetriableStreams = new HashSet<>();
1291

1292
    @GuardedBy("lock")
1293
    Status shutdownStatus;
1294

1295
    void onShutdown(Status reason) {
1296
      boolean shouldShutdownDelayedTransport = false;
1✔
1297
      synchronized (lock) {
1✔
1298
        if (shutdownStatus != null) {
1✔
1299
          return;
1✔
1300
        }
1301
        shutdownStatus = reason;
1✔
1302
        // Keep the delayedTransport open until there is no more uncommitted streams, b/c those
1303
        // retriable streams, which may be in backoff and not using any transport, are already
1304
        // started RPCs.
1305
        if (uncommittedRetriableStreams.isEmpty()) {
1✔
1306
          shouldShutdownDelayedTransport = true;
1✔
1307
        }
1308
      }
1✔
1309

1310
      if (shouldShutdownDelayedTransport) {
1✔
1311
        delayedTransport.shutdown(reason);
1✔
1312
      }
1313
    }
1✔
1314

1315
    void onShutdownNow(Status reason) {
1316
      onShutdown(reason);
1✔
1317
      Collection<ClientStream> streams;
1318

1319
      synchronized (lock) {
1✔
1320
        streams = new ArrayList<>(uncommittedRetriableStreams);
1✔
1321
      }
1✔
1322

1323
      for (ClientStream stream : streams) {
1✔
1324
        stream.cancel(reason);
1✔
1325
      }
1✔
1326
      delayedTransport.shutdownNow(reason);
1✔
1327
    }
1✔
1328

1329
    /**
1330
     * Registers a RetriableStream and return null if not shutdown, otherwise just returns the
1331
     * shutdown Status.
1332
     */
1333
    @Nullable
1334
    Status add(RetriableStream<?> retriableStream) {
1335
      synchronized (lock) {
1✔
1336
        if (shutdownStatus != null) {
1✔
1337
          return shutdownStatus;
1✔
1338
        }
1339
        uncommittedRetriableStreams.add(retriableStream);
1✔
1340
        return null;
1✔
1341
      }
1342
    }
1343

1344
    void remove(RetriableStream<?> retriableStream) {
1345
      Status shutdownStatusCopy = null;
1✔
1346

1347
      synchronized (lock) {
1✔
1348
        uncommittedRetriableStreams.remove(retriableStream);
1✔
1349
        if (uncommittedRetriableStreams.isEmpty()) {
1✔
1350
          shutdownStatusCopy = shutdownStatus;
1✔
1351
          // Because retriable transport is long-lived, we take this opportunity to down-size the
1352
          // hashmap.
1353
          uncommittedRetriableStreams = new HashSet<>();
1✔
1354
        }
1355
      }
1✔
1356

1357
      if (shutdownStatusCopy != null) {
1✔
1358
        delayedTransport.shutdown(shutdownStatusCopy);
1✔
1359
      }
1360
    }
1✔
1361
  }
1362

1363
  private final class LbHelperImpl extends LoadBalancer.Helper {
1✔
1364
    LoadBalancer lb;
1365

1366
    @Override
1367
    public AbstractSubchannel createSubchannel(CreateSubchannelArgs args) {
1368
      syncContext.throwIfNotInThisSynchronizationContext();
1✔
1369
      // No new subchannel should be created after load balancer has been shutdown.
1370
      checkState(!terminating, "Channel is being terminated");
1✔
1371
      return new SubchannelImpl(args);
1✔
1372
    }
1373

1374
    @Override
1375
    public void updateBalancingState(
1376
        final ConnectivityState newState, final SubchannelPicker newPicker) {
1377
      syncContext.throwIfNotInThisSynchronizationContext();
1✔
1378
      checkNotNull(newState, "newState");
1✔
1379
      checkNotNull(newPicker, "newPicker");
1✔
1380

1381
      if (LbHelperImpl.this != lbHelper || panicMode) {
1✔
1382
        return;
1✔
1383
      }
1384
      updateSubchannelPicker(newPicker);
1✔
1385
      // It's not appropriate to report SHUTDOWN state from lb.
1386
      // Ignore the case of newState == SHUTDOWN for now.
1387
      if (newState != SHUTDOWN) {
1✔
1388
        channelLogger.log(
1✔
1389
            ChannelLogLevel.INFO, "Entering {0} state with picker: {1}", newState, newPicker);
1390
        channelStateManager.gotoState(newState);
1✔
1391
      }
1392
    }
1✔
1393

1394
    @Override
1395
    public void refreshNameResolution() {
1396
      syncContext.throwIfNotInThisSynchronizationContext();
1✔
1397
      final class LoadBalancerRefreshNameResolution implements Runnable {
1✔
1398
        @Override
1399
        public void run() {
1400
          ManagedChannelImpl.this.refreshNameResolution();
1✔
1401
        }
1✔
1402
      }
1403

1404
      syncContext.execute(new LoadBalancerRefreshNameResolution());
1✔
1405
    }
1✔
1406

1407
    @Override
1408
    public ManagedChannel createOobChannel(EquivalentAddressGroup addressGroup, String authority) {
1409
      return createOobChannel(Collections.singletonList(addressGroup), authority);
×
1410
    }
1411

1412
    @Override
1413
    public ManagedChannel createOobChannel(List<EquivalentAddressGroup> addressGroup,
1414
        String authority) {
1415
      NameResolverRegistry nameResolverRegistry = new NameResolverRegistry();
1✔
1416
      OobNameResolverProvider resolverProvider =
1✔
1417
          new OobNameResolverProvider(authority, addressGroup, syncContext);
1418
      nameResolverRegistry.register(resolverProvider);
1✔
1419
      // We could use a hard-coded target, as the name resolver won't actually use this string.
1420
      // However, that would make debugging less clear, as we use the target to identify the
1421
      // channel.
1422
      String target;
1423
      try {
1424
        target = new URI("oob", "", "/" + authority, null, null).toString();
1✔
1425
      } catch (URISyntaxException ex) {
×
1426
        // Any special characters in the path will be percent encoded. So this should be impossible.
1427
        throw new AssertionError(ex);
×
1428
      }
1✔
1429
      ManagedChannel delegate = createResolvingOobChannelBuilder(
1✔
1430
          target, new DefaultChannelCreds(), nameResolverRegistry)
1431
          // TODO(zdapeng): executors should not outlive the parent channel.
1432
          .executor(balancerRpcExecutorHolder.getExecutor())
1✔
1433
          .idleTimeout(Integer.MAX_VALUE, TimeUnit.SECONDS)
1✔
1434
          .disableRetry()
1✔
1435
          .build();
1✔
1436
      return new OobChannel(delegate, resolverProvider);
1✔
1437
    }
1438

1439
    @Deprecated
1440
    @Override
1441
    public ManagedChannelBuilder<?> createResolvingOobChannelBuilder(String target) {
1442
      return createResolvingOobChannelBuilder(target, new DefaultChannelCreds())
1✔
1443
          // Override authority to keep the old behavior.
1444
          // createResolvingOobChannelBuilder(String target) will be deleted soon.
1445
          .overrideAuthority(getAuthority());
1✔
1446
    }
1447

1448
    @Override
1449
    public ManagedChannelBuilder<?> createResolvingOobChannelBuilder(
1450
        final String target, final ChannelCredentials channelCreds) {
1451
      return createResolvingOobChannelBuilder(target, channelCreds, nameResolverRegistry);
1✔
1452
    }
1453

1454
    // TODO(creamsoup) prevent main channel to shutdown if oob channel is not terminated
1455
    // TODO(zdapeng) register the channel as a subchannel of the parent channel in channelz.
1456
    private ManagedChannelBuilder<?> createResolvingOobChannelBuilder(
1457
        final String target, final ChannelCredentials channelCreds,
1458
        NameResolverRegistry nameResolverRegistry) {
1459
      checkNotNull(channelCreds, "channelCreds");
1✔
1460

1461
      final class ResolvingOobChannelBuilder
1462
          extends ForwardingChannelBuilder2<ResolvingOobChannelBuilder> {
1463
        final ManagedChannelBuilder<?> delegate;
1464

1465
        ResolvingOobChannelBuilder() {
1✔
1466
          final ClientTransportFactory transportFactory;
1467
          CallCredentials callCredentials;
1468
          if (channelCreds instanceof DefaultChannelCreds) {
1✔
1469
            // TODO(kannanjgithub) We should eventually refactor ManagedChannelImplBuilder so
1470
            // callCredentials can be resolved lazily at build() time, allowing transport factory
1471
            // retention to happen strictly inside buildClientTransportFactory().
1472
            transportFactory = originalTransportFactory.retain();
1✔
1473
            callCredentials = null;
1✔
1474
          } else {
1475
            SwapChannelCredentialsResult swapResult =
1✔
1476
                originalTransportFactory.swapChannelCredentials(channelCreds);
1✔
1477
            if (swapResult == null) {
1✔
1478
              delegate = Grpc.newChannelBuilder(target, channelCreds);
×
1479
              return;
×
1480
            } else {
1481
              transportFactory = swapResult.transportFactory;
1✔
1482
              callCredentials = swapResult.callCredentials;
1✔
1483
            }
1484
          }
1485
          ClientTransportFactoryBuilder transportFactoryBuilder =
1✔
1486
              new ClientTransportFactoryBuilder() {
1✔
1487
                @Override
1488
                public ClientTransportFactory buildClientTransportFactory() {
1489
                  return transportFactory;
1✔
1490
                }
1491
              };
1492
          delegate = new ManagedChannelImplBuilder(
1✔
1493
              target,
1494
              channelCreds,
1495
              callCredentials,
1496
              transportFactoryBuilder,
1497
              new FixedPortProvider(nameResolverArgs.getDefaultPort()))
1✔
1498
              .nameResolverRegistry(nameResolverRegistry);
1✔
1499
        }
1✔
1500

1501
        @Override
1502
        protected ManagedChannelBuilder<?> delegate() {
1503
          return delegate;
1✔
1504
        }
1505
      }
1506

1507
      checkState(!terminated, "Channel is terminated");
1✔
1508

1509
      ResolvingOobChannelBuilder builder = new ResolvingOobChannelBuilder();
1✔
1510

1511
      // Note that we follow the global configurator pattern and try to fuse the configurations as
1512
      // soon as the builder gets created
1513
      channelConfigurator.configureChannelBuilder(builder);
1✔
1514
      builder.childChannelConfigurator(channelConfigurator);
1✔
1515

1516
      return builder
1✔
1517
          // TODO(zdapeng): executors should not outlive the parent channel.
1518
          .executor(executor)
1✔
1519
          .offloadExecutor(offloadExecutorHolder.getExecutor())
1✔
1520
          .maxTraceEvents(maxTraceEvents)
1✔
1521
          .proxyDetector(nameResolverArgs.getProxyDetector())
1✔
1522
          .userAgent(userAgent);
1✔
1523
    }
1524

1525
    @Override
1526
    public ChannelCredentials getUnsafeChannelCredentials() {
1527
      if (originalChannelCreds == null) {
1✔
1528
        return new DefaultChannelCreds();
1✔
1529
      }
1530
      return originalChannelCreds;
×
1531
    }
1532

1533
    @Override
1534
    public void updateOobChannelAddresses(ManagedChannel channel, EquivalentAddressGroup eag) {
1535
      updateOobChannelAddresses(channel, Collections.singletonList(eag));
×
1536
    }
×
1537

1538
    @Override
1539
    public void updateOobChannelAddresses(ManagedChannel channel,
1540
        List<EquivalentAddressGroup> eag) {
1541
      checkArgument(channel instanceof OobChannel,
1✔
1542
          "channel must have been returned from createOobChannel");
1543
      ((OobChannel) channel).updateAddresses(eag);
1✔
1544
    }
1✔
1545

1546
    @Override
1547
    public String getAuthority() {
1548
      return ManagedChannelImpl.this.authority();
1✔
1549
    }
1550

1551
    @Override
1552
    public String getChannelTarget() {
1553
      return targetUri.toString();
1✔
1554
    }
1555

1556
    @Override
1557
    public SynchronizationContext getSynchronizationContext() {
1558
      return syncContext;
1✔
1559
    }
1560

1561
    @Override
1562
    public ScheduledExecutorService getScheduledExecutorService() {
1563
      return scheduledExecutor;
1✔
1564
    }
1565

1566
    @Override
1567
    public ChannelLogger getChannelLogger() {
1568
      return channelLogger;
1✔
1569
    }
1570

1571
    @Override
1572
    public NameResolver.Args getNameResolverArgs() {
1573
      return nameResolverArgs;
1✔
1574
    }
1575

1576
    @Override
1577
    public NameResolverRegistry getNameResolverRegistry() {
1578
      return nameResolverRegistry;
1✔
1579
    }
1580

1581
    @Override
1582
    public MetricRecorder getMetricRecorder() {
1583
      return metricRecorder;
1✔
1584
    }
1585

1586
    /**
1587
     * A placeholder for channel creds if user did not specify channel creds for the channel.
1588
     */
1589
    // TODO(zdapeng): get rid of this class and let all ChannelBuilders always provide a non-null
1590
    //     channel creds.
1591
    final class DefaultChannelCreds extends ChannelCredentials {
1✔
1592
      @Override
1593
      public ChannelCredentials withoutBearerTokens() {
1594
        return this;
×
1595
      }
1596
    }
1597
  }
1598

1599
  static final class OobChannel extends ForwardingManagedChannel {
1600
    private final OobNameResolverProvider resolverProvider;
1601

1602
    public OobChannel(ManagedChannel delegate, OobNameResolverProvider resolverProvider) {
1603
      super(delegate);
1✔
1604
      this.resolverProvider = checkNotNull(resolverProvider, "resolverProvider");
1✔
1605
    }
1✔
1606

1607
    public void updateAddresses(List<EquivalentAddressGroup> eags) {
1608
      resolverProvider.updateAddresses(eags);
1✔
1609
    }
1✔
1610
  }
1611

1612
  final class NameResolverListener extends NameResolver.Listener2 {
1613
    final LbHelperImpl helper;
1614
    final NameResolver resolver;
1615

1616
    NameResolverListener(LbHelperImpl helperImpl, NameResolver resolver) {
1✔
1617
      this.helper = checkNotNull(helperImpl, "helperImpl");
1✔
1618
      this.resolver = checkNotNull(resolver, "resolver");
1✔
1619
    }
1✔
1620

1621
    @Override
1622
    public void onResult(final ResolutionResult resolutionResult) {
1623
      syncContext.execute(() -> onResult2(resolutionResult));
×
1624
    }
×
1625

1626
    @SuppressWarnings("ReferenceEquality")
1627
    @Override
1628
    public Status onResult2(final ResolutionResult resolutionResult) {
1629
      syncContext.throwIfNotInThisSynchronizationContext();
1✔
1630
      if (ManagedChannelImpl.this.nameResolver != resolver) {
1✔
1631
        return Status.OK;
1✔
1632
      }
1633

1634
      StatusOr<List<EquivalentAddressGroup>> serversOrError =
1✔
1635
          resolutionResult.getAddressesOrError();
1✔
1636
      if (!serversOrError.hasValue()) {
1✔
1637
        handleErrorInSyncContext(serversOrError.getStatus());
1✔
1638
        return serversOrError.getStatus();
1✔
1639
      }
1640
      List<EquivalentAddressGroup> servers = serversOrError.getValue();
1✔
1641
      channelLogger.log(
1✔
1642
          ChannelLogLevel.DEBUG,
1643
          "Resolved address: {0}, config={1}",
1644
          servers,
1645
          resolutionResult.getAttributes());
1✔
1646

1647
      if (lastResolutionState != ResolutionState.SUCCESS) {
1✔
1648
        channelLogger.log(ChannelLogLevel.INFO, "Address resolved: {0}",
1✔
1649
            servers);
1650
        lastResolutionState = ResolutionState.SUCCESS;
1✔
1651
      }
1652
      ConfigOrError configOrError = resolutionResult.getServiceConfig();
1✔
1653
      InternalConfigSelector resolvedConfigSelector =
1✔
1654
          resolutionResult.getAttributes().get(InternalConfigSelector.KEY);
1✔
1655
      ManagedChannelServiceConfig validServiceConfig =
1656
          configOrError != null && configOrError.getConfig() != null
1✔
1657
              ? (ManagedChannelServiceConfig) configOrError.getConfig()
1✔
1658
              : null;
1✔
1659
      Status serviceConfigError = configOrError != null ? configOrError.getError() : null;
1✔
1660

1661
      ManagedChannelServiceConfig effectiveServiceConfig;
1662
      if (!lookUpServiceConfig) {
1✔
1663
        if (validServiceConfig != null) {
1✔
1664
          channelLogger.log(
1✔
1665
              ChannelLogLevel.INFO,
1666
              "Service config from name resolver discarded by channel settings");
1667
        }
1668
        effectiveServiceConfig =
1669
            defaultServiceConfig == null ? EMPTY_SERVICE_CONFIG : defaultServiceConfig;
1✔
1670
        if (resolvedConfigSelector != null) {
1✔
1671
          channelLogger.log(
1✔
1672
              ChannelLogLevel.INFO,
1673
              "Config selector from name resolver discarded by channel settings");
1674
        }
1675
        realChannel.updateConfigSelector(effectiveServiceConfig.getDefaultConfigSelector());
1✔
1676
      } else {
1677
        // Try to use config if returned from name resolver
1678
        // Otherwise, try to use the default config if available
1679
        if (validServiceConfig != null) {
1✔
1680
          effectiveServiceConfig = validServiceConfig;
1✔
1681
          if (resolvedConfigSelector != null) {
1✔
1682
            realChannel.updateConfigSelector(resolvedConfigSelector);
1✔
1683
            if (effectiveServiceConfig.getDefaultConfigSelector() != null) {
1✔
1684
              channelLogger.log(
×
1685
                  ChannelLogLevel.DEBUG,
1686
                  "Method configs in service config will be discarded due to presence of"
1687
                      + "config-selector");
1688
            }
1689
          } else {
1690
            realChannel.updateConfigSelector(effectiveServiceConfig.getDefaultConfigSelector());
1✔
1691
          }
1692
        } else if (defaultServiceConfig != null) {
1✔
1693
          effectiveServiceConfig = defaultServiceConfig;
1✔
1694
          realChannel.updateConfigSelector(effectiveServiceConfig.getDefaultConfigSelector());
1✔
1695
          channelLogger.log(
1✔
1696
              ChannelLogLevel.INFO,
1697
              "Received no service config, using default service config");
1698
        } else if (serviceConfigError != null) {
1✔
1699
          if (!serviceConfigUpdated) {
1✔
1700
            // First DNS lookup has invalid service config, and cannot fall back to default
1701
            channelLogger.log(
1✔
1702
                ChannelLogLevel.INFO,
1703
                "Fallback to error due to invalid first service config without default config");
1704
            // This error could be an "inappropriate" control plane error that should not bleed
1705
            // through to client code using gRPC. We let them flow through here to the LB as
1706
            // we later check for these error codes when investigating pick results in
1707
            // GrpcUtil.getTransportFromPickResult().
1708
            onError(configOrError.getError());
1✔
1709
            return configOrError.getError();
1✔
1710
          } else {
1711
            effectiveServiceConfig = lastServiceConfig;
1✔
1712
          }
1713
        } else {
1714
          effectiveServiceConfig = EMPTY_SERVICE_CONFIG;
1✔
1715
          realChannel.updateConfigSelector(null);
1✔
1716
        }
1717
        if (!effectiveServiceConfig.equals(lastServiceConfig)) {
1✔
1718
          channelLogger.log(
1✔
1719
              ChannelLogLevel.INFO,
1720
              "Service config changed{0}",
1721
              effectiveServiceConfig == EMPTY_SERVICE_CONFIG ? " to empty" : "");
1✔
1722
          lastServiceConfig = effectiveServiceConfig;
1✔
1723
          transportProvider.throttle = effectiveServiceConfig.getRetryThrottling();
1✔
1724
        }
1725

1726
        try {
1727
          // TODO(creamsoup): when `serversOrError` is empty and lastResolutionStateCopy == SUCCESS
1728
          //  and lbNeedAddress, it shouldn't call the handleServiceConfigUpdate. But,
1729
          //  lbNeedAddress is not deterministic
1730
          serviceConfigUpdated = true;
1✔
1731
        } catch (RuntimeException re) {
×
1732
          logger.log(
×
1733
              Level.WARNING,
1734
              "[" + getLogId() + "] Unexpected exception from parsing service config",
×
1735
              re);
1736
        }
1✔
1737
      }
1738

1739
      Attributes effectiveAttrs = resolutionResult.getAttributes();
1✔
1740
      // Call LB only if it's not shutdown.  If LB is shutdown, lbHelper won't match.
1741
      if (NameResolverListener.this.helper == ManagedChannelImpl.this.lbHelper) {
1✔
1742
        Attributes.Builder attrBuilder =
1✔
1743
            effectiveAttrs.toBuilder().discard(InternalConfigSelector.KEY);
1✔
1744
        Map<String, ?> healthCheckingConfig =
1✔
1745
            effectiveServiceConfig.getHealthCheckingConfig();
1✔
1746
        if (healthCheckingConfig != null) {
1✔
1747
          attrBuilder
1✔
1748
              .set(LoadBalancer.ATTR_HEALTH_CHECKING_CONFIG, healthCheckingConfig)
1✔
1749
              .build();
1✔
1750
        }
1751
        Attributes attributes = attrBuilder.build();
1✔
1752

1753
        ResolvedAddresses.Builder resolvedAddresses = ResolvedAddresses.newBuilder()
1✔
1754
            .setAddresses(serversOrError.getValue())
1✔
1755
            .setAttributes(attributes)
1✔
1756
            .setLoadBalancingPolicyConfig(effectiveServiceConfig.getLoadBalancingConfig());
1✔
1757
        Status addressAcceptanceStatus = helper.lb.acceptResolvedAddresses(
1✔
1758
            resolvedAddresses.build());
1✔
1759
        return addressAcceptanceStatus;
1✔
1760
      }
1761
      return Status.OK;
×
1762
    }
1763

1764
    @Override
1765
    public void onError(final Status error) {
1766
      checkArgument(!error.isOk(), "the error status must not be OK");
1✔
1767
      final class NameResolverErrorHandler implements Runnable {
1✔
1768
        @Override
1769
        public void run() {
1770
          handleErrorInSyncContext(error);
1✔
1771
        }
1✔
1772
      }
1773

1774
      syncContext.execute(new NameResolverErrorHandler());
1✔
1775
    }
1✔
1776

1777
    private void handleErrorInSyncContext(Status error) {
1778
      logger.log(Level.WARNING, "[{0}] Failed to resolve name. status={1}",
1✔
1779
          new Object[] {getLogId(), error});
1✔
1780
      realChannel.onConfigError();
1✔
1781
      if (lastResolutionState != ResolutionState.ERROR) {
1✔
1782
        channelLogger.log(ChannelLogLevel.WARNING, "Failed to resolve name: {0}", error);
1✔
1783
        lastResolutionState = ResolutionState.ERROR;
1✔
1784
      }
1785
      // Call LB only if it's not shutdown.  If LB is shutdown, lbHelper won't match.
1786
      if (NameResolverListener.this.helper != ManagedChannelImpl.this.lbHelper) {
1✔
1787
        return;
1✔
1788
      }
1789

1790
      helper.lb.handleNameResolutionError(error);
1✔
1791
    }
1✔
1792
  }
1793

1794
  private final class SubchannelImpl extends AbstractSubchannel {
1795
    final CreateSubchannelArgs args;
1796
    final InternalLogId subchannelLogId;
1797
    final ChannelLoggerImpl subchannelLogger;
1798
    final ChannelTracer subchannelTracer;
1799
    List<EquivalentAddressGroup> addressGroups;
1800
    InternalSubchannel subchannel;
1801
    boolean started;
1802
    boolean shutdown;
1803
    ScheduledHandle delayedShutdownTask;
1804

1805
    SubchannelImpl(CreateSubchannelArgs args) {
1✔
1806
      checkNotNull(args, "args");
1✔
1807
      addressGroups = args.getAddresses();
1✔
1808
      if (authorityOverride != null) {
1✔
1809
        List<EquivalentAddressGroup> eagsWithoutOverrideAttr =
1✔
1810
            stripOverrideAuthorityAttributes(args.getAddresses());
1✔
1811
        args = args.toBuilder().setAddresses(eagsWithoutOverrideAttr).build();
1✔
1812
      }
1813
      this.args = args;
1✔
1814
      subchannelLogId = InternalLogId.allocate("Subchannel", /*details=*/ authority());
1✔
1815
      subchannelTracer = new ChannelTracer(
1✔
1816
          subchannelLogId, maxTraceEvents, timeProvider.currentTimeNanos(),
1✔
1817
          "Subchannel for " + args.getAddresses());
1✔
1818
      subchannelLogger = new ChannelLoggerImpl(subchannelTracer, timeProvider);
1✔
1819
    }
1✔
1820

1821
    @Override
1822
    public void start(final SubchannelStateListener listener) {
1823
      syncContext.throwIfNotInThisSynchronizationContext();
1✔
1824
      checkState(!started, "already started");
1✔
1825
      checkState(!shutdown, "already shutdown");
1✔
1826
      checkState(!terminating, "Channel is being terminated");
1✔
1827
      started = true;
1✔
1828
      final class ManagedInternalSubchannelCallback extends InternalSubchannel.Callback {
1✔
1829
        // All callbacks are run in syncContext
1830
        @Override
1831
        void onTerminated(InternalSubchannel is) {
1832
          subchannels.remove(is);
1✔
1833
          channelz.removeSubchannel(is);
1✔
1834
          maybeTerminateChannel();
1✔
1835
        }
1✔
1836

1837
        @Override
1838
        void onStateChange(InternalSubchannel is, ConnectivityStateInfo newState) {
1839
          checkState(listener != null, "listener is null");
1✔
1840
          listener.onSubchannelState(newState);
1✔
1841
        }
1✔
1842

1843
        @Override
1844
        void onInUse(InternalSubchannel is) {
1845
          inUseStateAggregator.updateObjectInUse(is, true);
1✔
1846
        }
1✔
1847

1848
        @Override
1849
        void onNotInUse(InternalSubchannel is) {
1850
          inUseStateAggregator.updateObjectInUse(is, false);
1✔
1851
        }
1✔
1852
      }
1853

1854
      final InternalSubchannel internalSubchannel = new InternalSubchannel(
1✔
1855
          args,
1856
          authority(),
1✔
1857
          userAgent,
1✔
1858
          backoffPolicyProvider,
1✔
1859
          transportFactory,
1✔
1860
          transportFactory.getScheduledExecutorService(),
1✔
1861
          stopwatchSupplier,
1✔
1862
          syncContext,
1863
          new ManagedInternalSubchannelCallback(),
1864
          channelz,
1✔
1865
          callTracerFactory.create(),
1✔
1866
          subchannelTracer,
1867
          subchannelLogId,
1868
          subchannelLogger,
1869
          transportFilters, target,
1✔
1870
          lbHelper.getMetricRecorder());
1✔
1871

1872
      channelTracer.reportEvent(new ChannelTrace.Event.Builder()
1✔
1873
          .setDescription("Child Subchannel started")
1✔
1874
          .setSeverity(ChannelTrace.Event.Severity.CT_INFO)
1✔
1875
          .setTimestampNanos(timeProvider.currentTimeNanos())
1✔
1876
          .setSubchannelRef(internalSubchannel)
1✔
1877
          .build());
1✔
1878

1879
      this.subchannel = internalSubchannel;
1✔
1880
      channelz.addSubchannel(internalSubchannel);
1✔
1881
      subchannels.add(internalSubchannel);
1✔
1882
    }
1✔
1883

1884
    @Override
1885
    InternalInstrumented<ChannelStats> getInstrumentedInternalSubchannel() {
1886
      checkState(started, "not started");
1✔
1887
      return subchannel;
1✔
1888
    }
1889

1890
    @Override
1891
    public void shutdown() {
1892
      syncContext.throwIfNotInThisSynchronizationContext();
1✔
1893
      if (subchannel == null) {
1✔
1894
        // start() was not successful
1895
        shutdown = true;
×
1896
        return;
×
1897
      }
1898
      if (shutdown) {
1✔
1899
        if (terminating && delayedShutdownTask != null) {
1✔
1900
          // shutdown() was previously called when terminating == false, thus a delayed shutdown()
1901
          // was scheduled.  Now since terminating == true, We should expedite the shutdown.
1902
          delayedShutdownTask.cancel();
×
1903
          delayedShutdownTask = null;
×
1904
          // Will fall through to the subchannel.shutdown() at the end.
1905
        } else {
1906
          return;
1✔
1907
        }
1908
      } else {
1909
        shutdown = true;
1✔
1910
      }
1911
      // Add a delay to shutdown to deal with the race between 1) a transport being picked and
1912
      // newStream() being called on it, and 2) its Subchannel is shut down by LoadBalancer (e.g.,
1913
      // because of address change, or because LoadBalancer is shutdown by Channel entering idle
1914
      // mode). If (2) wins, the app will see a spurious error. We work around this by delaying
1915
      // shutdown of Subchannel for a few seconds here.
1916
      //
1917
      // TODO(zhangkun83): consider a better approach
1918
      // (https://github.com/grpc/grpc-java/issues/2562).
1919
      if (!terminating) {
1✔
1920
        final class ShutdownSubchannel implements Runnable {
1✔
1921
          @Override
1922
          public void run() {
1923
            subchannel.shutdown(SUBCHANNEL_SHUTDOWN_STATUS);
1✔
1924
          }
1✔
1925
        }
1926

1927
        delayedShutdownTask = syncContext.schedule(
1✔
1928
            new LogExceptionRunnable(new ShutdownSubchannel()),
1929
            SUBCHANNEL_SHUTDOWN_DELAY_SECONDS, TimeUnit.SECONDS,
1930
            transportFactory.getScheduledExecutorService());
1✔
1931
        return;
1✔
1932
      }
1933
      // When terminating == true, no more real streams will be created. It's safe and also
1934
      // desirable to shutdown timely.
1935
      subchannel.shutdown(SHUTDOWN_STATUS);
1✔
1936
    }
1✔
1937

1938
    @Override
1939
    public void requestConnection() {
1940
      syncContext.throwIfNotInThisSynchronizationContext();
1✔
1941
      checkState(started, "not started");
1✔
1942
      if (shutdown) {
1✔
1943
        return;
1✔
1944
      }
1945
      subchannel.obtainActiveTransport();
1✔
1946
    }
1✔
1947

1948
    @Override
1949
    public List<EquivalentAddressGroup> getAllAddresses() {
1950
      syncContext.throwIfNotInThisSynchronizationContext();
1✔
1951
      checkState(started, "not started");
1✔
1952
      return addressGroups;
1✔
1953
    }
1954

1955
    @Override
1956
    public Attributes getAttributes() {
1957
      return args.getAttributes();
1✔
1958
    }
1959

1960
    @Override
1961
    public String toString() {
1962
      return subchannelLogId.toString();
1✔
1963
    }
1964

1965
    @Override
1966
    public Channel asChannel() {
1967
      checkState(started, "not started");
1✔
1968
      return new SubchannelChannel(
1✔
1969
          subchannel, balancerRpcExecutorHolder.getExecutor(),
1✔
1970
          transportFactory.getScheduledExecutorService(),
1✔
1971
          callTracerFactory.create(),
1✔
1972
          new AtomicReference<InternalConfigSelector>(null));
1973
    }
1974

1975
    @Override
1976
    public Object getInternalSubchannel() {
1977
      checkState(started, "Subchannel is not started");
1✔
1978
      return subchannel;
1✔
1979
    }
1980

1981
    @Override
1982
    public ChannelLogger getChannelLogger() {
1983
      return subchannelLogger;
1✔
1984
    }
1985

1986
    @Override
1987
    public void updateAddresses(List<EquivalentAddressGroup> addrs) {
1988
      syncContext.throwIfNotInThisSynchronizationContext();
1✔
1989
      addressGroups = addrs;
1✔
1990
      if (authorityOverride != null) {
1✔
1991
        addrs = stripOverrideAuthorityAttributes(addrs);
1✔
1992
      }
1993
      subchannel.updateAddresses(addrs);
1✔
1994
    }
1✔
1995

1996
    @Override
1997
    public Attributes getConnectedAddressAttributes() {
1998
      return subchannel.getConnectedAddressAttributes();
1✔
1999
    }
2000

2001
    private List<EquivalentAddressGroup> stripOverrideAuthorityAttributes(
2002
        List<EquivalentAddressGroup> eags) {
2003
      List<EquivalentAddressGroup> eagsWithoutOverrideAttr = new ArrayList<>();
1✔
2004
      for (EquivalentAddressGroup eag : eags) {
1✔
2005
        EquivalentAddressGroup eagWithoutOverrideAttr = new EquivalentAddressGroup(
1✔
2006
            eag.getAddresses(),
1✔
2007
            eag.getAttributes().toBuilder().discard(ATTR_AUTHORITY_OVERRIDE).build());
1✔
2008
        eagsWithoutOverrideAttr.add(eagWithoutOverrideAttr);
1✔
2009
      }
1✔
2010
      return Collections.unmodifiableList(eagsWithoutOverrideAttr);
1✔
2011
    }
2012
  }
2013

2014
  @Override
2015
  public String toString() {
2016
    return MoreObjects.toStringHelper(this)
1✔
2017
        .add("logId", logId.getId())
1✔
2018
        .add("target", target)
1✔
2019
        .toString();
1✔
2020
  }
2021

2022
  /**
2023
   * Called from syncContext.
2024
   */
2025
  private final class DelayedTransportListener implements ManagedClientTransport.Listener {
1✔
2026
    @Override
2027
    public void transportShutdown(Status s, DisconnectError e) {
2028
      checkState(shutdown.get(), "Channel must have been shut down");
1✔
2029
    }
1✔
2030

2031
    @Override
2032
    public void transportReady() {
2033
      // Don't care
2034
    }
×
2035

2036
    @Override
2037
    public Attributes filterTransport(Attributes attributes) {
2038
      return attributes;
×
2039
    }
2040

2041
    @Override
2042
    public void transportInUse(final boolean inUse) {
2043
      inUseStateAggregator.updateObjectInUse(delayedTransport, inUse);
1✔
2044
      if (inUse) {
1✔
2045
        // It's possible to be in idle mode while inUseStateAggregator is in-use, if one of the
2046
        // subchannels is in use. But we should never be in idle mode when delayed transport is in
2047
        // use.
2048
        exitIdleMode();
1✔
2049
      }
2050
    }
1✔
2051

2052
    @Override
2053
    public void transportTerminated() {
2054
      checkState(shutdown.get(), "Channel must have been shut down");
1✔
2055
      terminating = true;
1✔
2056
      shutdownNameResolverAndLoadBalancer(false);
1✔
2057
      // No need to call channelStateManager since we are already in SHUTDOWN state.
2058
      // Until LoadBalancer is shutdown, it may still create new subchannels.  We catch them
2059
      // here.
2060
      maybeShutdownNowSubchannels();
1✔
2061
      maybeTerminateChannel();
1✔
2062
    }
1✔
2063
  }
2064

2065
  /**
2066
   * Must be accessed from syncContext.
2067
   */
2068
  private final class IdleModeStateAggregator extends InUseStateAggregator<Object> {
1✔
2069
    @Override
2070
    protected void handleInUse() {
2071
      exitIdleMode();
1✔
2072
    }
1✔
2073

2074
    @Override
2075
    protected void handleNotInUse() {
2076
      if (shutdown.get()) {
1✔
2077
        return;
1✔
2078
      }
2079
      rescheduleIdleTimer();
1✔
2080
    }
1✔
2081
  }
2082

2083
  /**
2084
   * Lazily request for Executor from an executor pool.
2085
   * Also act as an Executor directly to simply run a cmd
2086
   */
2087
  @VisibleForTesting
2088
  static final class ExecutorHolder implements Executor {
2089
    private final ObjectPool<? extends Executor> pool;
2090
    private Executor executor;
2091

2092
    ExecutorHolder(ObjectPool<? extends Executor> executorPool) {
1✔
2093
      this.pool = checkNotNull(executorPool, "executorPool");
1✔
2094
    }
1✔
2095

2096
    synchronized Executor getExecutor() {
2097
      if (executor == null) {
1✔
2098
        executor = checkNotNull(pool.getObject(), "%s.getObject()", executor);
1✔
2099
      }
2100
      return executor;
1✔
2101
    }
2102

2103
    synchronized void release() {
2104
      if (executor != null) {
1✔
2105
        executor = pool.returnObject(executor);
1✔
2106
      }
2107
    }
1✔
2108

2109
    @Override
2110
    public void execute(Runnable command) {
2111
      getExecutor().execute(command);
1✔
2112
    }
1✔
2113
  }
2114

2115
  private static final class RestrictedScheduledExecutor implements ScheduledExecutorService {
2116
    final ScheduledExecutorService delegate;
2117

2118
    private RestrictedScheduledExecutor(ScheduledExecutorService delegate) {
1✔
2119
      this.delegate = checkNotNull(delegate, "delegate");
1✔
2120
    }
1✔
2121

2122
    @Override
2123
    public <V> ScheduledFuture<V> schedule(Callable<V> callable, long delay, TimeUnit unit) {
2124
      return delegate.schedule(callable, delay, unit);
×
2125
    }
2126

2127
    @Override
2128
    public ScheduledFuture<?> schedule(Runnable cmd, long delay, TimeUnit unit) {
2129
      return delegate.schedule(cmd, delay, unit);
1✔
2130
    }
2131

2132
    @Override
2133
    public ScheduledFuture<?> scheduleAtFixedRate(
2134
        Runnable command, long initialDelay, long period, TimeUnit unit) {
2135
      return delegate.scheduleAtFixedRate(command, initialDelay, period, unit);
1✔
2136
    }
2137

2138
    @Override
2139
    public ScheduledFuture<?> scheduleWithFixedDelay(
2140
        Runnable command, long initialDelay, long delay, TimeUnit unit) {
2141
      return delegate.scheduleWithFixedDelay(command, initialDelay, delay, unit);
×
2142
    }
2143

2144
    @Override
2145
    public boolean awaitTermination(long timeout, TimeUnit unit)
2146
        throws InterruptedException {
2147
      return delegate.awaitTermination(timeout, unit);
×
2148
    }
2149

2150
    @Override
2151
    public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks)
2152
        throws InterruptedException {
2153
      return delegate.invokeAll(tasks);
×
2154
    }
2155

2156
    @Override
2157
    public <T> List<Future<T>> invokeAll(
2158
        Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit)
2159
        throws InterruptedException {
2160
      return delegate.invokeAll(tasks, timeout, unit);
×
2161
    }
2162

2163
    @Override
2164
    public <T> T invokeAny(Collection<? extends Callable<T>> tasks)
2165
        throws InterruptedException, ExecutionException {
2166
      return delegate.invokeAny(tasks);
×
2167
    }
2168

2169
    @Override
2170
    public <T> T invokeAny(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit)
2171
        throws InterruptedException, ExecutionException, TimeoutException {
2172
      return delegate.invokeAny(tasks, timeout, unit);
×
2173
    }
2174

2175
    @Override
2176
    public boolean isShutdown() {
2177
      return delegate.isShutdown();
×
2178
    }
2179

2180
    @Override
2181
    public boolean isTerminated() {
2182
      return delegate.isTerminated();
×
2183
    }
2184

2185
    @Override
2186
    public void shutdown() {
2187
      throw new UnsupportedOperationException("Restricted: shutdown() is not allowed");
1✔
2188
    }
2189

2190
    @Override
2191
    public List<Runnable> shutdownNow() {
2192
      throw new UnsupportedOperationException("Restricted: shutdownNow() is not allowed");
1✔
2193
    }
2194

2195
    @Override
2196
    public <T> Future<T> submit(Callable<T> task) {
2197
      return delegate.submit(task);
×
2198
    }
2199

2200
    @Override
2201
    public Future<?> submit(Runnable task) {
2202
      return delegate.submit(task);
×
2203
    }
2204

2205
    @Override
2206
    public <T> Future<T> submit(Runnable task, T result) {
2207
      return delegate.submit(task, result);
×
2208
    }
2209

2210
    @Override
2211
    public void execute(Runnable command) {
2212
      delegate.execute(command);
×
2213
    }
×
2214
  }
2215

2216
  /**
2217
   * A ResolutionState indicates the status of last name resolution.
2218
   */
2219
  enum ResolutionState {
1✔
2220
    NO_RESOLUTION,
1✔
2221
    SUCCESS,
1✔
2222
    ERROR
1✔
2223
  }
2224
}
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