Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 5 additions & 4 deletions rls/src/main/java/io/grpc/rls/CachingRlsLbClient.java
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -336,7 +336,8 @@ private CachedRouteLookupResponse asyncRlsCall(
}
final SettableFuture<RouteLookupResponse> 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)
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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());
}
}

Expand Down
12 changes: 8 additions & 4 deletions rls/src/main/java/io/grpc/rls/RlsProtoConverters.java
Original file line number Diff line number Diff line change
Expand Up @@ -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();
}
}

Expand Down
10 changes: 9 additions & 1 deletion rls/src/main/java/io/grpc/rls/RlsProtoData.java
Original file line number Diff line number Diff line change
Expand Up @@ -61,8 +61,16 @@ enum Reason {
/** Returns a map of key values extracted via key builders for the gRPC or HTTP request. */
abstract ImmutableMap<String, String> keyMap();

@Nullable
abstract String staleHeaderData();

static RouteLookupRequest create(
ImmutableMap<String, String> keyMap, Reason reason, @Nullable String staleHeaderData) {
return new AutoValue_RlsProtoData_RouteLookupRequest(reason, keyMap, staleHeaderData);
}

static RouteLookupRequest create(ImmutableMap<String, String> keyMap, Reason reason) {
return new AutoValue_RlsProtoData_RouteLookupRequest(reason, keyMap);
return create(keyMap, reason, null);
}
}

Expand Down
156 changes: 156 additions & 0 deletions rls/src/test/java/io/grpc/rls/CachingRlsLbClientTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -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<SubchannelPicker> 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<String, ?> routeLookupChannelServiceConfig =
Expand Down Expand Up @@ -1120,6 +1274,7 @@ private static final class StaticFixedDelayRlsServerImpl
private Map<RlsProtoData.RouteLookupRequestKey, RouteLookupResponse> lookupTable =
ImmutableMap.of();
io.grpc.lookup.v1.RouteLookupRequest.Reason routeLookupReason;
String routeLookupStaleHeaderData;

public StaticFixedDelayRlsServerImpl(
long responseDelayNano, ScheduledExecutorService scheduledExecutorService) {
Expand All @@ -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(
Expand Down
59 changes: 59 additions & 0 deletions rls/src/test/java/io/grpc/rls/RlsProtoConvertersTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,65 @@ public void convert_toRequestObject() {
assertThat(proto.getReason()).isEqualTo(RouteLookupRequest.Reason.REASON_MISS);
}

@Test
public void convert_toRequestProto_staleHeaderData() {
Converter<RouteLookupRequest, RlsProtoData.RouteLookupRequest> 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<RlsProtoData.RouteLookupRequest, RouteLookupRequest> 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<RouteLookupResponse, RlsProtoData.RouteLookupResponse> converter =
Expand Down
Loading