From 5a23fe5f5c8bcd7a91fa761b51d686f2114f61a6 Mon Sep 17 00:00:00 2001 From: annie-mac Date: Wed, 2 Nov 2022 00:00:42 -0700 Subject: [PATCH 1/8] improve on handling in rntbd when the request is cancelled --- .../RntbdTransportClient.java | 61 +------------------ .../rntbd/RntbdRequestManager.java | 25 ++++---- 2 files changed, 17 insertions(+), 69 deletions(-) diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java index d761cc5a227b..0207904392d7 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java @@ -234,7 +234,7 @@ public Mono invokeStoreAsync(final Uri addressUri, final RxDocume final Context reactorContext = Context.of(KEY_ON_ERROR_DROPPED, onErrorDropHookWithReduceLogLevel); - final Mono result = Mono.fromFuture(record.whenComplete((response, throwable) -> { + record.whenComplete((response, throwable) -> { record.stage(RntbdRequestRecord.Stage.COMPLETED); if (request.requestContext.cosmosDiagnostics == null) { @@ -254,8 +254,9 @@ public Mono invokeStoreAsync(final Uri addressUri, final RxDocume response.setChannelAcquisitionTimeline(record.getChannelAcquisitionTimeline()); } } + }); - })).onErrorMap(throwable -> { + return Mono.fromFuture(record).onErrorMap(throwable -> { Throwable error = throwable instanceof CompletionException ? throwable.getCause() : throwable; @@ -291,62 +292,6 @@ public Mono invokeStoreAsync(final Uri addressUri, final RxDocume } return cosmosException; - }); - - return result.doFinally(signalType -> { - - // This lambda ensures that a pending Direct TCP request in a reactive stream dropped by an end user or the - // HA layer completes without bubbling up to reactor.core.publisher.Hooks#onErrorDropped as a - // CompletionException error. Pending requests may be left outstanding when, for example, an end user calls - // CosmosAsyncClient#close or the HA layer detects that a partition split has occurred. This code guarantees - // that each pending Mono in the stream will run to completion with a new subscriber. - // Consequently the default Hooks#onErrorDropped method will not be called thus preventing distracting error - // messages. - // - // This lambda does not prevent requests that complete exceptionally before the call to this lambda from - // bubbling up to Hooks#onErrorDropped as CompletionException errors. We will still see some onErrorDropped - // messages due to CompletionException errors. Anecdotal evidence shows that this is more likely to be seen - // in low latency environments on Azure cloud. To avoid the onErrorDropped events to get logged in the - // default hook (which logs with level ERROR) we inject a local hook in the Reactor Context to just log it - // as DEBUG level for the lifecycle of this Mono (safe here because we know the onErrorDropped doesn't have - // any functional issues. - // - // One might be tempted to complete a pending request here, but that is ill advised. Testing and - // inspection of the reactor code shows that this does not prevent errors from bubbling up to - // reactor.core.publisher.Hooks#onErrorDropped. Worse than this it has been seen to cause failures in - // the HA layer: - // - // * Calling record.cancel or record.completeExceptionally causes failures in (low-latency) cloud - // environments and all errors bubble up Hooks#onErrorDropped. - // - // * Calling record.complete with a null value causes failures in all environments, depending on the - // operation being performed. In short: many of our tests fail. - - if (signalType != SignalType.CANCEL) { - return; - } - - result.subscribe( - response -> { - if (logger.isDebugEnabled()) { - logger.debug( - "received response to cancelled request: {\"request\":{},\"response\":{\"type\":{}," - + "\"value\":{}}}}", - RntbdObjectMapper.toJson(record), - response.getClass().getSimpleName(), - RntbdObjectMapper.toJson(response)); - } - }, - throwable -> { - if (logger.isDebugEnabled()) { - logger.debug( - "received response to cancelled request: {\"request\":{},\"response\":{\"type\":{}," - + "\"value\":{}}}", - RntbdObjectMapper.toJson(record), - throwable.getClass().getSimpleName(), - RntbdObjectMapper.toJson(throwable)); - } - }); }).contextWrite(reactorContext); } diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManager.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManager.java index e2c6ff2d1585..3c53d8119323 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManager.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManager.java @@ -66,6 +66,7 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; import static com.azure.cosmos.implementation.HttpConstants.StatusCodes; import static com.azure.cosmos.implementation.HttpConstants.SubStatusCodes; @@ -586,7 +587,6 @@ public void write(final ChannelHandlerContext context, final Object message, fin this.timestamps.channelWriteCompleted(); } }); - return; } @@ -661,30 +661,33 @@ Timestamps snapshotTimestamps() { private RntbdRequestRecord addPendingRequestRecord(final ChannelHandlerContext context, final RntbdRequestRecord record) { - return this.pendingRequests.compute(record.transportRequestId(), (id, current) -> { + AtomicReference pendingRequestTimeout = new AtomicReference<>(); + this.pendingRequests.compute(record.transportRequestId(), (id, current) -> { reportIssueUnless(current == null, context, "id: {}, current: {}, request: {}", record); record.pendingRequestQueueSize(pendingRequests.size()); - final Timeout pendingRequestTimeout = record.newTimeout(timeout -> { + pendingRequestTimeout.set(record.newTimeout(timeout -> { // We don't wish to complete on the timeout thread, but rather on a thread doled out by our executor requestExpirationExecutor.execute(record::expire); - }); - - record.whenComplete((response, error) -> { - this.pendingRequests.remove(id); - pendingRequestTimeout.cancel(); - }); + })); return record; + }); + record.whenComplete((response, error) -> { + this.pendingRequests.remove(record.transportRequestId()); + if (pendingRequestTimeout.get() != null) { + pendingRequestTimeout.get().cancel(); + } }); + + return record; } private void completeAllPendingRequestsExceptionally( - final ChannelHandlerContext context, final Throwable throwable - ) { + final ChannelHandlerContext context, final Throwable throwable) { reportIssueUnless(!this.closingExceptionally, context, "", throwable); this.closingExceptionally = true; From 2d794d33d40d937cb4d2a49f12fe75ea215fa5d0 Mon Sep 17 00:00:00 2001 From: annie-mac Date: Wed, 2 Nov 2022 00:02:02 -0700 Subject: [PATCH 2/8] only write to channel if the request record is not cancelled --- .../rntbd/RntbdRequestManager.java | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManager.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManager.java index 3c53d8119323..e4bee60d305d 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManager.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManager.java @@ -581,12 +581,15 @@ public void write(final ChannelHandlerContext context, final Object message, fin this.timestamps.channelWriteAttempted(); record.setSendingRequestHasStarted(); - context.write(this.addPendingRequestRecord(context, record), promise).addListener(completed -> { - record.stage(RntbdRequestRecord.Stage.SENT); - if (completed.isSuccess()) { - this.timestamps.channelWriteCompleted(); - } - }); + if (!record.isCancelled()) { + context.write(this.addPendingRequestRecord(context, record), promise).addListener(completed -> { + record.stage(RntbdRequestRecord.Stage.SENT); + if (completed.isSuccess()) { + this.timestamps.channelWriteCompleted(); + } + }); + } + return; } From 7de033d713ed7d6d522581c716b403f849e56733 Mon Sep 17 00:00:00 2001 From: annie-mac Date: Wed, 2 Nov 2022 00:28:24 -0700 Subject: [PATCH 3/8] changes --- .../RntbdTransportClient.java | 35 +++++++++++-------- 1 file changed, 20 insertions(+), 15 deletions(-) diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java index 0207904392d7..e4c4be1c3eac 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java @@ -234,29 +234,34 @@ public Mono invokeStoreAsync(final Uri addressUri, final RxDocume final Context reactorContext = Context.of(KEY_ON_ERROR_DROPPED, onErrorDropHookWithReduceLogLevel); - record.whenComplete((response, throwable) -> { + return Mono.fromFuture(record).map(storeResponse -> { record.stage(RntbdRequestRecord.Stage.COMPLETED); if (request.requestContext.cosmosDiagnostics == null) { request.requestContext.cosmosDiagnostics = request.createCosmosDiagnostics(); } - if (response != null) { - RequestTimeline timeline = record.takeTimelineSnapshot(); - response.setRequestTimeline(timeline); - response.setEndpointStatistics(record.serviceEndpointStatistics()); - response.setRntbdResponseLength(record.responseLength()); - response.setRntbdRequestLength(record.requestLength()); - response.setRequestPayloadLength(request.getContentLength()); - response.setRntbdChannelTaskQueueSize(record.channelTaskQueueLength()); - response.setRntbdPendingRequestSize(record.pendingRequestQueueSize()); - if(this.channelAcquisitionContextEnabled) { - response.setChannelAcquisitionTimeline(record.getChannelAcquisitionTimeline()); - } + RequestTimeline timeline = record.takeTimelineSnapshot(); + storeResponse.setRequestTimeline(timeline); + storeResponse.setEndpointStatistics(record.serviceEndpointStatistics()); + storeResponse.setRntbdResponseLength(record.responseLength()); + storeResponse.setRntbdRequestLength(record.requestLength()); + storeResponse.setRequestPayloadLength(request.getContentLength()); + storeResponse.setRntbdChannelTaskQueueSize(record.channelTaskQueueLength()); + storeResponse.setRntbdPendingRequestSize(record.pendingRequestQueueSize()); + if(this.channelAcquisitionContextEnabled) { + storeResponse.setChannelAcquisitionTimeline(record.getChannelAcquisitionTimeline()); } - }); - return Mono.fromFuture(record).onErrorMap(throwable -> { + return storeResponse; + + }).onErrorMap(throwable -> { + + record.stage(RntbdRequestRecord.Stage.COMPLETED); + + if (request.requestContext.cosmosDiagnostics == null) { + request.requestContext.cosmosDiagnostics = request.createCosmosDiagnostics(); + } Throwable error = throwable instanceof CompletionException ? throwable.getCause() : throwable; From 494e603e6e83a4fb4c294aee44aab4161b361aec Mon Sep 17 00:00:00 2001 From: annie-mac Date: Wed, 2 Nov 2022 22:15:54 -0700 Subject: [PATCH 4/8] add changelog --- sdk/cosmos/azure-cosmos/CHANGELOG.md | 1 + .../RntbdTransportClient.java | 8 +++++-- .../rntbd/RntbdRequestRecordTests.java | 24 +++++++++++++++++++ 3 files changed, 31 insertions(+), 2 deletions(-) diff --git a/sdk/cosmos/azure-cosmos/CHANGELOG.md b/sdk/cosmos/azure-cosmos/CHANGELOG.md index c379351eec3f..e81cb211e14b 100644 --- a/sdk/cosmos/azure-cosmos/CHANGELOG.md +++ b/sdk/cosmos/azure-cosmos/CHANGELOG.md @@ -14,6 +14,7 @@ #### Other Changes * Updated test dependency of apache `commons-text` to version 1.10.0 - CVE-2022-42889 - See [PR 31674](https://github.com/Azure/azure-sdk-for-java/pull/31674) * Updated `jackson-databind` dependency to 2.13.4.2 - CVE-2022-42003 - See [PR 31559](https://github.com/Azure/azure-sdk-for-java/pull/31559) +* Fixed issue on noisy `CancellationException` log - See [PR 31882](https://github.com/Azure/azure-sdk-for-java/pull/31882) ### 4.38.0 (2022-10-12) #### Features Added diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java index e4c4be1c3eac..b74a9d4a0068 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java @@ -229,11 +229,15 @@ public Mono invokeStoreAsync(final Uri addressUri, final RxDocume final URI address = addressUri.getURI(); final RntbdRequestArgs requestArgs = new RntbdRequestArgs(request, addressUri); + final RntbdEndpoint endpoint = this.endpointProvider.get(address); final RntbdRequestRecord record = endpoint.request(requestArgs); final Context reactorContext = Context.of(KEY_ON_ERROR_DROPPED, onErrorDropHookWithReduceLogLevel); + // Since reactor-core 3.4.23, if the Mono.fromCompletionStage is cancelled, then it will also cancel the internal future + // If SDK has not sent the request to server, then SDK will not send the request to server + // If the request has been sent to server, then SDK will discard the response once get from server return Mono.fromFuture(record).map(storeResponse -> { record.stage(RntbdRequestRecord.Stage.COMPLETED); @@ -249,7 +253,7 @@ public Mono invokeStoreAsync(final Uri addressUri, final RxDocume storeResponse.setRequestPayloadLength(request.getContentLength()); storeResponse.setRntbdChannelTaskQueueSize(record.channelTaskQueueLength()); storeResponse.setRntbdPendingRequestSize(record.pendingRequestQueueSize()); - if(this.channelAcquisitionContextEnabled) { + if (this.channelAcquisitionContextEnabled) { storeResponse.setChannelAcquisitionTimeline(record.getChannelAcquisitionTimeline()); } @@ -292,7 +296,7 @@ public Mono invokeStoreAsync(final Uri addressUri, final RxDocume BridgeInternal.setRntbdPendingRequestQueueSize(cosmosException, record.pendingRequestQueueSize()); BridgeInternal.setChannelTaskQueueSize(cosmosException, record.channelTaskQueueLength()); BridgeInternal.setSendingRequestStarted(cosmosException, record.hasSendingRequestStarted()); - if(this.channelAcquisitionContextEnabled) { + if (this.channelAcquisitionContextEnabled) { BridgeInternal.setChannelAcquisitionTimeline(cosmosException, record.getChannelAcquisitionTimeline()); } diff --git a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestRecordTests.java b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestRecordTests.java index 71ba14dd59e3..9b0f41ed94bf 100644 --- a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestRecordTests.java +++ b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestRecordTests.java @@ -8,12 +8,16 @@ import com.azure.cosmos.implementation.RequestTimeoutException; import com.azure.cosmos.implementation.ResourceType; import com.azure.cosmos.implementation.RxDocumentServiceRequest; +import com.azure.cosmos.implementation.directconnectivity.StoreResponse; import com.azure.cosmos.implementation.directconnectivity.Uri; import org.testng.annotations.DataProvider; import org.testng.annotations.Test; +import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; import java.net.URI; import java.net.URISyntaxException; +import java.time.Duration; import java.util.concurrent.ExecutionException; import static com.azure.cosmos.implementation.TestUtils.mockDiagnosticsClientContext; @@ -57,4 +61,24 @@ public void expireRecord(OperationType operationType, boolean requestSent, Class fail("Wrong exception"); } } + + @Test(groups = { "unit" }) + public void cancelRecord() throws URISyntaxException, InterruptedException { + + RntbdRequestArgs requestArgs = new RntbdRequestArgs( + RxDocumentServiceRequest.create(mockDiagnosticsClientContext(), OperationType.Read, ResourceType.Document), + new Uri(new URI("http://localhost/replica-path").toString()) + ); + + RntbdRequestTimer requestTimer = new RntbdRequestTimer(5000, 5000); + RntbdRequestRecord record = new AsyncRntbdRequestRecord(requestArgs, requestTimer); + Mono result = Mono.fromFuture(record) + .doOnNext(storeResponse -> fail("Record got cancelled should not reach here")) + .doOnError(throwable -> fail("Record got cancelled should not reach here")); + + result.cancelOn(Schedulers.boundedElastic()).subscribe().dispose(); + + Thread.sleep(100); + assertThat(record.isCancelled()).isTrue(); + } } From 61776e67136c3c552d3178742cfd6b51445efb03 Mon Sep 17 00:00:00 2001 From: annie-mac Date: Thu, 3 Nov 2022 07:52:24 -0700 Subject: [PATCH 5/8] update changelog --- sdk/cosmos/azure-cosmos/CHANGELOG.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/cosmos/azure-cosmos/CHANGELOG.md b/sdk/cosmos/azure-cosmos/CHANGELOG.md index 028372c05bf0..6ee63ab8163f 100644 --- a/sdk/cosmos/azure-cosmos/CHANGELOG.md +++ b/sdk/cosmos/azure-cosmos/CHANGELOG.md @@ -8,6 +8,7 @@ #### Bugs Fixed * Fixed a rare race condition for `query plan` cache exceeding the allowed size limit - See [PR 31859](https://github.com/Azure/azure-sdk-for-java/pull/31859) +* Fixed issue on noisy `CancellationException` log - See [PR 31882](https://github.com/Azure/azure-sdk-for-java/pull/31882) #### Other Changes @@ -15,7 +16,6 @@ #### Other Changes * Updated test dependency of apache `commons-text` to version 1.10.0 - CVE-2022-42889 - See [PR 31674](https://github.com/Azure/azure-sdk-for-java/pull/31674) * Updated `jackson-databind` dependency to 2.13.4.2 - CVE-2022-42003 - See [PR 31559](https://github.com/Azure/azure-sdk-for-java/pull/31559) -* Fixed issue on noisy `CancellationException` log - See [PR 31882](https://github.com/Azure/azure-sdk-for-java/pull/31882) ### 4.38.0 (2022-10-12) #### Features Added From d9801d437ac0e393fcacaf1756bad1b9508eea5d Mon Sep 17 00:00:00 2001 From: annie-mac Date: Thu, 17 Nov 2022 08:13:41 -0800 Subject: [PATCH 6/8] add tests --- sdk/cosmos/azure-cosmos/CHANGELOG.md | 2 +- .../RntbdTransportClient.java | 9 +++++ .../RntbdTransportClientTest.java | 33 +++++++++++++++++++ 3 files changed, 43 insertions(+), 1 deletion(-) diff --git a/sdk/cosmos/azure-cosmos/CHANGELOG.md b/sdk/cosmos/azure-cosmos/CHANGELOG.md index 1569a9952292..c436ba1ba849 100644 --- a/sdk/cosmos/azure-cosmos/CHANGELOG.md +++ b/sdk/cosmos/azure-cosmos/CHANGELOG.md @@ -9,6 +9,7 @@ #### Bugs Fixed #### Other Changes +* Fixed issue on noisy `CancellationException` log - See [PR 31882](https://github.com/Azure/azure-sdk-for-java/pull/31882) ### 4.39.0 (2022-11-16) @@ -18,7 +19,6 @@ * Fixed an issue in replica validation where addresses may have not sorted properly when replica validation is enabled. - See [PR 32022](https://github.com/Azure/azure-sdk-for-java/pull/32022) * Fixed unicode char handling in Uris in Cosmos Http Client. - See [PR 32058](https://github.com/Azure/azure-sdk-for-java/pull/32058) * Fixed an eager prefetch issue to lazily prefetch pages on a query - See [PR 32122](https://github.com/Azure/azure-sdk-for-java/pull/32122) -* Fixed issue on noisy `CancellationException` log - See [PR 31882](https://github.com/Azure/azure-sdk-for-java/pull/31882) #### Other Changes * Shaded `MurmurHash3` of apache `commons-codec` to enable removing of the `guava` dependency - CVE-2020-8908 - See [PR 31761](https://github.com/Azure/azure-sdk-for-java/pull/31761) diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java index f1eb6f3c8789..31c011262701 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java @@ -300,6 +300,15 @@ public Mono invokeStoreAsync(final Uri addressUri, final RxDocume } return cosmosException; + }).doFinally(signalType -> { + if (signalType != SignalType.CANCEL) { + return; + } + + // Since reactor-core 3.4.23, if the Mono.fromCompletionStage is cancelled, then it will also cancel the internal future + // But the stated behavior mat change in later versions. In order to keep consistent behavior, we internally will always cancel the future. + record.cancel(true); + }).contextWrite(reactorContext); } diff --git a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClientTest.java b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClientTest.java index 425557bbcfb7..7eca774dfa9e 100644 --- a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClientTest.java +++ b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClientTest.java @@ -53,6 +53,7 @@ import com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdRequestTimer; import com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdResponse; import com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdResponseDecoder; +import com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdServiceEndpoint; import com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdUUID; import com.azure.cosmos.implementation.guava25.base.Strings; import com.azure.cosmos.implementation.guava25.collect.ImmutableMap; @@ -67,9 +68,11 @@ import io.netty.handler.ssl.SslContextBuilder; import io.reactivex.subscribers.TestSubscriber; import org.apache.commons.lang3.StringUtils; +import org.mockito.Mockito; import org.testng.annotations.DataProvider; import org.testng.annotations.Test; import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; import java.io.IOException; import java.net.ConnectException; @@ -89,6 +92,7 @@ import static com.azure.cosmos.implementation.TestUtils.mockDiagnosticsClientContext; import static com.azure.cosmos.implementation.guava27.Strings.lenientFormat; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertTrue; import static org.testng.Assert.fail; @@ -792,6 +796,35 @@ public void sslHandshakeTimeoutTests() throws IOException { } } + @Test(groups = "unit") + public void cancelRequestMono() throws InterruptedException { + RxDocumentServiceRequest request = + RxDocumentServiceRequest.create(mockDiagnosticsClientContext(), OperationType.Read, ResourceType.Document); + RntbdRequestArgs requestArgs = new RntbdRequestArgs(request, physicalAddress); + RntbdRequestTimer requestTimer = new RntbdRequestTimer(5000, 5000); + RntbdRequestRecord rntbdRequestRecord = new AsyncRntbdRequestRecord(requestArgs, requestTimer); + + RntbdEndpoint rntbdEndpoint = Mockito.mock(RntbdServiceEndpoint.class); + Mockito.when(rntbdEndpoint.request(any())).thenReturn(rntbdRequestRecord); + + RntbdEndpoint.Provider endpointProvider = Mockito.mock(RntbdEndpoint.Provider.class); + Mockito.when(endpointProvider.get(physicalAddress.getURI())).thenReturn(rntbdEndpoint); + + RntbdTransportClient transportClient = new RntbdTransportClient(endpointProvider); + transportClient + .invokeStoreAsync( + physicalAddress, + RxDocumentServiceRequest.create(mockDiagnosticsClientContext(), OperationType.Read, ResourceType.Document)) + .cancelOn(Schedulers.boundedElastic()) + .subscribe() + .dispose(); + + // wait for the cancel signal to propagate + Thread.sleep(500); + + assertThat(rntbdRequestRecord.isCancelled()).isTrue(); + } + private static RntbdTransportClient getRntbdTransportClientUnderTest( final UserAgentContainer userAgent, final ConnectionPolicy connectionPolicy, From e5fa345d6e88ce8785da74f2022b9e4267e0a1c1 Mon Sep 17 00:00:00 2001 From: annie-mac Date: Thu, 17 Nov 2022 10:22:54 -0800 Subject: [PATCH 7/8] fix comments --- .../directconnectivity/RntbdTransportClient.java | 3 ++- .../directconnectivity/rntbd/RntbdMetrics.java | 7 +++++-- .../directconnectivity/RntbdTransportClientTest.java | 1 + 3 files changed, 8 insertions(+), 3 deletions(-) diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java index 31c011262701..cc770129ff60 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClient.java @@ -306,7 +306,8 @@ public Mono invokeStoreAsync(final Uri addressUri, final RxDocume } // Since reactor-core 3.4.23, if the Mono.fromCompletionStage is cancelled, then it will also cancel the internal future - // But the stated behavior mat change in later versions. In order to keep consistent behavior, we internally will always cancel the future. + // But the stated behavior may change in later versions (https://github.com/reactor/reactor-core/issues/3235). + // In order to keep consistent behavior, we internally will always cancel the future. record.cancel(true); }).contextWrite(reactorContext); diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdMetrics.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdMetrics.java index ddf8e88d0622..32355df37aa8 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdMetrics.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdMetrics.java @@ -241,8 +241,11 @@ public void markComplete(RntbdRequestRecord requestRecord) { requestRecord.stop(this.requests, requestRecord.isCompletedExceptionally() ? this.responseErrors : this.responseSuccesses); - this.requestSize.record(requestRecord.requestLength()); - this.responseSize.record(requestRecord.responseLength()); + + if (!requestRecord.isCancelled()) { + this.requestSize.record(requestRecord.requestLength()); + this.responseSize.record(requestRecord.responseLength()); + } } @Override diff --git a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClientTest.java b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClientTest.java index 7eca774dfa9e..2c01b8a8a54a 100644 --- a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClientTest.java +++ b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/RntbdTransportClientTest.java @@ -823,6 +823,7 @@ public void cancelRequestMono() throws InterruptedException { Thread.sleep(500); assertThat(rntbdRequestRecord.isCancelled()).isTrue(); + assertThat(rntbdRequestRecord.isCompletedExceptionally()).isTrue(); } private static RntbdTransportClient getRntbdTransportClientUnderTest( From c95ab51de2a3895880f2ef4e7fe21adc36a235c5 Mon Sep 17 00:00:00 2001 From: annie-mac Date: Thu, 17 Nov 2022 23:20:23 -0800 Subject: [PATCH 8/8] resolve comments --- .../directconnectivity/rntbd/RntbdRequestManager.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManager.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManager.java index 85f69f6f4719..d975bd883fda 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManager.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManager.java @@ -592,7 +592,7 @@ public void write(final ChannelHandlerContext context, final Object message, fin } }); } - + return; } @@ -682,6 +682,7 @@ private RntbdRequestRecord addPendingRequestRecord(final ChannelHandlerContext c return record; }); + // NOTE: please do not put the following logic inside the compute block. It may cause dead lock when a record is cancelled early record.whenComplete((response, error) -> { this.pendingRequests.remove(record.transportRequestId()); if (pendingRequestTimeout.get() != null) {