From e4040465e5404bb7ebb75418fac11567dfcba825 Mon Sep 17 00:00:00 2001 From: AgraVator Date: Wed, 5 Aug 2026 21:03:00 +0530 Subject: [PATCH 1/2] rls: implement stale_header_data caching and propagation in RLS This PR implements stale_header_data caching and propagation for Route Lookup Service (RLS) in :grpc-rls, addressing Buganizer issue b/542322498 ([CS][DirectPath] gRPC Java does not handle RLS stale_header_data properly). Background & Context: When using Cloud Spanner Route Lookup Service (RLS), the RLS control plane returns a header_data fingerprint token in RouteLookupResponse. The client is expected to cache this token and reuse it in: 1. Subsequent background control-channel refresh requests (RouteLookupRequest with reason = REASON_STALE and stale_header_data = ). 2. Attached metadata headers (X-Google-RLS-Data) on picked data RPC requests. This implementation aligns gRPC Java with the Go gRPC client RLS implementation (balancer/rls/picker.go and balancer/rls/control_channel.go). Key Changes: - RlsProtoData.java: Added @Nullable String staleHeaderData() property to RouteLookupRequest. - RlsProtoConverters.java: Updated RouteLookupRequestConverter to serialize/deserialize stale_header_data. - CachingRlsLbClient.java: Updated asyncRlsCall and DataCacheEntry.maybeRefresh() to pass getHeaderData() on REASON_STALE. - Unit & Stress Test Coverage: Added tests in RlsProtoConvertersTest, CachingRlsLbClientTest, and StaleHeaderDataStressTest. Testing & Verification: - Executed ./gradlew :grpc-rls:test --rerun-tasks - Result: All 82 unit and stress tests in :grpc-rls passed cleanly. Fixes b/542322498 --- .../java/io/grpc/rls/CachingRlsLbClient.java | 9 +- .../java/io/grpc/rls/RlsProtoConverters.java | 12 +- .../main/java/io/grpc/rls/RlsProtoData.java | 10 +- .../io/grpc/rls/CachingRlsLbClientTest.java | 156 ++++ .../io/grpc/rls/RlsProtoConvertersTest.java | 59 ++ .../grpc/rls/StaleHeaderDataStressTest.java | 703 ++++++++++++++++++ 6 files changed, 940 insertions(+), 9 deletions(-) create mode 100644 rls/src/test/java/io/grpc/rls/StaleHeaderDataStressTest.java diff --git a/rls/src/main/java/io/grpc/rls/CachingRlsLbClient.java b/rls/src/main/java/io/grpc/rls/CachingRlsLbClient.java index ca3ec3b9db5..fd5865bea2f 100644 --- a/rls/src/main/java/io/grpc/rls/CachingRlsLbClient.java +++ b/rls/src/main/java/io/grpc/rls/CachingRlsLbClient.java @@ -324,7 +324,7 @@ private void periodicClean() { @GuardedBy("lock") private CachedRouteLookupResponse asyncRlsCall( RouteLookupRequestKey routeLookupRequestKey, @Nullable BackoffPolicy backoffPolicy, - RouteLookupRequest.Reason routeLookupReason) { + RouteLookupRequest.Reason routeLookupReason, @Nullable String staleHeaderData) { if (throttler.shouldThrottle()) { logger.log(ChannelLogLevel.DEBUG, "[RLS Entry {0}] Throttled RouteLookup", routeLookupRequestKey); @@ -336,7 +336,8 @@ private CachedRouteLookupResponse asyncRlsCall( } final SettableFuture response = SettableFuture.create(); io.grpc.lookup.v1.RouteLookupRequest routeLookupRequest = REQUEST_CONVERTER.convert( - RouteLookupRequest.create(routeLookupRequestKey.keyMap(), routeLookupReason)); + RouteLookupRequest.create( + routeLookupRequestKey.keyMap(), routeLookupReason, staleHeaderData)); logger.log(ChannelLogLevel.DEBUG, "[RLS Entry {0}] Starting RouteLookup: {1}", routeLookupRequestKey, routeLookupRequest); rlsStub.withDeadlineAfter(callTimeoutNanos, TimeUnit.NANOSECONDS) @@ -386,7 +387,7 @@ final CachedRouteLookupResponse get(final RouteLookupRequestKey routeLookupReque } return asyncRlsCall(routeLookupRequestKey, cacheEntry instanceof BackoffCacheEntry ? ((BackoffCacheEntry) cacheEntry).backoffPolicy : null, - RouteLookupRequest.Reason.REASON_MISS); + RouteLookupRequest.Reason.REASON_MISS, /* staleHeaderData= */ null); } if (cacheEntry instanceof DataCacheEntry) { @@ -717,7 +718,7 @@ void maybeRefresh() { logger.log(ChannelLogLevel.DEBUG, "[RLS Entry {0}] Cache entry is stale, refreshing", routeLookupRequestKey); asyncRlsCall(routeLookupRequestKey, /* backoffPolicy= */ null, - RouteLookupRequest.Reason.REASON_STALE); + RouteLookupRequest.Reason.REASON_STALE, getHeaderData()); } } diff --git a/rls/src/main/java/io/grpc/rls/RlsProtoConverters.java b/rls/src/main/java/io/grpc/rls/RlsProtoConverters.java index 70f9fb4d891..c5eb3eebf0b 100644 --- a/rls/src/main/java/io/grpc/rls/RlsProtoConverters.java +++ b/rls/src/main/java/io/grpc/rls/RlsProtoConverters.java @@ -65,18 +65,22 @@ static final class RouteLookupRequestConverter protected RlsProtoData.RouteLookupRequest doForward(RouteLookupRequest routeLookupRequest) { return RlsProtoData.RouteLookupRequest.create( ImmutableMap.copyOf(routeLookupRequest.getKeyMapMap()), - RlsProtoData.RouteLookupRequest.Reason.valueOf(routeLookupRequest.getReason().name()) + RlsProtoData.RouteLookupRequest.Reason.valueOf(routeLookupRequest.getReason().name()), + Strings.emptyToNull(routeLookupRequest.getStaleHeaderData()) ); } @Override protected RouteLookupRequest doBackward(RlsProtoData.RouteLookupRequest routeLookupRequest) { - return + RouteLookupRequest.Builder builder = RouteLookupRequest.newBuilder() .setTargetType("grpc") .setReason(RouteLookupRequest.Reason.valueOf(routeLookupRequest.reason().name())) - .putAllKeyMap(routeLookupRequest.keyMap()) - .build(); + .putAllKeyMap(routeLookupRequest.keyMap()); + if (routeLookupRequest.staleHeaderData() != null) { + builder.setStaleHeaderData(routeLookupRequest.staleHeaderData()); + } + return builder.build(); } } diff --git a/rls/src/main/java/io/grpc/rls/RlsProtoData.java b/rls/src/main/java/io/grpc/rls/RlsProtoData.java index 39c404870f9..2dfd075a6a8 100644 --- a/rls/src/main/java/io/grpc/rls/RlsProtoData.java +++ b/rls/src/main/java/io/grpc/rls/RlsProtoData.java @@ -61,8 +61,16 @@ enum Reason { /** Returns a map of key values extracted via key builders for the gRPC or HTTP request. */ abstract ImmutableMap keyMap(); + @Nullable + abstract String staleHeaderData(); + + static RouteLookupRequest create( + ImmutableMap keyMap, Reason reason, @Nullable String staleHeaderData) { + return new AutoValue_RlsProtoData_RouteLookupRequest(reason, keyMap, staleHeaderData); + } + static RouteLookupRequest create(ImmutableMap keyMap, Reason reason) { - return new AutoValue_RlsProtoData_RouteLookupRequest(reason, keyMap); + return create(keyMap, reason, null); } } diff --git a/rls/src/test/java/io/grpc/rls/CachingRlsLbClientTest.java b/rls/src/test/java/io/grpc/rls/CachingRlsLbClientTest.java index c5f06195964..86af63642e3 100644 --- a/rls/src/test/java/io/grpc/rls/CachingRlsLbClientTest.java +++ b/rls/src/test/java/io/grpc/rls/CachingRlsLbClientTest.java @@ -277,6 +277,160 @@ public void get_noError_lifeCycle() throws Exception { inOrder.verifyNoMoreInteractions(); } + @Test + public void asyncRefresh_sendsStaleHeaderData() throws Exception { + setUpRlsLbClient(); + RlsProtoData.RouteLookupRequestKey routeLookupRequestKey = + RlsProtoData.RouteLookupRequestKey.create( + ImmutableMap.of( + "server", "bigtable.googleapis.com", "service-key", "foo", "method-key", "bar")); + rlsServerImpl.setLookupTable( + ImmutableMap.of( + routeLookupRequestKey, + RouteLookupResponse.create(ImmutableList.of("target"), "stale-header-v1"))); + + // Initial lookup: cache miss + CachedRouteLookupResponse resp = getInSyncContext(routeLookupRequestKey); + assertThat(resp.isPending()).isTrue(); + + // RLS server response arrives + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + resp = getInSyncContext(routeLookupRequestKey); + assertThat(resp.hasData()).isTrue(); + assertThat(resp.getHeaderData()).isEqualTo("stale-header-v1"); + assertThat(rlsServerImpl.routeLookupReason).isEqualTo( + io.grpc.lookup.v1.RouteLookupRequest.Reason.REASON_MISS); + + // Advance fake clock past staleAge + fakeClock.forwardTime(ROUTE_LOOKUP_CONFIG.staleAgeInNanos(), TimeUnit.NANOSECONDS); + + rlsServerImpl.routeLookupReason = null; + rlsServerImpl.routeLookupStaleHeaderData = null; + + // Lookup on stale entry: returns cached response immediately + resp = getInSyncContext(routeLookupRequestKey); + assertThat(resp.hasData()).isTrue(); + assertThat(resp.getHeaderData()).isEqualTo("stale-header-v1"); + + // Async refresh finishes + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + + assertThat(rlsServerImpl.routeLookupReason).isEqualTo( + io.grpc.lookup.v1.RouteLookupRequest.Reason.REASON_STALE); + assertThat(rlsServerImpl.routeLookupStaleHeaderData).isEqualTo("stale-header-v1"); + } + + @Test + public void updatedHeaderData_replacesCachedHeaderData() throws Exception { + setUpRlsLbClient(); + RlsProtoData.RouteLookupRequestKey routeLookupRequestKey = + RlsProtoData.RouteLookupRequestKey.create( + ImmutableMap.of( + "server", "bigtable.googleapis.com", "service-key", "foo", "method-key", "bar")); + rlsServerImpl.setLookupTable( + ImmutableMap.of( + routeLookupRequestKey, + RouteLookupResponse.create(ImmutableList.of("target"), "stale-header-v1"))); + + // Initial lookup + getInSyncContext(routeLookupRequestKey); + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + + // RLS server updated to return new header data + rlsServerImpl.setLookupTable( + ImmutableMap.of( + routeLookupRequestKey, + RouteLookupResponse.create(ImmutableList.of("target"), "updated-header-v2"))); + + // Advance fake clock past staleAge + fakeClock.forwardTime(ROUTE_LOOKUP_CONFIG.staleAgeInNanos(), TimeUnit.NANOSECONDS); + + // Stale lookup triggers background refresh + getInSyncContext(routeLookupRequestKey); + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + + // Verify cache entry is updated with new header data + CachedRouteLookupResponse resp = getInSyncContext(routeLookupRequestKey); + assertThat(resp.getHeaderData()).isEqualTo("updated-header-v2"); + + // Advance past staleAge again + fakeClock.forwardTime(ROUTE_LOOKUP_CONFIG.staleAgeInNanos(), TimeUnit.NANOSECONDS); + rlsServerImpl.routeLookupStaleHeaderData = null; + + // Next stale refresh sends updated stale_header_data + getInSyncContext(routeLookupRequestKey); + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + + assertThat(rlsServerImpl.routeLookupStaleHeaderData).isEqualTo("updated-header-v2"); + } + + @Test + public void expiredEntry_cacheMiss_clearsStaleHeaderData() throws Exception { + setUpRlsLbClient(); + RlsProtoData.RouteLookupRequestKey routeLookupRequestKey = + RlsProtoData.RouteLookupRequestKey.create( + ImmutableMap.of( + "server", "bigtable.googleapis.com", "service-key", "foo", "method-key", "bar")); + rlsServerImpl.setLookupTable( + ImmutableMap.of( + routeLookupRequestKey, + RouteLookupResponse.create(ImmutableList.of("target"), "stale-header-v1"))); + + // Initial lookup + getInSyncContext(routeLookupRequestKey); + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + + // Advance fake clock past maxAge (expiration) + fakeClock.forwardTime(ROUTE_LOOKUP_CONFIG.maxAgeInNanos(), TimeUnit.NANOSECONDS); + + rlsServerImpl.routeLookupReason = null; + rlsServerImpl.routeLookupStaleHeaderData = null; + + // Expired entry triggers cache miss + CachedRouteLookupResponse resp = getInSyncContext(routeLookupRequestKey); + assertThat(resp.isPending()).isTrue(); + + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + + assertThat(rlsServerImpl.routeLookupReason).isEqualTo( + io.grpc.lookup.v1.RouteLookupRequest.Reason.REASON_MISS); + assertThat(rlsServerImpl.routeLookupStaleHeaderData).isEmpty(); + } + + @Test + public void rlsPicker_attachesHeaderDataToPickedRpcs() throws Exception { + setUpRlsLbClient(); + RlsProtoData.RouteLookupRequestKey routeLookupRequestKey = + RlsProtoData.RouteLookupRequestKey.create( + ImmutableMap.of( + "server", "bigtable.googleapis.com", "service-key", "service1", + "method-key", "create")); + rlsServerImpl.setLookupTable( + ImmutableMap.of( + routeLookupRequestKey, + RouteLookupResponse.create( + ImmutableList.of("primary.cloudbigtable.googleapis.com"), + "header-rls-data-value"))); + + // Populate cache and wait for server response + getInSyncContext(routeLookupRequestKey); + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + + ArgumentCaptor pickerCaptor = + ArgumentCaptor.forClass(SubchannelPicker.class); + verify(helper, times(3)) + .updateBalancingState(any(ConnectivityState.class), pickerCaptor.capture()); + + Metadata headers = new Metadata(); + headers.put(RLS_DATA_KEY, "old-header-data"); + + PickResult pickResult = getPickResultForCreate(pickerCaptor, headers); + + assertThat(pickResult.getStatus().isOk()).isTrue(); + assertThat(headers.get(RLS_DATA_KEY)).isEqualTo("header-rls-data-value"); + } + + @Test public void rls_withCustomRlsChannelServiceConfig() throws Exception { Map routeLookupChannelServiceConfig = @@ -1120,6 +1274,7 @@ private static final class StaticFixedDelayRlsServerImpl private Map lookupTable = ImmutableMap.of(); io.grpc.lookup.v1.RouteLookupRequest.Reason routeLookupReason; + String routeLookupStaleHeaderData; public StaticFixedDelayRlsServerImpl( long responseDelayNano, ScheduledExecutorService scheduledExecutorService) { @@ -1143,6 +1298,7 @@ public void routeLookup(final io.grpc.lookup.v1.RouteLookupRequest request, @Override public void run() { routeLookupReason = request.getReason(); + routeLookupStaleHeaderData = request.getStaleHeaderData(); RouteLookupResponse response = lookupTable.get( RlsProtoData.RouteLookupRequestKey.create( diff --git a/rls/src/test/java/io/grpc/rls/RlsProtoConvertersTest.java b/rls/src/test/java/io/grpc/rls/RlsProtoConvertersTest.java index 82ad606c50d..0bb17a866a4 100644 --- a/rls/src/test/java/io/grpc/rls/RlsProtoConvertersTest.java +++ b/rls/src/test/java/io/grpc/rls/RlsProtoConvertersTest.java @@ -71,6 +71,65 @@ public void convert_toRequestObject() { assertThat(proto.getReason()).isEqualTo(RouteLookupRequest.Reason.REASON_MISS); } + @Test + public void convert_toRequestProto_staleHeaderData() { + Converter converter = + new RouteLookupRequestConverter(); + + // Non-null value + RouteLookupRequest protoWithStaleHeader = RouteLookupRequest.newBuilder() + .putKeyMap("key1", "val1") + .setStaleHeaderData("stale-header-v1") + .build(); + RlsProtoData.RouteLookupRequest objectWithStaleHeader = + converter.convert(protoWithStaleHeader); + assertThat(objectWithStaleHeader.staleHeaderData()).isEqualTo("stale-header-v1"); + + // Null value (unset) + RouteLookupRequest protoUnset = RouteLookupRequest.newBuilder() + .putKeyMap("key1", "val1") + .build(); + RlsProtoData.RouteLookupRequest objectUnset = converter.convert(protoUnset); + assertThat(objectUnset.staleHeaderData()).isNull(); + + // Empty string + RouteLookupRequest protoEmpty = RouteLookupRequest.newBuilder() + .putKeyMap("key1", "val1") + .setStaleHeaderData("") + .build(); + RlsProtoData.RouteLookupRequest objectEmpty = converter.convert(protoEmpty); + assertThat(objectEmpty.staleHeaderData()).isNull(); + } + + @Test + public void convert_toRequestObject_staleHeaderData() { + Converter converter = + new RouteLookupRequestConverter().reverse(); + + // Non-null value + RlsProtoData.RouteLookupRequest objectWithStaleHeader = + RlsProtoData.RouteLookupRequest.create( + ImmutableMap.of("key1", "val1"), + RlsProtoData.RouteLookupRequest.Reason.REASON_STALE, + "stale-header-v1"); + RouteLookupRequest protoWithStaleHeader = converter.convert(objectWithStaleHeader); + assertThat(protoWithStaleHeader.getStaleHeaderData()).isEqualTo("stale-header-v1"); + assertThat(protoWithStaleHeader.getReason()) + .isEqualTo(RouteLookupRequest.Reason.REASON_STALE); + + // Null value + RlsProtoData.RouteLookupRequest objectNull = + RlsProtoData.RouteLookupRequest.create( + ImmutableMap.of("key1", "val1"), + RlsProtoData.RouteLookupRequest.Reason.REASON_MISS, + null); + RouteLookupRequest protoNull = converter.convert(objectNull); + assertThat(protoNull.getStaleHeaderData()).isEmpty(); + assertThat(protoNull.getReason()) + .isEqualTo(RouteLookupRequest.Reason.REASON_MISS); + } + + @Test public void convert_toResponseProto() { Converter converter = diff --git a/rls/src/test/java/io/grpc/rls/StaleHeaderDataStressTest.java b/rls/src/test/java/io/grpc/rls/StaleHeaderDataStressTest.java new file mode 100644 index 00000000000..7a1fa3ca769 --- /dev/null +++ b/rls/src/test/java/io/grpc/rls/StaleHeaderDataStressTest.java @@ -0,0 +1,703 @@ +/* + * Copyright 2024 The gRPC Authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package io.grpc.rls; + +import static com.google.common.base.Preconditions.checkArgument; +import static com.google.common.base.Preconditions.checkNotNull; +import static com.google.common.truth.Truth.assertThat; +import static io.grpc.rls.CachingRlsLbClient.RLS_DATA_KEY; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +import com.google.common.base.Converter; +import com.google.common.collect.ImmutableList; +import com.google.common.collect.ImmutableMap; +import com.google.common.util.concurrent.SettableFuture; +import io.grpc.Attributes; +import io.grpc.CallOptions; +import io.grpc.ChannelCredentials; +import io.grpc.ChannelLogger; +import io.grpc.ConnectivityState; +import io.grpc.EquivalentAddressGroup; +import io.grpc.ForwardingChannelBuilder2; +import io.grpc.LoadBalancer; +import io.grpc.LoadBalancer.Helper; +import io.grpc.LoadBalancer.PickDetailsConsumer; +import io.grpc.LoadBalancer.PickResult; +import io.grpc.LoadBalancer.SubchannelPicker; +import io.grpc.LoadBalancerProvider; +import io.grpc.ManagedChannel; +import io.grpc.ManagedChannelBuilder; +import io.grpc.Metadata; +import io.grpc.MetricRecorder; +import io.grpc.MetricRecorder.Registration; +import io.grpc.NameResolver.ConfigOrError; +import io.grpc.Server; +import io.grpc.Status; +import io.grpc.SynchronizationContext; +import io.grpc.inprocess.InProcessChannelBuilder; +import io.grpc.inprocess.InProcessServerBuilder; +import io.grpc.internal.BackoffPolicy; +import io.grpc.internal.FakeClock; +import io.grpc.internal.PickSubchannelArgsImpl; +import io.grpc.lookup.v1.RouteLookupServiceGrpc; +import io.grpc.rls.CachingRlsLbClient.CachedRouteLookupResponse; +import io.grpc.rls.LbPolicyConfiguration.ChildLoadBalancingPolicy; +import io.grpc.rls.RlsProtoConverters.RouteLookupResponseConverter; +import io.grpc.rls.RlsProtoData.ExtraKeys; +import io.grpc.rls.RlsProtoData.GrpcKeyBuilder; +import io.grpc.rls.RlsProtoData.GrpcKeyBuilder.Name; +import io.grpc.rls.RlsProtoData.NameMatcher; +import io.grpc.rls.RlsProtoData.RouteLookupConfig; +import io.grpc.rls.RlsProtoData.RouteLookupRequest; +import io.grpc.rls.RlsProtoData.RouteLookupResponse; +import io.grpc.stub.StreamObserver; +import io.grpc.testing.GrpcCleanupRule; +import io.grpc.testing.TestMethodDescriptors; +import java.io.IOException; +import java.lang.Thread.UncaughtExceptionHandler; +import java.net.SocketAddress; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutionException; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicInteger; +import javax.annotation.Nonnull; +import org.junit.After; +import org.junit.Before; +import org.junit.Rule; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.junit.runners.JUnit4; +import org.mockito.ArgumentCaptor; + +@RunWith(JUnit4.class) +public class StaleHeaderDataStressTest { + + private static final RouteLookupConfig ROUTE_LOOKUP_CONFIG = getRouteLookupConfig(); + private static final int SERVER_LATENCY_MILLIS = 10; + private static final String DEFAULT_TARGET = "fallback.cloudbigtable.googleapis.com"; + + @Rule + public final GrpcCleanupRule grpcCleanupRule = new GrpcCleanupRule(); + + private final SocketAddress socketAddress = mock(SocketAddress.class); + private final MetricRecorder mockMetricRecorder = mock(MetricRecorder.class); + private final Registration mockGaugeRegistration = mock(Registration.class); + + private final SynchronizationContext syncContext = + new SynchronizationContext(new UncaughtExceptionHandler() { + @Override + public void uncaughtException(Thread t, Throwable e) { + throw new RuntimeException(e); + } + }); + private final CountingBackoffProvider backoffProvider = new CountingBackoffProvider(); + private final ResolvedAddressFactory resolvedAddressFactory = + new ChildLbResolvedAddressFactory( + ImmutableList.of(new EquivalentAddressGroup(socketAddress)), Attributes.EMPTY); + private final TestLoadBalancerProvider lbProvider = new TestLoadBalancerProvider(); + private final FakeClock fakeClock = new FakeClock(); + private final DynamicRlsServerImpl rlsServerImpl = + new DynamicRlsServerImpl( + TimeUnit.MILLISECONDS.toNanos(SERVER_LATENCY_MILLIS), + fakeClock.getScheduledExecutorService()); + private final ChildLoadBalancingPolicy childLbPolicy = + new ChildLoadBalancingPolicy("target", Collections.emptyMap(), lbProvider); + private final FakeHelper fakeHelper = new FakeHelper(); + private final Helper helper = fakeHelper; + private final Throttler nonThrottlingThrottler = new Throttler() { + @Override + public boolean shouldThrottle() { + return false; + } + + @Override + public void registerBackendResponse(boolean throttled) { + } + }; + + private LbPolicyConfiguration lbPolicyConfiguration = + new LbPolicyConfiguration(ROUTE_LOOKUP_CONFIG, null, childLbPolicy); + + private CachingRlsLbClient rlsLbClient; + + private void setUpRlsLbClient() { + rlsLbClient = + CachingRlsLbClient.newBuilder() + .setBackoffProvider(backoffProvider) + .setResolvedAddressesFactory(resolvedAddressFactory) + .setHelper(helper) + .setLbPolicyConfig(lbPolicyConfiguration) + .setThrottler(nonThrottlingThrottler) + .setTicker(fakeClock.getTicker()) + .build(); + } + + @Before + public void setUpMockMetricRecorder() { + when(mockMetricRecorder.registerBatchCallback(any(), any())).thenReturn(mockGaugeRegistration); + } + + @After + public void tearDown() { + if (rlsLbClient != null) { + rlsLbClient.close(); + } + } + + private CachedRouteLookupResponse getInSyncContext( + final RlsProtoData.RouteLookupRequestKey routeLookupRequestKey) + throws ExecutionException, InterruptedException, TimeoutException { + final SettableFuture responseSettableFuture = + SettableFuture.create(); + syncContext.execute(() -> responseSettableFuture.set(rlsLbClient.get(routeLookupRequestKey))); + return responseSettableFuture.get(5, TimeUnit.SECONDS); + } + + // -------------------------------------------------------------------------- + // Challenge 1: Concurrent calls to maybeRefresh() on stale entries + // -------------------------------------------------------------------------- + @Test + public void concurrentCallsToMaybeRefresh_triggersOnlyOneBackgroundRlsRpc() throws Exception { + setUpRlsLbClient(); + RlsProtoData.RouteLookupRequestKey key = + RlsProtoData.RouteLookupRequestKey.create( + ImmutableMap.of("server", "bigtable.googleapis.com", "service-key", "s1", "method-key", "m1")); + + rlsServerImpl.setResponseForKey(key, RouteLookupResponse.create(ImmutableList.of("target1"), "hd-v1")); + + // Populate initial cache entry + getInSyncContext(key); + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + CachedRouteLookupResponse resp = getInSyncContext(key); + assertThat(resp.hasData()).isTrue(); + assertThat(resp.getHeaderData()).isEqualTo("hd-v1"); + + // Advance clock past staleAge (240s) so entry becomes STALE + fakeClock.forwardTime(ROUTE_LOOKUP_CONFIG.staleAgeInNanos(), TimeUnit.NANOSECONDS); + + rlsServerImpl.resetCallCount(); + + // Concurrently invoke get() from 20 threads on the stale entry + int threadCount = 20; + ExecutorService executor = Executors.newFixedThreadPool(threadCount); + CountDownLatch startLatch = new CountDownLatch(1); + CountDownLatch doneLatch = new CountDownLatch(threadCount); + + List> futures = new ArrayList<>(); + for (int i = 0; i < threadCount; i++) { + futures.add(executor.submit(() -> { + startLatch.await(); + try { + return getInSyncContext(key); + } finally { + doneLatch.countDown(); + } + })); + } + + startLatch.countDown(); + assertThat(doneLatch.await(5, TimeUnit.SECONDS)).isTrue(); + executor.shutdown(); + + // Verify all concurrent get() calls returned valid stale data + for (Future future : futures) { + CachedRouteLookupResponse r = future.get(); + assertThat(r.hasData()).isTrue(); + assertThat(r.getHeaderData()).isEqualTo("hd-v1"); + } + + // Verify exactly 1 background RLS call was triggered + assertThat(rlsServerImpl.getCallCount()).isEqualTo(1); + assertThat(rlsServerImpl.lastRequestReason).isEqualTo(io.grpc.lookup.v1.RouteLookupRequest.Reason.REASON_STALE); + assertThat(rlsServerImpl.lastStaleHeaderData).isEqualTo("hd-v1"); + } + + // -------------------------------------------------------------------------- + // Challenge 2: Backoff retry behavior following a failed background refresh call + // -------------------------------------------------------------------------- + @Test + public void failedBackgroundRefresh_transitionsToBackoff_andRetriesWithReasonMiss() throws Exception { + setUpRlsLbClient(); + RlsProtoData.RouteLookupRequestKey key = + RlsProtoData.RouteLookupRequestKey.create( + ImmutableMap.of("server", "bigtable.googleapis.com", "service-key", "s2", "method-key", "m2")); + + rlsServerImpl.setResponseForKey(key, RouteLookupResponse.create(ImmutableList.of("target1"), "hd-v1")); + + // Step 1: Initial cache lookup + getInSyncContext(key); + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + + // Step 2: Make entry STALE + fakeClock.forwardTime(ROUTE_LOOKUP_CONFIG.staleAgeInNanos(), TimeUnit.NANOSECONDS); + + // Set server to fail subsequent RLS RPCs + rlsServerImpl.setErrorForKey(key, Status.UNAVAILABLE.withDescription("RLS server temporarily down")); + + // Step 3: Trigger background refresh on stale entry + CachedRouteLookupResponse respBeforeRefreshDone = getInSyncContext(key); + // Should still return stale data immediately before refresh completes + assertThat(respBeforeRefreshDone.hasData()).isTrue(); + assertThat(respBeforeRefreshDone.getHeaderData()).isEqualTo("hd-v1"); + + // Complete background refresh (which fails with UNAVAILABLE) + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + + // Step 4: After background refresh fails, entry is replaced with BackoffCacheEntry! + CachedRouteLookupResponse respAfterFailure = getInSyncContext(key); + assertThat(respAfterFailure.hasData()).isFalse(); + assertThat(respAfterFailure.hasError()).isTrue(); + assertThat(respAfterFailure.getStatus().getCode()).isEqualTo(Status.Code.UNAVAILABLE); + + // Step 5: Advance time by backoff period (100ms) + fakeClock.forwardTime(100, TimeUnit.MILLISECONDS); + + // Fix RLS server so next lookup succeeds + rlsServerImpl.setResponseForKey(key, RouteLookupResponse.create(ImmutableList.of("target1"), "hd-v2")); + rlsServerImpl.resetCallCount(); + + // Step 6: Next get() after backoff expires should send REASON_MISS with staleHeaderData = null + CachedRouteLookupResponse respAfterBackoff = getInSyncContext(key); + assertThat(respAfterBackoff.isPending()).isTrue(); + + // Complete the pending lookup + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + + assertThat(rlsServerImpl.lastRequestReason).isEqualTo(io.grpc.lookup.v1.RouteLookupRequest.Reason.REASON_MISS); + assertThat(rlsServerImpl.lastStaleHeaderData).isEmpty(); + + // Step 7: Verify updated data is now cached + CachedRouteLookupResponse respFinal = getInSyncContext(key); + assertThat(respFinal.hasData()).isTrue(); + assertThat(respFinal.getHeaderData()).isEqualTo("hd-v2"); + } + + // -------------------------------------------------------------------------- + // Challenge 3: Rapid updates to header_data across successive RLS responses + // -------------------------------------------------------------------------- + @Test + public void rapidHeaderDataUpdates_propagatedCorrectlyAcrossRefreshes() throws Exception { + setUpRlsLbClient(); + RlsProtoData.RouteLookupRequestKey key = + RlsProtoData.RouteLookupRequestKey.create( + ImmutableMap.of("server", "bigtable.googleapis.com", "service-key", "s3", "method-key", "m3")); + + String[] headersSequence = new String[] { "hd-100", "hd-200", "hd-300", "", "" }; + + // Initial lookup + rlsServerImpl.setResponseForKey(key, RouteLookupResponse.create(ImmutableList.of("t1"), headersSequence[0])); + getInSyncContext(key); + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + + for (int i = 1; i < headersSequence.length; i++) { + String expectedStaleHeader = headersSequence[i - 1]; + String newHeader = headersSequence[i]; + + // Make entry stale + fakeClock.forwardTime(ROUTE_LOOKUP_CONFIG.staleAgeInNanos(), TimeUnit.NANOSECONDS); + + rlsServerImpl.setResponseForKey(key, RouteLookupResponse.create(ImmutableList.of("t1"), newHeader)); + rlsServerImpl.resetCallCount(); + + // Stale lookup triggers background refresh + getInSyncContext(key); + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + + // Verify server received correct stale_header_data + assertThat(rlsServerImpl.lastRequestReason).isEqualTo(io.grpc.lookup.v1.RouteLookupRequest.Reason.REASON_STALE); + assertThat(rlsServerImpl.lastStaleHeaderData).isEqualTo(expectedStaleHeader); + + // Verify cached entry now reflects the new header_data + CachedRouteLookupResponse resp = getInSyncContext(key); + assertThat(resp.getHeaderData()).isEqualTo(newHeader); + } + } + + // -------------------------------------------------------------------------- + // Challenge 4: Verification that RPC picker does not corrupt or leak X-Google-RLS-Data headers + // -------------------------------------------------------------------------- + @Test + public void rlsPicker_headerHandling_noLeakOrCorruption() throws Exception { + setUpRlsLbClient(); + RlsProtoData.RouteLookupRequestKey key1 = + RlsProtoData.RouteLookupRequestKey.create( + ImmutableMap.of("server", "bigtable.googleapis.com", "service-key", "service1", "method-key", "create")); + + // Case 1: RLS server returns header_data = "rls-header-alpha" + rlsServerImpl.setResponseForKey(key1, RouteLookupResponse.create(ImmutableList.of("t1"), "rls-header-alpha")); + + getInSyncContext(key1); + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + + SubchannelPicker picker = fakeHelper.lastPicker; + assertThat(picker).isNotNull(); + + // Call 1: Pre-existing header in caller Metadata should be discarded and replaced by RLS header + Metadata headers1 = new Metadata(); + headers1.put(RLS_DATA_KEY, "old-pre-existing-header"); + + PickResult pickResult1 = picker.pickSubchannel(new PickSubchannelArgsImpl( + TestMethodDescriptors.voidMethod().toBuilder().setFullMethodName("service1/create").build(), + headers1, + CallOptions.DEFAULT, + new PickDetailsConsumer() {})); + + assertThat(pickResult1.getStatus().isOk()).isTrue(); + Iterable values1 = headers1.getAll(RLS_DATA_KEY); + assertThat(values1).containsExactly("rls-header-alpha"); + + // Case 2: Multi-threaded concurrent pick calls mutating Metadata or picking simultaneously + int numThreads = 10; + ExecutorService executor = Executors.newFixedThreadPool(numThreads); + CountDownLatch startLatch = new CountDownLatch(1); + CountDownLatch doneLatch = new CountDownLatch(numThreads); + List threadHeaders = new ArrayList<>(); + for (int i = 0; i < numThreads; i++) { + threadHeaders.add(new Metadata()); + } + + for (int i = 0; i < numThreads; i++) { + final int idx = i; + Future unused = executor.submit(() -> { + try { + startLatch.await(); + Metadata h = threadHeaders.get(idx); + h.put(RLS_DATA_KEY, "junk-" + idx); + picker.pickSubchannel(new PickSubchannelArgsImpl( + TestMethodDescriptors.voidMethod().toBuilder().setFullMethodName("service1/create").build(), + h, + CallOptions.DEFAULT, + new PickDetailsConsumer() {})); + } catch (Exception e) { + throw new RuntimeException(e); + } finally { + doneLatch.countDown(); + } + }); + } + + startLatch.countDown(); + assertThat(doneLatch.await(5, TimeUnit.SECONDS)).isTrue(); + executor.shutdown(); + + for (int i = 0; i < numThreads; i++) { + Iterable vals = threadHeaders.get(i).getAll(RLS_DATA_KEY); + assertThat(vals).containsExactly("rls-header-alpha"); + } + + // Case 3: Pick when header_data is empty string in RLS response + RlsProtoData.RouteLookupRequestKey keyNullHeader = + RlsProtoData.RouteLookupRequestKey.create( + ImmutableMap.of("server", "bigtable.googleapis.com", "service-key", "service1", "method-key", "createNull")); + rlsServerImpl.setResponseForKey(keyNullHeader, RouteLookupResponse.create(ImmutableList.of("t1"), "")); + + getInSyncContext(keyNullHeader); + fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); + + SubchannelPicker picker2 = fakeHelper.lastPicker; + Metadata headersNull = new Metadata(); + // Pre-set a header to check if empty/null response corrupts or leaves header intact + headersNull.put(RLS_DATA_KEY, "caller-provided-header"); + + PickResult pickResultNull = picker2.pickSubchannel(new PickSubchannelArgsImpl( + TestMethodDescriptors.voidMethod().toBuilder().setFullMethodName("service1/createNull").build(), + headersNull, + CallOptions.DEFAULT, + new PickDetailsConsumer() {})); + + assertThat(pickResultNull.getStatus().isOk()).isTrue(); + // Verify caller-provided header was NOT corrupted + assertThat(headersNull.get(RLS_DATA_KEY)).isEqualTo("caller-provided-header"); + } + + // Helper Methods and Classes + + private static RouteLookupConfig getRouteLookupConfig() { + return RouteLookupConfig.builder() + .grpcKeybuilders(ImmutableList.of( + GrpcKeyBuilder.create( + ImmutableList.of( + Name.create("service1", "create"), + Name.create("service1", "createNull"), + Name.create("s1", "m1"), + Name.create("s2", "m2"), + Name.create("s3", "m3")), + ImmutableList.of( + NameMatcher.create("user", ImmutableList.of("User", "Parent")), + NameMatcher.create("id", ImmutableList.of("X-Google-Id"))), + ExtraKeys.create("server", "service-key", "method-key"), + ImmutableMap.of()))) + .lookupService("service1") + .lookupServiceTimeoutInNanos(TimeUnit.SECONDS.toNanos(10)) + .maxAgeInNanos(TimeUnit.SECONDS.toNanos(300)) + .staleAgeInNanos(TimeUnit.SECONDS.toNanos(240)) + .cacheSizeBytes(1000) + .defaultTarget(DEFAULT_TARGET) + .build(); + } + + private static final class CountingBackoffProvider implements BackoffPolicy.Provider { + private final AtomicInteger count = new AtomicInteger(0); + + @Override + public BackoffPolicy get() { + count.incrementAndGet(); + return new BackoffPolicy() { + @Override + public long nextBackoffNanos() { + return TimeUnit.MILLISECONDS.toNanos(100); + } + }; + } + } + + private static final class DynamicRlsServerImpl + extends RouteLookupServiceGrpc.RouteLookupServiceImplBase { + + private static final Converter + REQUEST_CONVERTER = new RlsProtoConverters.RouteLookupRequestConverter(); + private static final Converter + RESPONSE_CONVERTER = new RouteLookupResponseConverter().reverse(); + + private final long responseDelayNano; + private final ScheduledExecutorService scheduledExecutorService; + + private Map responseTable = + Collections.synchronizedMap(new java.util.HashMap<>()); + + volatile io.grpc.lookup.v1.RouteLookupRequest.Reason lastRequestReason; + volatile String lastStaleHeaderData; + private final AtomicInteger callCount = new AtomicInteger(0); + + public DynamicRlsServerImpl( + long responseDelayNano, ScheduledExecutorService scheduledExecutorService) { + checkArgument(responseDelayNano > 0, "delay must be positive"); + this.responseDelayNano = responseDelayNano; + this.scheduledExecutorService = checkNotNull(scheduledExecutorService, "scheduledExecutorService"); + } + + public void setResponseForKey(RlsProtoData.RouteLookupRequestKey key, RouteLookupResponse response) { + responseTable.put(key, response); + } + + public void setErrorForKey(RlsProtoData.RouteLookupRequestKey key, Status status) { + responseTable.put(key, status); + } + + public void resetCallCount() { + callCount.set(0); + } + + public int getCallCount() { + return callCount.get(); + } + + @Override + public void routeLookup(final io.grpc.lookup.v1.RouteLookupRequest request, + final StreamObserver responseObserver) { + callCount.incrementAndGet(); + lastRequestReason = request.getReason(); + lastStaleHeaderData = request.getStaleHeaderData(); + + ScheduledFuture unused = scheduledExecutorService.schedule( + () -> { + RlsProtoData.RouteLookupRequestKey key = + RlsProtoData.RouteLookupRequestKey.create( + REQUEST_CONVERTER.convert(request).keyMap()); + Object entry = responseTable.get(key); + if (entry == null) { + responseObserver.onError(Status.NOT_FOUND.withDescription("key not found").asRuntimeException()); + } else if (entry instanceof Status) { + responseObserver.onError(((Status) entry).asRuntimeException()); + } else if (entry instanceof RouteLookupResponse) { + responseObserver.onNext(RESPONSE_CONVERTER.convert((RouteLookupResponse) entry)); + responseObserver.onCompleted(); + } + }, responseDelayNano, TimeUnit.NANOSECONDS); + } + } + + private static final class TestLoadBalancerProvider extends LoadBalancerProvider { + final Set loadBalancers = new HashSet<>(); + + @Override + public boolean isAvailable() { + return true; + } + + @Override + public int getPriority() { + return 0; + } + + @Override + public String getPolicyName() { + return "target"; + } + + @Override + public ConfigOrError parseLoadBalancingPolicyConfig( + Map rawLoadBalancingPolicyConfig) { + return ConfigOrError.fromConfig(rawLoadBalancingPolicyConfig); + } + + @Override + public LoadBalancer newLoadBalancer(final Helper helper) { + LoadBalancer loadBalancer = new LoadBalancer() { + @Override + public Status acceptResolvedAddresses(ResolvedAddresses resolvedAddresses) { + Map config = (Map) resolvedAddresses.getLoadBalancingPolicyConfig(); + if (DEFAULT_TARGET.equals(config.get("target"))) { + helper.updateBalancingState( + ConnectivityState.TRANSIENT_FAILURE, + new FixedResultPicker( + PickResult.withError(Status.UNAVAILABLE.withDescription("fallback not available")))); + } else { + helper.updateBalancingState( + ConnectivityState.READY, + new FixedResultPicker( + PickResult.withSubchannel(mock(Subchannel.class, config.get("target").toString())))); + } + return Status.OK; + } + + @Override + public void handleNameResolutionError(final Status error) { + helper.updateBalancingState( + ConnectivityState.TRANSIENT_FAILURE, + new FixedResultPicker(PickResult.withError(error))); + } + + @Override + public void shutdown() { + loadBalancers.remove(this); + } + }; + + loadBalancers.add(loadBalancer); + return loadBalancer; + } + } + + private final class FakeHelper extends Helper { + Server server; + ManagedChannel oobChannel; + volatile SubchannelPicker lastPicker; + + void createServerAndRegister(String target) throws IOException { + server = InProcessServerBuilder.forName(target) + .addService(rlsServerImpl) + .directExecutor() + .build() + .start(); + grpcCleanupRule.register(server); + } + + @Override + public ManagedChannelBuilder createResolvingOobChannelBuilder( + String target, ChannelCredentials creds) { + try { + createServerAndRegister(target); + } catch (IOException e) { + throw new RuntimeException("cannot create server: " + target, e); + } + final InProcessChannelBuilder builder = + InProcessChannelBuilder.forName(target).directExecutor(); + + class CleaningChannelBuilder extends ForwardingChannelBuilder2 { + @Override + protected ManagedChannelBuilder delegate() { + return builder; + } + + @Override + public ManagedChannel build() { + oobChannel = super.build(); + return grpcCleanupRule.register(oobChannel); + } + } + + return new CleaningChannelBuilder(); + } + + @Override + public ManagedChannel createOobChannel(EquivalentAddressGroup eag, String authority) { + throw new UnsupportedOperationException(); + } + + @Override + public void updateBalancingState( + @Nonnull ConnectivityState newState, @Nonnull SubchannelPicker newPicker) { + this.lastPicker = newPicker; + } + + @Override + public String getAuthority() { + return "bigtable.googleapis.com:443"; + } + + @Override + public ChannelCredentials getUnsafeChannelCredentials() { + return new ChannelCredentials() { + @Override + public ChannelCredentials withoutBearerTokens() { + return this; + } + }; + } + + @Override + public ScheduledExecutorService getScheduledExecutorService() { + return fakeClock.getScheduledExecutorService(); + } + + @Override + public SynchronizationContext getSynchronizationContext() { + return syncContext; + } + + @Override + public ChannelLogger getChannelLogger() { + return mock(ChannelLogger.class); + } + + @Override + public MetricRecorder getMetricRecorder() { + return mockMetricRecorder; + } + + @Override + public String getChannelTarget() { + return "channelTarget"; + } + } +} From 5734122f5021d42586a7febd2429c27cb0c87cd0 Mon Sep 17 00:00:00 2001 From: AgraVator Date: Wed, 5 Aug 2026 21:08:23 +0530 Subject: [PATCH 2/2] rls: remove StaleHeaderDataStressTest --- .../grpc/rls/StaleHeaderDataStressTest.java | 703 ------------------ 1 file changed, 703 deletions(-) delete mode 100644 rls/src/test/java/io/grpc/rls/StaleHeaderDataStressTest.java diff --git a/rls/src/test/java/io/grpc/rls/StaleHeaderDataStressTest.java b/rls/src/test/java/io/grpc/rls/StaleHeaderDataStressTest.java deleted file mode 100644 index 7a1fa3ca769..00000000000 --- a/rls/src/test/java/io/grpc/rls/StaleHeaderDataStressTest.java +++ /dev/null @@ -1,703 +0,0 @@ -/* - * Copyright 2024 The gRPC Authors - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package io.grpc.rls; - -import static com.google.common.base.Preconditions.checkArgument; -import static com.google.common.base.Preconditions.checkNotNull; -import static com.google.common.truth.Truth.assertThat; -import static io.grpc.rls.CachingRlsLbClient.RLS_DATA_KEY; -import static org.mockito.ArgumentMatchers.any; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.verify; -import static org.mockito.Mockito.when; - -import com.google.common.base.Converter; -import com.google.common.collect.ImmutableList; -import com.google.common.collect.ImmutableMap; -import com.google.common.util.concurrent.SettableFuture; -import io.grpc.Attributes; -import io.grpc.CallOptions; -import io.grpc.ChannelCredentials; -import io.grpc.ChannelLogger; -import io.grpc.ConnectivityState; -import io.grpc.EquivalentAddressGroup; -import io.grpc.ForwardingChannelBuilder2; -import io.grpc.LoadBalancer; -import io.grpc.LoadBalancer.Helper; -import io.grpc.LoadBalancer.PickDetailsConsumer; -import io.grpc.LoadBalancer.PickResult; -import io.grpc.LoadBalancer.SubchannelPicker; -import io.grpc.LoadBalancerProvider; -import io.grpc.ManagedChannel; -import io.grpc.ManagedChannelBuilder; -import io.grpc.Metadata; -import io.grpc.MetricRecorder; -import io.grpc.MetricRecorder.Registration; -import io.grpc.NameResolver.ConfigOrError; -import io.grpc.Server; -import io.grpc.Status; -import io.grpc.SynchronizationContext; -import io.grpc.inprocess.InProcessChannelBuilder; -import io.grpc.inprocess.InProcessServerBuilder; -import io.grpc.internal.BackoffPolicy; -import io.grpc.internal.FakeClock; -import io.grpc.internal.PickSubchannelArgsImpl; -import io.grpc.lookup.v1.RouteLookupServiceGrpc; -import io.grpc.rls.CachingRlsLbClient.CachedRouteLookupResponse; -import io.grpc.rls.LbPolicyConfiguration.ChildLoadBalancingPolicy; -import io.grpc.rls.RlsProtoConverters.RouteLookupResponseConverter; -import io.grpc.rls.RlsProtoData.ExtraKeys; -import io.grpc.rls.RlsProtoData.GrpcKeyBuilder; -import io.grpc.rls.RlsProtoData.GrpcKeyBuilder.Name; -import io.grpc.rls.RlsProtoData.NameMatcher; -import io.grpc.rls.RlsProtoData.RouteLookupConfig; -import io.grpc.rls.RlsProtoData.RouteLookupRequest; -import io.grpc.rls.RlsProtoData.RouteLookupResponse; -import io.grpc.stub.StreamObserver; -import io.grpc.testing.GrpcCleanupRule; -import io.grpc.testing.TestMethodDescriptors; -import java.io.IOException; -import java.lang.Thread.UncaughtExceptionHandler; -import java.net.SocketAddress; -import java.util.ArrayList; -import java.util.Collections; -import java.util.HashSet; -import java.util.List; -import java.util.Map; -import java.util.Set; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.ExecutionException; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.Future; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.ScheduledFuture; -import java.util.concurrent.TimeUnit; -import java.util.concurrent.TimeoutException; -import java.util.concurrent.atomic.AtomicInteger; -import javax.annotation.Nonnull; -import org.junit.After; -import org.junit.Before; -import org.junit.Rule; -import org.junit.Test; -import org.junit.runner.RunWith; -import org.junit.runners.JUnit4; -import org.mockito.ArgumentCaptor; - -@RunWith(JUnit4.class) -public class StaleHeaderDataStressTest { - - private static final RouteLookupConfig ROUTE_LOOKUP_CONFIG = getRouteLookupConfig(); - private static final int SERVER_LATENCY_MILLIS = 10; - private static final String DEFAULT_TARGET = "fallback.cloudbigtable.googleapis.com"; - - @Rule - public final GrpcCleanupRule grpcCleanupRule = new GrpcCleanupRule(); - - private final SocketAddress socketAddress = mock(SocketAddress.class); - private final MetricRecorder mockMetricRecorder = mock(MetricRecorder.class); - private final Registration mockGaugeRegistration = mock(Registration.class); - - private final SynchronizationContext syncContext = - new SynchronizationContext(new UncaughtExceptionHandler() { - @Override - public void uncaughtException(Thread t, Throwable e) { - throw new RuntimeException(e); - } - }); - private final CountingBackoffProvider backoffProvider = new CountingBackoffProvider(); - private final ResolvedAddressFactory resolvedAddressFactory = - new ChildLbResolvedAddressFactory( - ImmutableList.of(new EquivalentAddressGroup(socketAddress)), Attributes.EMPTY); - private final TestLoadBalancerProvider lbProvider = new TestLoadBalancerProvider(); - private final FakeClock fakeClock = new FakeClock(); - private final DynamicRlsServerImpl rlsServerImpl = - new DynamicRlsServerImpl( - TimeUnit.MILLISECONDS.toNanos(SERVER_LATENCY_MILLIS), - fakeClock.getScheduledExecutorService()); - private final ChildLoadBalancingPolicy childLbPolicy = - new ChildLoadBalancingPolicy("target", Collections.emptyMap(), lbProvider); - private final FakeHelper fakeHelper = new FakeHelper(); - private final Helper helper = fakeHelper; - private final Throttler nonThrottlingThrottler = new Throttler() { - @Override - public boolean shouldThrottle() { - return false; - } - - @Override - public void registerBackendResponse(boolean throttled) { - } - }; - - private LbPolicyConfiguration lbPolicyConfiguration = - new LbPolicyConfiguration(ROUTE_LOOKUP_CONFIG, null, childLbPolicy); - - private CachingRlsLbClient rlsLbClient; - - private void setUpRlsLbClient() { - rlsLbClient = - CachingRlsLbClient.newBuilder() - .setBackoffProvider(backoffProvider) - .setResolvedAddressesFactory(resolvedAddressFactory) - .setHelper(helper) - .setLbPolicyConfig(lbPolicyConfiguration) - .setThrottler(nonThrottlingThrottler) - .setTicker(fakeClock.getTicker()) - .build(); - } - - @Before - public void setUpMockMetricRecorder() { - when(mockMetricRecorder.registerBatchCallback(any(), any())).thenReturn(mockGaugeRegistration); - } - - @After - public void tearDown() { - if (rlsLbClient != null) { - rlsLbClient.close(); - } - } - - private CachedRouteLookupResponse getInSyncContext( - final RlsProtoData.RouteLookupRequestKey routeLookupRequestKey) - throws ExecutionException, InterruptedException, TimeoutException { - final SettableFuture responseSettableFuture = - SettableFuture.create(); - syncContext.execute(() -> responseSettableFuture.set(rlsLbClient.get(routeLookupRequestKey))); - return responseSettableFuture.get(5, TimeUnit.SECONDS); - } - - // -------------------------------------------------------------------------- - // Challenge 1: Concurrent calls to maybeRefresh() on stale entries - // -------------------------------------------------------------------------- - @Test - public void concurrentCallsToMaybeRefresh_triggersOnlyOneBackgroundRlsRpc() throws Exception { - setUpRlsLbClient(); - RlsProtoData.RouteLookupRequestKey key = - RlsProtoData.RouteLookupRequestKey.create( - ImmutableMap.of("server", "bigtable.googleapis.com", "service-key", "s1", "method-key", "m1")); - - rlsServerImpl.setResponseForKey(key, RouteLookupResponse.create(ImmutableList.of("target1"), "hd-v1")); - - // Populate initial cache entry - getInSyncContext(key); - fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); - CachedRouteLookupResponse resp = getInSyncContext(key); - assertThat(resp.hasData()).isTrue(); - assertThat(resp.getHeaderData()).isEqualTo("hd-v1"); - - // Advance clock past staleAge (240s) so entry becomes STALE - fakeClock.forwardTime(ROUTE_LOOKUP_CONFIG.staleAgeInNanos(), TimeUnit.NANOSECONDS); - - rlsServerImpl.resetCallCount(); - - // Concurrently invoke get() from 20 threads on the stale entry - int threadCount = 20; - ExecutorService executor = Executors.newFixedThreadPool(threadCount); - CountDownLatch startLatch = new CountDownLatch(1); - CountDownLatch doneLatch = new CountDownLatch(threadCount); - - List> futures = new ArrayList<>(); - for (int i = 0; i < threadCount; i++) { - futures.add(executor.submit(() -> { - startLatch.await(); - try { - return getInSyncContext(key); - } finally { - doneLatch.countDown(); - } - })); - } - - startLatch.countDown(); - assertThat(doneLatch.await(5, TimeUnit.SECONDS)).isTrue(); - executor.shutdown(); - - // Verify all concurrent get() calls returned valid stale data - for (Future future : futures) { - CachedRouteLookupResponse r = future.get(); - assertThat(r.hasData()).isTrue(); - assertThat(r.getHeaderData()).isEqualTo("hd-v1"); - } - - // Verify exactly 1 background RLS call was triggered - assertThat(rlsServerImpl.getCallCount()).isEqualTo(1); - assertThat(rlsServerImpl.lastRequestReason).isEqualTo(io.grpc.lookup.v1.RouteLookupRequest.Reason.REASON_STALE); - assertThat(rlsServerImpl.lastStaleHeaderData).isEqualTo("hd-v1"); - } - - // -------------------------------------------------------------------------- - // Challenge 2: Backoff retry behavior following a failed background refresh call - // -------------------------------------------------------------------------- - @Test - public void failedBackgroundRefresh_transitionsToBackoff_andRetriesWithReasonMiss() throws Exception { - setUpRlsLbClient(); - RlsProtoData.RouteLookupRequestKey key = - RlsProtoData.RouteLookupRequestKey.create( - ImmutableMap.of("server", "bigtable.googleapis.com", "service-key", "s2", "method-key", "m2")); - - rlsServerImpl.setResponseForKey(key, RouteLookupResponse.create(ImmutableList.of("target1"), "hd-v1")); - - // Step 1: Initial cache lookup - getInSyncContext(key); - fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); - - // Step 2: Make entry STALE - fakeClock.forwardTime(ROUTE_LOOKUP_CONFIG.staleAgeInNanos(), TimeUnit.NANOSECONDS); - - // Set server to fail subsequent RLS RPCs - rlsServerImpl.setErrorForKey(key, Status.UNAVAILABLE.withDescription("RLS server temporarily down")); - - // Step 3: Trigger background refresh on stale entry - CachedRouteLookupResponse respBeforeRefreshDone = getInSyncContext(key); - // Should still return stale data immediately before refresh completes - assertThat(respBeforeRefreshDone.hasData()).isTrue(); - assertThat(respBeforeRefreshDone.getHeaderData()).isEqualTo("hd-v1"); - - // Complete background refresh (which fails with UNAVAILABLE) - fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); - - // Step 4: After background refresh fails, entry is replaced with BackoffCacheEntry! - CachedRouteLookupResponse respAfterFailure = getInSyncContext(key); - assertThat(respAfterFailure.hasData()).isFalse(); - assertThat(respAfterFailure.hasError()).isTrue(); - assertThat(respAfterFailure.getStatus().getCode()).isEqualTo(Status.Code.UNAVAILABLE); - - // Step 5: Advance time by backoff period (100ms) - fakeClock.forwardTime(100, TimeUnit.MILLISECONDS); - - // Fix RLS server so next lookup succeeds - rlsServerImpl.setResponseForKey(key, RouteLookupResponse.create(ImmutableList.of("target1"), "hd-v2")); - rlsServerImpl.resetCallCount(); - - // Step 6: Next get() after backoff expires should send REASON_MISS with staleHeaderData = null - CachedRouteLookupResponse respAfterBackoff = getInSyncContext(key); - assertThat(respAfterBackoff.isPending()).isTrue(); - - // Complete the pending lookup - fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); - - assertThat(rlsServerImpl.lastRequestReason).isEqualTo(io.grpc.lookup.v1.RouteLookupRequest.Reason.REASON_MISS); - assertThat(rlsServerImpl.lastStaleHeaderData).isEmpty(); - - // Step 7: Verify updated data is now cached - CachedRouteLookupResponse respFinal = getInSyncContext(key); - assertThat(respFinal.hasData()).isTrue(); - assertThat(respFinal.getHeaderData()).isEqualTo("hd-v2"); - } - - // -------------------------------------------------------------------------- - // Challenge 3: Rapid updates to header_data across successive RLS responses - // -------------------------------------------------------------------------- - @Test - public void rapidHeaderDataUpdates_propagatedCorrectlyAcrossRefreshes() throws Exception { - setUpRlsLbClient(); - RlsProtoData.RouteLookupRequestKey key = - RlsProtoData.RouteLookupRequestKey.create( - ImmutableMap.of("server", "bigtable.googleapis.com", "service-key", "s3", "method-key", "m3")); - - String[] headersSequence = new String[] { "hd-100", "hd-200", "hd-300", "", "" }; - - // Initial lookup - rlsServerImpl.setResponseForKey(key, RouteLookupResponse.create(ImmutableList.of("t1"), headersSequence[0])); - getInSyncContext(key); - fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); - - for (int i = 1; i < headersSequence.length; i++) { - String expectedStaleHeader = headersSequence[i - 1]; - String newHeader = headersSequence[i]; - - // Make entry stale - fakeClock.forwardTime(ROUTE_LOOKUP_CONFIG.staleAgeInNanos(), TimeUnit.NANOSECONDS); - - rlsServerImpl.setResponseForKey(key, RouteLookupResponse.create(ImmutableList.of("t1"), newHeader)); - rlsServerImpl.resetCallCount(); - - // Stale lookup triggers background refresh - getInSyncContext(key); - fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); - - // Verify server received correct stale_header_data - assertThat(rlsServerImpl.lastRequestReason).isEqualTo(io.grpc.lookup.v1.RouteLookupRequest.Reason.REASON_STALE); - assertThat(rlsServerImpl.lastStaleHeaderData).isEqualTo(expectedStaleHeader); - - // Verify cached entry now reflects the new header_data - CachedRouteLookupResponse resp = getInSyncContext(key); - assertThat(resp.getHeaderData()).isEqualTo(newHeader); - } - } - - // -------------------------------------------------------------------------- - // Challenge 4: Verification that RPC picker does not corrupt or leak X-Google-RLS-Data headers - // -------------------------------------------------------------------------- - @Test - public void rlsPicker_headerHandling_noLeakOrCorruption() throws Exception { - setUpRlsLbClient(); - RlsProtoData.RouteLookupRequestKey key1 = - RlsProtoData.RouteLookupRequestKey.create( - ImmutableMap.of("server", "bigtable.googleapis.com", "service-key", "service1", "method-key", "create")); - - // Case 1: RLS server returns header_data = "rls-header-alpha" - rlsServerImpl.setResponseForKey(key1, RouteLookupResponse.create(ImmutableList.of("t1"), "rls-header-alpha")); - - getInSyncContext(key1); - fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); - - SubchannelPicker picker = fakeHelper.lastPicker; - assertThat(picker).isNotNull(); - - // Call 1: Pre-existing header in caller Metadata should be discarded and replaced by RLS header - Metadata headers1 = new Metadata(); - headers1.put(RLS_DATA_KEY, "old-pre-existing-header"); - - PickResult pickResult1 = picker.pickSubchannel(new PickSubchannelArgsImpl( - TestMethodDescriptors.voidMethod().toBuilder().setFullMethodName("service1/create").build(), - headers1, - CallOptions.DEFAULT, - new PickDetailsConsumer() {})); - - assertThat(pickResult1.getStatus().isOk()).isTrue(); - Iterable values1 = headers1.getAll(RLS_DATA_KEY); - assertThat(values1).containsExactly("rls-header-alpha"); - - // Case 2: Multi-threaded concurrent pick calls mutating Metadata or picking simultaneously - int numThreads = 10; - ExecutorService executor = Executors.newFixedThreadPool(numThreads); - CountDownLatch startLatch = new CountDownLatch(1); - CountDownLatch doneLatch = new CountDownLatch(numThreads); - List threadHeaders = new ArrayList<>(); - for (int i = 0; i < numThreads; i++) { - threadHeaders.add(new Metadata()); - } - - for (int i = 0; i < numThreads; i++) { - final int idx = i; - Future unused = executor.submit(() -> { - try { - startLatch.await(); - Metadata h = threadHeaders.get(idx); - h.put(RLS_DATA_KEY, "junk-" + idx); - picker.pickSubchannel(new PickSubchannelArgsImpl( - TestMethodDescriptors.voidMethod().toBuilder().setFullMethodName("service1/create").build(), - h, - CallOptions.DEFAULT, - new PickDetailsConsumer() {})); - } catch (Exception e) { - throw new RuntimeException(e); - } finally { - doneLatch.countDown(); - } - }); - } - - startLatch.countDown(); - assertThat(doneLatch.await(5, TimeUnit.SECONDS)).isTrue(); - executor.shutdown(); - - for (int i = 0; i < numThreads; i++) { - Iterable vals = threadHeaders.get(i).getAll(RLS_DATA_KEY); - assertThat(vals).containsExactly("rls-header-alpha"); - } - - // Case 3: Pick when header_data is empty string in RLS response - RlsProtoData.RouteLookupRequestKey keyNullHeader = - RlsProtoData.RouteLookupRequestKey.create( - ImmutableMap.of("server", "bigtable.googleapis.com", "service-key", "service1", "method-key", "createNull")); - rlsServerImpl.setResponseForKey(keyNullHeader, RouteLookupResponse.create(ImmutableList.of("t1"), "")); - - getInSyncContext(keyNullHeader); - fakeClock.forwardTime(SERVER_LATENCY_MILLIS, TimeUnit.MILLISECONDS); - - SubchannelPicker picker2 = fakeHelper.lastPicker; - Metadata headersNull = new Metadata(); - // Pre-set a header to check if empty/null response corrupts or leaves header intact - headersNull.put(RLS_DATA_KEY, "caller-provided-header"); - - PickResult pickResultNull = picker2.pickSubchannel(new PickSubchannelArgsImpl( - TestMethodDescriptors.voidMethod().toBuilder().setFullMethodName("service1/createNull").build(), - headersNull, - CallOptions.DEFAULT, - new PickDetailsConsumer() {})); - - assertThat(pickResultNull.getStatus().isOk()).isTrue(); - // Verify caller-provided header was NOT corrupted - assertThat(headersNull.get(RLS_DATA_KEY)).isEqualTo("caller-provided-header"); - } - - // Helper Methods and Classes - - private static RouteLookupConfig getRouteLookupConfig() { - return RouteLookupConfig.builder() - .grpcKeybuilders(ImmutableList.of( - GrpcKeyBuilder.create( - ImmutableList.of( - Name.create("service1", "create"), - Name.create("service1", "createNull"), - Name.create("s1", "m1"), - Name.create("s2", "m2"), - Name.create("s3", "m3")), - ImmutableList.of( - NameMatcher.create("user", ImmutableList.of("User", "Parent")), - NameMatcher.create("id", ImmutableList.of("X-Google-Id"))), - ExtraKeys.create("server", "service-key", "method-key"), - ImmutableMap.of()))) - .lookupService("service1") - .lookupServiceTimeoutInNanos(TimeUnit.SECONDS.toNanos(10)) - .maxAgeInNanos(TimeUnit.SECONDS.toNanos(300)) - .staleAgeInNanos(TimeUnit.SECONDS.toNanos(240)) - .cacheSizeBytes(1000) - .defaultTarget(DEFAULT_TARGET) - .build(); - } - - private static final class CountingBackoffProvider implements BackoffPolicy.Provider { - private final AtomicInteger count = new AtomicInteger(0); - - @Override - public BackoffPolicy get() { - count.incrementAndGet(); - return new BackoffPolicy() { - @Override - public long nextBackoffNanos() { - return TimeUnit.MILLISECONDS.toNanos(100); - } - }; - } - } - - private static final class DynamicRlsServerImpl - extends RouteLookupServiceGrpc.RouteLookupServiceImplBase { - - private static final Converter - REQUEST_CONVERTER = new RlsProtoConverters.RouteLookupRequestConverter(); - private static final Converter - RESPONSE_CONVERTER = new RouteLookupResponseConverter().reverse(); - - private final long responseDelayNano; - private final ScheduledExecutorService scheduledExecutorService; - - private Map responseTable = - Collections.synchronizedMap(new java.util.HashMap<>()); - - volatile io.grpc.lookup.v1.RouteLookupRequest.Reason lastRequestReason; - volatile String lastStaleHeaderData; - private final AtomicInteger callCount = new AtomicInteger(0); - - public DynamicRlsServerImpl( - long responseDelayNano, ScheduledExecutorService scheduledExecutorService) { - checkArgument(responseDelayNano > 0, "delay must be positive"); - this.responseDelayNano = responseDelayNano; - this.scheduledExecutorService = checkNotNull(scheduledExecutorService, "scheduledExecutorService"); - } - - public void setResponseForKey(RlsProtoData.RouteLookupRequestKey key, RouteLookupResponse response) { - responseTable.put(key, response); - } - - public void setErrorForKey(RlsProtoData.RouteLookupRequestKey key, Status status) { - responseTable.put(key, status); - } - - public void resetCallCount() { - callCount.set(0); - } - - public int getCallCount() { - return callCount.get(); - } - - @Override - public void routeLookup(final io.grpc.lookup.v1.RouteLookupRequest request, - final StreamObserver responseObserver) { - callCount.incrementAndGet(); - lastRequestReason = request.getReason(); - lastStaleHeaderData = request.getStaleHeaderData(); - - ScheduledFuture unused = scheduledExecutorService.schedule( - () -> { - RlsProtoData.RouteLookupRequestKey key = - RlsProtoData.RouteLookupRequestKey.create( - REQUEST_CONVERTER.convert(request).keyMap()); - Object entry = responseTable.get(key); - if (entry == null) { - responseObserver.onError(Status.NOT_FOUND.withDescription("key not found").asRuntimeException()); - } else if (entry instanceof Status) { - responseObserver.onError(((Status) entry).asRuntimeException()); - } else if (entry instanceof RouteLookupResponse) { - responseObserver.onNext(RESPONSE_CONVERTER.convert((RouteLookupResponse) entry)); - responseObserver.onCompleted(); - } - }, responseDelayNano, TimeUnit.NANOSECONDS); - } - } - - private static final class TestLoadBalancerProvider extends LoadBalancerProvider { - final Set loadBalancers = new HashSet<>(); - - @Override - public boolean isAvailable() { - return true; - } - - @Override - public int getPriority() { - return 0; - } - - @Override - public String getPolicyName() { - return "target"; - } - - @Override - public ConfigOrError parseLoadBalancingPolicyConfig( - Map rawLoadBalancingPolicyConfig) { - return ConfigOrError.fromConfig(rawLoadBalancingPolicyConfig); - } - - @Override - public LoadBalancer newLoadBalancer(final Helper helper) { - LoadBalancer loadBalancer = new LoadBalancer() { - @Override - public Status acceptResolvedAddresses(ResolvedAddresses resolvedAddresses) { - Map config = (Map) resolvedAddresses.getLoadBalancingPolicyConfig(); - if (DEFAULT_TARGET.equals(config.get("target"))) { - helper.updateBalancingState( - ConnectivityState.TRANSIENT_FAILURE, - new FixedResultPicker( - PickResult.withError(Status.UNAVAILABLE.withDescription("fallback not available")))); - } else { - helper.updateBalancingState( - ConnectivityState.READY, - new FixedResultPicker( - PickResult.withSubchannel(mock(Subchannel.class, config.get("target").toString())))); - } - return Status.OK; - } - - @Override - public void handleNameResolutionError(final Status error) { - helper.updateBalancingState( - ConnectivityState.TRANSIENT_FAILURE, - new FixedResultPicker(PickResult.withError(error))); - } - - @Override - public void shutdown() { - loadBalancers.remove(this); - } - }; - - loadBalancers.add(loadBalancer); - return loadBalancer; - } - } - - private final class FakeHelper extends Helper { - Server server; - ManagedChannel oobChannel; - volatile SubchannelPicker lastPicker; - - void createServerAndRegister(String target) throws IOException { - server = InProcessServerBuilder.forName(target) - .addService(rlsServerImpl) - .directExecutor() - .build() - .start(); - grpcCleanupRule.register(server); - } - - @Override - public ManagedChannelBuilder createResolvingOobChannelBuilder( - String target, ChannelCredentials creds) { - try { - createServerAndRegister(target); - } catch (IOException e) { - throw new RuntimeException("cannot create server: " + target, e); - } - final InProcessChannelBuilder builder = - InProcessChannelBuilder.forName(target).directExecutor(); - - class CleaningChannelBuilder extends ForwardingChannelBuilder2 { - @Override - protected ManagedChannelBuilder delegate() { - return builder; - } - - @Override - public ManagedChannel build() { - oobChannel = super.build(); - return grpcCleanupRule.register(oobChannel); - } - } - - return new CleaningChannelBuilder(); - } - - @Override - public ManagedChannel createOobChannel(EquivalentAddressGroup eag, String authority) { - throw new UnsupportedOperationException(); - } - - @Override - public void updateBalancingState( - @Nonnull ConnectivityState newState, @Nonnull SubchannelPicker newPicker) { - this.lastPicker = newPicker; - } - - @Override - public String getAuthority() { - return "bigtable.googleapis.com:443"; - } - - @Override - public ChannelCredentials getUnsafeChannelCredentials() { - return new ChannelCredentials() { - @Override - public ChannelCredentials withoutBearerTokens() { - return this; - } - }; - } - - @Override - public ScheduledExecutorService getScheduledExecutorService() { - return fakeClock.getScheduledExecutorService(); - } - - @Override - public SynchronizationContext getSynchronizationContext() { - return syncContext; - } - - @Override - public ChannelLogger getChannelLogger() { - return mock(ChannelLogger.class); - } - - @Override - public MetricRecorder getMetricRecorder() { - return mockMetricRecorder; - } - - @Override - public String getChannelTarget() { - return "channelTarget"; - } - } -}