diff --git a/.gitignore b/.gitignore index cc29f27cb..201acaa5f 100644 --- a/.gitignore +++ b/.gitignore @@ -6,6 +6,7 @@ .project .settings .classpath +.factorypath # Ignore all build/dist directories target diff --git a/docker-java-transport-httpclient5/src/main/java/com/github/dockerjava/httpclient5/ApacheDockerHttpClient.java b/docker-java-transport-httpclient5/src/main/java/com/github/dockerjava/httpclient5/ApacheDockerHttpClient.java index cf2b7300d..2c1890f80 100644 --- a/docker-java-transport-httpclient5/src/main/java/com/github/dockerjava/httpclient5/ApacheDockerHttpClient.java +++ b/docker-java-transport-httpclient5/src/main/java/com/github/dockerjava/httpclient5/ApacheDockerHttpClient.java @@ -13,6 +13,8 @@ public static final class Builder { private SSLConfig sslConfig = null; + private int maxConnections = Integer.MAX_VALUE; + public Builder dockerHost(URI value) { this.dockerHost = Objects.requireNonNull(value, "dockerHost"); return this; @@ -23,13 +25,18 @@ public Builder sslConfig(SSLConfig value) { return this; } + public Builder maxConnections(int value) { + this.maxConnections = value; + return this; + } + public ApacheDockerHttpClient build() { Objects.requireNonNull(dockerHost, "dockerHost"); - return new ApacheDockerHttpClient(dockerHost, sslConfig); + return new ApacheDockerHttpClient(dockerHost, sslConfig, maxConnections); } } - private ApacheDockerHttpClient(URI dockerHost, SSLConfig sslConfig) { - super(dockerHost, sslConfig); + private ApacheDockerHttpClient(URI dockerHost, SSLConfig sslConfig, int maxConnections) { + super(dockerHost, sslConfig, maxConnections); } } diff --git a/docker-java-transport-httpclient5/src/main/java/com/github/dockerjava/httpclient5/ApacheDockerHttpClientImpl.java b/docker-java-transport-httpclient5/src/main/java/com/github/dockerjava/httpclient5/ApacheDockerHttpClientImpl.java index d06bd81ab..40b13025f 100644 --- a/docker-java-transport-httpclient5/src/main/java/com/github/dockerjava/httpclient5/ApacheDockerHttpClientImpl.java +++ b/docker-java-transport-httpclient5/src/main/java/com/github/dockerjava/httpclient5/ApacheDockerHttpClientImpl.java @@ -41,12 +41,12 @@ class ApacheDockerHttpClientImpl implements DockerHttpClient { private final CloseableHttpClient httpClient; - private final HttpHost host; protected ApacheDockerHttpClientImpl( URI dockerHost, - SSLConfig sslConfig + SSLConfig sslConfig, + int maxConnections ) { Registry socketFactoryRegistry = createConnectionSocketFactoryRegistry(sslConfig, dockerHost); @@ -66,27 +66,30 @@ protected ApacheDockerHttpClientImpl( host = HttpHost.create(dockerHost); } + PoolingHttpClientConnectionManager connectionManager = new PoolingHttpClientConnectionManager( + socketFactoryRegistry, + new ManagedHttpClientConnectionFactory( + null, + null, + null, + null, + message -> { + Header transferEncodingHeader = message.getFirstHeader(HttpHeaders.TRANSFER_ENCODING); + if (transferEncodingHeader != null) { + if ("identity".equalsIgnoreCase(transferEncodingHeader.getValue())) { + return ContentLengthStrategy.UNDEFINED; + } + } + return DefaultContentLengthStrategy.INSTANCE.determineLength(message); + }, + null + ) + ); + connectionManager.setMaxTotal(maxConnections); + connectionManager.setDefaultMaxPerRoute(maxConnections); httpClient = HttpClients.custom() .setRequestExecutor(new HijackingHttpRequestExecutor(null)) - .setConnectionManager(new PoolingHttpClientConnectionManager( - socketFactoryRegistry, - new ManagedHttpClientConnectionFactory( - null, - null, - null, - null, - message -> { - Header transferEncodingHeader = message.getFirstHeader(HttpHeaders.TRANSFER_ENCODING); - if (transferEncodingHeader != null) { - if ("identity".equalsIgnoreCase(transferEncodingHeader.getValue())) { - return ContentLengthStrategy.UNDEFINED; - } - } - return DefaultContentLengthStrategy.INSTANCE.determineLength(message); - }, - null - ) - )) + .setConnectionManager(connectionManager) .build(); } diff --git a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyDockerHttpClient.java b/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyDockerHttpClient.java index 8eb3c2c6a..74ef77e4b 100644 --- a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyDockerHttpClient.java +++ b/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyDockerHttpClient.java @@ -54,9 +54,9 @@ public static final class Builder { private Integer connectTimeout = null; - private Integer maxTotalConnections = null; + private Integer maxTotalConnections = Integer.MAX_VALUE; - private Integer maxPerRouteConnections = null; + private Integer maxPerRouteConnections = Integer.MAX_VALUE; private Integer connectionRequestTimeout = null; diff --git a/docker-java-transport-zerodep/src/main/java/com/github/dockerjava/httpclient5/ZerodepDockerHttpClient.java b/docker-java-transport-zerodep/src/main/java/com/github/dockerjava/httpclient5/ZerodepDockerHttpClient.java index 2298da816..a0d2abaaf 100644 --- a/docker-java-transport-zerodep/src/main/java/com/github/dockerjava/httpclient5/ZerodepDockerHttpClient.java +++ b/docker-java-transport-zerodep/src/main/java/com/github/dockerjava/httpclient5/ZerodepDockerHttpClient.java @@ -14,6 +14,8 @@ public static final class Builder { private SSLConfig sslConfig = null; + private int maxConnections = Integer.MAX_VALUE; + public Builder dockerHost(URI value) { this.dockerHost = Objects.requireNonNull(value, "dockerHost"); return this; @@ -24,13 +26,18 @@ public Builder sslConfig(SSLConfig value) { return this; } + public Builder maxConnections(int value) { + this.maxConnections = value; + return this; + } + public ZerodepDockerHttpClient build() { Objects.requireNonNull(dockerHost, "dockerHost"); - return new ZerodepDockerHttpClient(dockerHost, sslConfig); + return new ZerodepDockerHttpClient(dockerHost, sslConfig, maxConnections); } } - protected ZerodepDockerHttpClient(URI dockerHost, SSLConfig sslConfig) { - super(dockerHost, sslConfig); + protected ZerodepDockerHttpClient(URI dockerHost, SSLConfig sslConfig, int maxConnections) { + super(dockerHost, sslConfig, maxConnections); } } diff --git a/docker-java/src/test/java/com/github/dockerjava/cmd/LogContainerCmdIT.java b/docker-java/src/test/java/com/github/dockerjava/cmd/LogContainerCmdIT.java index 37bf5f393..b0de380db 100644 --- a/docker-java/src/test/java/com/github/dockerjava/cmd/LogContainerCmdIT.java +++ b/docker-java/src/test/java/com/github/dockerjava/cmd/LogContainerCmdIT.java @@ -1,7 +1,10 @@ package com.github.dockerjava.cmd; +import com.github.dockerjava.api.DockerClient; +import com.github.dockerjava.api.async.ResultCallback; import com.github.dockerjava.api.command.CreateContainerResponse; import com.github.dockerjava.api.exception.NotFoundException; +import com.github.dockerjava.api.model.Frame; import com.github.dockerjava.api.model.StreamType; import com.github.dockerjava.utils.LogContainerTestCallback; import org.junit.Test; @@ -9,13 +12,25 @@ import org.slf4j.LoggerFactory; import java.io.IOException; +import java.util.List; +import java.util.concurrent.Callable; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.stream.Collectors; +import java.util.stream.LongStream; +import static org.awaitility.Awaitility.await; +import static org.hamcrest.CoreMatchers.everyItem; import static org.hamcrest.CoreMatchers.is; import static org.hamcrest.MatcherAssert.assertThat; import static org.hamcrest.Matchers.containsString; -import static org.hamcrest.Matchers.equalTo; import static org.hamcrest.Matchers.emptyString; +import static org.hamcrest.Matchers.equalTo; +import static org.hamcrest.Matchers.hasSize; +import static org.hamcrest.Matchers.hasToString; import static org.hamcrest.Matchers.not; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; @@ -197,4 +212,52 @@ public void asyncLogContainerWithSince() throws Exception { assertThat(loggingCallback.toString(), containsString(snippet)); } + + @Test(timeout = 10_000) + public void simultaneousCommands() throws Exception { + // Create a new client to not affect other tests + DockerClient client = dockerRule.newClient(); + CreateContainerResponse container = client.createContainerCmd("busybox") + .withCmd("/bin/sh", "-c", "echo hello world; sleep infinity") + .exec(); + + client.startContainerCmd(container.getId()).exec(); + + // Simulate 100 simultaneous connections + int connections = 100; + + ExecutorService executor = Executors.newFixedThreadPool(connections); + try { + List firstFrames = new CopyOnWriteArrayList<>(); + executor.invokeAll( + LongStream.range(0, connections).>mapToObj(__ -> { + return () -> { + return client.logContainerCmd(container.getId()) + .withStdOut(true) + .withFollowStream(true) + .exec(new ResultCallback.Adapter() { + + final AtomicBoolean first = new AtomicBoolean(true); + + @Override + public void onNext(Frame object) { + if (first.compareAndSet(true, false)) { + firstFrames.add(object); + } + super.onNext(object); + } + }); + }; + }).collect(Collectors.toList()) + ); + + await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> { + assertThat(firstFrames, hasSize(connections)); + }); + + assertThat(firstFrames, everyItem(hasToString("STDOUT: hello world"))); + } finally { + executor.shutdownNow(); + } + } } diff --git a/docker-java/src/test/java/com/github/dockerjava/core/DockerRule.java b/docker-java/src/test/java/com/github/dockerjava/core/DockerRule.java index c7ca2c0d9..e6e922bc6 100644 --- a/docker-java/src/test/java/com/github/dockerjava/core/DockerRule.java +++ b/docker-java/src/test/java/com/github/dockerjava/core/DockerRule.java @@ -46,11 +46,7 @@ public DockerRule(CmdIT cmdIT) { } - public DockerClient getClient() { - if (dockerClient != null) { - return dockerClient; - } - + public DockerClient newClient() { DockerClientImpl dockerClient = cmdIT.getFactoryType().createDockerClient(config()); DockerHttpClient dockerHttpClient = dockerClient.getHttpClient(); @@ -88,7 +84,7 @@ public CreateVolumeCmd.Exec createCreateVolumeCmdExec() { } ); - return this.dockerClient = new DockerClientDelegate(dockerClient) { + return new DockerClientDelegate(dockerClient) { @Override public DockerHttpClient getHttpClient() { return dockerHttpClient; @@ -96,6 +92,13 @@ public DockerHttpClient getHttpClient() { }; } + public DockerClient getClient() { + if (dockerClient != null) { + return dockerClient; + } + return this.dockerClient = newClient(); + } + @Override public Statement apply(Statement base, Description description) { return super.apply(base, description);