From 2f9b870979ca254c4cfce92858c8775e76e34e54 Mon Sep 17 00:00:00 2001 From: annie-mac Date: Tue, 15 Nov 2022 09:14:24 -0800 Subject: [PATCH 1/3] always populate start time --- .../changefeed/common/ChangeFeedStateV1.java | 17 +++ .../implementation/ChangeFeedStateTest.java | 137 +++++++++++++++++- 2 files changed, 146 insertions(+), 8 deletions(-) diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/changefeed/common/ChangeFeedStateV1.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/changefeed/common/ChangeFeedStateV1.java index e48d58dad5e0..c83a9938d7ab 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/changefeed/common/ChangeFeedStateV1.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/changefeed/common/ChangeFeedStateV1.java @@ -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); effectiveStartFrom.populateRequest(request, this.mode); } diff --git a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/ChangeFeedStateTest.java b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/ChangeFeedStateTest.java index fa7571a4e664..b4aa48be8fd1 100644 --- a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/ChangeFeedStateTest.java +++ b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/ChangeFeedStateTest.java @@ -5,22 +5,47 @@ 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.BEGINNING, true }, + { ChangeFeedMode.FULL_FIDELITY, ChangeFeedStartFromTypes.NOW, true }, + { ChangeFeedMode.FULL_FIDELITY, ChangeFeedStartFromTypes.POINT_IN_TIME, true }, + { ChangeFeedMode.FULL_FIDELITY, ChangeFeedStartFromTypes.BEGINNING, false }, + { ChangeFeedMode.FULL_FIDELITY, ChangeFeedStartFromTypes.NOW, false }, + { ChangeFeedMode.FULL_FIDELITY, ChangeFeedStartFromTypes.POINT_IN_TIME, false }, + }; + } + @Test(groups = "unit") public void changeFeedState_incrementalMode_startFromNow_PKRangeId_toJsonFromJson() { String containerRid = "/cols/" + UUID.randomUUID().toString(); @@ -187,12 +212,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\"," + @@ -204,13 +246,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); } @@ -221,7 +263,7 @@ public void changeFeedState_extractContinuationTokens() { String continuationCCToEE = UUID.randomUUID().toString(); List tokens = this - .createStateWithContinuation(continuationAAToCC, continuationCCToEE) + .createDefaultStateWithContinuation(continuationAAToCC, continuationCCToEE) .extractForEffectiveRange(new Range<>("AA", "CC", true, false)) .extractContinuationTokens(); @@ -239,7 +281,7 @@ public void changeFeedState_extractContinuationTokens() { tokens = this - .createStateWithContinuation(continuationAAToCC, continuationCCToEE) + .createDefaultStateWithContinuation(continuationAAToCC, continuationCCToEE) .extractForEffectiveRange(new Range<>("BB", "DD", true, false)) .extractContinuationTokens(); @@ -269,7 +311,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( @@ -325,4 +367,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 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 headers = serviceRequest.getHeaders(); + + for (String key : expectedHeaders.keySet()) { + assertThat(headers.containsKey(key)).isTrue(); + assertThat(headers.get(key)).isEqualTo(expectedHeaders.get(key)); + } + } } From 2ca4883cdb61d257fc98f57e385a23fd4409fb69 Mon Sep 17 00:00:00 2001 From: annie-mac Date: Wed, 16 Nov 2022 12:28:39 -0800 Subject: [PATCH 2/3] fix --- .../changefeed/epkversion/ServiceItemLeaseV1.java | 11 ++--------- 1 file changed, 2 insertions(+), 9 deletions(-) diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/changefeed/epkversion/ServiceItemLeaseV1.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/changefeed/epkversion/ServiceItemLeaseV1.java index a0b125755513..fcd7851d52d2 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/changefeed/epkversion/ServiceItemLeaseV1.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/changefeed/epkversion/ServiceItemLeaseV1.java @@ -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; @@ -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(), changeFeedState.getContinuation()); } From 4e61517481bfb9c0717ff2f072c572f3560794fb Mon Sep 17 00:00:00 2001 From: annie-mac Date: Mon, 28 Nov 2022 09:08:35 -0800 Subject: [PATCH 3/3] resolve comments --- .../azure/cosmos/implementation/ChangeFeedStateTest.java | 6 +----- 1 file changed, 1 insertion(+), 5 deletions(-) diff --git a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/ChangeFeedStateTest.java b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/ChangeFeedStateTest.java index b4aa48be8fd1..90786ca76dc1 100644 --- a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/ChangeFeedStateTest.java +++ b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/ChangeFeedStateTest.java @@ -37,12 +37,8 @@ public Object[][] populateRequestArgProvider() { { ChangeFeedMode.INCREMENTAL, ChangeFeedStartFromTypes.BEGINNING, false }, { ChangeFeedMode.INCREMENTAL, ChangeFeedStartFromTypes.NOW, false }, { ChangeFeedMode.INCREMENTAL, ChangeFeedStartFromTypes.POINT_IN_TIME, false }, - { ChangeFeedMode.FULL_FIDELITY, ChangeFeedStartFromTypes.BEGINNING, true }, { ChangeFeedMode.FULL_FIDELITY, ChangeFeedStartFromTypes.NOW, true }, - { ChangeFeedMode.FULL_FIDELITY, ChangeFeedStartFromTypes.POINT_IN_TIME, true }, - { ChangeFeedMode.FULL_FIDELITY, ChangeFeedStartFromTypes.BEGINNING, false }, - { ChangeFeedMode.FULL_FIDELITY, ChangeFeedStartFromTypes.NOW, false }, - { ChangeFeedMode.FULL_FIDELITY, ChangeFeedStartFromTypes.POINT_IN_TIME, false }, + { ChangeFeedMode.FULL_FIDELITY, ChangeFeedStartFromTypes.NOW, false } }; }