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 =