From 0c6ed32702d81d564cdd3072ec7a5f6a58b20a02 Mon Sep 17 00:00:00 2001 From: annie-mac Date: Mon, 17 Oct 2022 18:05:50 -0700 Subject: [PATCH 1/7] add transitTimeout check during healthCheck --- .../RntbdTransportClient.java | 11 + .../RntbdClientChannelHealthChecker.java | 221 ++++++++++-------- .../rntbd/RntbdEndpoint.java | 5 + .../rntbd/RntbdRequestManager.java | 7 +- .../rntbd/RntbdRequestRecord.java | 10 + 5 files changed, 153 insertions(+), 101 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..7c95005afec4 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 @@ -474,6 +474,8 @@ public static final class Options { @JsonProperty() private final Duration sslHandshakeTimeoutMinDuration; + @JsonProperty() + private final int transientTimeoutDetectionThreshold; // endregion // region Constructors @@ -508,6 +510,7 @@ private Options(final Builder builder) { this.tcpKeepIdle = builder.tcpKeepIdle; this.preferTcpNative = builder.preferTcpNative; this.sslHandshakeTimeoutMinDuration = builder.sslHandshakeTimeoutMinDuration; + this.transientTimeoutDetectionThreshold = builder.transientTimeoutDetectionThreshold; this.connectTimeout = builder.connectTimeout == null ? builder.tcpNetworkRequestTimeout @@ -541,6 +544,7 @@ private Options(final ConnectionPolicy connectionPolicy) { this.tcpKeepIntvl = 1; // Configuration for EpollChannelOption.TCP_KEEPINTVL this.tcpKeepIdle = 30; // Configuration for EpollChannelOption.TCP_KEEPIDLE this.sslHandshakeTimeoutMinDuration = Duration.ofSeconds(5); + this.transientTimeoutDetectionThreshold = 5; this.preferTcpNative = true; } @@ -646,6 +650,11 @@ public long sslHandshakeTimeoutInMillis() { return Math.max(this.sslHandshakeTimeoutMinDuration.toMillis(), this.connectTimeout.toMillis()); } + public int transientTimeoutDetectionThreshold() { + return this.transientTimeoutDetectionThreshold; + } + + // endregion // region Methods @@ -804,6 +813,7 @@ public static class Builder { private int tcpKeepIdle; private boolean preferTcpNative; private Duration sslHandshakeTimeoutMinDuration; + private int transientTimeoutDetectionThreshold; // endregion @@ -838,6 +848,7 @@ public Builder(ConnectionPolicy connectionPolicy) { this.tcpKeepIdle = DEFAULT_OPTIONS.tcpKeepIdle; this.preferTcpNative = DEFAULT_OPTIONS.preferTcpNative; this.sslHandshakeTimeoutMinDuration = DEFAULT_OPTIONS.sslHandshakeTimeoutMinDuration; + this.transientTimeoutDetectionThreshold = DEFAULT_OPTIONS.transientTimeoutDetectionThreshold; } // endregion diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthChecker.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthChecker.java index 992de40635b3..e34ecd08196e 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthChecker.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthChecker.java @@ -3,6 +3,7 @@ package com.azure.cosmos.implementation.directconnectivity.rntbd; +import com.azure.cosmos.implementation.apachecommons.lang.StringUtils; import com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdEndpoint.Config; import com.fasterxml.jackson.annotation.JsonProperty; import io.netty.channel.Channel; @@ -14,6 +15,7 @@ import java.text.MessageFormat; import java.util.Optional; +import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; import java.util.concurrent.atomic.AtomicLongFieldUpdater; import static com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdReporter.reportIssueUnless; @@ -54,6 +56,9 @@ public final class RntbdClientChannelHealthChecker implements ChannelHealthCheck @JsonProperty private final long writeDelayLimitInNanos; + @JsonProperty + private final int transitTimeoutDetectionThreshold; + // endregion // region Constructors @@ -73,6 +78,7 @@ public RntbdClientChannelHealthChecker(final Config config) { this.idleConnectionTimeoutInNanos = config.idleConnectionTimeoutInNanos(); this.readDelayLimitInNanos = config.receiveHangDetectionTimeInNanos(); this.writeDelayLimitInNanos = config.sendHangDetectionTimeInNanos(); + this.transitTimeoutDetectionThreshold = config.transitTimeoutDetectionThreshold(); } @@ -144,54 +150,19 @@ public Future isHealthy(final Channel channel) { return promise.setSuccess(Boolean.TRUE); // because we recently received data } - // Black hole detection, part 1: - // Treat the channel as unhealthy if the gap between the last attempted write and the last successful write - // grew beyond acceptable limits, unless a write was attempted recently. This is a sign of a nonresponding write. - - final long writeDelayInNanos = - timestamps.lastChannelWriteAttemptNanoTime() - timestamps.lastChannelWriteNanoTime(); - - final long writeHangDurationInNanos = - currentTime - timestamps.lastChannelWriteAttemptNanoTime(); - - if (writeDelayInNanos > this.writeDelayLimitInNanos && writeHangDurationInNanos > writeHangGracePeriodInNanos) { - - final Optional rntbdContext = requestManager.rntbdContext(); - final int pendingRequestCount = requestManager.pendingRequestCount(); - - logger.warn("{} health check failed due to nonresponding write: {lastChannelWriteAttemptNanoTime: {}, " + - "lastChannelWriteNanoTime: {}, writeDelayInNanos: {}, writeDelayLimitInNanos: {}, " + - "rntbdContext: {}, pendingRequestCount: {}}", - channel, timestamps.lastChannelWriteAttemptNanoTime(), timestamps.lastChannelWriteNanoTime(), - writeDelayInNanos, this.writeDelayLimitInNanos, rntbdContext, pendingRequestCount); - + String writeIsHangMessage = this.isWriteHang(timestamps, currentTime, requestManager, channel); + if (StringUtils.isNotEmpty(writeIsHangMessage)) { return promise.setSuccess(Boolean.FALSE); } - // Black hole detection, part 2: - // Treat the connection as unhealthy if the gap between the last successful write and the last successful read - // grew beyond acceptable limits, unless a write succeeded recently. This is a sign of a nonresponding read. - - final long readDelay = timestamps.lastChannelWriteNanoTime() - timestamps.lastChannelReadNanoTime(); - final long readHangDuration = currentTime - timestamps.lastChannelWriteNanoTime(); - - if (readDelay > this.readDelayLimitInNanos && readHangDuration > readHangGracePeriodInNanos) { - - final Optional rntbdContext = requestManager.rntbdContext(); - final int pendingRequestCount = requestManager.pendingRequestCount(); - - logger.warn("{} health check failed due to nonresponding read: {lastChannelWrite: {}, lastChannelRead: {}, " - + "readDelay: {}, readDelayLimit: {}, rntbdContext: {}, pendingRequestCount: {}}", channel, - timestamps.lastChannelWriteNanoTime(), timestamps.lastChannelReadNanoTime(), readDelay, - this.readDelayLimitInNanos, rntbdContext, pendingRequestCount); - + String readIsHangMessage = this.isReadHang(timestamps, currentTime, requestManager, channel); + if (StringUtils.isNotEmpty(readIsHangMessage)) { return promise.setSuccess(Boolean.FALSE); } - if (this.idleConnectionTimeoutInNanos > 0L) { - if (currentTime - timestamps.lastChannelReadNanoTime() > this.idleConnectionTimeoutInNanos) { - return promise.setSuccess(Boolean.FALSE); - } + String idleConnectionValidationMessage = this.idleConnectionValidation(timestamps, currentTime, channel); + if(StringUtils.isNotEmpty(idleConnectionValidationMessage)) { + return promise.setSuccess(Boolean.FALSE); } channel.writeAndFlush(RntbdHealthCheckRequest.MESSAGE).addListener(completed -> { @@ -232,94 +203,127 @@ public Future isHealthyWithFailureReason(final Channel channel) { return promise.setSuccess(RntbdConstants.RntbdHealthCheckResults.SuccessValue); } - // Black hole detection, part 1: - // Treat the channel as unhealthy if the gap between the last attempted write and the last successful write - // grew beyond acceptable limits, unless a write was attempted recently. This is a sign of a nonresponding write. + String writeIsHangMessage = this.isWriteHang(timestamps, currentTime, requestManager, channel); + if (StringUtils.isNotEmpty(writeIsHangMessage)) { + return promise.setSuccess(writeIsHangMessage); + } + + String readIsHangMessage = this.isReadHang(timestamps, currentTime, requestManager, channel); + if (StringUtils.isNotEmpty(readIsHangMessage)) { + return promise.setSuccess(readIsHangMessage); + } + + String idleConnectionValidationMessage = this.idleConnectionValidation(timestamps, currentTime, channel); + if(StringUtils.isNotEmpty(idleConnectionValidationMessage)) { + return promise.setSuccess(idleConnectionValidationMessage); + } + + channel.writeAndFlush(RntbdHealthCheckRequest.MESSAGE).addListener(completed -> { + if (completed.isSuccess()) { + promise.setSuccess(RntbdConstants.RntbdHealthCheckResults.SuccessValue); + } else { + String msg = MessageFormat.format( + "{0} health check request failed due to: {1}", + channel, + completed.cause().toString() + ); + + logger.warn(msg); + promise.setSuccess(msg); + } + }); + + return promise; + } + + private String isWriteHang(Timestamps timestamps, long currentTime, RntbdRequestManager requestManager, Channel channel) { + // Treat the channel as unhealthy if the gap between the last attempted to write and the last successful write + // grew beyond acceptable limits, unless a write was attempted recently. This is a sign of a non-responding write. final long writeDelayInNanos = - timestamps.lastChannelWriteAttemptNanoTime() - timestamps.lastChannelWriteNanoTime(); + timestamps.lastChannelWriteAttemptNanoTime() - timestamps.lastChannelWriteNanoTime(); final long writeHangDurationInNanos = - currentTime - timestamps.lastChannelWriteAttemptNanoTime(); + currentTime - timestamps.lastChannelWriteAttemptNanoTime(); + + String writeHangMessage = StringUtils.EMPTY; if (writeDelayInNanos > this.writeDelayLimitInNanos && writeHangDurationInNanos > writeHangGracePeriodInNanos) { final Optional rntbdContext = requestManager.rntbdContext(); final int pendingRequestCount = requestManager.pendingRequestCount(); - logger.warn("{} health check failed due to nonresponding write: {lastChannelWriteAttemptNanoTime: {}, " + - "lastChannelWriteNanoTime: {}, writeDelayInNanos: {}, writeDelayLimitInNanos: {}, " + - "rntbdContext: {}, pendingRequestCount: {}}", - channel, timestamps.lastChannelWriteAttemptNanoTime(), timestamps.lastChannelWriteNanoTime(), - writeDelayInNanos, this.writeDelayLimitInNanos, rntbdContext, pendingRequestCount); - - String msg = MessageFormat.format( - "{0} health check failed due to nonresponding write: (lastChannelWriteAttemptNanoTime: {1}, " + - "lastChannelWriteNanoTime: {2}, writeDelayInNanos: {3}, writeDelayLimitInNanos: {4}, " + - "rntbdContext: {5}, pendingRequestCount: {6})", - channel, timestamps.lastChannelWriteAttemptNanoTime(), timestamps.lastChannelWriteNanoTime(), - writeDelayInNanos, this.writeDelayLimitInNanos, rntbdContext, pendingRequestCount - ); - - return promise.setSuccess(msg); + writeHangMessage = MessageFormat.format( + "{0} health check failed due to non-responding write: [lastChannelWriteAttemptNanoTime: {1}, " + + "lastChannelWriteNanoTime: {2}, writeDelayInNanos: {3}, writeDelayLimitInNanos: {4}, " + + "rntbdContext: {5}, pendingRequestCount: {6}]", + channel, + timestamps.lastChannelWriteAttemptNanoTime(), + timestamps.lastChannelWriteNanoTime(), + writeDelayInNanos, + this.writeDelayLimitInNanos, + rntbdContext, + pendingRequestCount); + + logger.warn(writeHangMessage); } - // Black hole detection, part 2: + return writeHangMessage; + } + + private String isReadHang(Timestamps timestamps, long currentTime, RntbdRequestManager requestManager, Channel channel) { // Treat the connection as unhealthy if the gap between the last successful write and the last successful read - // grew beyond acceptable limits, unless a write succeeded recently. This is a sign of a nonresponding read. + // grew beyond acceptable limits, unless a write succeeded recently or transitTimeout is below threshold. This is a sign of a non-responding read. final long readDelay = timestamps.lastChannelWriteNanoTime() - timestamps.lastChannelReadNanoTime(); final long readHangDuration = currentTime - timestamps.lastChannelWriteNanoTime(); - if (readDelay > this.readDelayLimitInNanos && readHangDuration > readHangGracePeriodInNanos) { + String readHangMessage = StringUtils.EMPTY; + + if (readDelay > this.readDelayLimitInNanos && + (readHangDuration > readHangGracePeriodInNanos || timestamps.transitTimeoutCount() > this.transitTimeoutDetectionThreshold)) { final Optional rntbdContext = requestManager.rntbdContext(); final int pendingRequestCount = requestManager.pendingRequestCount(); - logger.warn("{} health check failed due to nonresponding read: {lastChannelWrite: {}, lastChannelRead: {}, " - + "readDelay: {}, readDelayLimit: {}, rntbdContext: {}, pendingRequestCount: {}}", channel, - timestamps.lastChannelWriteNanoTime(), timestamps.lastChannelReadNanoTime(), readDelay, - this.readDelayLimitInNanos, rntbdContext, pendingRequestCount); + readHangMessage = MessageFormat.format( + "{0} health check failed due to non-responding read: [lastChannelWrite: {1}, lastChannelRead: {2}, " + + "readDelay: {3}, readDelayLimit: {4}, rntbdContext: {5}, pendingRequestCount: {6}, transitTimeoutCount: {7}]", + channel, + timestamps.lastChannelWriteNanoTime(), + timestamps.lastChannelReadNanoTime(), + readDelay, + this.readDelayLimitInNanos, + rntbdContext, + pendingRequestCount, + timestamps.transitTimeoutCount()); - String msg = MessageFormat.format( - "{0} health check failed due to nonresponding read: (lastChannelWrite: {1}, lastChannelRead: {2}, " - + "readDelay: {3}, readDelayLimit: {4}, rntbdContext: {5}, pendingRequestCount: {6})", channel, - timestamps.lastChannelWriteNanoTime(), timestamps.lastChannelReadNanoTime(), readDelay, - this.readDelayLimitInNanos, rntbdContext, pendingRequestCount - ); - return promise.setSuccess(msg); + logger.warn(readHangMessage); } + return readHangMessage; + } + + private String idleConnectionValidation(Timestamps timestamps, long currentTime, Channel channel) { + String errorMessage = StringUtils.EMPTY; + if (this.idleConnectionTimeoutInNanos > 0L) { if (currentTime - timestamps.lastChannelReadNanoTime() > this.idleConnectionTimeoutInNanos) { - String msg = MessageFormat.format( - "{0} health check failed due to idle connection timeout: (lastChannelWrite: {1}, lastChannelRead: {2}, " - + "idleConnectionTimeout: {3}, currentTime: {4}", channel, - timestamps.lastChannelWriteNanoTime(), timestamps.lastChannelReadNanoTime(), - idleConnectionTimeoutInNanos, currentTime - ); - return promise.setSuccess(msg); + errorMessage = MessageFormat.format( + "{0} health check failed due to idle connection timeout: [lastChannelWrite: {1}, lastChannelRead: {2}, " + + "idleConnectionTimeout: {3}, currentTime: {4}]", + channel, + timestamps.lastChannelWriteNanoTime(), + timestamps.lastChannelReadNanoTime(), + idleConnectionTimeoutInNanos, + currentTime); + + logger.warn(errorMessage); } } - channel.writeAndFlush(RntbdHealthCheckRequest.MESSAGE).addListener(completed -> { - if (completed.isSuccess()) { - promise.setSuccess(RntbdConstants.RntbdHealthCheckResults.SuccessValue); - } else { - logger.warn("{} health check request failed due to:", channel, completed.cause()); - - String msg = MessageFormat.format( - "{0} health check request failed due to: {1}", - channel, - completed.cause().toString() - ); - - promise.setSuccess(msg); - } - }); - - return promise; + return errorMessage; } @Override @@ -345,10 +349,14 @@ static final class Timestamps { private static final AtomicLongFieldUpdater lastWriteAttemptUpdater = newUpdater(Timestamps.class, "lastWriteAttemptNanoTime"); + private static final AtomicIntegerFieldUpdater transitTimeoutCountUpdater = + AtomicIntegerFieldUpdater.newUpdater(Timestamps.class, "transitTimeoutCount"); + private volatile long lastPingNanoTime; private volatile long lastReadNanoTime; private volatile long lastWriteNanoTime; private volatile long lastWriteAttemptNanoTime; + private volatile int transitTimeoutCount; public Timestamps() { } @@ -360,6 +368,7 @@ public Timestamps(Timestamps other) { this.lastReadNanoTime = lastReadUpdater.get(other); this.lastWriteNanoTime = lastWriteUpdater.get(other); this.lastWriteAttemptNanoTime = lastWriteAttemptUpdater.get(other); + this.transitTimeoutCount = transitTimeoutCountUpdater.get(other); } public void channelPingCompleted() { @@ -378,6 +387,13 @@ public void channelWriteCompleted() { lastWriteAttemptUpdater.set(this, System.nanoTime()); } + public void transitTimeout() { + transitTimeoutCountUpdater.incrementAndGet(this); + } + public void resetTransitTimeout() { + transitTimeoutCountUpdater.set(this, 0); + } + @JsonProperty public long lastChannelPingNanoTime() { return lastPingUpdater.get(this); @@ -398,6 +414,11 @@ public long lastChannelWriteAttemptNanoTime() { return lastWriteAttemptUpdater.get(this); } + @JsonProperty + public long transitTimeoutCount() { + return transitTimeoutCountUpdater.get(this); + } + @Override public String toString() { return RntbdObjectMapper.toString(this); diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdEndpoint.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdEndpoint.java index ddff0c5998a3..8cf8029d7ffb 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdEndpoint.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdEndpoint.java @@ -247,6 +247,11 @@ public long sslHandshakeTimeoutInMillis() { return this.options.sslHandshakeTimeoutInMillis(); } + @JsonProperty + public int transitTimeoutDetectionThreshold() { + return this.options.transientTimeoutDetectionThreshold(); + } + @Override public String toString() { return RntbdObjectMapper.toString(this); 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..f5c4279742a4 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 @@ -187,6 +187,9 @@ public void channelRead(final ChannelHandlerContext context, final Object messag this.traceOperation(context, "channelRead"); + this.timestamps.channelReadCompleted(); + this.timestamps.resetTransitTimeout(); // we have got a successful read, so reset the transitTimeout count. + try { if (message.getClass() == RntbdResponse.class) { @@ -230,7 +233,6 @@ public void channelRead(final ChannelHandlerContext context, final Object messag @Override public void channelReadComplete(final ChannelHandlerContext context) { this.traceOperation(context, "channelReadComplete"); - this.timestamps.channelReadCompleted(); context.fireChannelReadComplete(); } @@ -577,7 +579,10 @@ public void write(final ChannelHandlerContext context, final Object message, fin if (message instanceof RntbdRequestRecord) { final RntbdRequestRecord record = (RntbdRequestRecord) message; + this.timestamps.channelWriteAttempted(); + record.setTimestamps(this.timestamps); + record.setSendingRequestHasStarted(); context.write(this.addPendingRequestRecord(context, record), promise).addListener(completed -> { diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestRecord.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestRecord.java index 5847c8f4dfe4..db675142b6a2 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestRecord.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestRecord.java @@ -67,6 +67,7 @@ public abstract class RntbdRequestRecord extends CompletableFuture Date: Mon, 17 Oct 2022 22:10:34 -0700 Subject: [PATCH 2/7] add unit tests --- .../RntbdTransportClient.java | 2 +- .../RntbdClientChannelHealthChecker.java | 9 +- .../directconnectivity/ReflectionUtils.java | 6 + .../RntbdTransportClientTest.java | 5 +- .../RntbdClientChannelHealthCheckerTests.java | 156 ++++++++++++++++++ .../rntbd/RntbdRequestManagerTests.java | 78 +++++++++ 6 files changed, 249 insertions(+), 7 deletions(-) create mode 100644 sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthCheckerTests.java create mode 100644 sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManagerTests.java 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 7c95005afec4..82c799ed7b4b 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 @@ -544,7 +544,7 @@ private Options(final ConnectionPolicy connectionPolicy) { this.tcpKeepIntvl = 1; // Configuration for EpollChannelOption.TCP_KEEPINTVL this.tcpKeepIdle = 30; // Configuration for EpollChannelOption.TCP_KEEPIDLE this.sslHandshakeTimeoutMinDuration = Duration.ofSeconds(5); - this.transientTimeoutDetectionThreshold = 5; + this.transientTimeoutDetectionThreshold = 3; this.preferTcpNative = true; } diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthChecker.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthChecker.java index e34ecd08196e..a8e443dcd75e 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthChecker.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthChecker.java @@ -15,7 +15,6 @@ import java.text.MessageFormat; import java.util.Optional; -import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; import java.util.concurrent.atomic.AtomicLongFieldUpdater; import static com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdReporter.reportIssueUnless; @@ -335,7 +334,7 @@ public String toString() { // region Types - static final class Timestamps { + public static final class Timestamps { private static final AtomicLongFieldUpdater lastPingUpdater = newUpdater(Timestamps.class, "lastPingNanoTime"); @@ -349,14 +348,14 @@ static final class Timestamps { private static final AtomicLongFieldUpdater lastWriteAttemptUpdater = newUpdater(Timestamps.class, "lastWriteAttemptNanoTime"); - private static final AtomicIntegerFieldUpdater transitTimeoutCountUpdater = - AtomicIntegerFieldUpdater.newUpdater(Timestamps.class, "transitTimeoutCount"); + private static final AtomicLongFieldUpdater transitTimeoutCountUpdater = + newUpdater(Timestamps.class, "transitTimeoutCount"); private volatile long lastPingNanoTime; private volatile long lastReadNanoTime; private volatile long lastWriteNanoTime; private volatile long lastWriteAttemptNanoTime; - private volatile int transitTimeoutCount; + private volatile long transitTimeoutCount; public Timestamps() { } diff --git a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/ReflectionUtils.java b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/ReflectionUtils.java index eceeea23b727..8c9a61a2822c 100644 --- a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/ReflectionUtils.java +++ b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/ReflectionUtils.java @@ -11,6 +11,8 @@ import com.azure.cosmos.implementation.ApiType; import com.azure.cosmos.implementation.AsyncDocumentClient; import com.azure.cosmos.implementation.ClientSideRequestStatistics; +import com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdClientChannelHealthChecker; +import com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdRequestManager; import com.azure.cosmos.models.CosmosClientTelemetryConfig; import com.azure.cosmos.implementation.ConnectionPolicy; import com.azure.cosmos.implementation.DocumentCollection; @@ -402,4 +404,8 @@ public static AtomicReference getHealthStatus(Uri uri) { public static Set getReplicaValidationScopes(GatewayAddressCache gatewayAddressCache) { return get(Set.class, gatewayAddressCache, "replicaValidationScopes"); } + + public static RntbdClientChannelHealthChecker.Timestamps getTimestamps(RntbdRequestManager rntbdRequestManager) { + return get(RntbdClientChannelHealthChecker.Timestamps.class, rntbdRequestManager, "timestamps"); + } } 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 d354542b69c7..425557bbcfb7 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 @@ -101,6 +101,7 @@ public final class RntbdTransportClientTest { private static final Uri physicalAddress = new Uri("rntbd://host:10251/replica-path/"); private static final Duration requestTimeout = Duration.ofSeconds(1000); private static final int sslHandshakeTimeoutInMillis = 5000; + private static final int transitTimeoutDetectionThreshold = 3; @DataProvider(name = "fromMockedNetworkFailureToExpectedDocumentClientException") public Object[][] fromMockedNetworkFailureToExpectedDocumentClientException() { @@ -737,6 +738,7 @@ public void transportClientDefaultOptionsTests() { .build(); assertEquals(options.sslHandshakeTimeoutInMillis(), sslHandshakeTimeoutInMillis); + assertEquals(options.transientTimeoutDetectionThreshold(), transitTimeoutDetectionThreshold); } // TODO: add validations for other properties @@ -744,7 +746,7 @@ public void transportClientDefaultOptionsTests() { @Test(enabled = false, groups = "unit") public void transportClientCustomizedOptionsTests() { try { - System.setProperty("azure.cosmos.directTcp.defaultOptions", "{\"sslHandshakeTimeoutMinDuration\":\"PT15S\"}"); + System.setProperty("azure.cosmos.directTcp.defaultOptions", "{\"sslHandshakeTimeoutMinDuration\":\"PT15S\",\"transitTimeoutDetectionThreshold\":\"10\" }"); ConnectionPolicy connectionPolicy = new ConnectionPolicy(DirectConnectionConfig.getDefaultConfig()); UserAgentContainer userAgentContainer = new UserAgentContainer(); @@ -754,6 +756,7 @@ public void transportClientCustomizedOptionsTests() { .build(); assertEquals(options.sslHandshakeTimeoutInMillis(), Duration.ofSeconds(15).toMillis()); + assertEquals(options.transientTimeoutDetectionThreshold(), 10); } finally { System.clearProperty("azure.cosmos.directTcp.defaultOptions"); diff --git a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthCheckerTests.java b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthCheckerTests.java new file mode 100644 index 000000000000..2bf089b47508 --- /dev/null +++ b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthCheckerTests.java @@ -0,0 +1,156 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. + +package com.azure.cosmos.implementation.directconnectivity.rntbd; + +import com.azure.cosmos.implementation.ConnectionPolicy; +import com.azure.cosmos.implementation.directconnectivity.RntbdTransportClient; +import io.netty.channel.Channel; +import io.netty.channel.ChannelPipeline; +import io.netty.channel.DefaultEventLoop; +import io.netty.channel.SingleThreadEventLoop; +import io.netty.handler.logging.LogLevel; +import io.netty.handler.ssl.SslContext; +import io.netty.util.concurrent.Future; +import org.mockito.Mockito; +import org.testng.annotations.DataProvider; +import org.testng.annotations.Test; + +import static org.assertj.core.api.AssertionsForClassTypes.assertThat; + +public class RntbdClientChannelHealthCheckerTests { + private static final long writeHangGracePeriodInNanos = 2L * 1_000_000_000L; + private static final long readHangGracePeriodInNanos = (45L + 10L) * 1_000_000_000L; + + @DataProvider + public static Object[][] isHealthyWithReasonArgs() { + return new Object[][]{ + // expect failureReason + { false }, + { true } + }; + } + + @Test(groups = { "unit" }, dataProvider = "isHealthyWithReasonArgs") + public void isHealthyForWriteHangTests(boolean withFailureReason) { + SslContext sslContextMock = Mockito.mock(SslContext.class); + + RntbdEndpoint.Config config = new RntbdEndpoint.Config( + new RntbdTransportClient.Options.Builder(ConnectionPolicy.getDefaultPolicy()).build(), + sslContextMock, + LogLevel.INFO); + + RntbdClientChannelHealthChecker healthChecker = new RntbdClientChannelHealthChecker(config); + Channel channelMock = Mockito.mock(Channel.class); + ChannelPipeline channelPipelineMock = Mockito.mock(ChannelPipeline.class); + RntbdRequestManager rntbdRequestManagerMock = Mockito.mock(RntbdRequestManager.class); + SingleThreadEventLoop eventLoopMock = new DefaultEventLoop(); + RntbdClientChannelHealthChecker.Timestamps timestampsMock = Mockito.mock(RntbdClientChannelHealthChecker.Timestamps.class); + + Mockito.when(channelMock.pipeline()).thenReturn(channelPipelineMock); + Mockito.when(channelPipelineMock.get(RntbdRequestManager.class)).thenReturn(rntbdRequestManagerMock); + Mockito.when(channelMock.eventLoop()).thenReturn(eventLoopMock); + Mockito.when(rntbdRequestManagerMock.snapshotTimestamps()).thenReturn(timestampsMock); + + long lastChannelWriteAttemptNanoTime = System.nanoTime() - writeHangGracePeriodInNanos - 10; + long lastChannelWriteNanoTime = lastChannelWriteAttemptNanoTime - config.sendHangDetectionTimeInNanos() - 10; + + Mockito.when(timestampsMock.lastChannelWriteAttemptNanoTime()).thenReturn(lastChannelWriteAttemptNanoTime); + Mockito.when(timestampsMock.lastChannelWriteNanoTime()).thenReturn(lastChannelWriteNanoTime); + + if (withFailureReason) { + Future healthyResult = healthChecker.isHealthyWithFailureReason(channelMock); + assertThat(healthyResult.isSuccess()).isTrue(); + assertThat(healthyResult.getNow()).isNotEqualTo(RntbdConstants.RntbdHealthCheckResults.SuccessValue); + assertThat(healthyResult.getNow().contains("health check failed due to non-responding write")); + } else { + Future healthyResult = healthChecker.isHealthy(channelMock); + assertThat(healthyResult.isSuccess()).isTrue(); + assertThat(healthyResult.getNow()).isFalse(); + } + } + + @Test(groups = { "unit" }, dataProvider = "isHealthyWithReasonArgs") + public void isHealthyForReadHangTests(boolean withFailureReason) { + SslContext sslContextMock = Mockito.mock(SslContext.class); + + RntbdEndpoint.Config config = new RntbdEndpoint.Config( + new RntbdTransportClient.Options.Builder(ConnectionPolicy.getDefaultPolicy()).build(), + sslContextMock, + LogLevel.INFO); + + RntbdClientChannelHealthChecker healthChecker = new RntbdClientChannelHealthChecker(config); + Channel channelMock = Mockito.mock(Channel.class); + ChannelPipeline channelPipelineMock = Mockito.mock(ChannelPipeline.class); + RntbdRequestManager rntbdRequestManagerMock = Mockito.mock(RntbdRequestManager.class); + SingleThreadEventLoop eventLoopMock = new DefaultEventLoop(); + RntbdClientChannelHealthChecker.Timestamps timestampsMock = Mockito.mock(RntbdClientChannelHealthChecker.Timestamps.class); + + Mockito.when(channelMock.pipeline()).thenReturn(channelPipelineMock); + Mockito.when(channelPipelineMock.get(RntbdRequestManager.class)).thenReturn(rntbdRequestManagerMock); + Mockito.when(channelMock.eventLoop()).thenReturn(eventLoopMock); + Mockito.when(rntbdRequestManagerMock.snapshotTimestamps()).thenReturn(timestampsMock); + + long lastChannelWriteNanoTime = System.nanoTime() - readHangGracePeriodInNanos - 10; + long lastChannelWriteAttemptNanoTime = lastChannelWriteNanoTime; + long lastChannelReadNanoTime = lastChannelWriteNanoTime - config.receiveHangDetectionTimeInNanos() - 10; + + Mockito.when(timestampsMock.lastChannelWriteAttemptNanoTime()).thenReturn(lastChannelWriteAttemptNanoTime); + Mockito.when(timestampsMock.lastChannelWriteNanoTime()).thenReturn(lastChannelWriteNanoTime); + Mockito.when(timestampsMock.lastChannelReadNanoTime()).thenReturn(lastChannelReadNanoTime); + + if (withFailureReason) { + Future healthyResult = healthChecker.isHealthyWithFailureReason(channelMock); + assertThat(healthyResult.isSuccess()).isTrue(); + assertThat(healthyResult.getNow()).isNotEqualTo(RntbdConstants.RntbdHealthCheckResults.SuccessValue); + assertThat(healthyResult.getNow().contains("health check failed due to non-responding read")); + } else { + Future healthyResult = healthChecker.isHealthy(channelMock); + assertThat(healthyResult.isSuccess()).isTrue(); + assertThat(healthyResult.getNow()).isFalse(); + } + } + + @Test(groups = { "unit" }, dataProvider = "isHealthyWithReasonArgs") + public void isHealthyForReadHangWithTransitTimeoutTests(boolean withFailureReason) { + SslContext sslContextMock = Mockito.mock(SslContext.class); + + RntbdEndpoint.Config config = new RntbdEndpoint.Config( + new RntbdTransportClient.Options.Builder(ConnectionPolicy.getDefaultPolicy()).build(), + sslContextMock, + LogLevel.INFO); + + RntbdClientChannelHealthChecker healthChecker = new RntbdClientChannelHealthChecker(config); + Channel channelMock = Mockito.mock(Channel.class); + ChannelPipeline channelPipelineMock = Mockito.mock(ChannelPipeline.class); + RntbdRequestManager rntbdRequestManagerMock = Mockito.mock(RntbdRequestManager.class); + SingleThreadEventLoop eventLoopMock = new DefaultEventLoop(); + RntbdClientChannelHealthChecker.Timestamps timestampsMock = Mockito.mock(RntbdClientChannelHealthChecker.Timestamps.class); + + Mockito.when(channelMock.pipeline()).thenReturn(channelPipelineMock); + Mockito.when(channelPipelineMock.get(RntbdRequestManager.class)).thenReturn(rntbdRequestManagerMock); + Mockito.when(channelMock.eventLoop()).thenReturn(eventLoopMock); + Mockito.when(rntbdRequestManagerMock.snapshotTimestamps()).thenReturn(timestampsMock); + + long lastChannelWriteNanoTime = System.nanoTime(); + long lastChannelWriteAttemptNanoTime = lastChannelWriteNanoTime; + long lastChannelReadNanoTime = lastChannelWriteNanoTime - config.receiveHangDetectionTimeInNanos() - 10; + long transitTimeoutCount = config.transitTimeoutDetectionThreshold() + 1; + + Mockito.when(timestampsMock.lastChannelWriteAttemptNanoTime()).thenReturn(lastChannelWriteAttemptNanoTime); + Mockito.when(timestampsMock.lastChannelWriteNanoTime()).thenReturn(lastChannelWriteNanoTime); + Mockito.when(timestampsMock.lastChannelReadNanoTime()).thenReturn(lastChannelReadNanoTime); + Mockito.when(timestampsMock.transitTimeoutCount()).thenReturn(transitTimeoutCount); + + if (withFailureReason) { + Future healthyResult = healthChecker.isHealthyWithFailureReason(channelMock); + assertThat(healthyResult.isSuccess()).isTrue(); + assertThat(healthyResult.getNow()).isNotEqualTo(RntbdConstants.RntbdHealthCheckResults.SuccessValue); + assertThat(healthyResult.getNow().contains("health check failed due to non-responding read")); + } else { + Future healthyResult = healthChecker.isHealthy(channelMock); + assertThat(healthyResult.isSuccess()).isTrue(); + assertThat(healthyResult.getNow()).isFalse(); + } + } +} diff --git a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManagerTests.java b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManagerTests.java new file mode 100644 index 000000000000..fc722620d0e8 --- /dev/null +++ b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdRequestManagerTests.java @@ -0,0 +1,78 @@ +// Copyright (c) Microsoft Corporation. All rights reserved. +// Licensed under the MIT License. + +package com.azure.cosmos.implementation.directconnectivity.rntbd; + +import com.azure.cosmos.implementation.ConnectionPolicy; +import com.azure.cosmos.implementation.OperationType; +import com.azure.cosmos.implementation.ResourceType; +import com.azure.cosmos.implementation.RxDocumentServiceRequest; +import com.azure.cosmos.implementation.directconnectivity.ReflectionUtils; +import com.azure.cosmos.implementation.directconnectivity.RntbdTransportClient; +import com.azure.cosmos.implementation.directconnectivity.Uri; +import io.netty.channel.ChannelFuture; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelPromise; +import io.netty.handler.logging.LogLevel; +import io.netty.handler.ssl.SslContext; +import org.mockito.Mockito; +import org.testng.annotations.Test; + +import java.net.URI; +import java.net.URISyntaxException; +import java.time.Duration; + +import static com.azure.cosmos.implementation.TestUtils.mockDiagnosticsClientContext; +import static org.assertj.core.api.AssertionsForClassTypes.assertThat; + +public class RntbdRequestManagerTests { + + @Test + public void transitTimeoutTimestampTests() throws URISyntaxException { + SslContext sslContextMock = Mockito.mock(SslContext.class); + RntbdEndpoint.Config config = new RntbdEndpoint.Config( + new RntbdTransportClient.Options.Builder(ConnectionPolicy.getDefaultPolicy()).build(), + sslContextMock, + LogLevel.INFO); + RntbdClientChannelHealthChecker healthChecker = new RntbdClientChannelHealthChecker(config); + + RntbdConnectionStateListener connectionStateListener = Mockito.mock(RntbdConnectionStateListener.class); + + RntbdRequestManager rntbdRequestManager = new RntbdRequestManager( + healthChecker, + 30, + connectionStateListener, + Duration.ofSeconds(1).toNanos()); + RntbdClientChannelHealthChecker.Timestamps timestamps = ReflectionUtils.getTimestamps(rntbdRequestManager); + + ChannelHandlerContext channelHandlerContext = Mockito.mock(ChannelHandlerContext.class); + ChannelFuture channelFuture = Mockito.mock(ChannelFuture.class); + Mockito.when(channelHandlerContext.write(Mockito.any(), Mockito.any())).thenReturn(channelFuture); + Mockito.when(channelFuture.addListener(Mockito.any())).thenReturn(channelFuture); + + RntbdRequestArgs requestArgs = new RntbdRequestArgs( + RxDocumentServiceRequest.create(mockDiagnosticsClientContext(), OperationType.Read, ResourceType.Document), + new Uri(new URI("http://localhost/replica-path").toString()) + ); + long requestTimeoutInNanos = Duration.ofMinutes(5).toNanos(); + RntbdRequestTimer requestTimer = new RntbdRequestTimer(requestTimeoutInNanos, requestTimeoutInNanos); + RntbdRequestRecord rntbdRequestRecord = new AsyncRntbdRequestRecord(requestArgs, requestTimer); + + ChannelPromise promise = Mockito.mock(ChannelPromise.class); + + // Test transitTimeout is 0 at start point + rntbdRequestManager.write(channelHandlerContext, rntbdRequestRecord, promise); + assertThat(timestamps.transitTimeoutCount()).isZero(); + + // Test when a transit timeout happens, the transitTimeoutCount is increased + rntbdRequestRecord.expire(); + assertThat(timestamps.transitTimeoutCount()).isOne(); + + // Test when there is channelRead, transitTimeout is cleared out + Mockito.when(channelHandlerContext.flush()).thenReturn(channelHandlerContext); + ChannelFuture closeChannelFuture = Mockito.mock(ChannelFuture.class); + Mockito.when(channelHandlerContext.close()).thenReturn(closeChannelFuture); + rntbdRequestManager.channelRead(channelHandlerContext, rntbdRequestRecord); + assertThat(timestamps.transitTimeoutCount()).isZero(); + } +} From 840d61022d52342214e11e1bcadbed7b39600ef7 Mon Sep 17 00:00:00 2001 From: annie-mac Date: Mon, 17 Oct 2022 22:14:34 -0700 Subject: [PATCH 3/7] add changelog --- sdk/cosmos/azure-cosmos/CHANGELOG.md | 1 + 1 file changed, 1 insertion(+) diff --git a/sdk/cosmos/azure-cosmos/CHANGELOG.md b/sdk/cosmos/azure-cosmos/CHANGELOG.md index 9e6ab05af170..4d82a0463b4f 100644 --- a/sdk/cosmos/azure-cosmos/CHANGELOG.md +++ b/sdk/cosmos/azure-cosmos/CHANGELOG.md @@ -9,6 +9,7 @@ #### Bugs Fixed #### Other Changes +* Added improvement in `RntbdClientChannelHealthChecker` for continuous transit timeout. - See [PR 31544](https://github.com/Azure/azure-sdk-for-java/pull/31544) ### 4.38.0 (2022-10-12) #### Features Added From b59227f9e4e910e8b8c14443306a8615ffe60927 Mon Sep 17 00:00:00 2001 From: annie-mac Date: Tue, 18 Oct 2022 11:16:32 -0700 Subject: [PATCH 4/7] update javadoc for connectionEndpointRediscoveryEnabled --- .../main/java/com/azure/cosmos/DirectConnectionConfig.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/DirectConnectionConfig.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/DirectConnectionConfig.java index d9178298f4af..17729e6c7ca7 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/DirectConnectionConfig.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/DirectConnectionConfig.java @@ -60,7 +60,7 @@ public DirectConnectionConfig() { *

* The connection endpoint rediscovery feature is designed to reduce and spread-out latency spikes that may occur during maintenance operations. * - * By default, connection endpoint rediscovery is disabled. + * By default, connection endpoint rediscovery is enabled. * * @return {@code true} if Direct TCP connection endpoint rediscovery is enabled; {@code false} otherwise. */ @@ -73,7 +73,7 @@ public boolean isConnectionEndpointRediscoveryEnabled() { *

* The connection endpoint rediscovery feature is designed to reduce and spread-out latency spikes that may occur during maintenance operations. * - * By default, connection endpoint rediscovery is disabled. + * By default, connection endpoint rediscovery is enabled. * * @param connectionEndpointRediscoveryEnabled {@code true} if connection endpoint rediscovery is enabled; {@code * false} otherwise. From b3211f981dc55813e3290e701b2f67416bb68a06 Mon Sep 17 00:00:00 2001 From: annie-mac Date: Tue, 18 Oct 2022 15:21:28 -0700 Subject: [PATCH 5/7] update tests --- .../rntbd/RntbdClientChannelHealthChecker.java | 2 +- .../rntbd/RntbdClientChannelHealthCheckerTests.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthChecker.java b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthChecker.java index a8e443dcd75e..99e85728a4a1 100644 --- a/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthChecker.java +++ b/sdk/cosmos/azure-cosmos/src/main/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthChecker.java @@ -280,7 +280,7 @@ private String isReadHang(Timestamps timestamps, long currentTime, RntbdRequestM String readHangMessage = StringUtils.EMPTY; if (readDelay > this.readDelayLimitInNanos && - (readHangDuration > readHangGracePeriodInNanos || timestamps.transitTimeoutCount() > this.transitTimeoutDetectionThreshold)) { + (readHangDuration > readHangGracePeriodInNanos || timestamps.transitTimeoutCount() >= this.transitTimeoutDetectionThreshold)) { final Optional rntbdContext = requestManager.rntbdContext(); final int pendingRequestCount = requestManager.pendingRequestCount(); diff --git a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthCheckerTests.java b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthCheckerTests.java index 2bf089b47508..70d73cb9592f 100644 --- a/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthCheckerTests.java +++ b/sdk/cosmos/azure-cosmos/src/test/java/com/azure/cosmos/implementation/directconnectivity/rntbd/RntbdClientChannelHealthCheckerTests.java @@ -135,7 +135,7 @@ public void isHealthyForReadHangWithTransitTimeoutTests(boolean withFailureReaso long lastChannelWriteNanoTime = System.nanoTime(); long lastChannelWriteAttemptNanoTime = lastChannelWriteNanoTime; long lastChannelReadNanoTime = lastChannelWriteNanoTime - config.receiveHangDetectionTimeInNanos() - 10; - long transitTimeoutCount = config.transitTimeoutDetectionThreshold() + 1; + long transitTimeoutCount = config.transitTimeoutDetectionThreshold(); Mockito.when(timestampsMock.lastChannelWriteAttemptNanoTime()).thenReturn(lastChannelWriteAttemptNanoTime); Mockito.when(timestampsMock.lastChannelWriteNanoTime()).thenReturn(lastChannelWriteNanoTime); From 8032d107c1075ee46bdf24b064f2baaaaa017013 Mon Sep 17 00:00:00 2001 From: annie-mac Date: Wed, 2 Nov 2022 22:35:50 -0700 Subject: [PATCH 6/7] resolve comments --- .../directconnectivity/RntbdTransportClient.java | 14 ++++++++------ .../rntbd/RntbdRequestManager.java | 4 +--- 2 files changed, 9 insertions(+), 9 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 82c799ed7b4b..d139c6a92332 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 @@ -14,11 +14,7 @@ import com.azure.cosmos.implementation.RxDocumentServiceRequest; import com.azure.cosmos.implementation.UserAgentContainer; import com.azure.cosmos.implementation.clienttelemetry.ClientTelemetry; -import com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdEndpoint; -import com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdObjectMapper; -import com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdRequestArgs; -import com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdRequestRecord; -import com.azure.cosmos.implementation.directconnectivity.rntbd.RntbdServiceEndpoint; +import com.azure.cosmos.implementation.directconnectivity.rntbd.*; import com.azure.cosmos.implementation.guava25.base.Strings; import com.fasterxml.jackson.annotation.JsonCreator; import com.fasterxml.jackson.annotation.JsonIgnore; @@ -474,6 +470,11 @@ public static final class Options { @JsonProperty() private final Duration sslHandshakeTimeoutMinDuration; + /** + * This property will be used in {@link RntbdClientChannelHealthChecker} to determine whether there is a readHang. + * If there is no successful reads for up to receiveHangDetectionTime, and the number of consecutive timeout has also reached this config, + * then SDK is going to treat the channel as unhealthy and close it. + */ @JsonProperty() private final int transientTimeoutDetectionThreshold; // endregion @@ -715,7 +716,8 @@ public String toDiagnosticsString() { * "requestTimerResolution": "PT100MS", * "sendHangDetectionTime": "PT10S", * "shutdownTimeout": "PT15S", - * "threadCount": 16 + * "threadCount": 16, + * "transientTimeoutDetectionThreshold": 3 * }} * * 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 f5c4279742a4..140165d9c4b5 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 @@ -186,13 +186,10 @@ public void channelInactive(final ChannelHandlerContext context) { public void channelRead(final ChannelHandlerContext context, final Object message) { this.traceOperation(context, "channelRead"); - - this.timestamps.channelReadCompleted(); this.timestamps.resetTransitTimeout(); // we have got a successful read, so reset the transitTimeout count. try { if (message.getClass() == RntbdResponse.class) { - try { this.messageReceived(context, (RntbdResponse) message); } catch (CorruptedFrameException error) { @@ -233,6 +230,7 @@ public void channelRead(final ChannelHandlerContext context, final Object messag @Override public void channelReadComplete(final ChannelHandlerContext context) { this.traceOperation(context, "channelReadComplete"); + this.timestamps.channelReadCompleted(); context.fireChannelReadComplete(); } From 9bebffea7fa20c3daf09a345512cee7f00e22c18 Mon Sep 17 00:00:00 2001 From: annie-mac Date: Thu, 3 Nov 2022 07:56:26 -0700 Subject: [PATCH 7/7] 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 ff7b48db31c4..e3ffe7479843 100644 --- a/sdk/cosmos/azure-cosmos/CHANGELOG.md +++ b/sdk/cosmos/azure-cosmos/CHANGELOG.md @@ -8,9 +8,9 @@ #### 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) +* Added improvement in `RntbdClientChannelHealthChecker` for detecting continuous transit timeout. - See [PR 31544](https://github.com/Azure/azure-sdk-for-java/pull/31544) #### Other Changes -* Added improvement in `RntbdClientChannelHealthChecker` for continuous transit timeout. - See [PR 31544](https://github.com/Azure/azure-sdk-for-java/pull/31544) ### 4.38.1 (2022-10-21) #### Other Changes