From 0c87505ac9f85353a9508129212bd379eac3fd22 Mon Sep 17 00:00:00 2001 From: Sergei Egorov Date: Tue, 23 Jun 2020 18:18:52 +0200 Subject: [PATCH] Detect leaked `DockerHttpClient.Response` objects --- .../core/DefaultInvocationBuilder.java | 31 ++++++- .../core/exec/ResizeContainerCmdExec.java | 8 +- .../core/exec/ResizeExecCmdExec.java | 8 +- .../dockerjava/okhttp/OkDockerHttpClient.java | 13 ++- .../java/com/github/dockerjava/cmd/CmdIT.java | 36 +++++--- .../cmd/DockerHttpClientLeakDetector.java | 31 +++++++ .../github/dockerjava/cmd/SaveImageCmdIT.java | 20 ++-- .../dockerjava/cmd/SaveImagesCmdIT.java | 25 +++-- .../cmd/TrackingDockerHttpClient.java | 91 +++++++++++++++++++ 9 files changed, 225 insertions(+), 38 deletions(-) create mode 100644 docker-java/src/test/java/com/github/dockerjava/cmd/DockerHttpClientLeakDetector.java create mode 100644 docker-java/src/test/java/com/github/dockerjava/cmd/TrackingDockerHttpClient.java diff --git a/docker-java-core/src/main/java/com/github/dockerjava/core/DefaultInvocationBuilder.java b/docker-java-core/src/main/java/com/github/dockerjava/core/DefaultInvocationBuilder.java index 100904d7f..c63daf42d 100644 --- a/docker-java-core/src/main/java/com/github/dockerjava/core/DefaultInvocationBuilder.java +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/DefaultInvocationBuilder.java @@ -18,6 +18,7 @@ import org.apache.commons.io.IOUtils; import java.io.ByteArrayInputStream; +import java.io.FilterInputStream; import java.io.IOException; import java.io.InputStream; import java.nio.charset.StandardCharsets; @@ -99,7 +100,17 @@ public InputStream post(Object entity) { .body(encode(entity)) .build(); - return execute(request).getBody(); + DockerHttpClient.Response response = execute(request); + return new FilterInputStream(response.getBody()) { + @Override + public void close() throws IOException { + try { + super.close(); + } finally { + response.close(); + } + } + }; } @Override @@ -188,7 +199,17 @@ public InputStream get() { .method(DockerHttpClient.Request.Method.GET) .build(); - return execute(request).getBody(); + DockerHttpClient.Response response = execute(request); + return new FilterInputStream(response.getBody()) { + @Override + public void close() throws IOException { + try { + super.close(); + } finally { + response.close(); + } + } + }; } @Override @@ -244,8 +265,12 @@ protected void executeAndStream( Consumer sourceConsumer ) { Thread thread = new Thread(() -> { + Thread streamingThread = Thread.currentThread(); try (DockerHttpClient.Response response = execute(request)) { - callback.onStart(response); + callback.onStart(() -> { + streamingThread.interrupt(); + response.close(); + }); sourceConsumer.accept(response); callback.onComplete(); diff --git a/docker-java-core/src/main/java/com/github/dockerjava/core/exec/ResizeContainerCmdExec.java b/docker-java-core/src/main/java/com/github/dockerjava/core/exec/ResizeContainerCmdExec.java index f13fc582a..4913bde79 100644 --- a/docker-java-core/src/main/java/com/github/dockerjava/core/exec/ResizeContainerCmdExec.java +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/exec/ResizeContainerCmdExec.java @@ -7,6 +7,8 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.io.IOException; + public class ResizeContainerCmdExec extends AbstrSyncDockerCmdExec implements ResizeContainerCmd.Exec { private static final Logger LOGGER = LoggerFactory.getLogger(ResizeContainerCmdExec.class); @@ -23,7 +25,11 @@ protected Void execute(ResizeContainerCmd command) { LOGGER.trace("POST: {}", webResource); - webResource.request().accept(MediaType.APPLICATION_JSON).post(command); + try { + webResource.request().accept(MediaType.APPLICATION_JSON).post(command).close(); + } catch (IOException e) { + throw new RuntimeException(e); + } return null; } diff --git a/docker-java-core/src/main/java/com/github/dockerjava/core/exec/ResizeExecCmdExec.java b/docker-java-core/src/main/java/com/github/dockerjava/core/exec/ResizeExecCmdExec.java index ba01c5ffe..e799a95d5 100644 --- a/docker-java-core/src/main/java/com/github/dockerjava/core/exec/ResizeExecCmdExec.java +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/exec/ResizeExecCmdExec.java @@ -7,6 +7,8 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.io.IOException; + public class ResizeExecCmdExec extends AbstrSyncDockerCmdExec implements ResizeExecCmd.Exec { @@ -23,7 +25,11 @@ protected Void execute(ResizeExecCmd command) { LOGGER.trace("POST: {}", webResource); - webResource.request().accept(MediaType.APPLICATION_JSON).post(null); + try { + webResource.request().accept(MediaType.APPLICATION_JSON).post(null).close(); + } catch (IOException e) { + throw new RuntimeException(e); + } return null; } diff --git a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkDockerHttpClient.java b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkDockerHttpClient.java index dad9e478d..5bb943bef 100644 --- a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkDockerHttpClient.java +++ b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkDockerHttpClient.java @@ -2,6 +2,7 @@ import com.github.dockerjava.transport.DockerHttpClient; import com.github.dockerjava.transport.SSLConfig; +import okhttp3.Call; import okhttp3.ConnectionPool; import okhttp3.Dns; import okhttp3.HttpUrl; @@ -219,9 +220,11 @@ public Response execute(Request request) { clientToUse = streamingClient; } + Call call = clientToUse.newCall(requestBuilder.build()); try { - return new OkResponse(clientToUse.newCall(requestBuilder.build()).execute()); + return new OkResponse(call); } catch (IOException e) { + call.cancel(); throw new UncheckedIOException("Error while executing " + request, e); } } @@ -239,10 +242,13 @@ static class OkResponse implements Response { static final ThreadLocal CLOSING = ThreadLocal.withInitial(() -> false); + private final Call call; + private final okhttp3.Response response; - OkResponse(okhttp3.Response response) { - this.response = response; + OkResponse(Call call) throws IOException { + this.call = call; + this.response = call.execute(); } @Override @@ -270,6 +276,7 @@ public void close() { boolean previous = CLOSING.get(); CLOSING.set(true); try { + call.cancel(); response.close(); } catch (Exception | AssertionError e) { LOGGER.debug("Failed to close the response", e); diff --git a/docker-java/src/test/java/com/github/dockerjava/cmd/CmdIT.java b/docker-java/src/test/java/com/github/dockerjava/cmd/CmdIT.java index 8dab7f96c..12664c4e5 100644 --- a/docker-java/src/test/java/com/github/dockerjava/cmd/CmdIT.java +++ b/docker-java/src/test/java/com/github/dockerjava/cmd/CmdIT.java @@ -39,11 +39,13 @@ public DockerClientImpl createDockerClient(DockerClientConfig config) { public DockerClientImpl createDockerClient(DockerClientConfig config) { return (DockerClientImpl) DockerClientBuilder.getInstance(config) .withDockerHttpClient( - new JerseyDockerHttpClient.Builder() - .dockerHost(config.getDockerHost()) - .sslConfig(config.getSSLConfig()) - .connectTimeout(30 * 1000) - .build() + new TrackingDockerHttpClient( + new JerseyDockerHttpClient.Builder() + .dockerHost(config.getDockerHost()) + .sslConfig(config.getSSLConfig()) + .connectTimeout(30 * 1000) + .build() + ) ) .build(); } @@ -53,11 +55,13 @@ public DockerClientImpl createDockerClient(DockerClientConfig config) { public DockerClientImpl createDockerClient(DockerClientConfig config) { return (DockerClientImpl) DockerClientBuilder.getInstance(config) .withDockerHttpClient( - new OkDockerHttpClient.Builder() - .dockerHost(config.getDockerHost()) - .sslConfig(config.getSSLConfig()) - .connectTimeout(30 * 100) - .build() + new TrackingDockerHttpClient( + new OkDockerHttpClient.Builder() + .dockerHost(config.getDockerHost()) + .sslConfig(config.getSSLConfig()) + .connectTimeout(30 * 100) + .build() + ) ) .build(); } @@ -67,10 +71,12 @@ public DockerClientImpl createDockerClient(DockerClientConfig config) { public DockerClientImpl createDockerClient(DockerClientConfig config) { return (DockerClientImpl) DockerClientBuilder.getInstance(config) .withDockerHttpClient( - new ApacheDockerHttpClient.Builder() - .dockerHost(config.getDockerHost()) - .sslConfig(config.getSSLConfig()) - .build() + new TrackingDockerHttpClient( + new ApacheDockerHttpClient.Builder() + .dockerHost(config.getDockerHost()) + .sslConfig(config.getSSLConfig()) + .build() + ) ) .build(); } @@ -110,4 +116,6 @@ public FactoryType getFactoryType() { @Rule public DockerRule dockerRule = new DockerRule( this); + @Rule + public DockerHttpClientLeakDetector leakDetector = new DockerHttpClientLeakDetector(); } diff --git a/docker-java/src/test/java/com/github/dockerjava/cmd/DockerHttpClientLeakDetector.java b/docker-java/src/test/java/com/github/dockerjava/cmd/DockerHttpClientLeakDetector.java new file mode 100644 index 000000000..1da12f3e0 --- /dev/null +++ b/docker-java/src/test/java/com/github/dockerjava/cmd/DockerHttpClientLeakDetector.java @@ -0,0 +1,31 @@ +package com.github.dockerjava.cmd; + +import org.junit.rules.ExternalResource; + +public class DockerHttpClientLeakDetector extends ExternalResource { + + @Override + protected void before() { + synchronized (TrackingDockerHttpClient.ACTIVE_RESPONSES) { + TrackingDockerHttpClient.ACTIVE_RESPONSES.clear(); + } + } + + @Override + protected void after() { + synchronized (TrackingDockerHttpClient.ACTIVE_RESPONSES) { + if (TrackingDockerHttpClient.ACTIVE_RESPONSES.isEmpty()) { + return; + } + + System.out.println("Leaked responses:"); + IllegalStateException exception = new IllegalStateException("Leaked responses!"); + exception.setStackTrace(new StackTraceElement[0]); + + TrackingDockerHttpClient.ACTIVE_RESPONSES.forEach(response -> { + exception.addSuppressed(response.allocatedAt); + }); + throw exception; + } + } +} diff --git a/docker-java/src/test/java/com/github/dockerjava/cmd/SaveImageCmdIT.java b/docker-java/src/test/java/com/github/dockerjava/cmd/SaveImageCmdIT.java index ab2f3aca9..cb5a4666c 100644 --- a/docker-java/src/test/java/com/github/dockerjava/cmd/SaveImageCmdIT.java +++ b/docker-java/src/test/java/com/github/dockerjava/cmd/SaveImageCmdIT.java @@ -16,13 +16,19 @@ public class SaveImageCmdIT extends CmdIT { @Test public void saveImage() throws Exception { - InputStream image = IOUtils.toBufferedInputStream(dockerRule.getClient().saveImageCmd("busybox").exec()); - assertThat(image.read(), not(-1)); - - InputStream image2 = IOUtils.toBufferedInputStream(dockerRule.getClient().saveImageCmd("busybox").withTag("latest").exec()); - assertThat(image2.read(), not(-1)); - - + try ( + InputStream inputStream = dockerRule.getClient().saveImageCmd("busybox").exec(); + InputStream image = IOUtils.toBufferedInputStream(inputStream) + ) { + assertThat(image.read(), not(-1)); + } + + try ( + InputStream inputStream = dockerRule.getClient().saveImageCmd("busybox").withTag("latest").exec(); + InputStream image2 = IOUtils.toBufferedInputStream(inputStream) + ) { + assertThat(image2.read(), not(-1)); + } } } diff --git a/docker-java/src/test/java/com/github/dockerjava/cmd/SaveImagesCmdIT.java b/docker-java/src/test/java/com/github/dockerjava/cmd/SaveImagesCmdIT.java index 2b5305d68..86b246029 100644 --- a/docker-java/src/test/java/com/github/dockerjava/cmd/SaveImagesCmdIT.java +++ b/docker-java/src/test/java/com/github/dockerjava/cmd/SaveImagesCmdIT.java @@ -15,7 +15,10 @@ public class SaveImagesCmdIT extends CmdIT { @Test public void saveNoImages() throws Exception { - try(final InputStream image = IOUtils.toBufferedInputStream(dockerRule.getClient().saveImagesCmd().exec())){ + try ( + InputStream inputStream = dockerRule.getClient().saveImagesCmd().exec(); + InputStream image = IOUtils.toBufferedInputStream(inputStream) + ){ assertThat(image.read(), not(-1)); } @@ -23,8 +26,10 @@ public void saveNoImages() throws Exception { @Test public void saveImagesWithNameAndTag() throws Exception { - - try(final InputStream image = IOUtils.toBufferedInputStream(dockerRule.getClient().saveImagesCmd().withImage("busybox", "latest").exec())) { + try ( + InputStream inputStream = dockerRule.getClient().saveImagesCmd().withImage("busybox", "latest").exec(); + InputStream image = IOUtils.toBufferedInputStream(inputStream) + ) { assertThat(image.read(), not(-1)); } @@ -32,12 +37,14 @@ public void saveImagesWithNameAndTag() throws Exception { @Test public void saveMultipleImages() throws Exception { - - try(final InputStream image = IOUtils.toBufferedInputStream(dockerRule.getClient().saveImagesCmd() - // Not a real life use-case but "busybox" is the only one I dare to assume is really there. - .withImage("busybox", "latest") - .withImage("busybox", "latest") - .exec())) { + try ( + InputStream inputStream = dockerRule.getClient().saveImagesCmd() + // Not a real life use-case but "busybox" is the only one I dare to assume is really there. + .withImage("busybox", "latest") + .withImage("busybox", "latest") + .exec(); + InputStream image = IOUtils.toBufferedInputStream(inputStream) + ) { assertThat(image.read(), not(-1)); } } diff --git a/docker-java/src/test/java/com/github/dockerjava/cmd/TrackingDockerHttpClient.java b/docker-java/src/test/java/com/github/dockerjava/cmd/TrackingDockerHttpClient.java new file mode 100644 index 000000000..3c991a8f1 --- /dev/null +++ b/docker-java/src/test/java/com/github/dockerjava/cmd/TrackingDockerHttpClient.java @@ -0,0 +1,91 @@ +package com.github.dockerjava.cmd; + +import com.github.dockerjava.transport.DockerHttpClient; + +import javax.annotation.Nonnull; +import javax.annotation.Nullable; +import java.io.IOException; +import java.io.InputStream; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; + +class TrackingDockerHttpClient implements DockerHttpClient { + + static final Set ACTIVE_RESPONSES = Collections.newSetFromMap(new ConcurrentHashMap<>()); + + private final DockerHttpClient delegate; + + TrackingDockerHttpClient(DockerHttpClient delegate) { + this.delegate = delegate; + } + + @Override + public Response execute(Request request) { + return new TrackedResponse(delegate.execute(request)) { + { + synchronized (ACTIVE_RESPONSES) { + ACTIVE_RESPONSES.add(this); + } + } + + @Override + public void close() { + synchronized (ACTIVE_RESPONSES) { + ACTIVE_RESPONSES.remove(this); + } + super.close(); + } + }; + } + + @Override + public void close() throws IOException { + delegate.close(); + } + + static class TrackedResponse implements Response { + + private static class AllocatedAt extends Exception { + public AllocatedAt(String message) { + super(message); + } + } + + final Exception allocatedAt = new AllocatedAt(this.toString()); + + private final Response delegate; + + TrackedResponse(Response delegate) { + this.delegate = delegate; + } + + @Override + public int getStatusCode() { + return delegate.getStatusCode(); + } + + @Override + public Map> getHeaders() { + return delegate.getHeaders(); + } + + @Override + public InputStream getBody() { + return delegate.getBody(); + } + + @Override + public void close() { + delegate.close(); + } + + @Override + @Nullable + public String getHeader(@Nonnull String name) { + return delegate.getHeader(name); + } + } +}