From b1b1956d1be55242a0dde07489fce7140d1adfbe Mon Sep 17 00:00:00 2001 From: Fabian Meiswinkel Date: Wed, 12 Nov 2025 02:05:50 +0000 Subject: [PATCH] Test changes --- .../cosmos/benchmark/ReadMyWriteWorkflow.java | 9 + .../implementation/ConsistencyTests2.java | 1 + .../DocumentQuerySpyWireContentTest.java | 24 +- .../RequestHeadersSpyWireTest.java | 4 +- .../cosmos/implementation/SessionTest.java | 223 +++++++++--------- .../DCDocumentCrudTest.java | 6 +- .../directconnectivity/ReflectionUtils.java | 9 + .../RntbdTransportClientTest.java | 5 + .../com/azure/cosmos/rx/OfferQueryTest.java | 84 ++++--- .../azure/cosmos/rx/OfferReadReplaceTest.java | 7 +- .../azure/cosmos/rx/ReadFeedOffersTest.java | 77 +++--- .../azure/cosmos/rx/ResourceTokenTest.java | 8 +- 12 files changed, 234 insertions(+), 223 deletions(-) diff --git a/sdk/cosmos/azure-cosmos-benchmark/src/main/java/com/azure/cosmos/benchmark/ReadMyWriteWorkflow.java b/sdk/cosmos/azure-cosmos-benchmark/src/main/java/com/azure/cosmos/benchmark/ReadMyWriteWorkflow.java index a3d5b79ad4ab..5584f243c419 100644 --- a/sdk/cosmos/azure-cosmos-benchmark/src/main/java/com/azure/cosmos/benchmark/ReadMyWriteWorkflow.java +++ b/sdk/cosmos/azure-cosmos-benchmark/src/main/java/com/azure/cosmos/benchmark/ReadMyWriteWorkflow.java @@ -498,4 +498,13 @@ protected String getDocumentLink(Document doc) { return doc.getSelfLink(); } } + + @Override + void shutdown() { + if (this.client != null) { + this.client.close(); + } + + super.shutdown(); + } } diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/ConsistencyTests2.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/ConsistencyTests2.java index 38d8c941efa7..1143dfcba0e8 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/ConsistencyTests2.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/ConsistencyTests2.java @@ -233,6 +233,7 @@ public void validateNoChargeOnFailedSessionRead() throws Exception { new CosmosClientTelemetryConfig() .sendClientTelemetryToService(ClientTelemetry.DEFAULT_CLIENT_TELEMETRY_ENABLED)) .build(); + QueryFeedOperationState dummyState = null; try { // CREATE collection diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/DocumentQuerySpyWireContentTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/DocumentQuerySpyWireContentTest.java index fcbee5b7f3e9..5540662212d8 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/DocumentQuerySpyWireContentTest.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/DocumentQuerySpyWireContentTest.java @@ -95,8 +95,7 @@ public void queryWithContinuationTokenLimit(CosmosQueryRequestOptions options, S client.clearCapturedRequests(); - QueryFeedOperationState dummyState = TestUtils - .createDummyQueryFeedOperationState(ResourceType.Document, OperationType.Query, options, client); + QueryFeedOperationState dummyState = TestUtils.createDummyQueryFeedOperationState(ResourceType.Document, OperationType.Query, options, client); try { Flux> queryObservable = client .queryDocuments( @@ -148,6 +147,11 @@ public Document createDocument(AsyncDocumentClient client, String collectionLink @BeforeClass(groups = { "fast" }, timeOut = SETUP_TIMEOUT) public void before_DocumentQuerySpyWireContentTest() throws Exception { + SpyClientUnderTestFactory.ClientUnderTest oldSnapshot = client; + if (oldSnapshot != null) { + oldSnapshot.close(); + } + client = new SpyClientBuilder(this.clientBuilder()).build(); createdDatabase = SHARED_DATABASE; @@ -177,18 +181,14 @@ public void before_DocumentQuerySpyWireContentTest() throws Exception { options, client ); - try { - // do the query once to ensure the collection is cached. - client.queryDocuments(getMultiPartitionCollectionLink(), "select * from root", state, Document.class) - .then().block(); + // do the query once to ensure the collection is cached. + client.queryDocuments(getMultiPartitionCollectionLink(), "select * from root", state, Document.class) + .then().block(); - // do the query once to ensure the collection is cached. - client.queryDocuments(getSinglePartitionCollectionLink(), "select * from root", state, Document.class) - .then().block(); - } finally { - safeClose(state); - } + // do the query once to ensure the collection is cached. + client.queryDocuments(getSinglePartitionCollectionLink(), "select * from root", state, Document.class) + .then().block(); } @AfterClass(groups = { "fast" }, timeOut = SHUTDOWN_TIMEOUT, alwaysRun = true) diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/RequestHeadersSpyWireTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/RequestHeadersSpyWireTest.java index b9cbde1fc7fd..9f5e783d1f67 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/RequestHeadersSpyWireTest.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/RequestHeadersSpyWireTest.java @@ -137,8 +137,7 @@ public void queryWithMaxIntegratedCacheStaleness(CosmosQueryRequestOptions optio client.clearCapturedRequests(); - QueryFeedOperationState dummyState = TestUtils - .createDummyQueryFeedOperationState(ResourceType.Document, OperationType.Query, options, client); + QueryFeedOperationState dummyState = TestUtils.createDummyQueryFeedOperationState(ResourceType.Document, OperationType.Query, options, client); try { client.queryDocuments( collectionLink, @@ -175,7 +174,6 @@ public void queryWithMaxIntegratedCacheStalenessInNanoseconds() { ); try { - assertThatThrownBy(() -> client .queryDocuments(collectionLink, query, state, Document.class) .blockLast()) diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/SessionTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/SessionTest.java index eb1ff945e297..ee03e5e34414 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/SessionTest.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/SessionTest.java @@ -203,132 +203,129 @@ public void partitionedSessionToken(boolean isNameBased) throws NoSuchMethodExce spyClient ); - try { + spyClient.queryDocuments(getCollectionLink(isNameBased), query, dummyState, Document.class).blockFirst(); + assertThat(getSessionTokensInRequests()).hasSize(1); + assertThat(getSessionTokensInRequests().get(0)).isNotEmpty(); + assertThat(getSessionTokensInRequests().get(0)).doesNotContain(","); // making sure we have only one scope session token - spyClient.queryDocuments(getCollectionLink(isNameBased), query, dummyState, Document.class).blockFirst(); - assertThat(getSessionTokensInRequests()).hasSize(1); - assertThat(getSessionTokensInRequests().get(0)).isNotEmpty(); - assertThat(getSessionTokensInRequests().get(0)).doesNotContain(","); // making sure we have only one scope session token + // Session token validation for cross partition query + spyClient.clearCapturedRequests(); + queryRequestOptions = new CosmosQueryRequestOptions(); - // Session token validation for cross partition query - spyClient.clearCapturedRequests(); - queryRequestOptions = new CosmosQueryRequestOptions(); - safeClose(dummyState); - dummyState = TestUtils.createDummyQueryFeedOperationState( - ResourceType.Document, - OperationType.Query, - queryRequestOptions, - spyClient - ); - spyClient.queryDocuments(getCollectionLink(isNameBased), query, dummyState, Document.class).blockFirst(); - assertThat(getSessionTokensInRequests().size()).isGreaterThanOrEqualTo(1); - assertThat(getSessionTokensInRequests().get(0)).isNotEmpty(); - assertThat(getSessionTokensInRequests().get(0)).doesNotContain(","); // making sure we have only one scope session token + safeClose(dummyState); + dummyState = TestUtils.createDummyQueryFeedOperationState( + ResourceType.Document, + OperationType.Query, + queryRequestOptions, + spyClient + ); + spyClient.queryDocuments(getCollectionLink(isNameBased), query, dummyState, Document.class).blockFirst(); + assertThat(getSessionTokensInRequests().size()).isGreaterThanOrEqualTo(1); + assertThat(getSessionTokensInRequests().get(0)).isNotEmpty(); + assertThat(getSessionTokensInRequests().get(0)).doesNotContain(","); // making sure we have only one scope session token - // Session token validation for feed ranges query - spyClient.clearCapturedRequests(); - List feedRanges = spyClient.getFeedRanges(getCollectionLink(isNameBased), true).block(); - queryRequestOptions = new CosmosQueryRequestOptions(); - queryRequestOptions.setFeedRange(feedRanges.get(0)); - safeClose(dummyState); - dummyState = TestUtils.createDummyQueryFeedOperationState( - ResourceType.Document, - OperationType.Query, - queryRequestOptions, - spyClient - ); - spyClient.queryDocuments(getCollectionLink(isNameBased), query, dummyState, Document.class).blockFirst(); - assertThat(getSessionTokensInRequests().size()).isGreaterThanOrEqualTo(1); - assertThat(getSessionTokensInRequests().get(0)).isNotEmpty(); - assertThat(getSessionTokensInRequests().get(0)).doesNotContain(","); // making sure we have only one scope session token + // Session token validation for feed ranges query + spyClient.clearCapturedRequests(); + List feedRanges = spyClient.getFeedRanges(getCollectionLink(isNameBased), true).block(); + queryRequestOptions = new CosmosQueryRequestOptions(); + queryRequestOptions.setFeedRange(feedRanges.get(0)); + safeClose(dummyState); + dummyState = TestUtils.createDummyQueryFeedOperationState( + ResourceType.Document, + OperationType.Query, + queryRequestOptions, + spyClient + ); + spyClient.queryDocuments(getCollectionLink(isNameBased), query, dummyState, Document.class).blockFirst(); + assertThat(getSessionTokensInRequests().size()).isGreaterThanOrEqualTo(1); + assertThat(getSessionTokensInRequests().get(0)).isNotEmpty(); + assertThat(getSessionTokensInRequests().get(0)).doesNotContain(","); // making sure we have only one scope session token - // Session token validation for readAll with partition query - spyClient.clearCapturedRequests(); - queryRequestOptions = new CosmosQueryRequestOptions(); - safeClose(dummyState); - dummyState = TestUtils.createDummyQueryFeedOperationState( - ResourceType.Document, - OperationType.ReadFeed, - queryRequestOptions, - spyClient - ); - spyClient.readAllDocuments( - getCollectionLink(isNameBased), - new PartitionKey(documentCreated.getId()), - dummyState, - Document.class).blockFirst(); - assertThat(getSessionTokensInRequests().size()).isEqualTo(1); - assertThat(getSessionTokensInRequests().get(0)).isNotEmpty(); - assertThat(getSessionTokensInRequests().get(0)).doesNotContain(","); // making sure we have only one scope session token + // Session token validation for readAll with partition query + spyClient.clearCapturedRequests(); + queryRequestOptions = new CosmosQueryRequestOptions(); + safeClose(dummyState); + dummyState = TestUtils.createDummyQueryFeedOperationState( + ResourceType.Document, + OperationType.ReadFeed, + queryRequestOptions, + spyClient + ); + spyClient.readAllDocuments( + getCollectionLink(isNameBased), + new PartitionKey(documentCreated.getId()), + dummyState, + Document.class).blockFirst(); + assertThat(getSessionTokensInRequests().size()).isEqualTo(1); + assertThat(getSessionTokensInRequests().get(0)).isNotEmpty(); + assertThat(getSessionTokensInRequests().get(0)).doesNotContain(","); // making sure we have only one scope session token - // Session token validation for readAll with cross partition - spyClient.clearCapturedRequests(); - queryRequestOptions = new CosmosQueryRequestOptions(); + // Session token validation for readAll with cross partition + spyClient.clearCapturedRequests(); + queryRequestOptions = new CosmosQueryRequestOptions(); - safeClose(dummyState); - dummyState = TestUtils.createDummyQueryFeedOperationState( - ResourceType.Document, - OperationType.ReadFeed, - queryRequestOptions, - spyClient - ); + safeClose(dummyState); + dummyState = TestUtils.createDummyQueryFeedOperationState( + ResourceType.Document, + OperationType.ReadFeed, + queryRequestOptions, + spyClient + ); - spyClient.readDocuments(getCollectionLink(isNameBased), dummyState, Document.class).blockFirst(); - assertThat(getSessionTokensInRequests().size()).isGreaterThanOrEqualTo(1); - assertThat(getSessionTokensInRequests().get(0)).isNotEmpty(); - assertThat(getSessionTokensInRequests().get(0)).doesNotContain(","); // making sure we have only one scope session token + spyClient.readDocuments(getCollectionLink(isNameBased), dummyState, Document.class).blockFirst(); + assertThat(getSessionTokensInRequests().size()).isGreaterThanOrEqualTo(1); + assertThat(getSessionTokensInRequests().get(0)).isNotEmpty(); + assertThat(getSessionTokensInRequests().get(0)).doesNotContain(","); // making sure we have only one scope session token - // Session token validation for readMany with cross partition + // Session token validation for readMany with cross partition + spyClient.clearCapturedRequests(); + queryRequestOptions = new CosmosQueryRequestOptions(); + CosmosItemIdentity cosmosItemIdentity = new CosmosItemIdentity(new PartitionKey(documentCreated.getId()), documentCreated.getId()); + List cosmosItemIdentities = new ArrayList<>(); + cosmosItemIdentities.add(cosmosItemIdentity); + safeClose(dummyState); + dummyState = TestUtils.createDummyQueryFeedOperationState(ResourceType.Document, OperationType.Query, queryRequestOptions, spyClient); + spyClient.readMany( + cosmosItemIdentities, + getCollectionLink(isNameBased), + dummyState, + InternalObjectNode.class).block(); + assertThat(getSessionTokensInRequests().size()).isEqualTo(1); + assertThat(getSessionTokensInRequests().get(0)).isNotEmpty(); + assertThat(getSessionTokensInRequests().get(0)).doesNotContain(","); // making sure we have only one scope session token + // session token + + // Session token validation for create in Batch + if(isNameBased) { // Batch only work with name based url spyClient.clearCapturedRequests(); - queryRequestOptions = new CosmosQueryRequestOptions(); - CosmosItemIdentity cosmosItemIdentity = new CosmosItemIdentity(new PartitionKey(documentCreated.getId()), documentCreated.getId()); - List cosmosItemIdentities = new ArrayList<>(); - cosmosItemIdentities.add(cosmosItemIdentity); - safeClose(dummyState); - dummyState = TestUtils - .createDummyQueryFeedOperationState(ResourceType.Document, OperationType.Query, queryRequestOptions, spyClient); - spyClient.readMany( - cosmosItemIdentities, - getCollectionLink(isNameBased), - dummyState, - InternalObjectNode.class).block(); + Document document = newDocument(); + document.set("mypk", document.getId()); + ItemBatchOperation itemBatchOperation = new ItemBatchOperation(CosmosItemOperationType.CREATE, + documentCreated.getId(), new PartitionKey(documentCreated.getId()), new RequestOptions(), document); + List> itemBatchOperations = new ArrayList<>(); + itemBatchOperations.add(itemBatchOperation); + + Method method = SinglePartitionKeyServerBatchRequest.class.getDeclaredMethod("createBatchRequest", + PartitionKey.class, + List.class); + method.setAccessible(true); + SinglePartitionKeyServerBatchRequest serverBatchRequest = + (SinglePartitionKeyServerBatchRequest) method.invoke(SinglePartitionKeyServerBatchRequest.class, new PartitionKey(document.getId()), + itemBatchOperations); + spyClient + .executeBatchRequest( + getCollectionLink(isNameBased), + serverBatchRequest, + new RequestOptions(), + false, + true) + .block(); assertThat(getSessionTokensInRequests().size()).isEqualTo(1); assertThat(getSessionTokensInRequests().get(0)).isNotEmpty(); assertThat(getSessionTokensInRequests().get(0)).doesNotContain(","); // making sure we have only one scope session token - // session token - - // Session token validation for create in Batch - if (isNameBased) { // Batch only work with name based url - spyClient.clearCapturedRequests(); - Document document = newDocument(); - document.set("mypk", document.getId()); - ItemBatchOperation itemBatchOperation = new ItemBatchOperation(CosmosItemOperationType.CREATE, - documentCreated.getId(), new PartitionKey(documentCreated.getId()), new RequestOptions(), document); - List> itemBatchOperations = new ArrayList<>(); - itemBatchOperations.add(itemBatchOperation); - - Method method = SinglePartitionKeyServerBatchRequest.class.getDeclaredMethod("createBatchRequest", - PartitionKey.class, - List.class); - method.setAccessible(true); - SinglePartitionKeyServerBatchRequest serverBatchRequest = - (SinglePartitionKeyServerBatchRequest) method.invoke(SinglePartitionKeyServerBatchRequest.class, new PartitionKey(document.getId()), - itemBatchOperations); - spyClient - .executeBatchRequest( - getCollectionLink(isNameBased), - serverBatchRequest, - new RequestOptions(), - false, - true) - .block(); - assertThat(getSessionTokensInRequests().size()).isEqualTo(1); - assertThat(getSessionTokensInRequests().get(0)).isNotEmpty(); - assertThat(getSessionTokensInRequests().get(0)).doesNotContain(","); // making sure we have only one scope session token - } - } finally { - safeClose(dummyState); } + + safeClose(dummyState); } @Test(groups = { "fast" }, timeOut = TIMEOUT, dataProvider = "sessionTestArgProvider") diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/directconnectivity/DCDocumentCrudTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/directconnectivity/DCDocumentCrudTest.java index 27969a39fd7c..7e7410836eb5 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/directconnectivity/DCDocumentCrudTest.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/directconnectivity/DCDocumentCrudTest.java @@ -229,9 +229,8 @@ public void crossPartitionQuery() { options.setMaxDegreeOfParallelism(-1); ModelBridgeInternal.setQueryRequestOptionsMaxItemCount(options, 100); - QueryFeedOperationState dummyState = TestUtils - .createDummyQueryFeedOperationState(ResourceType.Document, OperationType.Query, options, client); - + QueryFeedOperationState dummyState = + TestUtils.createDummyQueryFeedOperationState(ResourceType.Document, OperationType.Query, options, client); try { Flux> results = client.queryDocuments( getCollectionLink(), @@ -245,6 +244,7 @@ public void crossPartitionQuery() { validateQuerySuccess(results, validator, QUERY_TIMEOUT); validateNoDocumentQueryOperationThroughGateway(); + // validates only the first query for fetching query plan goes to gateway. assertThat(client.getCapturedRequests().stream().filter(r -> r.getResourceType() == ResourceType.Document)).hasSize(1); } finally { diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/directconnectivity/ReflectionUtils.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/directconnectivity/ReflectionUtils.java index cf126f74b3fe..cf1110aec87f 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/directconnectivity/ReflectionUtils.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/directconnectivity/ReflectionUtils.java @@ -317,6 +317,15 @@ public static void setTransportClient(StoreReader storeReader, TransportClient t set(storeReader, transportClient, "transportClient"); } + public static void setTransportClient(CosmosClient client, TransportClient transportClient) { + StoreClient storeClient = getStoreClient((RxDocumentClientImpl) CosmosBridgeInternal.getAsyncDocumentClient(client)); + set(storeClient, transportClient, "transportClient"); + ReplicatedResourceClient replicatedResClient = getReplicatedResourceClient(storeClient); + ConsistencyWriter writer = getConsistencyWriter(replicatedResClient); + set(replicatedResClient, transportClient, "transportClient"); + set(writer, transportClient, "transportClient"); + } + public static TransportClient getTransportClient(ReplicatedResourceClient replicatedResourceClient) { return get(TransportClient.class, replicatedResourceClient, "transportClient"); } diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClientTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClientTest.java index 07af54b8c73c..b1c93b2ee7a3 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClientTest.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClientTest.java @@ -1138,6 +1138,11 @@ public URI serviceEndpoint() { return null; } + @Override + public URI serverKeyUsedAsActualRemoteAddress() { + return this.remoteURI; + } + @Override public void injectConnectionErrors(String ruleId, double threshold, Class eventType) { throw new NotImplementedException("injectConnectionErrors is not supported in FakeEndpoint"); diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/OfferQueryTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/OfferQueryTest.java index d52803c0336f..fb814862e6ea 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/OfferQueryTest.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/OfferQueryTest.java @@ -61,19 +61,20 @@ public void queryOffersWithFilter() throws Exception { String query = String.format("SELECT * from c where c.offerResourceId = '%s'", collectionResourceId); CosmosQueryRequestOptions options = new CosmosQueryRequestOptions(); - QueryFeedOperationState dummyState = TestUtils + ModelBridgeInternal.setQueryRequestOptionsMaxItemCount(options, 2); + + QueryFeedOperationState queryDummyState = TestUtils .createDummyQueryFeedOperationState(ResourceType.Offer, OperationType.Query, options, client); - QueryFeedOperationState dummyState2 = null; + QueryFeedOperationState offerDummyState = TestUtils + .createDummyQueryFeedOperationState(ResourceType.Offer, OperationType.ReadFeed, options, client); + try { - ModelBridgeInternal.setQueryRequestOptionsMaxItemCount(options, 2); Flux> queryObservable = client.queryOffers( query, - dummyState); + queryDummyState); - dummyState2 = TestUtils - .createDummyQueryFeedOperationState(ResourceType.Offer, OperationType.ReadFeed, options, client); List allOffers = client - .readOffers(dummyState2) + .readOffers(offerDummyState) .flatMap(f -> Flux.fromIterable(f.getResults())).collectList().single().block(); List expectedOffers = allOffers.stream().filter(o -> collectionResourceId.equals(o.getString("offerResourceId"))).collect(Collectors.toList()); @@ -92,8 +93,8 @@ public void queryOffersWithFilter() throws Exception { validateQuerySuccess(queryObservable, validator, 10000); } finally { - safeClose(dummyState); - safeClose(dummyState2); + safeClose(queryDummyState); + safeClose(offerDummyState); } } @@ -107,18 +108,20 @@ public void queryOffersFilterMorePages() throws Exception { CosmosQueryRequestOptions options = new CosmosQueryRequestOptions(); ModelBridgeInternal.setQueryRequestOptionsMaxItemCount(options, 1); - QueryFeedOperationState dummyState = TestUtils + QueryFeedOperationState queryDummyState = TestUtils .createDummyQueryFeedOperationState(ResourceType.Offer, OperationType.Query, options, client); - QueryFeedOperationState dummyState2 = null; + + QueryFeedOperationState offerDummyState = TestUtils + .createDummyQueryFeedOperationState(ResourceType.Offer, OperationType.ReadFeed, new CosmosQueryRequestOptions(), client); + try { + Flux> queryObservable = client.queryOffers( query, - dummyState); + queryDummyState); - dummyState2 = TestUtils - .createDummyQueryFeedOperationState(ResourceType.Offer, OperationType.ReadFeed, new CosmosQueryRequestOptions(), client); List expectedOffers = client - .readOffers(dummyState2) + .readOffers(offerDummyState) .flatMap(f -> Flux.fromIterable(f.getResults())) .collectList() .single().block() @@ -140,8 +143,8 @@ public void queryOffersFilterMorePages() throws Exception { validateQuerySuccess(queryObservable, validator, 10000); } finally { - safeClose(dummyState); - safeClose(dummyState2); + safeClose(queryDummyState); + safeClose(offerDummyState); } } @@ -150,35 +153,30 @@ public void queryCollections_NoResults() throws Exception { String query = "SELECT * from root r where r.id = '2'"; CosmosQueryRequestOptions options = new CosmosQueryRequestOptions(); - try(CosmosAsyncClient cosmosClient = new CosmosClientBuilder() + CosmosAsyncClient cosmosClient = new CosmosClientBuilder() .key(TestConfigurations.MASTER_KEY) .endpoint(TestConfigurations.HOST) - .buildAsyncClient()) { - QueryFeedOperationState dummyState = new QueryFeedOperationState( - cosmosClient, - "SomeSpanName", - "SomeDBName", - "SomeContainerName", - ResourceType.Document, - OperationType.Query, - null, - options, - new CosmosPagedFluxOptions() - ); - try { - Flux> queryObservable = client.queryCollections(getDatabaseLink(), query, dummyState); - - FeedResponseListValidator validator = new FeedResponseListValidator.Builder() - .containsExactly(new ArrayList<>()) - .numberOfPages(1) - .pageSatisfy(0, new FeedResponseValidator.Builder() + .buildAsyncClient(); + QueryFeedOperationState dummyState = new QueryFeedOperationState( + cosmosClient, + "SomeSpanName", + "SomeDBName", + "SomeContainerName", + ResourceType.Document, + OperationType.Query, + null, + options, + new CosmosPagedFluxOptions() + ); + Flux> queryObservable = client.queryCollections(getDatabaseLink(), query, dummyState); + + FeedResponseListValidator validator = new FeedResponseListValidator.Builder() + .containsExactly(new ArrayList<>()) + .numberOfPages(1) + .pageSatisfy(0, new FeedResponseValidator.Builder() .requestChargeGreaterThanOrEqualTo(1.0).build()) - .build(); - validateQuerySuccess(queryObservable, validator); - } finally { - safeClose(dummyState); - } - } + .build(); + validateQuerySuccess(queryObservable, validator); } @BeforeClass(groups = { "query" }, timeOut = SETUP_TIMEOUT) diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/OfferReadReplaceTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/OfferReadReplaceTest.java index c3ac8eb1e5a4..ed4f1a523e71 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/OfferReadReplaceTest.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/OfferReadReplaceTest.java @@ -43,14 +43,15 @@ public OfferReadReplaceTest(AsyncDocumentClient.Builder clientBuilder) { @Test(groups = { "emulator" }, timeOut = TIMEOUT) public void readAndReplaceOffer() { - QueryFeedOperationState dummyState = TestUtils.createDummyQueryFeedOperationState( + QueryFeedOperationState offerDummyState = TestUtils.createDummyQueryFeedOperationState( ResourceType.Offer, OperationType.ReadFeed, new CosmosQueryRequestOptions(), client); + try { List offers = client - .readOffers(dummyState) + .readOffers(offerDummyState) .map(FeedResponse::getResults) .flatMap(list -> Flux.fromIterable(list)).collectList().block(); @@ -87,7 +88,7 @@ public void readAndReplaceOffer() { validateSuccess(replaceObservable, validatorForReplace); } finally { - safeClose(dummyState); + safeClose(offerDummyState); } } diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/ReadFeedOffersTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/ReadFeedOffersTest.java index 6a2bccf66250..6bc22bd03a2d 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/ReadFeedOffersTest.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/ReadFeedOffersTest.java @@ -59,41 +59,35 @@ public void readOffers() throws Exception { CosmosQueryRequestOptions options = new CosmosQueryRequestOptions(); ModelBridgeInternal.setQueryRequestOptionsMaxItemCount(options, 2); - try (CosmosAsyncClient cosmosClient = new CosmosClientBuilder() + CosmosAsyncClient cosmosClient = new CosmosClientBuilder() .key(TestConfigurations.MASTER_KEY) .endpoint(TestConfigurations.HOST) - .buildAsyncClient()) { - QueryFeedOperationState dummyState = new QueryFeedOperationState( - cosmosClient, - "SomeSpanName", - "SomeDBName", - "SomeContainerName", - ResourceType.Document, - OperationType.Query, - null, - options, - new CosmosPagedFluxOptions() - ); - - - try { - Flux> feedObservable = client.readOffers(dummyState); - - int maxItemCount = ModelBridgeInternal.getMaxItemCountFromQueryRequestOptions(options); - int expectedPageSize = (allOffers.size() + maxItemCount - 1) / maxItemCount; - - FeedResponseListValidator validator = new FeedResponseListValidator.Builder() - .totalSize(allOffers.size()) - .exactlyContainsInAnyOrder(allOffers.stream().map(d -> d.getResourceId()).collect(Collectors.toList())) - .numberOfPages(expectedPageSize) - .pageSatisfy(0, new FeedResponseValidator.Builder() + .buildAsyncClient(); + QueryFeedOperationState dummyState = new QueryFeedOperationState( + cosmosClient, + "SomeSpanName", + "SomeDBName", + "SomeContainerName", + ResourceType.Document, + OperationType.Query, + null, + options, + new CosmosPagedFluxOptions() + ); + + Flux> feedObservable = client.readOffers(dummyState); + + int maxItemCount = ModelBridgeInternal.getMaxItemCountFromQueryRequestOptions(options); + int expectedPageSize = (allOffers.size() + maxItemCount - 1) / maxItemCount; + + FeedResponseListValidator validator = new FeedResponseListValidator.Builder() + .totalSize(allOffers.size()) + .exactlyContainsInAnyOrder(allOffers.stream().map(d -> d.getResourceId()).collect(Collectors.toList())) + .numberOfPages(expectedPageSize) + .pageSatisfy(0, new FeedResponseValidator.Builder() .requestChargeGreaterThanOrEqualTo(1.0).build()) - .build(); - validateQuerySuccess(feedObservable, validator, FEED_TIMEOUT); - } finally { - safeClose(dummyState); - } - } + .build(); + validateQuerySuccess(feedObservable, validator, FEED_TIMEOUT); } @BeforeClass(groups = { "query" }, timeOut = SETUP_TIMEOUT) @@ -105,20 +99,23 @@ public void before_ReadFeedOffersTest() { createCollections(client); } - QueryFeedOperationState dummyState = TestUtils.createDummyQueryFeedOperationState( + QueryFeedOperationState offerDummyState = TestUtils.createDummyQueryFeedOperationState( ResourceType.Offer, OperationType.ReadFeed, new CosmosQueryRequestOptions(), client ); - allOffers = client.readOffers(dummyState) - .map(FeedResponse::getResults) - .collectList() - .map(list -> list.stream().flatMap(Collection::stream).collect(Collectors.toList())) - .single() - .doFinally(signal -> safeClose(dummyState)) - .block(); + try { + allOffers = client.readOffers(offerDummyState) + .map(FeedResponse::getResults) + .collectList() + .map(list -> list.stream().flatMap(Collection::stream).collect(Collectors.toList())) + .single() + .block(); + } finally { + safeClose(offerDummyState); + } } @AfterClass(groups = { "query" }, timeOut = SHUTDOWN_TIMEOUT, alwaysRun = true) diff --git a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/ResourceTokenTest.java b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/ResourceTokenTest.java index d42367a5ef49..ead3cdc9f785 100644 --- a/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/ResourceTokenTest.java +++ b/sdk/cosmos/azure-cosmos-tests/src/test/java/com/azure/cosmos/rx/ResourceTokenTest.java @@ -483,11 +483,7 @@ public void queryItemFromResourceToken(DocumentCollection documentCollection, Pe queryRequestOptions.setPartitionKey(partitionKey); dummyState = TestUtils - .createDummyQueryFeedOperationState( - ResourceType.Document, - OperationType.Query, - queryRequestOptions, - asyncClientResourceToken); + .createDummyQueryFeedOperationState(ResourceType.Document, OperationType.Query, queryRequestOptions, asyncClientResourceToken); Flux> queryObservable = asyncClientResourceToken.queryDocuments( @@ -503,8 +499,8 @@ public void queryItemFromResourceToken(DocumentCollection documentCollection, Pe validateQuerySuccess(queryObservable, validator); } finally { - safeClose(dummyState); safeClose(asyncClientResourceToken); + safeClose(dummyState); } }