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

grpc / grpc-java / #20406

14 Aug 2026 01:27PM UTC coverage: 89.164% (-0.04%) from 89.205%
#20406

push

github

web-flow
rls: implement stale_header_data caching and propagation in RLS (#12972)

This PR implements stale_header_data caching and propagation for Route Lookup Service (RLS) in :grpc-rls, addressing gRPC Java does not handle RLS stale_header_data properly).

38535 of 43218 relevant lines covered (89.16%)

0.89 hits per line

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

88.99
/../rls/src/main/java/io/grpc/rls/CachingRlsLbClient.java
1
/*
2
 * Copyright 2020 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.rls;
18

19
import static com.google.common.base.Preconditions.checkNotNull;
20
import static com.google.common.base.Preconditions.checkState;
21

22
import com.google.common.annotations.VisibleForTesting;
23
import com.google.common.base.Converter;
24
import com.google.common.base.MoreObjects;
25
import com.google.common.base.MoreObjects.ToStringHelper;
26
import com.google.common.base.Ticker;
27
import com.google.common.util.concurrent.Futures;
28
import com.google.common.util.concurrent.ListenableFuture;
29
import com.google.common.util.concurrent.MoreExecutors;
30
import com.google.common.util.concurrent.SettableFuture;
31
import com.google.errorprone.annotations.CheckReturnValue;
32
import com.google.errorprone.annotations.concurrent.GuardedBy;
33
import io.grpc.ChannelLogger;
34
import io.grpc.ChannelLogger.ChannelLogLevel;
35
import io.grpc.ConnectivityState;
36
import io.grpc.Grpc;
37
import io.grpc.LoadBalancer.Helper;
38
import io.grpc.LoadBalancer.PickResult;
39
import io.grpc.LoadBalancer.PickSubchannelArgs;
40
import io.grpc.LoadBalancer.ResolvedAddresses;
41
import io.grpc.LoadBalancer.SubchannelPicker;
42
import io.grpc.LongCounterMetricInstrument;
43
import io.grpc.LongGaugeMetricInstrument;
44
import io.grpc.ManagedChannel;
45
import io.grpc.ManagedChannelBuilder;
46
import io.grpc.Metadata;
47
import io.grpc.MetricInstrumentRegistry;
48
import io.grpc.MetricRecorder.BatchCallback;
49
import io.grpc.MetricRecorder.BatchRecorder;
50
import io.grpc.MetricRecorder.Registration;
51
import io.grpc.Status;
52
import io.grpc.internal.BackoffPolicy;
53
import io.grpc.internal.ExponentialBackoffPolicy;
54
import io.grpc.lookup.v1.RouteLookupServiceGrpc;
55
import io.grpc.lookup.v1.RouteLookupServiceGrpc.RouteLookupServiceStub;
56
import io.grpc.rls.ChildLoadBalancerHelper.ChildLoadBalancerHelperProvider;
57
import io.grpc.rls.LbPolicyConfiguration.ChildPolicyWrapper;
58
import io.grpc.rls.LbPolicyConfiguration.RefCountedChildPolicyWrapperFactory;
59
import io.grpc.rls.LruCache.EvictionListener;
60
import io.grpc.rls.LruCache.EvictionType;
61
import io.grpc.rls.RlsProtoConverters.RouteLookupResponseConverter;
62
import io.grpc.rls.RlsProtoData.RouteLookupConfig;
63
import io.grpc.rls.RlsProtoData.RouteLookupRequest;
64
import io.grpc.rls.RlsProtoData.RouteLookupRequestKey;
65
import io.grpc.rls.RlsProtoData.RouteLookupResponse;
66
import io.grpc.stub.StreamObserver;
67
import io.grpc.util.ForwardingLoadBalancerHelper;
68
import java.net.URI;
69
import java.net.URISyntaxException;
70
import java.util.Arrays;
71
import java.util.Collections;
72
import java.util.HashMap;
73
import java.util.List;
74
import java.util.Map;
75
import java.util.UUID;
76
import java.util.concurrent.Future;
77
import java.util.concurrent.ScheduledExecutorService;
78
import java.util.concurrent.TimeUnit;
79
import javax.annotation.Nullable;
80
import javax.annotation.concurrent.ThreadSafe;
81

82
/**
83
 * A CachingRlsLbClient is a core implementation of RLS loadbalancer supports dynamic request
84
 * routing by fetching the decision from route lookup server. Every single request is routed by
85
 * the server's decision. To reduce the performance penalty, {@link LruCache} is used.
86
 */
87
@ThreadSafe
88
final class CachingRlsLbClient {
89

90
  private static final Converter<RouteLookupRequest, io.grpc.lookup.v1.RouteLookupRequest>
91
      REQUEST_CONVERTER = new RlsProtoConverters.RouteLookupRequestConverter().reverse();
1 ✔
92
  private static final Converter<RouteLookupResponse, io.grpc.lookup.v1.RouteLookupResponse>
93
      RESPONSE_CONVERTER = new RouteLookupResponseConverter().reverse();
1 ✔
94
  public static final long MIN_EVICTION_TIME_DELTA_NANOS = TimeUnit.SECONDS.toNanos(5);
1 ✔
95
  public static final int BYTES_PER_CHAR = 2;
96
  public static final int STRING_OVERHEAD_BYTES = 38;
97
  /** Minimum bytes for a Java Object. */
98
  public static final int OBJ_OVERHEAD_B = 16;
99

100
  private static final LongCounterMetricInstrument DEFAULT_TARGET_PICKS_COUNTER;
101
  private static final LongCounterMetricInstrument TARGET_PICKS_COUNTER;
102
  private static final LongCounterMetricInstrument FAILED_PICKS_COUNTER;
103
  private static final LongGaugeMetricInstrument CACHE_ENTRIES_GAUGE;
104
  private static final LongGaugeMetricInstrument CACHE_SIZE_GAUGE;
105
  private final Registration gaugeRegistration;
106
  private final String metricsInstanceUuid = UUID.randomUUID().toString();
1 ✔
107

108
  // All cache status changes (pending, backoff, success) must be under this lock
109
  private final Object lock = new Object();
1 ✔
110
  // LRU cache based on access order (BACKOFF and actual data will be here)
111
  @GuardedBy("lock")
112
  private final RlsAsyncLruCache linkedHashLruCache;
113
  private final Future<?> periodicCleaner;
114
  // any RPC on the fly will cached in this map
115
  @GuardedBy("lock")
1 ✔
116
  private final Map<RouteLookupRequestKey, PendingCacheEntry> pendingCallCache = new HashMap<>();
117

118
  private final ScheduledExecutorService scheduledExecutorService;
119
  private final Ticker ticker;
120
  private final Throttler throttler;
121

122
  private final LbPolicyConfiguration lbPolicyConfig;
123
  private final BackoffPolicy.Provider backoffProvider;
124
  private final long maxAgeNanos;
125
  private final long staleAgeNanos;
126
  private final long callTimeoutNanos;
127

128
  private final RlsLbHelper helper;
129
  private final ManagedChannel rlsChannel;
130
  private final RouteLookupServiceStub rlsStub;
131
  private final RlsPicker rlsPicker;
132
  private final ResolvedAddressFactory childLbResolvedAddressFactory;
133
  @GuardedBy("lock")
134
  private final RefCountedChildPolicyWrapperFactory refCountedChildPolicyWrapperFactory;
135
  private final ChannelLogger logger;
136
  private final ChildPolicyWrapper fallbackChildPolicyWrapper;
137

138
  static {
139
    MetricInstrumentRegistry metricInstrumentRegistry
140
        = MetricInstrumentRegistry.getDefaultRegistry();
1 ✔
141
    DEFAULT_TARGET_PICKS_COUNTER = metricInstrumentRegistry.registerLongCounter(
1 ✔
142
        "grpc.lb.rls.default_target_picks",
143
        "EXPERIMENTAL. Number of LB picks sent to the default target", "{pick}",
144
        Arrays.asList("grpc.target", "grpc.lb.rls.server_target",
1 ✔
145
            "grpc.lb.rls.data_plane_target", "grpc.lb.pick_result"),
146
        Arrays.asList("grpc.client.call.custom"),
1 ✔
147
        false);
148
    TARGET_PICKS_COUNTER = metricInstrumentRegistry.registerLongCounter("grpc.lb.rls.target_picks",
1 ✔
149
        "EXPERIMENTAL. Number of LB picks sent to each RLS target. Note that if the default "
150
            + "target is also returned by the RLS server, RPCs sent to that target from the cache "
151
            + "will be counted in this metric, not in grpc.rls.default_target_picks.", "{pick}",
152
        Arrays.asList("grpc.target", "grpc.lb.rls.server_target", "grpc.lb.rls.data_plane_target",
1 ✔
153
            "grpc.lb.pick_result"),
154
        Arrays.asList("grpc.client.call.custom"),
1 ✔
155
        false);
156
    FAILED_PICKS_COUNTER = metricInstrumentRegistry.registerLongCounter("grpc.lb.rls.failed_picks",
1 ✔
157
        "EXPERIMENTAL. Number of LB picks failed due to either a failed RLS request or the "
158
            + "RLS channel being throttled", "{pick}",
159
        Arrays.asList("grpc.target", "grpc.lb.rls.server_target"),
1 ✔
160
        Arrays.asList("grpc.client.call.custom"), false);
1 ✔
161
    CACHE_ENTRIES_GAUGE = metricInstrumentRegistry.registerLongGauge("grpc.lb.rls.cache_entries",
1 ✔
162
        "EXPERIMENTAL. Number of entries in the RLS cache", "{entry}",
163
        Arrays.asList("grpc.target", "grpc.lb.rls.server_target", "grpc.lb.rls.instance_uuid"),
1 ✔
164
        Collections.emptyList(), false);
1 ✔
165
    CACHE_SIZE_GAUGE = metricInstrumentRegistry.registerLongGauge("grpc.lb.rls.cache_size",
1 ✔
166
        "EXPERIMENTAL. The current size of the RLS cache", "By",
167
        Arrays.asList("grpc.target", "grpc.lb.rls.server_target", "grpc.lb.rls.instance_uuid"),
1 ✔
168
        Collections.emptyList(), false);
1 ✔
169
  }
170

171
  private CachingRlsLbClient(Builder builder) {
1 ✔
172
    helper = new RlsLbHelper(checkNotNull(builder.helper, "helper"));
1 ✔
173
    scheduledExecutorService = helper.getScheduledExecutorService();
1 ✔
174
    lbPolicyConfig = checkNotNull(builder.lbPolicyConfig, "lbPolicyConfig");
1 ✔
175
    RouteLookupConfig rlsConfig = lbPolicyConfig.getRouteLookupConfig();
1 ✔
176
    maxAgeNanos = rlsConfig.maxAgeInNanos();
1 ✔
177
    staleAgeNanos = rlsConfig.staleAgeInNanos();
1 ✔
178
    callTimeoutNanos = rlsConfig.lookupServiceTimeoutInNanos();
1 ✔
179
    ticker = checkNotNull(builder.ticker, "ticker");
1 ✔
180
    throttler = checkNotNull(builder.throttler, "throttler");
1 ✔
181
    linkedHashLruCache =
1 ✔
182
        new RlsAsyncLruCache(
183
            rlsConfig.cacheSizeBytes(),
1 ✔
184
            new AutoCleaningEvictionListener(builder.evictionListener),
1 ✔
185
            ticker,
186
            helper);
187
    periodicCleaner =
1 ✔
188
        scheduledExecutorService.scheduleAtFixedRate(this::periodicClean, 1, 1, TimeUnit.MINUTES);
1 ✔
189
    logger = helper.getChannelLogger();
1 ✔
190
    String serverHost = null;
1 ✔
191
    try {
192
      serverHost = new URI(null, helper.getAuthority(), null, null, null).getHost();
1 ✔
193
    } catch (URISyntaxException ignore) {
×
194
      // handled by the following null check
195
    }
1 ✔
196
    if (serverHost == null) {
1 ✔
197
      logger.log(
×
198
          ChannelLogLevel.DEBUG, "Can not get hostname from authority: {0}", helper.getAuthority());
×
199
      serverHost = helper.getAuthority();
×
200
    }
201
    RlsRequestFactory requestFactory = new RlsRequestFactory(
1 ✔
202
        lbPolicyConfig.getRouteLookupConfig(), serverHost);
1 ✔
203
    rlsPicker = new RlsPicker(requestFactory, rlsConfig.lookupService());
1 ✔
204
    // It is safe to use helper.getUnsafeChannelCredentials() because the client authenticates the
205
    // RLS server using the same authority as the backends, even though the RLS server’s addresses
206
    // will be looked up differently than the backends; overrideAuthority(helper.getAuthority()) is
207
    // called to impose the authority security restrictions.
208
    ManagedChannelBuilder<?> rlsChannelBuilder = helper.createResolvingOobChannelBuilder(
1 ✔
209
        rlsConfig.lookupService(), helper.getUnsafeChannelCredentials());
1 ✔
210
    rlsChannelBuilder.overrideAuthority(helper.getAuthority());
1 ✔
211
    Map<String, ?> routeLookupChannelServiceConfig =
1 ✔
212
        lbPolicyConfig.getRouteLookupChannelServiceConfig();
1 ✔
213
    if (routeLookupChannelServiceConfig != null) {
1 ✔
214
      logger.log(
1 ✔
215
          ChannelLogLevel.DEBUG,
216
          "RLS channel service config: {0}",
217
          routeLookupChannelServiceConfig);
218
      rlsChannelBuilder.defaultServiceConfig(routeLookupChannelServiceConfig);
1 ✔
219
      rlsChannelBuilder.disableServiceConfigLookUp();
1 ✔
220
    }
221
    rlsChannel = rlsChannelBuilder.build();
1 ✔
222
    Runnable rlsServerConnectivityStateChangeHandler = new Runnable() {
1 ✔
223
      private boolean wasInTransientFailure;
224
      @Override
225
      public void run() {
226
        ConnectivityState currentState = rlsChannel.getState(false);
1 ✔
227
        if (currentState == ConnectivityState.TRANSIENT_FAILURE) {
1 ✔
228
          wasInTransientFailure = true;
1 ✔
229
        } else if (wasInTransientFailure && currentState == ConnectivityState.READY) {
1 ✔
230
          wasInTransientFailure = false;
1 ✔
231
          synchronized (lock) {
1 ✔
232
            boolean anyBackoffsCanceled = false;
1 ✔
233
            for (CacheEntry value : linkedHashLruCache.values()) {
1 ✔
234
              if (value instanceof BackoffCacheEntry) {
1 ✔
235
                if (((BackoffCacheEntry) value).scheduledFuture.cancel(false)) {
1 ✔
236
                  anyBackoffsCanceled = true;
1 ✔
237
                }
238
              }
239
            }
1 ✔
240
            if (anyBackoffsCanceled) {
1 ✔
241
              // Cache updated. updateBalancingState() to reattempt picks
242
              helper.triggerPendingRpcProcessing();
1 ✔
243
            }
244
          }
1 ✔
245
        }
246
        rlsChannel.notifyWhenStateChanged(currentState, this);
1 ✔
247
      }
1 ✔
248
    };
249
    rlsChannel.notifyWhenStateChanged(
1 ✔
250
        ConnectivityState.IDLE, rlsServerConnectivityStateChangeHandler);
251
    rlsStub = RouteLookupServiceGrpc.newStub(rlsChannel);
1 ✔
252
    childLbResolvedAddressFactory =
1 ✔
253
        checkNotNull(builder.resolvedAddressFactory, "resolvedAddressFactory");
1 ✔
254
    backoffProvider = builder.backoffProvider;
1 ✔
255
    ChildLoadBalancerHelperProvider childLbHelperProvider =
1 ✔
256
        new ChildLoadBalancerHelperProvider(helper, new SubchannelStateManagerImpl(), rlsPicker);
257
    refCountedChildPolicyWrapperFactory =
1 ✔
258
        new RefCountedChildPolicyWrapperFactory(
259
            lbPolicyConfig.getLoadBalancingPolicy(), childLbResolvedAddressFactory,
1 ✔
260
            childLbHelperProvider);
261
    // TODO(creamsoup) wait until lb is ready
262
    String defaultTarget = lbPolicyConfig.getRouteLookupConfig().defaultTarget();
1 ✔
263
    if (defaultTarget != null && !defaultTarget.isEmpty()) {
1 ✔
264
      fallbackChildPolicyWrapper = refCountedChildPolicyWrapperFactory.createOrGet(defaultTarget);
1 ✔
265
    } else {
266
      fallbackChildPolicyWrapper = null;
1 ✔
267
    }
268

269
    gaugeRegistration = helper.getMetricRecorder()
1 ✔
270
        .registerBatchCallback(new BatchCallback() {
1 ✔
271
          @Override
272
          public void accept(BatchRecorder recorder) {
273
            int estimatedSize;
274
            long estimatedSizeBytes;
275
            synchronized (lock) {
1 ✔
276
              estimatedSize = linkedHashLruCache.estimatedSize();
1 ✔
277
              estimatedSizeBytes = linkedHashLruCache.estimatedSizeBytes();
1 ✔
278
            }
1 ✔
279
            recorder.recordLongGauge(CACHE_ENTRIES_GAUGE, estimatedSize,
1 ✔
280
                Arrays.asList(helper.getChannelTarget(), rlsConfig.lookupService(),
1 ✔
281
                    metricsInstanceUuid), Collections.emptyList());
1 ✔
282
            recorder.recordLongGauge(CACHE_SIZE_GAUGE, estimatedSizeBytes,
1 ✔
283
                Arrays.asList(helper.getChannelTarget(), rlsConfig.lookupService(),
1 ✔
284
                    metricsInstanceUuid), Collections.emptyList());
1 ✔
285
          }
1 ✔
286
        }, CACHE_ENTRIES_GAUGE, CACHE_SIZE_GAUGE);
287

288
    logger.log(ChannelLogLevel.DEBUG, "CachingRlsLbClient created");
1 ✔
289
  }
1 ✔
290

291
  void init() {
292
    synchronized (lock) {
1 ✔
293
      refCountedChildPolicyWrapperFactory.init();
1 ✔
294
    }
1 ✔
295
  }
1 ✔
296

297
  Status acceptResolvedAddressFactory(ResolvedAddressFactory childLbResolvedAddressFactory) {
298
    synchronized (lock) {
1 ✔
299
      return refCountedChildPolicyWrapperFactory.acceptResolvedAddressFactory(
1 ✔
300
          childLbResolvedAddressFactory);
301
    }
302
  }
303

304
  /**
305
   * Convert the status to UNAVAILABLE and enhance the error message.
306
   * @param status status as provided by server
307
   * @param serverName Used for error description
308
   * @return Transformed status
309
   */
310
  static Status convertRlsServerStatus(Status status, String serverName) {
311
    return Status.UNAVAILABLE.withCause(status.getCause()).withDescription(
1 ✔
312
        String.format("Unable to retrieve RLS targets from RLS server %s.  "
1 ✔
313
                + "RLS server returned: %s: %s",
314
            serverName, status.getCode(), status.getDescription()));
1 ✔
315
  }
316

317
  private void periodicClean() {
318
    synchronized (lock) {
1 ✔
319
      linkedHashLruCache.cleanupExpiredEntries();
1 ✔
320
    }
1 ✔
321
  }
1 ✔
322

323
  /** Populates async cache entry for new request. */
324
  @GuardedBy("lock")
325
  private CachedRouteLookupResponse asyncRlsCall(
326
      RouteLookupRequestKey routeLookupRequestKey, @Nullable BackoffPolicy backoffPolicy,
327
      RouteLookupRequest.Reason routeLookupReason, @Nullable String staleHeaderData) {
328
    if (throttler.shouldThrottle()) {
1 ✔
329
      logger.log(ChannelLogLevel.DEBUG, "[RLS Entry {0}] Throttled RouteLookup",
1 ✔
330
          routeLookupRequestKey);
331
      // Cache updated, but no need to call updateBalancingState because no RPCs were queued waiting
332
      // on this result
333
      return CachedRouteLookupResponse.backoffEntry(createBackOffEntry(
1 ✔
334
          routeLookupRequestKey, Status.RESOURCE_EXHAUSTED.withDescription("RLS throttled"),
1 ✔
335
          backoffPolicy));
336
    }
337
    final SettableFuture<RouteLookupResponse> response = SettableFuture.create();
1 ✔
338
    io.grpc.lookup.v1.RouteLookupRequest routeLookupRequest = REQUEST_CONVERTER.convert(
1 ✔
339
        RouteLookupRequest.create(
1 ✔
340
            routeLookupRequestKey.keyMap(), routeLookupReason, staleHeaderData));
1 ✔
341
    logger.log(ChannelLogLevel.DEBUG,
1 ✔
342
        "[RLS Entry {0}] Starting RouteLookup: {1}", routeLookupRequestKey, routeLookupRequest);
343
    rlsStub.withDeadlineAfter(callTimeoutNanos, TimeUnit.NANOSECONDS)
1 ✔
344
        .routeLookup(
1 ✔
345
            routeLookupRequest,
346
            new StreamObserver<io.grpc.lookup.v1.RouteLookupResponse>() {
1 ✔
347
              @Override
348
              public void onNext(io.grpc.lookup.v1.RouteLookupResponse value) {
349
                logger.log(ChannelLogLevel.DEBUG,
1 ✔
350
                    "[RLS Entry {0}] RouteLookup succeeded: {1}", routeLookupRequestKey, value);
351
                response.set(RESPONSE_CONVERTER.reverse().convert(value));
1 ✔
352
              }
1 ✔
353

354
              @Override
355
              public void onError(Throwable t) {
356
                logger.log(ChannelLogLevel.DEBUG,
1 ✔
357
                    "[RLS Entry {0}] RouteLookup failed: {1}", routeLookupRequestKey, t);
358
                response.setException(t);
1 ✔
359
                throttler.registerBackendResponse(true);
1 ✔
360
              }
1 ✔
361

362
              @Override
363
              public void onCompleted() {
364
                throttler.registerBackendResponse(false);
1 ✔
365
              }
1 ✔
366
            });
367
    return CachedRouteLookupResponse.pendingResponse(
1 ✔
368
        createPendingEntry(routeLookupRequestKey, response, backoffPolicy));
1 ✔
369
  }
370

371
  /**
372
   * Returns async response of the {@code request}. The returned value can be in 3 different states;
373
   * cached, pending and backed-off due to error. The result remains same even if the status is
374
   * changed after the return.
375
   */
376
  @CheckReturnValue
377
  final CachedRouteLookupResponse get(final RouteLookupRequestKey routeLookupRequestKey) {
378
    synchronized (lock) {
1 ✔
379
      final CacheEntry cacheEntry;
380
      cacheEntry = linkedHashLruCache.read(routeLookupRequestKey);
1 ✔
381
      if (cacheEntry == null
1 ✔
382
          || (cacheEntry instanceof BackoffCacheEntry
383
          && !((BackoffCacheEntry) cacheEntry).isInBackoffPeriod())) {
1 ✔
384
        PendingCacheEntry pendingEntry = pendingCallCache.get(routeLookupRequestKey);
1 ✔
385
        if (pendingEntry != null) {
1 ✔
386
          return CachedRouteLookupResponse.pendingResponse(pendingEntry);
×
387
        }
388
        return asyncRlsCall(routeLookupRequestKey, cacheEntry instanceof BackoffCacheEntry
1 ✔
389
            ? ((BackoffCacheEntry) cacheEntry).backoffPolicy : null,
1 ✔
390
            RouteLookupRequest.Reason.REASON_MISS, /* staleHeaderData= */ null);
391
      }
392

393
      if (cacheEntry instanceof DataCacheEntry) {
1 ✔
394
        // cache hit, initiate async-refresh if entry is staled
395
        DataCacheEntry dataEntry = ((DataCacheEntry) cacheEntry);
1 ✔
396
        if (dataEntry.isStaled(ticker.read())) {
1 ✔
397
          dataEntry.maybeRefresh();
1 ✔
398
        }
399
        return CachedRouteLookupResponse.dataEntry((DataCacheEntry) cacheEntry);
1 ✔
400
      }
401
      return CachedRouteLookupResponse.backoffEntry((BackoffCacheEntry) cacheEntry);
1 ✔
402
    }
403
  }
404

405
  /** Performs any pending maintenance operations needed by the cache. */
406
  void close() {
407
    logger.log(ChannelLogLevel.DEBUG, "CachingRlsLbClient closed");
1 ✔
408
    synchronized (lock) {
1 ✔
409
      periodicCleaner.cancel(false);
1 ✔
410
      // all childPolicyWrapper will be returned via AutoCleaningEvictionListener
411
      linkedHashLruCache.close();
1 ✔
412
      // TODO(creamsoup) maybe cancel all pending requests
413
      pendingCallCache.clear();
1 ✔
414
      rlsChannel.shutdownNow();
1 ✔
415
      rlsPicker.close();
1 ✔
416
      gaugeRegistration.close();
1 ✔
417
    }
1 ✔
418
  }
1 ✔
419

420
  void requestConnection() {
421
    rlsChannel.getState(true);
×
422
  }
×
423

424
  @GuardedBy("lock")
425
  private PendingCacheEntry createPendingEntry(
426
      RouteLookupRequestKey routeLookupRequestKey,
427
      ListenableFuture<RouteLookupResponse> pendingCall,
428
      @Nullable BackoffPolicy backoffPolicy) {
429
    PendingCacheEntry entry = new PendingCacheEntry(routeLookupRequestKey, pendingCall,
1 ✔
430
        backoffPolicy);
431
    // Add the entry to the map before adding the Listener, because the listener removes the
432
    // entry from the map
433
    pendingCallCache.put(routeLookupRequestKey, entry);
1 ✔
434
    // Beware that the listener can run immediately on the current thread
435
    pendingCall.addListener(() -> pendingRpcComplete(entry), MoreExecutors.directExecutor());
1 ✔
436
    return entry;
1 ✔
437
  }
438

439
  private void pendingRpcComplete(PendingCacheEntry entry) {
440
    synchronized (lock) {
1 ✔
441
      boolean clientClosed = pendingCallCache.remove(entry.routeLookupRequestKey) == null;
1 ✔
442
      if (clientClosed) {
1 ✔
443
        return;
1 ✔
444
      }
445

446
      try {
447
        createDataEntry(entry.routeLookupRequestKey, Futures.getDone(entry.pendingCall));
1 ✔
448
        // Cache updated. DataCacheEntry constructor indirectly calls updateBalancingState() to
449
        // reattempt picks when the child LB is done connecting
450
      } catch (Exception e) {
1 ✔
451
        createBackOffEntry(entry.routeLookupRequestKey, Status.fromThrowable(e),
1 ✔
452
            entry.backoffPolicy);
1 ✔
453
        // Cache updated. updateBalancingState() to reattempt picks
454
        helper.triggerPendingRpcProcessing();
1 ✔
455
      }
1 ✔
456
    }
1 ✔
457
  }
1 ✔
458

459
  @GuardedBy("lock")
460
  private DataCacheEntry createDataEntry(
461
      RouteLookupRequestKey routeLookupRequestKey, RouteLookupResponse routeLookupResponse) {
462
    logger.log(
1 ✔
463
        ChannelLogLevel.DEBUG,
464
        "[RLS Entry {0}] Transition to data cache: routeLookupResponse={1}",
465
        routeLookupRequestKey, routeLookupResponse);
466
    DataCacheEntry entry = new DataCacheEntry(routeLookupRequestKey, routeLookupResponse);
1 ✔
467
    // Constructor for DataCacheEntry causes updateBalancingState, but the picks can't happen until
468
    // this cache update because the lock is held
469
    linkedHashLruCache.cacheAndClean(routeLookupRequestKey, entry);
1 ✔
470
    return entry;
1 ✔
471
  }
472

473
  @GuardedBy("lock")
474
  private BackoffCacheEntry createBackOffEntry(RouteLookupRequestKey routeLookupRequestKey,
475
      Status status, @Nullable BackoffPolicy backoffPolicy) {
476
    if (backoffPolicy == null) {
1 ✔
477
      backoffPolicy = backoffProvider.get();
1 ✔
478
    }
479
    long delayNanos = backoffPolicy.nextBackoffNanos();
1 ✔
480
    logger.log(
1 ✔
481
        ChannelLogLevel.DEBUG,
482
        "[RLS Entry {0}] Transition to back off: status={1}, delayNanos={2}",
483
        routeLookupRequestKey, status, delayNanos);
1 ✔
484
    BackoffCacheEntry entry = new BackoffCacheEntry(routeLookupRequestKey, status, backoffPolicy,
1 ✔
485
        ticker.read() + delayNanos * 2);
1 ✔
486
    // Lock is held, so the task can't execute before the assignment
487
    entry.scheduledFuture = scheduledExecutorService.schedule(
1 ✔
488
        () -> refreshBackoffEntry(entry), delayNanos, TimeUnit.NANOSECONDS);
1 ✔
489
    linkedHashLruCache.cacheAndClean(routeLookupRequestKey, entry);
1 ✔
490
    return entry;
1 ✔
491
  }
492

493
  private void refreshBackoffEntry(BackoffCacheEntry entry) {
494
    synchronized (lock) {
1 ✔
495
      // This checks whether the task has been cancelled and prevents a second execution.
496
      if (!entry.scheduledFuture.cancel(false)) {
1 ✔
497
        // Future was previously cancelled
498
        return;
×
499
      }
500
      // Cache updated. updateBalancingState() to reattempt picks
501
      helper.triggerPendingRpcProcessing();
1 ✔
502
    }
1 ✔
503
  }
1 ✔
504

505
  private static final class RlsLbHelper extends ForwardingLoadBalancerHelper {
506

507
    final Helper helper;
508
    private ConnectivityState state;
509
    private SubchannelPicker picker;
510

511
    RlsLbHelper(Helper helper) {
1 ✔
512
      this.helper = helper;
1 ✔
513
    }
1 ✔
514

515
    @Override
516
    protected Helper delegate() {
517
      return helper;
1 ✔
518
    }
519

520
    @Override
521
    public void updateBalancingState(ConnectivityState newState, SubchannelPicker newPicker) {
522
      state = newState;
1 ✔
523
      picker = newPicker;
1 ✔
524
      super.updateBalancingState(newState, newPicker);
1 ✔
525
    }
1 ✔
526

527
    void triggerPendingRpcProcessing() {
528
      checkState(state != null, "updateBalancingState hasn't yet been called");
1 ✔
529
      helper.getSynchronizationContext().execute(
1 ✔
530
          () -> super.updateBalancingState(state, picker));
1 ✔
531
    }
1 ✔
532
  }
533

534
  /**
535
   * Viewer class for cached {@link RouteLookupResponse} and associated {@link ChildPolicyWrapper}.
536
   */
537
  static final class CachedRouteLookupResponse {
538
    // Should only have 1 of following 3 cache entries
539
    @Nullable
540
    private final DataCacheEntry dataCacheEntry;
541
    @Nullable
542
    private final PendingCacheEntry pendingCacheEntry;
543
    @Nullable
544
    private final BackoffCacheEntry backoffCacheEntry;
545

546
    CachedRouteLookupResponse(
547
        DataCacheEntry dataCacheEntry,
548
        PendingCacheEntry pendingCacheEntry,
549
        BackoffCacheEntry backoffCacheEntry) {
1 ✔
550
      this.dataCacheEntry = dataCacheEntry;
1 ✔
551
      this.pendingCacheEntry = pendingCacheEntry;
1 ✔
552
      this.backoffCacheEntry = backoffCacheEntry;
1 ✔
553
      checkState((dataCacheEntry != null ^ pendingCacheEntry != null ^ backoffCacheEntry != null)
1 ✔
554
          && !(dataCacheEntry != null && pendingCacheEntry != null && backoffCacheEntry != null),
555
          "Expected only 1 cache entry value provided");
556
    }
1 ✔
557

558
    static CachedRouteLookupResponse pendingResponse(PendingCacheEntry pendingEntry) {
559
      return new CachedRouteLookupResponse(null, pendingEntry, null);
1 ✔
560
    }
561

562
    static CachedRouteLookupResponse backoffEntry(BackoffCacheEntry backoffEntry) {
563
      return new CachedRouteLookupResponse(null, null, backoffEntry);
1 ✔
564
    }
565

566
    static CachedRouteLookupResponse dataEntry(DataCacheEntry dataEntry) {
567
      return new CachedRouteLookupResponse(dataEntry, null, null);
1 ✔
568
    }
569

570
    boolean hasData() {
571
      return dataCacheEntry != null;
1 ✔
572
    }
573

574
    @Nullable
575
    ChildPolicyWrapper getChildPolicyWrapper() {
576
      if (!hasData()) {
1 ✔
577
        return null;
×
578
      }
579
      return dataCacheEntry.getChildPolicyWrapper();
1 ✔
580
    }
581

582
    @VisibleForTesting
583
    @Nullable
584
    ChildPolicyWrapper getChildPolicyWrapper(String target) {
585
      if (!hasData()) {
1 ✔
586
        return null;
×
587
      }
588
      return dataCacheEntry.getChildPolicyWrapper(target);
1 ✔
589
    }
590

591
    @Nullable
592
    String getHeaderData() {
593
      if (!hasData()) {
1 ✔
594
        return null;
1 ✔
595
      }
596
      return dataCacheEntry.getHeaderData();
1 ✔
597
    }
598

599
    boolean hasError() {
600
      return backoffCacheEntry != null;
1 ✔
601
    }
602

603
    boolean isPending() {
604
      return pendingCacheEntry != null;
1 ✔
605
    }
606

607
    @Nullable
608
    Status getStatus() {
609
      if (!hasError()) {
1 ✔
610
        return null;
×
611
      }
612
      return backoffCacheEntry.getStatus();
1 ✔
613
    }
614

615
    @Override
616
    public String toString() {
617
      ToStringHelper toStringHelper = MoreObjects.toStringHelper(this);
×
618
      if (dataCacheEntry != null) {
×
619
        toStringHelper.add("dataCacheEntry", dataCacheEntry);
×
620
      }
621
      if (pendingCacheEntry != null) {
×
622
        toStringHelper.add("pendingCacheEntry", pendingCacheEntry);
×
623
      }
624
      if (backoffCacheEntry != null) {
×
625
        toStringHelper.add("backoffCacheEntry", backoffCacheEntry);
×
626
      }
627
      return toStringHelper.toString();
×
628
    }
629
  }
630

631
  /** A pending cache entry when the async RouteLookup RPC is still on the fly. */
632
  static final class PendingCacheEntry {
633
    private final ListenableFuture<RouteLookupResponse> pendingCall;
634
    private final RouteLookupRequestKey routeLookupRequestKey;
635
    @Nullable
636
    private final BackoffPolicy backoffPolicy;
637

638
    PendingCacheEntry(
639
        RouteLookupRequestKey routeLookupRequestKey,
640
        ListenableFuture<RouteLookupResponse> pendingCall,
641
        @Nullable BackoffPolicy backoffPolicy) {
1 ✔
642
      this.routeLookupRequestKey = checkNotNull(routeLookupRequestKey, "request");
1 ✔
643
      this.pendingCall = checkNotNull(pendingCall, "pendingCall");
1 ✔
644
      this.backoffPolicy = backoffPolicy;
1 ✔
645
    }
1 ✔
646

647
    @Override
648
    public String toString() {
649
      return MoreObjects.toStringHelper(this)
×
650
          .add("routeLookupRequestKey", routeLookupRequestKey)
×
651
          .toString();
×
652
    }
653
  }
654

655
  /** Common cache entry data for {@link RlsAsyncLruCache}. */
656
  abstract static class CacheEntry {
657

658
    protected final RouteLookupRequestKey routeLookupRequestKey;
659

660
    CacheEntry(RouteLookupRequestKey routeLookupRequestKey) {
1 ✔
661
      this.routeLookupRequestKey = checkNotNull(routeLookupRequestKey, "request");
1 ✔
662
    }
1 ✔
663

664
    abstract int getSizeBytes();
665

666
    abstract boolean isExpired(long now);
667

668
    abstract void cleanup();
669

670
    protected boolean isOldEnoughToBeEvicted(long now) {
671
      return true;
×
672
    }
673
  }
674

675
  /** Implementation of {@link CacheEntry} contains valid data. */
676
  final class DataCacheEntry extends CacheEntry {
677
    private final RouteLookupResponse response;
678
    private final long minEvictionTime;
679
    private final long expireTime;
680
    private final long staleTime;
681
    private final List<ChildPolicyWrapper> childPolicyWrappers;
682

683
    // GuardedBy CachingRlsLbClient.lock
684
    DataCacheEntry(RouteLookupRequestKey routeLookupRequestKey,
685
        final RouteLookupResponse response) {
1 ✔
686
      super(routeLookupRequestKey);
1 ✔
687
      this.response = checkNotNull(response, "response");
1 ✔
688
      checkState(!response.targets().isEmpty(), "No targets returned by RLS");
1 ✔
689
      childPolicyWrappers =
1 ✔
690
          refCountedChildPolicyWrapperFactory
1 ✔
691
              .createOrGet(response.targets());
1 ✔
692
      long now = ticker.read();
1 ✔
693
      minEvictionTime = now + MIN_EVICTION_TIME_DELTA_NANOS;
1 ✔
694
      expireTime = now + maxAgeNanos;
1 ✔
695
      staleTime = now + staleAgeNanos;
1 ✔
696
    }
1 ✔
697

698
    /**
699
     * Refreshes cache entry by creating {@link PendingCacheEntry}. When the {@code
700
     * PendingCacheEntry} received data from RLS server, it will replace the data entry if valid
701
     * data still exists. Flow looks like following.
702
     *
703
     * <pre>
704
     * Timeline                       | async refresh
705
     *                                V put new cache (entry2)
706
     * entry1: Pending | hasValue | staled  |
707
     * entry2:                        | OV* | pending | hasValue | staled |
708
     *
709
     * OV: old value
710
     * </pre>
711
     */
712
    void maybeRefresh() {
713
      synchronized (lock) { // Lock is already held, but ErrorProne can't tell
1 ✔
714
        if (pendingCallCache.containsKey(routeLookupRequestKey)) {
1 ✔
715
          // pending already requested
716
          return;
×
717
        }
718
        logger.log(ChannelLogLevel.DEBUG,
1 ✔
719
            "[RLS Entry {0}] Cache entry is stale, refreshing", routeLookupRequestKey);
720
        asyncRlsCall(routeLookupRequestKey, /* backoffPolicy= */ null,
1 ✔
721
            RouteLookupRequest.Reason.REASON_STALE, getHeaderData());
1 ✔
722
      }
1 ✔
723
    }
1 ✔
724

725
    @VisibleForTesting
726
    ChildPolicyWrapper getChildPolicyWrapper(String target) {
727
      for (ChildPolicyWrapper childPolicyWrapper : childPolicyWrappers) {
1 ✔
728
        if (childPolicyWrapper.getTarget().equals(target)) {
1 ✔
729
          return childPolicyWrapper;
1 ✔
730
        }
731
      }
1 ✔
732

733
      throw new RuntimeException("Target not found:" + target);
×
734
    }
735

736
    @Nullable
737
    ChildPolicyWrapper getChildPolicyWrapper() {
738
      for (ChildPolicyWrapper childPolicyWrapper : childPolicyWrappers) {
1 ✔
739
        if (childPolicyWrapper.getState() != ConnectivityState.TRANSIENT_FAILURE) {
1 ✔
740
          return childPolicyWrapper;
1 ✔
741
        }
742
      }
1 ✔
743
      return childPolicyWrappers.get(0);
1 ✔
744
    }
745

746
    String getHeaderData() {
747
      return response.getHeaderData();
1 ✔
748
    }
749

750
    // Assume UTF-16 (2 bytes) and overhead of a String object is 38 bytes
751
    int calcStringSize(String target) {
752
      return target.length() * BYTES_PER_CHAR + STRING_OVERHEAD_BYTES;
1 ✔
753
    }
754

755
    @Override
756
    int getSizeBytes() {
757
      int targetSize = 0;
1 ✔
758
      for (String target : response.targets()) {
1 ✔
759
        targetSize += calcStringSize(target);
1 ✔
760
      }
1 ✔
761
      return targetSize + calcStringSize(response.getHeaderData()) + OBJ_OVERHEAD_B // response size
1 ✔
762
          + Long.SIZE * 2 + OBJ_OVERHEAD_B; // Other fields
763
    }
764

765
    @Override
766
    boolean isExpired(long now) {
767
      return expireTime - now <= 0;
1 ✔
768
    }
769

770
    boolean isStaled(long now) {
771
      return staleTime - now <= 0;
1 ✔
772
    }
773

774
    @Override
775
    protected boolean isOldEnoughToBeEvicted(long now) {
776
      return minEvictionTime - now <= 0;
×
777
    }
778

779
    @Override
780
    void cleanup() {
781
      synchronized (lock) {
1 ✔
782
        for (ChildPolicyWrapper policyWrapper : childPolicyWrappers) {
1 ✔
783
          refCountedChildPolicyWrapperFactory.release(policyWrapper);
1 ✔
784
        }
1 ✔
785
      }
1 ✔
786
    }
1 ✔
787

788
    @Override
789
    public String toString() {
790
      return MoreObjects.toStringHelper(this)
×
791
          .add("request", routeLookupRequestKey)
×
792
          .add("response", response)
×
793
          .add("expireTime", expireTime)
×
794
          .add("staleTime", staleTime)
×
795
          .add("childPolicyWrappers", childPolicyWrappers)
×
796
          .toString();
×
797
    }
798
  }
799

800
  /**
801
   * Implementation of {@link CacheEntry} contains error. This entry will transition to pending
802
   * status when the backoff time is expired.
803
   */
804
  private static final class BackoffCacheEntry extends CacheEntry {
805

806
    private final Status status;
807
    private final BackoffPolicy backoffPolicy;
808
    private final long expiryTimeNanos;
809
    private Future<?> scheduledFuture;
810

811
    BackoffCacheEntry(RouteLookupRequestKey routeLookupRequestKey, Status status,
812
        BackoffPolicy backoffPolicy, long expiryTimeNanos) {
813
      super(routeLookupRequestKey);
1 ✔
814
      this.status = checkNotNull(status, "status");
1 ✔
815
      this.backoffPolicy = checkNotNull(backoffPolicy, "backoffPolicy");
1 ✔
816
      this.expiryTimeNanos = expiryTimeNanos;
1 ✔
817
    }
1 ✔
818

819
    Status getStatus() {
820
      return status;
1 ✔
821
    }
822

823
    @Override
824
    int getSizeBytes() {
825
      return OBJ_OVERHEAD_B * 3 + Long.SIZE + 8; // 3 java objects, 1 long and a boolean
1 ✔
826
    }
827

828
    boolean isInBackoffPeriod() {
829
      return !scheduledFuture.isDone();
1 ✔
830
    }
831

832
    @Override
833
    boolean isExpired(long nowNanos) {
834
      return nowNanos > expiryTimeNanos;
1 ✔
835
    }
836

837
    @Override
838
    void cleanup() {
839
      scheduledFuture.cancel(false);
1 ✔
840
    }
1 ✔
841

842
    @Override
843
    public String toString() {
844
      return MoreObjects.toStringHelper(this)
×
845
          .add("request", routeLookupRequestKey)
×
846
          .add("status", status)
×
847
          .toString();
×
848
    }
849
  }
850

851
  /** Returns a Builder for {@link CachingRlsLbClient}. */
852
  static Builder newBuilder() {
853
    return new Builder();
1 ✔
854
  }
855

856
  /** A Builder for {@link CachingRlsLbClient}. */
857
  static final class Builder {
1 ✔
858

859
    private Helper helper;
860
    private LbPolicyConfiguration lbPolicyConfig;
861
    private Throttler throttler = new HappyThrottler();
1 ✔
862
    private ResolvedAddressFactory resolvedAddressFactory;
863
    private Ticker ticker = Ticker.systemTicker();
1 ✔
864
    private EvictionListener<RouteLookupRequestKey, CacheEntry> evictionListener;
865
    private BackoffPolicy.Provider backoffProvider = new ExponentialBackoffPolicy.Provider();
1 ✔
866

867
    Builder setHelper(Helper helper) {
868
      this.helper = checkNotNull(helper, "helper");
1 ✔
869
      return this;
1 ✔
870
    }
871

872
    Builder setLbPolicyConfig(LbPolicyConfiguration lbPolicyConfig) {
873
      this.lbPolicyConfig = checkNotNull(lbPolicyConfig, "lbPolicyConfig");
1 ✔
874
      return this;
1 ✔
875
    }
876

877
    Builder setThrottler(Throttler throttler) {
878
      this.throttler = checkNotNull(throttler, "throttler");
1 ✔
879
      return this;
1 ✔
880
    }
881

882
    /**
883
     * Sets a factory to create {@link ResolvedAddresses} for child load balancer.
884
     */
885
    Builder setResolvedAddressesFactory(
886
        ResolvedAddressFactory resolvedAddressFactory) {
887
      this.resolvedAddressFactory =
1 ✔
888
          checkNotNull(resolvedAddressFactory, "resolvedAddressFactory");
1 ✔
889
      return this;
1 ✔
890
    }
891

892
    Builder setTicker(Ticker ticker) {
893
      this.ticker = checkNotNull(ticker, "ticker");
1 ✔
894
      return this;
1 ✔
895
    }
896

897
    Builder setEvictionListener(
898
        @Nullable EvictionListener<RouteLookupRequestKey, CacheEntry> evictionListener) {
899
      this.evictionListener = evictionListener;
1 ✔
900
      return this;
1 ✔
901
    }
902

903
    Builder setBackoffProvider(BackoffPolicy.Provider provider) {
904
      this.backoffProvider = checkNotNull(provider, "provider");
1 ✔
905
      return this;
1 ✔
906
    }
907

908
    CachingRlsLbClient build() {
909
      CachingRlsLbClient client = new CachingRlsLbClient(this);
1 ✔
910
      client.init();
1 ✔
911
      return client;
1 ✔
912
    }
913
  }
914

915
  /**
916
   * When any {@link CacheEntry} is evicted from {@link LruCache}, it performs {@link
917
   * CacheEntry#cleanup()} after original {@link EvictionListener} is finished.
918
   */
919
  private static final class AutoCleaningEvictionListener
920
      implements EvictionListener<RouteLookupRequestKey, CacheEntry> {
921

922
    private final EvictionListener<RouteLookupRequestKey, CacheEntry> delegate;
923

924
    AutoCleaningEvictionListener(
925
        @Nullable EvictionListener<RouteLookupRequestKey, CacheEntry> delegate) {
1 ✔
926
      this.delegate = delegate;
1 ✔
927
    }
1 ✔
928

929
    @Override
930
    public void onEviction(RouteLookupRequestKey key, CacheEntry value, EvictionType cause) {
931
      if (delegate != null) {
1 ✔
932
        delegate.onEviction(key, value, cause);
1 ✔
933
      }
934
      // performs cleanup after delegation
935
      value.cleanup();
1 ✔
936
    }
1 ✔
937
  }
938

939
  /** A Throttler never throttles. */
940
  private static final class HappyThrottler implements Throttler {
941

942
    @Override
943
    public boolean shouldThrottle() {
944
      return false;
×
945
    }
946

947
    @Override
948
    public void registerBackendResponse(boolean throttled) {
949
      // no-op
950
    }
×
951
  }
952

953
  /** Implementation of {@link LinkedHashLruCache} for RLS. */
954
  private static final class RlsAsyncLruCache
955
      extends LinkedHashLruCache<RouteLookupRequestKey, CacheEntry> {
956
    private final RlsLbHelper helper;
957

958
    RlsAsyncLruCache(long maxEstimatedSizeBytes,
959
        @Nullable EvictionListener<RouteLookupRequestKey, CacheEntry> evictionListener,
960
        Ticker ticker, RlsLbHelper helper) {
961
      super(maxEstimatedSizeBytes, evictionListener, ticker);
1 ✔
962
      this.helper = checkNotNull(helper, "helper");
1 ✔
963
    }
1 ✔
964

965
    @Override
966
    protected boolean isExpired(RouteLookupRequestKey key, CacheEntry value, long nowNanos) {
967
      return value.isExpired(nowNanos);
1 ✔
968
    }
969

970
    @Override
971
    protected int estimateSizeOf(RouteLookupRequestKey key, CacheEntry value) {
972
      return value.getSizeBytes();
1 ✔
973
    }
974

975
    @Override
976
    protected boolean shouldInvalidateEldestEntry(
977
        RouteLookupRequestKey eldestKey, CacheEntry eldestValue, long now) {
978
      if (!eldestValue.isOldEnoughToBeEvicted(now)) {
×
979
        return false;
×
980
      }
981

982
      // eldest entry should be evicted if size limit exceeded
983
      return this.estimatedSizeBytes() > this.estimatedMaxSizeBytes();
×
984
    }
985

986
    public CacheEntry cacheAndClean(RouteLookupRequestKey key, CacheEntry value) {
987
      CacheEntry newEntry = cache(key, value);
1 ✔
988

989
      // force cleanup if new entry pushed cache over max size (in bytes)
990
      if (fitToLimit()) {
1 ✔
991
        helper.triggerPendingRpcProcessing();
×
992
      }
993
      return newEntry;
1 ✔
994
    }
995
  }
996

997
  /** A header will be added when RLS server respond with additional header data. */
998
  @VisibleForTesting
999
  static final Metadata.Key<String> RLS_DATA_KEY =
1 ✔
1000
      Metadata.Key.of("X-Google-RLS-Data", Metadata.ASCII_STRING_MARSHALLER);
1 ✔
1001

1002
  final class RlsPicker extends SubchannelPicker {
1003

1004
    private final RlsRequestFactory requestFactory;
1005
    private final String lookupService;
1006

1007
    RlsPicker(RlsRequestFactory requestFactory, String lookupService) {
1 ✔
1008
      this.requestFactory = checkNotNull(requestFactory, "requestFactory");
1 ✔
1009
      this.lookupService = checkNotNull(lookupService, "rlsConfig");
1 ✔
1010
    }
1 ✔
1011

1012
    @Override
1013
    public PickResult pickSubchannel(PickSubchannelArgs args) {
1014
      String serviceName = args.getMethodDescriptor().getServiceName();
1 ✔
1015
      String methodName = args.getMethodDescriptor().getBareMethodName();
1 ✔
1016
      RlsProtoData.RouteLookupRequestKey lookupRequestKey =
1 ✔
1017
          requestFactory.create(serviceName, methodName, args.getHeaders());
1 ✔
1018
      final CachedRouteLookupResponse response = CachingRlsLbClient.this.get(lookupRequestKey);
1 ✔
1019

1020
      if (response.getHeaderData() != null && !response.getHeaderData().isEmpty()) {
1 ✔
1021
        Metadata headers = args.getHeaders();
1 ✔
1022
        headers.discardAll(RLS_DATA_KEY);
1 ✔
1023
        headers.put(RLS_DATA_KEY, response.getHeaderData());
1 ✔
1024
      }
1025
      String defaultTarget = lbPolicyConfig.getRouteLookupConfig().defaultTarget();
1 ✔
1026
      boolean hasFallback = defaultTarget != null && !defaultTarget.isEmpty();
1 ✔
1027
      if (response.hasData()) {
1 ✔
1028
        ChildPolicyWrapper childPolicyWrapper = response.getChildPolicyWrapper();
1 ✔
1029
        SubchannelPicker picker =
1030
            (childPolicyWrapper != null) ? childPolicyWrapper.getPicker() : null;
1 ✔
1031
        if (picker == null) {
1 ✔
1032
          // Child policy is connecting. Preserve leaf delay type.
1033
          return PickResult.withNoResult(
×
1034
              "connecting", "RLS child policy connecting");
1035
        }
1036
        // Happy path
1037
        PickResult pickResult = picker.pickSubchannel(args);
1 ✔
1038
        if (pickResult.hasResult()) {
1 ✔
1039
          helper.getMetricRecorder().addLongCounter(TARGET_PICKS_COUNTER, 1,
1 ✔
1040
              Arrays.asList(helper.getChannelTarget(), lookupService,
1 ✔
1041
                  childPolicyWrapper.getTarget(), determineMetricsPickResult(pickResult)),
1 ✔
1042
              Arrays.asList(determineCustomLabel(args)));
1 ✔
1043
        } else if (pickResult.getDelayType() != null) {
1 ✔
1044
          return PickResult.withNoResult(
1 ✔
1045
              pickResult.getDelayType(),
1 ✔
1046
              "RLS child (" + childPolicyWrapper.getTarget() + ") delayed: "
1 ✔
1047
                  + pickResult.getDelayReason());
1 ✔
1048
        }
1049
        return pickResult;
1 ✔
1050
      } else if (response.hasError()) {
1 ✔
1051
        if (hasFallback) {
1 ✔
1052
          return useFallback(args);
1 ✔
1053
        }
1054
        helper.getMetricRecorder().addLongCounter(FAILED_PICKS_COUNTER, 1,
1 ✔
1055
            Arrays.asList(helper.getChannelTarget(), lookupService),
1 ✔
1056
            Arrays.asList(determineCustomLabel(args)));
1 ✔
1057
        return PickResult.withError(
1 ✔
1058
            convertRlsServerStatus(response.getStatus(),
1 ✔
1059
                lbPolicyConfig.getRouteLookupConfig().lookupService()));
1 ✔
1060
      } else {
1061
        // RLS control-plane query is pending.
1062
        return PickResult.withNoResult(
1 ✔
1063
            "rls_lookup_pending",
1064
            "Route Lookup Service query pending on " + lookupService);
1065
      }
1066
    }
1067

1068
    /** Uses Subchannel connected to default target. */
1069
    private PickResult useFallback(PickSubchannelArgs args) {
1070
      SubchannelPicker picker = fallbackChildPolicyWrapper.getPicker();
1 ✔
1071
      if (picker == null) {
1 ✔
1072
        return PickResult.withNoResult(
×
1073
            "connecting", "RLS fallback child policy connecting");
1074
      }
1075
      PickResult pickResult = picker.pickSubchannel(args);
1 ✔
1076
      if (pickResult.hasResult()) {
1 ✔
1077
        helper.getMetricRecorder().addLongCounter(DEFAULT_TARGET_PICKS_COUNTER, 1,
1 ✔
1078
            Arrays.asList(helper.getChannelTarget(), lookupService,
1 ✔
1079
                fallbackChildPolicyWrapper.getTarget(), determineMetricsPickResult(pickResult)),
1 ✔
1080
            Arrays.asList(determineCustomLabel(args)));
1 ✔
1081
      } else if (pickResult.getDelayType() != null) {
1 ✔
1082
        return PickResult.withNoResult(
1 ✔
1083
            pickResult.getDelayType(),
1 ✔
1084
            "RLS fallback (" + fallbackChildPolicyWrapper.getTarget() + ") delayed: "
1 ✔
1085
                + pickResult.getDelayReason());
1 ✔
1086
      }
1087
      return pickResult;
1 ✔
1088
    }
1089

1090
    private String determineMetricsPickResult(PickResult pickResult) {
1091
      if (pickResult.getStatus().isOk()) {
1 ✔
1092
        return "complete";
1 ✔
1093
      } else if (pickResult.isDrop()) {
1 ✔
1094
        return "drop";
×
1095
      } else {
1096
        return "fail";
1 ✔
1097
      }
1098
    }
1099

1100
    private String determineCustomLabel(PickSubchannelArgs args) {
1101
      return args.getCallOptions().getOption(Grpc.CALL_OPTION_CUSTOM_LABEL);
1 ✔
1102
    }
1103

1104
    // GuardedBy CachingRlsLbClient.lock
1105
    void close() {
1106
      synchronized (lock) { // Lock is already held, but ErrorProne can't tell
1 ✔
1107
        if (fallbackChildPolicyWrapper != null) {
1 ✔
1108
          refCountedChildPolicyWrapperFactory.release(fallbackChildPolicyWrapper);
1 ✔
1109
        }
1110
      }
1 ✔
1111
    }
1 ✔
1112

1113
    @Override
1114
    public String toString() {
1115
      return MoreObjects.toStringHelper(this)
×
1116
          .add("target", lbPolicyConfig.getRouteLookupConfig().lookupService())
×
1117
          .toString();
×
1118
    }
1119
  }
1120

1121
}
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