Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,23 @@ private void populateEffectiveRangeAndStartFromSettingsToRequest(RxDocumentServi
new FeedRangeEpkImpl(continuationToken.getRange()));
}

this.populateStartFrom(this.startFromSettings, effectiveStartFrom, request);
}

private void populateStartFrom(
ChangeFeedStartFromInternal initialStartFrom,
ChangeFeedStartFromInternal effectiveStartFrom,
RxDocumentServiceRequest request) {

checkNotNull(initialStartFrom, "Argument 'initialStartFrom' should not be null");
checkNotNull(effectiveStartFrom, "Argument 'effectiveStartFromSettings' should not be null");
checkNotNull(request, "Argument 'request' should not be null");

// When a merge happens, the child partition will contain documents ordered by LSN but the _ts/creation time
// of the documents may not be sequential. So when reading the changeFeed by LSN, it is possible to encounter documents with lower _ts.
// In order to guarantee we always get the documents after customer's point start time, we will need to always pass the start time in the header.
// NOTE: the sequence of calling populate request order matters here, as both can try to populate the same header, effective ones will win
initialStartFrom.populateRequest(request, this.mode);
Comment thread
xinlian12 marked this conversation as resolved.
effectiveStartFrom.populateRequest(request, this.mode);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -6,12 +6,11 @@
import com.azure.cosmos.implementation.InternalObjectNode;
import com.azure.cosmos.implementation.Utils;
import com.azure.cosmos.implementation.changefeed.Lease;
import com.azure.cosmos.implementation.changefeed.common.ChangeFeedStartFromInternal;
import com.azure.cosmos.implementation.changefeed.common.ChangeFeedMode;
import com.azure.cosmos.implementation.changefeed.common.ChangeFeedState;
import com.azure.cosmos.implementation.changefeed.common.ChangeFeedStateV1;
import com.azure.cosmos.implementation.changefeed.common.LeaseVersion;
import com.azure.cosmos.implementation.feedranges.FeedRangeInternal;
import com.azure.cosmos.implementation.changefeed.common.ChangeFeedMode;
import com.azure.cosmos.models.ModelBridgeInternal;
import com.fasterxml.jackson.core.JsonGenerator;
import com.fasterxml.jackson.databind.JsonNode;
Expand Down Expand Up @@ -150,17 +149,11 @@ public ChangeFeedState getContinuationState(String containerRid, ChangeFeedMode

// Lease token are stored in Base64 encoded json - and contains the complete ChangeFeedState
ChangeFeedState changeFeedState = ChangeFeedStateV1.fromString(this.continuationToken);
// Calculating this token from epk based lease format
// This token is then used to pass as lsn in form of etag.
String token = changeFeedState.getContinuation().getCurrentContinuationToken().getToken();
// This token has extra quotes
token = token.replace("\"", "");

return new ChangeFeedStateV1(
containerRid,
this.feedRangeInternal,
changeFeedMode,
ChangeFeedStartFromInternal.createFromETagAndFeedRange(token, this.feedRangeInternal),
changeFeedState.getStartFromSettings(),
Comment thread
xinlian12 marked this conversation as resolved.
changeFeedState.getContinuation());
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,22 +5,43 @@

import com.azure.cosmos.implementation.changefeed.common.ChangeFeedMode;
import com.azure.cosmos.implementation.changefeed.common.ChangeFeedStartFromInternal;
import com.azure.cosmos.implementation.changefeed.common.ChangeFeedStartFromTypes;
import com.azure.cosmos.implementation.changefeed.common.ChangeFeedState;
import com.azure.cosmos.implementation.changefeed.common.ChangeFeedStateV1;
import com.azure.cosmos.implementation.feedranges.FeedRangeContinuation;
import com.azure.cosmos.implementation.feedranges.FeedRangePartitionKeyRangeImpl;
import com.azure.cosmos.implementation.query.CompositeContinuationToken;
import com.azure.cosmos.implementation.routing.Range;
import org.testng.annotations.DataProvider;
import org.testng.annotations.Test;

import java.nio.charset.StandardCharsets;
import java.time.Instant;
import java.util.Base64;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.UUID;

import static com.azure.cosmos.implementation.TestUtils.mockDiagnosticsClientContext;
import static org.assertj.core.api.Assertions.assertThat;

public class ChangeFeedStateTest {
@DataProvider(name = "populateRequestArgProvider")
public Object[][] populateRequestArgProvider() {
return new Object[][] {
// changeFeed mode, changeFeed startFrom type, use continuation
{ ChangeFeedMode.INCREMENTAL, ChangeFeedStartFromTypes.BEGINNING, true },
{ ChangeFeedMode.INCREMENTAL, ChangeFeedStartFromTypes.NOW, true },
{ ChangeFeedMode.INCREMENTAL, ChangeFeedStartFromTypes.POINT_IN_TIME, true },
{ ChangeFeedMode.INCREMENTAL, ChangeFeedStartFromTypes.BEGINNING, false },
{ ChangeFeedMode.INCREMENTAL, ChangeFeedStartFromTypes.NOW, false },
{ ChangeFeedMode.INCREMENTAL, ChangeFeedStartFromTypes.POINT_IN_TIME, false },
{ ChangeFeedMode.FULL_FIDELITY, ChangeFeedStartFromTypes.NOW, true },
{ ChangeFeedMode.FULL_FIDELITY, ChangeFeedStartFromTypes.NOW, false }
};
}

@Test(groups = "unit")
public void changeFeedState_incrementalMode_startFromNow_PKRangeId_toJsonFromJson() {
String containerRid = "/cols/" + UUID.randomUUID().toString();
Expand Down Expand Up @@ -187,12 +208,29 @@ public void changeFeedState_fullFidelityMode_startFromNow_PKRangeId_toJsonFromJs
assertThat(representationAfterDeserialization).isEqualTo(base64EncodedJsonRepresentation);
}

private ChangeFeedState createStateWithContinuation(String continuationAAToCC, String continuationCCToEE)
private ChangeFeedState createDefaultStateWithContinuation(String continuationAAToCC, String continuationCCToEE)
{
String containerRid = "/cols/" + UUID.randomUUID().toString();
String containerRid = "/cols/" + UUID.randomUUID();
String pkRangeId = UUID.randomUUID().toString();
FeedRangePartitionKeyRangeImpl feedRange = new FeedRangePartitionKeyRangeImpl(pkRangeId);
ChangeFeedStartFromInternal startFromSettings = ChangeFeedStartFromInternal.createFromNow();

return this.createStateWithContinuation(
containerRid,
feedRange,
continuationAAToCC,
continuationCCToEE,
ChangeFeedMode.INCREMENTAL,
ChangeFeedStartFromInternal.createFromNow());
}

private ChangeFeedState createStateWithContinuation(
String containerRid,
FeedRangePartitionKeyRangeImpl feedRange,
String continuationAAToCC,
String continuationCCToEE,
ChangeFeedMode changeFeedMode,
ChangeFeedStartFromInternal startFromSettings)
{
String continuationJson = String.format(
"{\"V\":1," +
"\"Rid\":\"%s\"," +
Expand All @@ -204,13 +242,13 @@ private ChangeFeedState createStateWithContinuation(String continuationAAToCC, S
containerRid,
continuationAAToCC,
continuationCCToEE,
pkRangeId);
feedRange.getPartitionKeyRangeId());

FeedRangeContinuation continuation = FeedRangeContinuation.convert(continuationJson);
return new ChangeFeedStateV1(
containerRid,
feedRange,
ChangeFeedMode.INCREMENTAL,
changeFeedMode,
startFromSettings,
continuation);
}
Expand All @@ -221,7 +259,7 @@ public void changeFeedState_extractContinuationTokens() {
String continuationCCToEE = UUID.randomUUID().toString();
List<CompositeContinuationToken> tokens =
this
.createStateWithContinuation(continuationAAToCC, continuationCCToEE)
.createDefaultStateWithContinuation(continuationAAToCC, continuationCCToEE)
.extractForEffectiveRange(new Range<>("AA", "CC", true, false))
.extractContinuationTokens();

Expand All @@ -239,7 +277,7 @@ public void changeFeedState_extractContinuationTokens() {

tokens =
this
.createStateWithContinuation(continuationAAToCC, continuationCCToEE)
.createDefaultStateWithContinuation(continuationAAToCC, continuationCCToEE)
.extractForEffectiveRange(new Range<>("BB", "DD", true, false))
.extractContinuationTokens();

Expand Down Expand Up @@ -269,7 +307,7 @@ public void changeFeedState_merge() {
String continuationAAToCC = UUID.randomUUID().toString();
String continuationCCToEE = UUID.randomUUID().toString();
ChangeFeedState original = this
.createStateWithContinuation(continuationAAToCC, continuationCCToEE);
.createDefaultStateWithContinuation(continuationAAToCC, continuationCCToEE);
ChangeFeedState stateAAToBB = original.extractForEffectiveRange(
new Range<>("AA", "BB", true, false));
ChangeFeedState stateBBToDD = original.extractForEffectiveRange(
Expand Down Expand Up @@ -325,4 +363,83 @@ public void changeFeedState_merge() {
.isNotNull()
.isEqualTo(continuationCCToEE);
}

@Test(dataProvider = "populateRequestArgProvider", groups = "unit")
public void changeFeedState_populateRequest(
ChangeFeedMode changeFeedMode,
ChangeFeedStartFromTypes initialChangeFeedStartFromTypes,
boolean useContinuationToken) {

String containerRid = "/cols/" + UUID.randomUUID();
String pkRangeId = UUID.randomUUID().toString();
FeedRangePartitionKeyRangeImpl feedRange = new FeedRangePartitionKeyRangeImpl(pkRangeId);
ChangeFeedStartFromInternal changeFeedStartFromInternal;
Map<String, String> expectedHeaders = new HashMap<>();

switch (initialChangeFeedStartFromTypes) {
case BEGINNING:
changeFeedStartFromInternal = ChangeFeedStartFromInternal.createFromBeginning();
break;
case NOW:
changeFeedStartFromInternal = ChangeFeedStartFromInternal.createFromNow();
expectedHeaders.put(HttpConstants.HttpHeaders.IF_NONE_MATCH, HttpConstants.HeaderValues.IF_NONE_MATCH_ALL);
break;
case POINT_IN_TIME:
Instant startTime = Instant.now();
changeFeedStartFromInternal = ChangeFeedStartFromInternal.createFromPointInTime(startTime);
expectedHeaders.put(HttpConstants.HttpHeaders.IF_MODIFIED_SINCE, Utils.instantAsUTCRFC1123(startTime));
break;
default:
throw new IllegalStateException("Invalid initialChangeFeedStartFromTypes " + initialChangeFeedStartFromTypes);
}

ChangeFeedState changeFeedState;
if (useContinuationToken) {
String continuationAAToCC = UUID.randomUUID().toString();
String continuationCCToEE = UUID.randomUUID().toString();
changeFeedState = this.createStateWithContinuation(
containerRid,
feedRange,
continuationAAToCC,
continuationCCToEE,
changeFeedMode,
changeFeedStartFromInternal);

expectedHeaders.put(
HttpConstants.HttpHeaders.IF_NONE_MATCH,
changeFeedState.getContinuation().getCurrentContinuationToken().getToken());
} else {
changeFeedState = new ChangeFeedStateV1(
containerRid,
feedRange,
changeFeedMode,
changeFeedStartFromInternal,
null);
}

int maxItemCount = 1;
expectedHeaders.put(HttpConstants.HttpHeaders.PAGE_SIZE, String.valueOf(maxItemCount));
expectedHeaders.put(HttpConstants.HttpHeaders.POPULATE_QUERY_METRICS, "true");
if (changeFeedMode == ChangeFeedMode.INCREMENTAL) {
expectedHeaders.put(HttpConstants.HttpHeaders.A_IM, HttpConstants.A_IMHeaderValues.INCREMENTAL_FEED);
} else {
expectedHeaders.put(HttpConstants.HttpHeaders.A_IM, HttpConstants.A_IMHeaderValues.FULL_FIDELITY_FEED);
expectedHeaders.put(
HttpConstants.HttpHeaders.CHANGE_FEED_WIRE_FORMAT_VERSION,
HttpConstants.ChangeFeedWireFormatVersions.SEPARATE_METADATA_WITH_CRTS);
}

RxDocumentServiceRequest serviceRequest =
RxDocumentServiceRequest.create(
mockDiagnosticsClientContext(),
OperationType.Read,
ResourceType.Document);
changeFeedState.populateRequest(serviceRequest, maxItemCount);
Map<String, String> headers = serviceRequest.getHeaders();

for (String key : expectedHeaders.keySet()) {
assertThat(headers.containsKey(key)).isTrue();
assertThat(headers.get(key)).isEqualTo(expectedHeaders.get(key));
}
}
}