Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
.project
.settings
.classpath
.factorypath

# Ignore all build/dist directories
target
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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);
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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<ConnectionSocketFactory> socketFactoryRegistry = createConnectionSocketFactoryRegistry(sslConfig, dockerHost);

Expand All @@ -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();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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);
}
}
Original file line number Diff line number Diff line change
@@ -1,21 +1,36 @@
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;
import org.slf4j.Logger;
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;
Expand Down Expand Up @@ -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<Frame> firstFrames = new CopyOnWriteArrayList<>();
executor.invokeAll(
LongStream.range(0, connections).<Callable<Object>>mapToObj(__ -> {
return () -> {
return client.logContainerCmd(container.getId())
.withStdOut(true)
.withFollowStream(true)
.exec(new ResultCallback.Adapter<Frame>() {

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();
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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();

Expand Down Expand Up @@ -88,14 +84,21 @@ public CreateVolumeCmd.Exec createCreateVolumeCmdExec() {
}
);

return this.dockerClient = new DockerClientDelegate(dockerClient) {
return new DockerClientDelegate(dockerClient) {
@Override
public DockerHttpClient getHttpClient() {
return dockerHttpClient;
}
};
}

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);
Expand Down