From 0cf93fb463be5dfd486144e3226dad1e0db58947 Mon Sep 17 00:00:00 2001 From: Kushagra Thapar Date: Sat, 4 Jan 2020 10:18:45 -0800 Subject: [PATCH] Added ability to pass connection policy details to http transport client to better configure it --- .../com/azure/data/cosmos/internal/Configs.java | 5 ----- .../cosmos/internal/RxDocumentClientImpl.java | 2 +- .../directconnectivity/HttpTransportClient.java | 15 +++++++++------ .../directconnectivity/StoreClientFactory.java | 7 ++++--- .../data/cosmos/internal/http/HttpClient.java | 2 +- .../HttpTransportClientTest.java | 5 +++-- 6 files changed, 18 insertions(+), 18 deletions(-) diff --git a/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/Configs.java b/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/Configs.java index a78db3824f1e..d881633d35bb 100644 --- a/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/Configs.java +++ b/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/Configs.java @@ -58,7 +58,6 @@ public class Configs { // Reactor Netty Constants private static final int MAX_IDLE_CONNECTION_TIMEOUT_IN_MILLIS = 60 * 1000; private static final int CONNECTION_ACQUIRE_TIMEOUT_IN_MILLIS = 45 * 1000; - private static final int REACTOR_NETTY_MAX_CONNECTION_POOL_SIZE = 1000; private static final String REACTOR_NETTY_CONNECTION_POOL_NAME = "reactor-netty-connection-pool"; public Configs() { @@ -159,10 +158,6 @@ public int getConnectionAcquireTimeoutInMillis() { return CONNECTION_ACQUIRE_TIMEOUT_IN_MILLIS; } - public int getReactorNettyMaxConnectionPoolSize() { - return REACTOR_NETTY_MAX_CONNECTION_POOL_SIZE; - } - private static String getJVMConfigAsString(String propName, String defaultValue) { String propValue = System.getProperty(propName); return StringUtils.defaultString(propValue, defaultValue); diff --git a/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/RxDocumentClientImpl.java b/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/RxDocumentClientImpl.java index 5dfc4e41122c..87484013b567 100644 --- a/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/RxDocumentClientImpl.java +++ b/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/RxDocumentClientImpl.java @@ -281,7 +281,7 @@ private void initializeDirectConnectivity() { this.storeClientFactory = new StoreClientFactory( this.configs, - this.connectionPolicy.requestTimeoutInMillis() / 1000, + this.connectionPolicy, // this.maxConcurrentConnectionOpenRequests, 0, this.userAgentContainer diff --git a/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/directconnectivity/HttpTransportClient.java b/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/directconnectivity/HttpTransportClient.java index d33d6792e55a..7d0a4a9f3d41 100644 --- a/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/directconnectivity/HttpTransportClient.java +++ b/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/directconnectivity/HttpTransportClient.java @@ -6,6 +6,7 @@ import com.azure.data.cosmos.BadRequestException; import com.azure.data.cosmos.BridgeInternal; import com.azure.data.cosmos.ConflictException; +import com.azure.data.cosmos.ConnectionPolicy; import com.azure.data.cosmos.CosmosClientException; import com.azure.data.cosmos.ForbiddenException; import com.azure.data.cosmos.GoneException; @@ -67,18 +68,20 @@ public class HttpTransportClient extends TransportClient { private final Map defaultHeaders; private final Configs configs; - HttpClient createHttpClient(int requestTimeout) { + HttpClient createHttpClient(ConnectionPolicy connectionPolicy) { // TODO: use one instance of SSL context everywhere - HttpClientConfig httpClientConfig = new HttpClientConfig(this.configs); - httpClientConfig.withRequestTimeoutInMillis(requestTimeout * 1000); - httpClientConfig.withPoolSize(configs.getDirectHttpsMaxConnectionLimit()); + HttpClientConfig httpClientConfig = new HttpClientConfig(this.configs) + .withMaxIdleConnectionTimeoutInMillis(connectionPolicy.idleConnectionTimeoutInMillis()) + .withPoolSize(connectionPolicy.maxPoolSize()) + .withHttpProxy(connectionPolicy.proxy()) + .withRequestTimeoutInMillis(connectionPolicy.requestTimeoutInMillis()); return HttpClient.createFixed(httpClientConfig); } - public HttpTransportClient(Configs configs, int requestTimeout, UserAgentContainer userAgent) { + public HttpTransportClient(Configs configs, ConnectionPolicy connectionPolicy, UserAgentContainer userAgent) { this.configs = configs; - this.httpClient = createHttpClient(requestTimeout); + this.httpClient = createHttpClient(connectionPolicy); this.defaultHeaders = new HashMap<>(); diff --git a/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/directconnectivity/StoreClientFactory.java b/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/directconnectivity/StoreClientFactory.java index dddad8cf95d5..66aef3a1a22b 100644 --- a/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/directconnectivity/StoreClientFactory.java +++ b/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/directconnectivity/StoreClientFactory.java @@ -3,6 +3,7 @@ package com.azure.data.cosmos.internal.directconnectivity; +import com.azure.data.cosmos.ConnectionPolicy; import com.azure.data.cosmos.internal.Configs; import com.azure.data.cosmos.internal.IAuthorizationTokenProvider; import com.azure.data.cosmos.internal.SessionContainer; @@ -22,17 +23,17 @@ public class StoreClientFactory implements AutoCloseable { public StoreClientFactory( Configs configs, - int requestTimeoutInSeconds, + ConnectionPolicy connectionPolicy, int maxConcurrentConnectionOpenRequests, UserAgentContainer userAgent) { this.configs = configs; this.protocol = configs.getProtocol(); - this.requestTimeoutInSeconds = requestTimeoutInSeconds; + this.requestTimeoutInSeconds = connectionPolicy.requestTimeoutInMillis() / 1000; this.maxConcurrentConnectionOpenRequests = maxConcurrentConnectionOpenRequests; if (protocol == Protocol.HTTPS) { - this.transportClient = new HttpTransportClient(configs, requestTimeoutInSeconds, userAgent); + this.transportClient = new HttpTransportClient(configs, connectionPolicy, userAgent); } else if (protocol == Protocol.TCP){ this.transportClient = new RntbdTransportClient(configs, requestTimeoutInSeconds, userAgent); } else { diff --git a/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/http/HttpClient.java b/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/http/HttpClient.java index 9ea150006cf7..d92b1ac74d13 100644 --- a/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/http/HttpClient.java +++ b/sdk/cosmos/microsoft-azure-cosmos/src/main/java/com/azure/data/cosmos/internal/http/HttpClient.java @@ -36,7 +36,7 @@ static HttpClient createFixed(HttpClientConfig httpClientConfig) { } // Default pool size - Integer maxPoolSize = httpClientConfig.getConfigs().getReactorNettyMaxConnectionPoolSize(); + Integer maxPoolSize = httpClientConfig.getConfigs().getDirectHttpsMaxConnectionLimit(); if (httpClientConfig.getMaxPoolSize() != null) { maxPoolSize = httpClientConfig.getMaxPoolSize(); } diff --git a/sdk/cosmos/microsoft-azure-cosmos/src/test/java/com/azure/data/cosmos/internal/directconnectivity/HttpTransportClientTest.java b/sdk/cosmos/microsoft-azure-cosmos/src/test/java/com/azure/data/cosmos/internal/directconnectivity/HttpTransportClientTest.java index bfebc6844a51..f1148845de5e 100644 --- a/sdk/cosmos/microsoft-azure-cosmos/src/test/java/com/azure/data/cosmos/internal/directconnectivity/HttpTransportClientTest.java +++ b/sdk/cosmos/microsoft-azure-cosmos/src/test/java/com/azure/data/cosmos/internal/directconnectivity/HttpTransportClientTest.java @@ -5,6 +5,7 @@ import com.azure.data.cosmos.BadRequestException; import com.azure.data.cosmos.ConflictException; +import com.azure.data.cosmos.ConnectionPolicy; import com.azure.data.cosmos.ForbiddenException; import com.azure.data.cosmos.GoneException; import com.azure.data.cosmos.InternalServerErrorException; @@ -109,11 +110,11 @@ public static HttpTransportClient getHttpTransportClientUnderTest(int requestTim HttpClient httpClient) { class HttpTransportClientUnderTest extends HttpTransportClient { public HttpTransportClientUnderTest(int requestTimeout, UserAgentContainer userAgent) { - super(configs, requestTimeout, userAgent); + super(configs, ConnectionPolicy.defaultPolicy().requestTimeoutInMillis(requestTimeout * 1000), userAgent); } @Override - HttpClient createHttpClient(int requestTimeout) { + HttpClient createHttpClient(ConnectionPolicy connectionPolicy) { return httpClient; } }