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

grpc / grpc-java / #20434

04 Sep 2026 05:20AM UTC coverage: 89.211% (+0.009%) from 89.202%
#20434

push

github

web-flow
core: Reset configSelector on realChannel while entering IDLE state (#12832)

ManagedChannel will be stuck in IDLE state when xDS control plane
doesn't have a resource anymore.
The scenario is following:
1. Channel is open for a xds resource.
2. XdsNameResolver subscribes to the resource on xDS control plane. 
3. The resource is removed from xDS control plane for extended period of
time (unhealthy for more than the idle timeout on Channel)
4. Channel enters into TRANSIENT_FAILURE
5. The Idle timeout triggers and Channel shutdowns XdsNameResolver and
other resources. xDS watchers are removed.
6. The resource comes back online on xDS control plane.
7. A new GRPC call is executed targeting the channel.
8. The channel stays in IDLE state and reports:
"io.grpc.StatusRuntimeException: UNAVAILABLE: LDS resource xxxx does not
exist nodeID: yyyy" because realChannel.configSelector still points to
old state.

```    
public <ReqT, RespT> ClientCall<ReqT, RespT> newCall(
        MethodDescriptor<ReqT, RespT> method, CallOptions callOptions) {
      if (configSelector.get() != INITIAL_PENDING_SELECTOR) {
        return newClientCall(method, callOptions);
      }
...
```

38641 of 43314 relevant lines covered (89.21%)

0.89 hits per line

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

93.41
/../core/src/main/java/io/grpc/internal/ManagedChannelImpl.java
1
/*
2
 * Copyright 2016 The gRPC Authors
3
 *
4
 * Licensed under the Apache License, Version 2.0 (the "License");
5
 * you may not use this file except in compliance with the License.
6
 * You may obtain a copy of the License at
7
 *
8
 *     http://www.apache.org/licenses/LICENSE-2.0
9
 *
10
 * Unless required by applicable law or agreed to in writing, software
11
 * distributed under the License is distributed on an "AS IS" BASIS,
12
 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13
 * See the License for the specific language governing permissions and
14
 * limitations under the License.
15
 */
16

17
package io.grpc.internal;
18

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

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

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

131
  static final long IDLE_TIMEOUT_MILLIS_DISABLE = -1;
132

133
  static final long SUBCHANNEL_SHUTDOWN_DELAY_SECONDS = 5;
134

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

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

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

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

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

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

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

206
  private boolean fullStreamDecompression;
207

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

1099
    private ClientCall<ReqT, RespT> delegate;
1100

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

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

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

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

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

1159
      callExecutor.execute(new CloseInContext());
1✔
1160
    }
1✔
1161

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

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

1174
    @Override
1175
    public void request(int numMessages) {}
1✔
1176

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

1180
    @Override
1181
    public void halfClose() {}
×
1182

1183
    @Override
1184
    public void sendMessage(Object message) {}
×
1185

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

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

1210
      terminated = true;
1✔
1211
      terminatedLatch.countDown();
1✔
1212
    }
1213
  }
1✔
1214

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

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

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

1243
    syncContext.execute(new NotifyStateChanged());
1✔
1244
  }
1✔
1245

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

1263
    syncContext.execute(new ResetConnectBackoff());
1✔
1264
  }
1✔
1265

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

1279
    syncContext.execute(new PrepareToLoseNetworkRunnable());
1✔
1280
  }
1✔
1281

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

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

1294
    @GuardedBy("lock")
1295
    Status shutdownStatus;
1296

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

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

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

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

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

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

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

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

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

1365
  private final class LbHelperImpl extends LoadBalancer.Helper {
1✔
1366
    LoadBalancer lb;
1367

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

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

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

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

1406
      syncContext.execute(new LoadBalancerRefreshNameResolution());
1✔
1407
    }
1✔
1408

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

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

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

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

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

1463
      final class ResolvingOobChannelBuilder
1464
          extends ForwardingChannelBuilder2<ResolvingOobChannelBuilder> {
1465
        final ManagedChannelBuilder<?> delegate;
1466

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

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

1509
      checkState(!terminated, "Channel is terminated");
1✔
1510

1511
      ResolvingOobChannelBuilder builder = new ResolvingOobChannelBuilder();
1✔
1512

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

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

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

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

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

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

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

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

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

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

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

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

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

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

1601
  static final class OobChannel extends ForwardingManagedChannel {
1602
    private final OobNameResolverProvider resolverProvider;
1603

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

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

1614
  final class NameResolverListener extends NameResolver.Listener2 {
1615
    final LbHelperImpl helper;
1616
    final NameResolver resolver;
1617

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

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

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

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

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

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

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

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

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

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

1776
      syncContext.execute(new NameResolverErrorHandler());
1✔
1777
    }
1✔
1778

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

1792
      helper.lb.handleNameResolutionError(error);
1✔
1793
    }
1✔
1794
  }
1795

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

2117
  private static final class RestrictedScheduledExecutor implements ScheduledExecutorService {
2118
    final ScheduledExecutorService delegate;
2119

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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