From a3d2b7ca7adf05187fba62bc2bd95a8eb5e63252 Mon Sep 17 00:00:00 2001 From: Sergei Egorov Date: Thu, 19 Mar 2020 20:39:15 +0100 Subject: [PATCH 01/13] Draft DockerHttpClient abstraction --- .../DelegatingDockerCmdExecFactory.java | 377 +++++++++++++++ docker-java-core/pom.xml | 7 + .../core/DefaultDockerCmdExecFactory.java | 154 ++++++ .../core/DefaultInvocationBuilder.java | 296 ++++++++++++ .../dockerjava/core/DockerClientImpl.java | 23 + .../dockerjava/core/DockerHttpClient.java | 67 +++ .../core/FramedInputStreamConsumer.java | 87 ++++ .../github/dockerjava/okhttp/FramedSink.java | 81 ---- .../okhttp/HijackingInterceptor.java | 61 +++ .../okhttp/NamedPipeSocketFactory.java | 7 +- .../dockerjava/okhttp/OkDockerHttpClient.java | 227 +++++++++ .../okhttp/OkHttpDockerCmdExecFactory.java | 165 +------ .../okhttp/OkHttpInvocationBuilder.java | 372 --------------- .../dockerjava/okhttp/OkHttpWebTarget.java | 148 ------ .../dockerjava/okhttp/UnixSocketFactory.java | 4 +- .../junit/DockerCmdExecFactoryDelegate.java | 448 +----------------- 16 files changed, 1332 insertions(+), 1192 deletions(-) create mode 100644 docker-java-api/src/main/java/com/github/dockerjava/api/command/DelegatingDockerCmdExecFactory.java create mode 100644 docker-java-core/src/main/java/com/github/dockerjava/core/DefaultDockerCmdExecFactory.java create mode 100644 docker-java-core/src/main/java/com/github/dockerjava/core/DefaultInvocationBuilder.java create mode 100644 docker-java-core/src/main/java/com/github/dockerjava/core/DockerHttpClient.java create mode 100644 docker-java-core/src/main/java/com/github/dockerjava/core/FramedInputStreamConsumer.java delete mode 100644 docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/FramedSink.java create mode 100644 docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/HijackingInterceptor.java create mode 100644 docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkDockerHttpClient.java delete mode 100644 docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpInvocationBuilder.java delete mode 100644 docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpWebTarget.java diff --git a/docker-java-api/src/main/java/com/github/dockerjava/api/command/DelegatingDockerCmdExecFactory.java b/docker-java-api/src/main/java/com/github/dockerjava/api/command/DelegatingDockerCmdExecFactory.java new file mode 100644 index 000000000..dbd691974 --- /dev/null +++ b/docker-java-api/src/main/java/com/github/dockerjava/api/command/DelegatingDockerCmdExecFactory.java @@ -0,0 +1,377 @@ +package com.github.dockerjava.api.command; + +import java.io.IOException; + +public class DelegatingDockerCmdExecFactory implements DockerCmdExecFactory { + + // We're not using abstract class because we want + // the compiler to force us to implement new DockerCmdExecFactory when added + public DockerCmdExecFactory getDockerCmdExecFactory() { + throw new IllegalStateException("Implement me!"); + } + + @Override + public AuthCmd.Exec createAuthCmdExec() { + return getDockerCmdExecFactory().createAuthCmdExec(); + } + + @Override + public InfoCmd.Exec createInfoCmdExec() { + return getDockerCmdExecFactory().createInfoCmdExec(); + } + + @Override + public PingCmd.Exec createPingCmdExec() { + return getDockerCmdExecFactory().createPingCmdExec(); + } + + @Override + public ExecCreateCmd.Exec createExecCmdExec() { + return getDockerCmdExecFactory().createExecCmdExec(); + } + + @Override + public VersionCmd.Exec createVersionCmdExec() { + return getDockerCmdExecFactory().createVersionCmdExec(); + } + + @Override + public PullImageCmd.Exec createPullImageCmdExec() { + return getDockerCmdExecFactory().createPullImageCmdExec(); + } + + @Override + public PushImageCmd.Exec createPushImageCmdExec() { + return getDockerCmdExecFactory().createPushImageCmdExec(); + } + + @Override + public SaveImageCmd.Exec createSaveImageCmdExec() { + return getDockerCmdExecFactory().createSaveImageCmdExec(); + } + + @Override + public SaveImagesCmd.Exec createSaveImagesCmdExec() { + return getDockerCmdExecFactory().createSaveImagesCmdExec(); + } + + @Override + public CreateImageCmd.Exec createCreateImageCmdExec() { + return getDockerCmdExecFactory().createCreateImageCmdExec(); + } + + @Override + public LoadImageCmd.Exec createLoadImageCmdExec() { + return getDockerCmdExecFactory().createLoadImageCmdExec(); + } + + @Override + public SearchImagesCmd.Exec createSearchImagesCmdExec() { + return getDockerCmdExecFactory().createSearchImagesCmdExec(); + } + + @Override + public RemoveImageCmd.Exec createRemoveImageCmdExec() { + return getDockerCmdExecFactory().createRemoveImageCmdExec(); + } + + @Override + public ListImagesCmd.Exec createListImagesCmdExec() { + return getDockerCmdExecFactory().createListImagesCmdExec(); + } + + @Override + public InspectImageCmd.Exec createInspectImageCmdExec() { + return getDockerCmdExecFactory().createInspectImageCmdExec(); + } + + @Override + public ListContainersCmd.Exec createListContainersCmdExec() { + return getDockerCmdExecFactory().createListContainersCmdExec(); + } + + @Override + public CreateContainerCmd.Exec createCreateContainerCmdExec() { + return getDockerCmdExecFactory().createCreateContainerCmdExec(); + } + + @Override + public StartContainerCmd.Exec createStartContainerCmdExec() { + return getDockerCmdExecFactory().createStartContainerCmdExec(); + } + + @Override + public InspectContainerCmd.Exec createInspectContainerCmdExec() { + return getDockerCmdExecFactory().createInspectContainerCmdExec(); + } + + @Override + public RemoveContainerCmd.Exec createRemoveContainerCmdExec() { + return getDockerCmdExecFactory().createRemoveContainerCmdExec(); + } + + @Override + public WaitContainerCmd.Exec createWaitContainerCmdExec() { + return getDockerCmdExecFactory().createWaitContainerCmdExec(); + } + + @Override + public AttachContainerCmd.Exec createAttachContainerCmdExec() { + return getDockerCmdExecFactory().createAttachContainerCmdExec(); + } + + @Override + public ExecStartCmd.Exec createExecStartCmdExec() { + return getDockerCmdExecFactory().createExecStartCmdExec(); + } + + @Override + public InspectExecCmd.Exec createInspectExecCmdExec() { + return getDockerCmdExecFactory().createInspectExecCmdExec(); + } + + @Override + public LogContainerCmd.Exec createLogContainerCmdExec() { + return getDockerCmdExecFactory().createLogContainerCmdExec(); + } + + @Override + public CopyFileFromContainerCmd.Exec createCopyFileFromContainerCmdExec() { + return getDockerCmdExecFactory().createCopyFileFromContainerCmdExec(); + } + + @Override + public CopyArchiveFromContainerCmd.Exec createCopyArchiveFromContainerCmdExec() { + return getDockerCmdExecFactory().createCopyArchiveFromContainerCmdExec(); + } + + @Override + public CopyArchiveToContainerCmd.Exec createCopyArchiveToContainerCmdExec() { + return getDockerCmdExecFactory().createCopyArchiveToContainerCmdExec(); + } + + @Override + public StopContainerCmd.Exec createStopContainerCmdExec() { + return getDockerCmdExecFactory().createStopContainerCmdExec(); + } + + @Override + public ContainerDiffCmd.Exec createContainerDiffCmdExec() { + return getDockerCmdExecFactory().createContainerDiffCmdExec(); + } + + @Override + public KillContainerCmd.Exec createKillContainerCmdExec() { + return getDockerCmdExecFactory().createKillContainerCmdExec(); + } + + @Override + public UpdateContainerCmd.Exec createUpdateContainerCmdExec() { + return getDockerCmdExecFactory().createUpdateContainerCmdExec(); + } + + @Override + public RenameContainerCmd.Exec createRenameContainerCmdExec() { + return getDockerCmdExecFactory().createRenameContainerCmdExec(); + } + + @Override + public RestartContainerCmd.Exec createRestartContainerCmdExec() { + return getDockerCmdExecFactory().createRestartContainerCmdExec(); + } + + @Override + public CommitCmd.Exec createCommitCmdExec() { + return getDockerCmdExecFactory().createCommitCmdExec(); + } + + @Override + public BuildImageCmd.Exec createBuildImageCmdExec() { + return getDockerCmdExecFactory().createBuildImageCmdExec(); + } + + @Override + public TopContainerCmd.Exec createTopContainerCmdExec() { + return getDockerCmdExecFactory().createTopContainerCmdExec(); + } + + @Override + public TagImageCmd.Exec createTagImageCmdExec() { + return getDockerCmdExecFactory().createTagImageCmdExec(); + } + + @Override + public PauseContainerCmd.Exec createPauseContainerCmdExec() { + return getDockerCmdExecFactory().createPauseContainerCmdExec(); + } + + @Override + public UnpauseContainerCmd.Exec createUnpauseContainerCmdExec() { + return getDockerCmdExecFactory().createUnpauseContainerCmdExec(); + } + + @Override + public EventsCmd.Exec createEventsCmdExec() { + return getDockerCmdExecFactory().createEventsCmdExec(); + } + + @Override + public StatsCmd.Exec createStatsCmdExec() { + return getDockerCmdExecFactory().createStatsCmdExec(); + } + + @Override + public CreateVolumeCmd.Exec createCreateVolumeCmdExec() { + return getDockerCmdExecFactory().createCreateVolumeCmdExec(); + } + + @Override + public InspectVolumeCmd.Exec createInspectVolumeCmdExec() { + return getDockerCmdExecFactory().createInspectVolumeCmdExec(); + } + + @Override + public RemoveVolumeCmd.Exec createRemoveVolumeCmdExec() { + return getDockerCmdExecFactory().createRemoveVolumeCmdExec(); + } + + @Override + public ListVolumesCmd.Exec createListVolumesCmdExec() { + return getDockerCmdExecFactory().createListVolumesCmdExec(); + } + + @Override + public ListNetworksCmd.Exec createListNetworksCmdExec() { + return getDockerCmdExecFactory().createListNetworksCmdExec(); + } + + @Override + public InspectNetworkCmd.Exec createInspectNetworkCmdExec() { + return getDockerCmdExecFactory().createInspectNetworkCmdExec(); + } + + @Override + public CreateNetworkCmd.Exec createCreateNetworkCmdExec() { + return getDockerCmdExecFactory().createCreateNetworkCmdExec(); + } + + @Override + public RemoveNetworkCmd.Exec createRemoveNetworkCmdExec() { + return getDockerCmdExecFactory().createRemoveNetworkCmdExec(); + } + + @Override + public ConnectToNetworkCmd.Exec createConnectToNetworkCmdExec() { + return getDockerCmdExecFactory().createConnectToNetworkCmdExec(); + } + + @Override + public DisconnectFromNetworkCmd.Exec createDisconnectFromNetworkCmdExec() { + return getDockerCmdExecFactory().createDisconnectFromNetworkCmdExec(); + } + + @Override + public InitializeSwarmCmd.Exec createInitializeSwarmCmdExec() { + return getDockerCmdExecFactory().createInitializeSwarmCmdExec(); + } + + @Override + public InspectSwarmCmd.Exec createInspectSwarmCmdExec() { + return getDockerCmdExecFactory().createInspectSwarmCmdExec(); + } + + @Override + public JoinSwarmCmd.Exec createJoinSwarmCmdExec() { + return getDockerCmdExecFactory().createJoinSwarmCmdExec(); + } + + @Override + public LeaveSwarmCmd.Exec createLeaveSwarmCmdExec() { + return getDockerCmdExecFactory().createLeaveSwarmCmdExec(); + } + + @Override + public UpdateSwarmCmd.Exec createUpdateSwarmCmdExec() { + return getDockerCmdExecFactory().createUpdateSwarmCmdExec(); + } + + @Override + public ListServicesCmd.Exec createListServicesCmdExec() { + return getDockerCmdExecFactory().createListServicesCmdExec(); + } + + @Override + public CreateServiceCmd.Exec createCreateServiceCmdExec() { + return getDockerCmdExecFactory().createCreateServiceCmdExec(); + } + + @Override + public InspectServiceCmd.Exec createInspectServiceCmdExec() { + return getDockerCmdExecFactory().createInspectServiceCmdExec(); + } + + @Override + public UpdateServiceCmd.Exec createUpdateServiceCmdExec() { + return getDockerCmdExecFactory().createUpdateServiceCmdExec(); + } + + @Override + public RemoveServiceCmd.Exec createRemoveServiceCmdExec() { + return getDockerCmdExecFactory().createRemoveServiceCmdExec(); + } + + @Override + public LogSwarmObjectCmd.Exec logSwarmObjectExec(String endpoint) { + return getDockerCmdExecFactory().logSwarmObjectExec(endpoint); + } + + @Override + public ListSwarmNodesCmd.Exec listSwarmNodeCmdExec() { + return getDockerCmdExecFactory().listSwarmNodeCmdExec(); + } + + @Override + public InspectSwarmNodeCmd.Exec inspectSwarmNodeCmdExec() { + return getDockerCmdExecFactory().inspectSwarmNodeCmdExec(); + } + + @Override + public RemoveSwarmNodeCmd.Exec removeSwarmNodeCmdExec() { + return getDockerCmdExecFactory().removeSwarmNodeCmdExec(); + } + + @Override + public UpdateSwarmNodeCmd.Exec updateSwarmNodeCmdExec() { + return getDockerCmdExecFactory().updateSwarmNodeCmdExec(); + } + + @Override + public ListTasksCmd.Exec listTasksCmdExec() { + return getDockerCmdExecFactory().listTasksCmdExec(); + } + + @Override + public PruneCmd.Exec pruneCmdExec() { + return getDockerCmdExecFactory().pruneCmdExec(); + } + + @Override + public ListSecretsCmd.Exec createListSecretsCmdExec() { + return getDockerCmdExecFactory().createListSecretsCmdExec(); + } + + @Override + public CreateSecretCmd.Exec createCreateSecretCmdExec() { + return getDockerCmdExecFactory().createCreateSecretCmdExec(); + } + + @Override + public RemoveSecretCmd.Exec createRemoveSecretCmdExec() { + return getDockerCmdExecFactory().createRemoveSecretCmdExec(); + } + + @Override + public void close() throws IOException { + getDockerCmdExecFactory().close(); + } +} diff --git a/docker-java-core/pom.xml b/docker-java-core/pom.xml index 5238f816a..ed77c9db7 100644 --- a/docker-java-core/pom.xml +++ b/docker-java-core/pom.xml @@ -70,6 +70,13 @@ 3.0.1u2 provided + + + org.immutables + value + 2.8.2 + provided + diff --git a/docker-java-core/src/main/java/com/github/dockerjava/core/DefaultDockerCmdExecFactory.java b/docker-java-core/src/main/java/com/github/dockerjava/core/DefaultDockerCmdExecFactory.java new file mode 100644 index 000000000..8f3dc3b04 --- /dev/null +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/DefaultDockerCmdExecFactory.java @@ -0,0 +1,154 @@ +package com.github.dockerjava.core; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.google.common.collect.HashMultimap; +import com.google.common.collect.ImmutableList; +import com.google.common.collect.MultimapBuilder; +import com.google.common.collect.SetMultimap; +import com.google.common.escape.Escaper; +import com.google.common.net.UrlEscapers; +import org.apache.commons.lang.StringUtils; + +import java.io.IOException; +import java.util.Map; +import java.util.Objects; +import java.util.Set; +import java.util.stream.Collectors; + +public final class DefaultDockerCmdExecFactory extends AbstractDockerCmdExecFactory { + + private final DockerHttpClient dockerHttpClient; + + private final ObjectMapper objectMapper; + + public DefaultDockerCmdExecFactory( + DockerHttpClient dockerHttpClient, + ObjectMapper objectMapper + ) { + this.dockerHttpClient = dockerHttpClient; + this.objectMapper = objectMapper; + } + + public DockerHttpClient getDockerHttpClient() { + return dockerHttpClient; + } + + @Override + protected WebTarget getBaseResource() { + return new DefaultWebTarget(); + } + + @Override + public void close() throws IOException { + dockerHttpClient.close(); + } + + private class DefaultWebTarget implements WebTarget { + + final ImmutableList path; + + final SetMultimap queryParams; + + public DefaultWebTarget() { + this( + ImmutableList.of(), + MultimapBuilder.hashKeys().hashSetValues().build() + ); + } + + DefaultWebTarget( + ImmutableList path, + SetMultimap queryParams + ) { + this.path = path; + this.queryParams = queryParams; + } + + @Override + public String toString() { + return String.format("DefaultWebTarget{path=%s, queryParams=%s}", path, queryParams); + } + + @Override + public InvocationBuilder request() { + String resource = StringUtils.join(path, "/"); + + if (!resource.startsWith("/")) { + resource = "/" + resource; + } + + if (!queryParams.isEmpty()) { + Escaper urlFormParameterEscaper = UrlEscapers.urlFormParameterEscaper(); + resource = queryParams.asMap().entrySet().stream() + .flatMap(entry -> { + return entry.getValue().stream().map(s -> { + return entry.getKey() + "=" + urlFormParameterEscaper.escape(s); + }); + }) + .collect(Collectors.joining("&", resource + "?", "")); + } + + return new DefaultInvocationBuilder( + dockerHttpClient, objectMapper, resource + ); + } + + @Override + public DefaultWebTarget path(String... components) { + ImmutableList newPath = ImmutableList.builder() + .addAll(path) + .add(components) + .build(); + return new DefaultWebTarget(newPath, queryParams); + } + + @Override + public DefaultWebTarget resolveTemplate(String name, Object value) { + ImmutableList.Builder newPath = ImmutableList.builder(); + for (String component : path) { + component = component.replaceAll( + "\\{" + name + "\\}", + UrlEscapers.urlPathSegmentEscaper().escape(value.toString()) + ); + newPath.add(component); + } + + return new DefaultWebTarget(newPath.build(), queryParams); + } + + @Override + public DefaultWebTarget queryParam(String name, Object value) { + if (value == null) { + return this; + } + + SetMultimap newQueryParams = HashMultimap.create(queryParams); + newQueryParams.put(name, value.toString()); + + return new DefaultWebTarget(path, newQueryParams); + } + + @Override + public DefaultWebTarget queryParamsSet(String name, Set values) { + SetMultimap newQueryParams = HashMultimap.create(queryParams); + newQueryParams.replaceValues(name, values.stream().filter(Objects::nonNull).map(Object::toString).collect(Collectors.toSet())); + + return new DefaultWebTarget(path, newQueryParams); + } + + @Override + public DefaultWebTarget queryParamsJsonMap(String name, Map values) { + if (values == null || values.isEmpty()) { + return this; + } + + // when param value is JSON string + try { + return queryParam(name, objectMapper.writeValueAsString(values)); + } catch (JsonProcessingException e) { + throw new RuntimeException(e); + } + } + } +} 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 new file mode 100644 index 000000000..f94ac19c7 --- /dev/null +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/DefaultInvocationBuilder.java @@ -0,0 +1,296 @@ +package com.github.dockerjava.core; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.core.type.TypeReference; +import com.fasterxml.jackson.databind.MappingIterator; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.github.dockerjava.api.async.ResultCallback; +import com.github.dockerjava.api.exception.BadRequestException; +import com.github.dockerjava.api.exception.ConflictException; +import com.github.dockerjava.api.exception.DockerException; +import com.github.dockerjava.api.exception.InternalServerErrorException; +import com.github.dockerjava.api.exception.NotAcceptableException; +import com.github.dockerjava.api.exception.NotFoundException; +import com.github.dockerjava.api.exception.NotModifiedException; +import com.github.dockerjava.api.exception.UnauthorizedException; +import com.github.dockerjava.api.model.Frame; +import org.apache.commons.io.IOUtils; + +import java.io.ByteArrayInputStream; +import java.io.IOException; +import java.io.InputStream; +import java.nio.charset.StandardCharsets; +import java.util.Objects; +import java.util.function.Consumer; + +class DefaultInvocationBuilder implements InvocationBuilder { + + private final DockerHttpClient.Request.Builder requestBuilder; + private final DockerHttpClient dockerHttpClient; + private final ObjectMapper objectMapper; + + public DefaultInvocationBuilder(DockerHttpClient dockerHttpClient, ObjectMapper objectMapper, String path) { + this.requestBuilder = DockerHttpClient.Request.builder().path(path); + this.dockerHttpClient = dockerHttpClient; + this.objectMapper = objectMapper; + } + + @Override + public DefaultInvocationBuilder accept(MediaType mediaType) { + return header("accept", mediaType.getMediaType()); + } + + @Override + public DefaultInvocationBuilder header(String name, String value) { + requestBuilder.putHeader(name, value); + return this; + } + + @Override + public void delete() { + DockerHttpClient.Request request = requestBuilder + .method(DockerHttpClient.Request.Method.DELETE) + .build(); + + execute(request).close(); + } + + @Override + public void get(ResultCallback resultCallback) { + DockerHttpClient.Request request = requestBuilder + .method(DockerHttpClient.Request.Method.GET) + .build(); + + executeAndStream( + request, + resultCallback, + new FramedInputStreamConsumer(resultCallback) + ); + } + + @Override + public T get(TypeReference typeReference) { + try (InputStream inputStream = get()) { + return objectMapper.readValue(inputStream, typeReference); + } catch (IOException e) { + throw new RuntimeException(e); + } + } + + @Override + public void get(TypeReference typeReference, ResultCallback resultCallback) { + DockerHttpClient.Request request = requestBuilder + .method(DockerHttpClient.Request.Method.GET) + .build(); + + executeAndStream( + request, + resultCallback, + new JsonSink<>(typeReference, resultCallback) + ); + } + + @Override + public InputStream post(Object entity) { + DockerHttpClient.Request request = requestBuilder + .method(DockerHttpClient.Request.Method.POST) + .putHeader("content-type", "application/json") + .body(encode(entity)) + .build(); + + return execute(request).getBody(); + } + + @Override + public T post(Object entity, TypeReference typeReference) { + try { + DockerHttpClient.Request request = requestBuilder + .method(DockerHttpClient.Request.Method.POST) + .putHeader("content-type", "application/json") + .body(new ByteArrayInputStream(objectMapper.writeValueAsBytes(entity))) + .build(); + + try (DockerHttpClient.Response response = execute(request)) { + return objectMapper.readValue(response.getBody(), typeReference); + } + } catch (IOException e) { + throw new RuntimeException(e); + } + } + + @Override + public void post(Object entity, TypeReference typeReference, ResultCallback resultCallback) { + try { + post(typeReference, resultCallback, new ByteArrayInputStream(objectMapper.writeValueAsBytes(entity))); + } catch (JsonProcessingException e) { + throw new RuntimeException(e); + } + } + + @Override + public T post(TypeReference typeReference, InputStream body) { + try (InputStream inputStream = post(body)) { + return objectMapper.readValue(inputStream, typeReference); + } catch (IOException e) { + throw new RuntimeException(e); + } + } + + @Override + public void post(Object entity, InputStream stdin, ResultCallback resultCallback) { + final DockerHttpClient.Request request; + try { + request = requestBuilder + .method(DockerHttpClient.Request.Method.POST) + .putHeader("content-type", "application/json") + .body(new ByteArrayInputStream(objectMapper.writeValueAsBytes(entity))) + .hijackedInput(stdin) + .build(); + } catch (JsonProcessingException e) { + throw new RuntimeException(e); + } + + executeAndStream( + request, + resultCallback, + new FramedInputStreamConsumer(resultCallback) + ); + } + + @Override + public void post(TypeReference typeReference, ResultCallback resultCallback, InputStream body) { + DockerHttpClient.Request request = requestBuilder + .method(DockerHttpClient.Request.Method.POST) + .body(body) + .build(); + + executeAndStream( + request, + resultCallback, + new JsonSink<>(typeReference, resultCallback) + ); + } + + @Override + public void postStream(InputStream body) { + DockerHttpClient.Request request = requestBuilder + .method(DockerHttpClient.Request.Method.POST) + .body(body) + .build(); + + execute(request).close(); + } + + @Override + public InputStream get() { + DockerHttpClient.Request request = requestBuilder + .method(DockerHttpClient.Request.Method.GET) + .build(); + + return execute(request).getBody(); + } + + @Override + public void put(InputStream body, MediaType mediaType) { + DockerHttpClient.Request request = requestBuilder + .method(DockerHttpClient.Request.Method.PUT) + .putHeader("content-type", mediaType.toString()) + .body(body) + .build(); + + execute(request).close(); + } + + protected DockerHttpClient.Response execute(DockerHttpClient.Request request) { + try { + DockerHttpClient.Response response = dockerHttpClient.execute(request); + int statusCode = response.getStatusCode(); + if (statusCode < 200 || statusCode > 299) { + try { + String body = IOUtils.toString(response.getBody(), StandardCharsets.UTF_8); + switch (statusCode) { + case 304: + throw new NotModifiedException(body); + case 400: + throw new BadRequestException(body); + case 401: + throw new UnauthorizedException(body); + case 404: + throw new NotFoundException(body); + case 406: + throw new NotAcceptableException(body); + case 409: + throw new ConflictException(body); + case 500: + throw new InternalServerErrorException(body); + default: + throw new DockerException(body, statusCode); + } + } finally { + response.close(); + } + } else { + return response; + } + } catch (IOException e) { + throw new RuntimeException(e); + } + } + + protected void executeAndStream( + DockerHttpClient.Request request, + ResultCallback callback, + Consumer sourceConsumer + ) { + Thread thread = new Thread(() -> { + try (DockerHttpClient.Response response = execute(request)) { + callback.onStart(response); + + sourceConsumer.accept(response); + callback.onComplete(); + } catch (Exception e) { + callback.onError(e); + } + }, "docker-java-okhttp-stream-" + Objects.hashCode(request)); + thread.setDaemon(true); + + thread.start(); + } + + private InputStream encode(Object entity) { + if (entity == null) { + return null; + } + + try { + return new ByteArrayInputStream(objectMapper.writeValueAsBytes(entity)); + } catch (JsonProcessingException e) { + throw new RuntimeException(e); + } + } + + private class JsonSink implements Consumer { + + private final TypeReference typeReference; + + private final ResultCallback resultCallback; + + JsonSink(TypeReference typeReference, ResultCallback resultCallback) { + this.typeReference = typeReference; + this.resultCallback = resultCallback; + } + + @Override + public void accept(DockerHttpClient.Response response) { + try { + InputStream body = response.getBody(); + MappingIterator iterator = objectMapper.readerFor(typeReference).readValues(body); + while (iterator.hasNextValue()) { + resultCallback.onNext((T) iterator.nextValue()); + } + } catch (Exception e) { + resultCallback.onError(e); + } + } + } +} diff --git a/docker-java-core/src/main/java/com/github/dockerjava/core/DockerClientImpl.java b/docker-java-core/src/main/java/com/github/dockerjava/core/DockerClientImpl.java index 85f21d476..3e1b59b1e 100644 --- a/docker-java-core/src/main/java/com/github/dockerjava/core/DockerClientImpl.java +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/DockerClientImpl.java @@ -150,6 +150,7 @@ import com.github.dockerjava.core.command.WaitContainerCmdImpl; import javax.annotation.Nonnull; +import javax.annotation.Nullable; import java.io.Closeable; import java.io.File; import java.io.IOException; @@ -196,6 +197,27 @@ public static DockerClientImpl getInstance(String serverUrl) { return new DockerClientImpl(serverUrl); } + public DockerClientImpl withHttpClient(DockerHttpClient httpClient) { + return withDockerCmdExecFactory(new DefaultDockerCmdExecFactory(httpClient, dockerClientConfig.getObjectMapper())); + } + + /** + * + * @return {@link DockerHttpClient} or null if not set + */ + @Nullable + public DockerHttpClient getHttpClient() { + if (dockerCmdExecFactory instanceof DefaultDockerCmdExecFactory) { + return ((DefaultDockerCmdExecFactory) dockerCmdExecFactory).getDockerHttpClient(); + } else { + return null; + } + } + + /** + * @deprecated use {{@link #withHttpClient(DockerHttpClient)}} + */ + @Deprecated public DockerClientImpl withDockerCmdExecFactory(DockerCmdExecFactory dockerCmdExecFactory) { checkNotNull(dockerCmdExecFactory, "dockerCmdExecFactory was not specified"); this.dockerCmdExecFactory = dockerCmdExecFactory; @@ -205,6 +227,7 @@ public DockerClientImpl withDockerCmdExecFactory(DockerCmdExecFactory dockerCmdE return this; } + @Deprecated private DockerCmdExecFactory getDockerCmdExecFactory() { checkNotNull(dockerCmdExecFactory, "dockerCmdExecFactory was not specified"); return dockerCmdExecFactory; diff --git a/docker-java-core/src/main/java/com/github/dockerjava/core/DockerHttpClient.java b/docker-java-core/src/main/java/com/github/dockerjava/core/DockerHttpClient.java new file mode 100644 index 000000000..7c412fd78 --- /dev/null +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/DockerHttpClient.java @@ -0,0 +1,67 @@ +package com.github.dockerjava.core; + +import org.immutables.value.Value; + +import javax.annotation.Nullable; +import java.io.Closeable; +import java.io.InputStream; +import java.util.List; +import java.util.Map; + +public interface DockerHttpClient extends Closeable { + + Response execute(Request request); + + interface Response extends Closeable { + + int getStatusCode(); + + Map> getHeaders(); + + InputStream getBody(); + + @Override + void close(); + } + + @Value.Immutable + @Value.Style( + visibility = Value.Style.ImplementationVisibility.PACKAGE, + overshadowImplementation = true, + depluralize = true + ) + abstract class Request { + + public enum Method { + GET, + POST, + PUT, + DELETE, + OPTIONS, + PATCH, + } + + public static class Builder extends ImmutableRequest.Builder { + + public Builder method(Method method) { + return method(method.name()); + } + } + + public static Builder builder() { + return new Builder(); + } + + public abstract String method(); + + public abstract String path(); + + @Nullable + public abstract InputStream body(); + + @Nullable + public abstract InputStream hijackedInput(); + + public abstract Map headers(); + } +} diff --git a/docker-java-core/src/main/java/com/github/dockerjava/core/FramedInputStreamConsumer.java b/docker-java-core/src/main/java/com/github/dockerjava/core/FramedInputStreamConsumer.java new file mode 100644 index 000000000..ce707eb58 --- /dev/null +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/FramedInputStreamConsumer.java @@ -0,0 +1,87 @@ +package com.github.dockerjava.core; + +import com.github.dockerjava.api.async.ResultCallback; +import com.github.dockerjava.api.model.Frame; +import com.github.dockerjava.api.model.StreamType; + +import java.io.InputStream; +import java.nio.ByteBuffer; +import java.util.Arrays; +import java.util.function.Consumer; + +class FramedInputStreamConsumer implements Consumer { + + private static final int HEADER_SIZE = 8; + + private final ResultCallback resultCallback; + + FramedInputStreamConsumer(ResultCallback resultCallback) { + this.resultCallback = resultCallback; + } + + @Override + public void accept(DockerHttpClient.Response response) { + try { + InputStream body = response.getBody(); + while (true) { + // See https://docs.docker.com/engine/api/v1.37/#operation/ContainerAttach + // [8]byte{STREAM_TYPE, 0, 0, 0, SIZE1, SIZE2, SIZE3, SIZE4}[]byte{OUTPUT} + + byte[] header = new byte[HEADER_SIZE]; + if (body.read(header) < 0) { + // TODO log? + return; + } + + int streamTypeByte = header[0]; + + StreamType streamType = streamType(streamTypeByte); + + byte[] buffer = new byte[1024]; + + if (streamType == StreamType.RAW) { + // FIXME should check the header instead + resultCallback.onNext(new Frame(StreamType.RAW, header)); + + int readBytes; + while ((readBytes = body.read(buffer)) >= 0) { + resultCallback.onNext(new Frame(StreamType.RAW, Arrays.copyOf(buffer, readBytes))); + } + return; + } + + int bytesToRead = ByteBuffer.wrap(header, 4, 4).getInt(); + + do { + int readBytes = body.read(buffer); + if (readBytes < 0) { + // TODO log? + return; + } + + if (readBytes == buffer.length) { + resultCallback.onNext(new Frame(streamType, buffer)); + } else { + resultCallback.onNext(new Frame(streamType, Arrays.copyOf(buffer, readBytes))); + } + bytesToRead -= buffer.length; + } while (bytesToRead > 0); + } + } catch (Exception e) { + resultCallback.onError(e); + } + } + + private static StreamType streamType(int streamType) { + switch (streamType) { + case 0: + return StreamType.STDIN; + case 1: + return StreamType.STDOUT; + case 2: + return StreamType.STDERR; + default: + return StreamType.RAW; + } + } +} diff --git a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/FramedSink.java b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/FramedSink.java deleted file mode 100644 index f49fd9247..000000000 --- a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/FramedSink.java +++ /dev/null @@ -1,81 +0,0 @@ -package com.github.dockerjava.okhttp; - -import com.github.dockerjava.api.async.ResultCallback; -import com.github.dockerjava.api.model.Frame; -import com.github.dockerjava.api.model.StreamType; -import okio.BufferedSource; - -import java.io.IOException; -import java.nio.ByteBuffer; -import java.util.Arrays; -import java.util.function.Consumer; - -class FramedSink implements Consumer { - - private static final int HEADER_SIZE = 8; - - private final ResultCallback resultCallback; - - FramedSink(ResultCallback resultCallback) { - this.resultCallback = resultCallback; - } - - @Override - public void accept(BufferedSource source) { - try { - while (true) { - try { - if (source.exhausted()) { - break; - } - } catch (IOException e) { - break; - } - // See https://docs.docker.com/engine/api/v1.37/#operation/ContainerAttach - // [8]byte{STREAM_TYPE, 0, 0, 0, SIZE1, SIZE2, SIZE3, SIZE4}[]byte{OUTPUT} - - if (!source.request(HEADER_SIZE)) { - return; - } - byte[] bytes = source.readByteArray(HEADER_SIZE); - - StreamType streamType = streamType(bytes[0]); - - if (streamType == StreamType.RAW) { - resultCallback.onNext(new Frame(StreamType.RAW, bytes)); - byte[] buffer = new byte[1024]; - while (!source.exhausted()) { - int readBytes = source.read(buffer); - if (readBytes != -1) { - resultCallback.onNext(new Frame(StreamType.RAW, Arrays.copyOf(buffer, readBytes))); - } - } - return; - } - - int payloadSize = ByteBuffer.wrap(bytes, 4, 4).getInt(); - if (!source.request(payloadSize)) { - return; - } - byte[] payload = source.readByteArray(payloadSize); - - resultCallback.onNext(new Frame(streamType, payload)); - } - } catch (Exception e) { - resultCallback.onError(e); - } - } - - private static StreamType streamType(byte streamType) { - switch (streamType) { - case 0: - return StreamType.STDIN; - case 1: - return StreamType.STDOUT; - case 2: - return StreamType.STDERR; - default: - return StreamType.RAW; - } - } -} diff --git a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/HijackingInterceptor.java b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/HijackingInterceptor.java new file mode 100644 index 000000000..29aa524e7 --- /dev/null +++ b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/HijackingInterceptor.java @@ -0,0 +1,61 @@ +package com.github.dockerjava.okhttp; + +import com.github.dockerjava.core.DockerHttpClient; +import okhttp3.Interceptor; +import okhttp3.Request; +import okhttp3.Response; +import okhttp3.internal.connection.Exchange; +import okhttp3.internal.http.RealInterceptorChain; +import okhttp3.internal.ws.RealWebSocket; +import okio.BufferedSink; + +import java.io.IOException; +import java.io.InputStream; + +class HijackingInterceptor implements Interceptor { + + @Override + public Response intercept(Chain chain) throws IOException { + Request request = chain.request(); + Response response = chain.proceed(request); + if (!response.isSuccessful()) { + return response; + } + + DockerHttpClient.Request originalRequest = request.tag(DockerHttpClient.Request.class); + + if (originalRequest == null) { + // WTF? + return response; + } + + InputStream stdin = originalRequest.hijackedInput(); + + if (stdin == null) { + return response; + } + + chain.call().timeout().clearTimeout().clearDeadline(); + + Exchange exchange = ((RealInterceptorChain) chain).exchange(); + RealWebSocket.Streams streams = exchange.newWebSocketStreams(); + Thread thread = new Thread(() -> { + try (BufferedSink sink = streams.sink) { + while (sink.isOpen()) { + int aByte = stdin.read(); + if (aByte < 0) { + break; + } + sink.writeByte(aByte); + sink.emit(); + } + } catch (Exception e) { + throw new RuntimeException(e); + } + }); + thread.setName("okhttp-hijack-streaming-" + System.identityHashCode(request)); + thread.setDaemon(true); + thread.start(); + return response; + } +} diff --git a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/NamedPipeSocketFactory.java b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/NamedPipeSocketFactory.java index 149816ad8..d993a9c52 100644 --- a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/NamedPipeSocketFactory.java +++ b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/NamedPipeSocketFactory.java @@ -1,5 +1,6 @@ package com.github.dockerjava.okhttp; +import com.github.dockerjava.okhttp.OkDockerHttpClient.OkResponse; import com.sun.jna.platform.win32.Kernel32; import javax.net.SocketFactory; @@ -61,7 +62,7 @@ public void connect(SocketAddress endpoint, int timeout) { is = new InputStream() { @Override public int read(byte[] bytes, int off, int len) throws IOException { - if (OkHttpInvocationBuilder.CLOSING.get()) { + if (OkResponse.CLOSING.get()) { return 0; } return file.read(bytes, off, len); @@ -69,7 +70,7 @@ public int read(byte[] bytes, int off, int len) throws IOException { @Override public int read() throws IOException { - if (OkHttpInvocationBuilder.CLOSING.get()) { + if (OkResponse.CLOSING.get()) { return 0; } return file.read(); @@ -77,7 +78,7 @@ public int read() throws IOException { @Override public int read(byte[] bytes) throws IOException { - if (OkHttpInvocationBuilder.CLOSING.get()) { + if (OkResponse.CLOSING.get()) { return 0; } return file.read(bytes); 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 new file mode 100644 index 000000000..cbe849d01 --- /dev/null +++ b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkDockerHttpClient.java @@ -0,0 +1,227 @@ +package com.github.dockerjava.okhttp; + +import com.github.dockerjava.core.DockerClientConfig; +import com.github.dockerjava.core.DockerHttpClient; +import com.github.dockerjava.core.SSLConfig; +import okhttp3.ConnectionPool; +import okhttp3.Dns; +import okhttp3.HttpUrl; +import okhttp3.MediaType; +import okhttp3.OkHttpClient; +import okhttp3.RequestBody; +import okhttp3.ResponseBody; +import okio.BufferedSink; +import okio.Okio; + +import javax.net.ssl.SSLContext; +import javax.net.ssl.X509TrustManager; +import java.io.IOException; +import java.io.InputStream; +import java.io.UncheckedIOException; +import java.net.InetAddress; +import java.net.URI; +import java.security.cert.X509Certificate; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.concurrent.TimeUnit; + +public class OkDockerHttpClient implements DockerHttpClient { + + private static final String SOCKET_SUFFIX = ".socket"; + + final OkHttpClient client; + + final OkHttpClient streamingClient; + + private final HttpUrl baseUrl; + + public OkDockerHttpClient(DockerClientConfig dockerClientConfig) { + okhttp3.OkHttpClient.Builder clientBuilder = new okhttp3.OkHttpClient.Builder() + .addNetworkInterceptor(new HijackingInterceptor()) + .readTimeout(0, TimeUnit.MILLISECONDS) + .retryOnConnectionFailure(true); + + URI dockerHost = dockerClientConfig.getDockerHost(); + switch (dockerHost.getScheme()) { + case "unix": + case "npipe": + String socketPath = dockerHost.getPath(); + + if ("unix".equals(dockerHost.getScheme())) { + clientBuilder.socketFactory(new UnixSocketFactory(socketPath)); + } else { + clientBuilder.socketFactory(new NamedPipeSocketFactory(socketPath)); + } + + clientBuilder + .connectionPool(new ConnectionPool(0, 1, TimeUnit.SECONDS)) + .dns(hostname -> { + if (hostname.endsWith(SOCKET_SUFFIX)) { + return Collections.singletonList(InetAddress.getByAddress(hostname, new byte[]{0, 0, 0, 0})); + } else { + return Dns.SYSTEM.lookup(hostname); + } + }); + break; + default: + } + + boolean isSSL = false; + SSLConfig sslConfig = dockerClientConfig.getSSLConfig(); + if (sslConfig != null) { + try { + SSLContext sslContext = sslConfig.getSSLContext(); + if (sslContext != null) { + isSSL = true; + clientBuilder.sslSocketFactory(sslContext.getSocketFactory(), new TrustAllX509TrustManager()); + } + } catch (Exception e) { + throw new RuntimeException(e); + } + } + + client = clientBuilder.build(); + + streamingClient = client.newBuilder().build(); + + HttpUrl.Builder baseUrlBuilder; + + switch (dockerHost.getScheme()) { + case "unix": + case "npipe": + baseUrlBuilder = new HttpUrl.Builder() + .scheme("http") + .host("docker" + SOCKET_SUFFIX); + break; + case "tcp": + baseUrlBuilder = new HttpUrl.Builder() + .scheme(isSSL ? "https" : "http") + .host(dockerHost.getHost()) + .port(dockerHost.getPort()); + break; + default: + baseUrlBuilder = HttpUrl.get(dockerHost.toString()).newBuilder(); + } + baseUrl = baseUrlBuilder.build(); + } + + private RequestBody toRequestBody(Request request) { + InputStream body = request.body(); + if (body != null) { + return new RequestBody() { + @Override + public MediaType contentType() { + return null; + } + + @Override + public void writeTo(BufferedSink sink) throws IOException { + sink.writeAll(Okio.source(body)); + } + }; + } + switch (request.method()) { + case "POST": + return RequestBody.create(null, ""); + default: + return null; + } + } + + @Override + public Response execute(Request request) { + String url = baseUrl.toString(); + if (url.endsWith("/") && request.path().startsWith("/")) { + url = url.substring(0, url.length() - 1); + } + okhttp3.Request.Builder requestBuilder = new okhttp3.Request.Builder() + .url(url + request.path()) + .tag(Request.class, request) + .method(request.method(), toRequestBody(request)); + + request.headers().forEach(requestBuilder::header); + + final OkHttpClient clientToUse; + + if (request.hijackedInput() == null) { + clientToUse = client; + } else { + clientToUse = streamingClient; + } + + try { + return new OkResponse(clientToUse.newCall(requestBuilder.build()).execute()); + } catch (IOException e) { + throw new UncheckedIOException("Error while executing " + request, e); + } + } + + @Override + public void close() throws IOException { + for (OkHttpClient client : new OkHttpClient[]{client, streamingClient}) { + client.dispatcher().cancelAll(); + client.dispatcher().executorService().shutdown(); + client.connectionPool().evictAll(); + } + } + + static class OkResponse implements Response { + + static final ThreadLocal CLOSING = ThreadLocal.withInitial(() -> false); + + private final okhttp3.Response response; + + public OkResponse(okhttp3.Response response) { + this.response = response; + } + + @Override + public int getStatusCode() { + return response.code(); + } + + @Override + public Map> getHeaders() { + return response.headers().toMultimap(); + } + + @Override + public InputStream getBody() { + ResponseBody body = response.body(); + if (body == null) { + return null; + } + + return body.source().inputStream(); + } + + @Override + public void close() { + boolean previous = CLOSING.get(); + CLOSING.set(true); + try { + response.close(); + } finally { + CLOSING.set(previous); + } + } + } + + static class TrustAllX509TrustManager implements X509TrustManager { + @Override + public void checkClientTrusted(X509Certificate[] x509Certificates, String s) { + + } + + @Override + public void checkServerTrusted(X509Certificate[] x509Certificates, String s) { + + } + + @Override + public X509Certificate[] getAcceptedIssuers() { + return new X509Certificate[0]; + } + } +} diff --git a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpDockerCmdExecFactory.java b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpDockerCmdExecFactory.java index 7b0030e40..e586fb145 100644 --- a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpDockerCmdExecFactory.java +++ b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpDockerCmdExecFactory.java @@ -1,159 +1,32 @@ package com.github.dockerjava.okhttp; -import com.fasterxml.jackson.databind.ObjectMapper; -import com.github.dockerjava.core.AbstractDockerCmdExecFactory; +import com.github.dockerjava.api.command.DelegatingDockerCmdExecFactory; +import com.github.dockerjava.api.command.DockerCmdExecFactory; +import com.github.dockerjava.core.DefaultDockerCmdExecFactory; import com.github.dockerjava.core.DockerClientConfig; -import com.github.dockerjava.core.SSLConfig; -import com.google.common.collect.ImmutableList; -import com.google.common.collect.MultimapBuilder; -import okhttp3.ConnectionPool; -import okhttp3.Dns; -import okhttp3.HttpUrl; -import okhttp3.OkHttpClient; +import com.github.dockerjava.core.DockerClientConfigAware; +import com.github.dockerjava.core.DockerClientImpl; +import com.github.dockerjava.core.DockerHttpClient; -import javax.net.ssl.SSLContext; -import javax.net.ssl.X509TrustManager; -import java.io.IOException; -import java.net.InetAddress; -import java.net.URI; -import java.security.cert.X509Certificate; -import java.util.Collections; -import java.util.concurrent.TimeUnit; +/** + * @deprecated use {@link OkDockerHttpClient} with {@link DockerClientImpl#withHttpClient(DockerHttpClient)} + */ +@Deprecated +public class OkHttpDockerCmdExecFactory extends DelegatingDockerCmdExecFactory implements DockerClientConfigAware { -import static java.util.Objects.nonNull; - -public class OkHttpDockerCmdExecFactory extends AbstractDockerCmdExecFactory { - - private static final String SOCKET_SUFFIX = ".socket"; - - private ObjectMapper objectMapper; - - private OkHttpClient okHttpClient; - private Boolean retryOnConnectionFailure; - - private HttpUrl baseUrl; - - public OkHttpDockerCmdExecFactory setRetryOnConnectionFailure(Boolean retryOnConnectionFailure) { - this.retryOnConnectionFailure = retryOnConnectionFailure; - return this; - } + private DefaultDockerCmdExecFactory dockerCmdExecFactory; @Override - public void init(DockerClientConfig dockerClientConfig) { - super.init(dockerClientConfig); - - OkHttpClient.Builder clientBuilder = new OkHttpClient.Builder(); - if (nonNull(readTimeout)) { - clientBuilder.readTimeout(readTimeout, TimeUnit.MILLISECONDS); - } else { - // default is too small for most docker commands, set default like in jersey/netty - clientBuilder.readTimeout(0, TimeUnit.MILLISECONDS); - } - - if (nonNull(connectTimeout)) { - clientBuilder.connectTimeout(connectTimeout, TimeUnit.MILLISECONDS); - } - - if (nonNull(retryOnConnectionFailure)) { - clientBuilder.retryOnConnectionFailure(retryOnConnectionFailure); - } else { - clientBuilder.retryOnConnectionFailure(true); - } - - URI dockerHost = dockerClientConfig.getDockerHost(); - switch (dockerHost.getScheme()) { - case "unix": - case "npipe": - String socketPath = dockerHost.getPath(); - - if ("unix".equals(dockerHost.getScheme())) { - clientBuilder.socketFactory(new UnixSocketFactory(socketPath)); - } else { - clientBuilder.socketFactory(new NamedPipeSocketFactory(socketPath)); - } - - clientBuilder - // Disable pooling - .connectionPool(new ConnectionPool(0, 1, TimeUnit.SECONDS)) - .dns(hostname -> { - if (hostname.endsWith(SOCKET_SUFFIX)) { - return Collections.singletonList(InetAddress.getByAddress(hostname, new byte[]{0, 0, 0, 0})); - } else { - return Dns.SYSTEM.lookup(hostname); - } - }); - default: - } - - SSLConfig sslConfig = dockerClientConfig.getSSLConfig(); - boolean isSSL = false; - if (sslConfig != null) { - try { - SSLContext sslContext = sslConfig.getSSLContext(); - if (sslContext != null) { - isSSL = true; - clientBuilder.sslSocketFactory(sslContext.getSocketFactory(), new TrustAllX509TrustManager()); - } - } catch (Exception e) { - throw new RuntimeException(e); - } - } - - okHttpClient = clientBuilder.build(); - - HttpUrl.Builder baseUrlBuilder; - - switch (dockerHost.getScheme()) { - case "unix": - case "npipe": - baseUrlBuilder = new HttpUrl.Builder() - .scheme("http") - .host("docker" + SOCKET_SUFFIX); - break; - case "tcp": - baseUrlBuilder = new HttpUrl.Builder() - .scheme(isSSL ? "https" : "http") - .host(dockerHost.getHost()) - .port(dockerHost.getPort()); - break; - default: - baseUrlBuilder = HttpUrl.get(dockerHost.toString()).newBuilder(); - } - baseUrl = baseUrlBuilder.build(); - - objectMapper = dockerClientConfig.getObjectMapper(); + public final DockerCmdExecFactory getDockerCmdExecFactory() { + return dockerCmdExecFactory; } @Override - protected OkHttpWebTarget getBaseResource() { - return new OkHttpWebTarget( - objectMapper, - okHttpClient, - baseUrl, - ImmutableList.of(), - MultimapBuilder.hashKeys().hashSetValues().build() + public void init(DockerClientConfig dockerClientConfig) { + dockerCmdExecFactory = new DefaultDockerCmdExecFactory( + new OkDockerHttpClient(dockerClientConfig), + dockerClientConfig.getObjectMapper() ); - } - - @Override - public void close() throws IOException { - - } - - private static class TrustAllX509TrustManager implements X509TrustManager { - @Override - public void checkClientTrusted(X509Certificate[] x509Certificates, String s) { - - } - - @Override - public void checkServerTrusted(X509Certificate[] x509Certificates, String s) { - - } - - @Override - public X509Certificate[] getAcceptedIssuers() { - return new X509Certificate[0]; - } + dockerCmdExecFactory.init(dockerClientConfig); } } diff --git a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpInvocationBuilder.java b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpInvocationBuilder.java deleted file mode 100644 index 31c45be23..000000000 --- a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpInvocationBuilder.java +++ /dev/null @@ -1,372 +0,0 @@ -package com.github.dockerjava.okhttp; - -import com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.core.type.TypeReference; -import com.fasterxml.jackson.databind.ObjectMapper; -import com.github.dockerjava.api.async.ResultCallback; -import com.github.dockerjava.api.exception.BadRequestException; -import com.github.dockerjava.api.exception.ConflictException; -import com.github.dockerjava.api.exception.DockerException; -import com.github.dockerjava.api.exception.InternalServerErrorException; -import com.github.dockerjava.api.exception.NotAcceptableException; -import com.github.dockerjava.api.exception.NotFoundException; -import com.github.dockerjava.api.exception.NotModifiedException; -import com.github.dockerjava.api.exception.UnauthorizedException; -import com.github.dockerjava.api.model.Frame; -import com.github.dockerjava.core.InvocationBuilder; -import okhttp3.HttpUrl; -import okhttp3.MediaType; -import okhttp3.OkHttpClient; -import okhttp3.Request; -import okhttp3.RequestBody; -import okhttp3.Response; -import okhttp3.internal.connection.RealConnection; -import okio.BufferedSink; -import okio.BufferedSource; -import okio.Okio; -import okio.Source; - -import java.io.ByteArrayInputStream; -import java.io.IOException; -import java.io.InputStream; -import java.lang.reflect.Field; -import java.util.Objects; -import java.util.function.Consumer; - -class OkHttpInvocationBuilder implements InvocationBuilder { - - static final ThreadLocal CLOSING = ThreadLocal.withInitial(() -> false); - - private final ObjectMapper objectMapper; - - private final OkHttpClient okHttpClient; - - private final Request.Builder requestBuilder; - - OkHttpInvocationBuilder(ObjectMapper objectMapper, OkHttpClient okHttpClient, HttpUrl httpUrl) { - this.objectMapper = objectMapper; - this.okHttpClient = okHttpClient; - - requestBuilder = new Request.Builder() - .url(httpUrl); - } - - @Override - public OkHttpInvocationBuilder accept(com.github.dockerjava.core.MediaType mediaType) { - return header("accept", mediaType.getMediaType()); - } - - @Override - public OkHttpInvocationBuilder header(String name, String value) { - requestBuilder.header(name, value); - return this; - } - - @Override - public void delete() { - Request request = requestBuilder - .delete() - .build(); - - execute(request).close(); - } - - @Override - public void get(ResultCallback resultCallback) { - Request request = requestBuilder - .get() - .build(); - - executeAndStream( - request, - resultCallback, - new FramedSink(resultCallback) - ); - } - - @Override - public T get(TypeReference typeReference) { - try (InputStream inputStream = get()) { - return objectMapper.readValue(inputStream, typeReference); - } catch (IOException e) { - throw new RuntimeException(e); - } - } - - @Override - public void get(TypeReference typeReference, ResultCallback resultCallback) { - Request request = requestBuilder - .get() - .build(); - - executeAndStream( - request, - resultCallback, - new JsonSink<>(objectMapper, typeReference, resultCallback) - ); - } - - @Override - public InputStream post(Object entity) { - try { - Request request = requestBuilder - .post(RequestBody.create(MediaType.parse("application/json"), objectMapper.writeValueAsBytes(entity))) - .build(); - - return execute(request).body().byteStream(); - } catch (JsonProcessingException e) { - throw new RuntimeException(e); - } - } - - @Override - public T post(Object entity, TypeReference typeReference) { - try { - Request request = requestBuilder - .post(RequestBody.create(MediaType.parse("application/json"), objectMapper.writeValueAsBytes(entity))) - .build(); - - try (Response response = execute(request)) { - String inputStream = response.body().string(); - return objectMapper.readValue(inputStream, typeReference); - } - } catch (IOException e) { - throw new RuntimeException(e); - } - } - - @Override - public void post(Object entity, TypeReference typeReference, ResultCallback resultCallback) { - try { - post(typeReference, resultCallback, new ByteArrayInputStream(objectMapper.writeValueAsBytes(entity))); - } catch (JsonProcessingException e) { - throw new RuntimeException(e); - } - } - - @Override - public T post(TypeReference typeReference, InputStream body) { - try (InputStream inputStream = post(body)) { - return objectMapper.readValue(inputStream, typeReference); - } catch (IOException e) { - throw new RuntimeException(e); - } - } - - @Override - public void post(Object entity, InputStream stdin, ResultCallback resultCallback) { - final Request request; - try { - request = requestBuilder - .post(RequestBody.create(MediaType.parse("application/json"), objectMapper.writeValueAsBytes(entity))) - .build(); - } catch (JsonProcessingException e) { - throw new RuntimeException(e); - } - - OkHttpClient clientToUse = this.okHttpClient; - - if (stdin != null) { - // FIXME there must be a better way of handling it - clientToUse = clientToUse.newBuilder() - .addNetworkInterceptor(chain -> { - Response response = chain.proceed(chain.request()); - if (response.isSuccessful()) { - Thread thread = new Thread(() -> { - try { - Field sinkField = RealConnection.class.getDeclaredField("sink"); - sinkField.setAccessible(true); - - try ( - BufferedSink sink = (BufferedSink) sinkField.get(chain.connection()); - Source source = Okio.source(stdin); - ) { - while (sink.isOpen()) { - int available = stdin.available(); - if (available > 0) { - sink.write(source, available); - sink.emit(); - } - } - } - } catch (Exception e) { - throw new RuntimeException(e); - } - }); - thread.start(); - } - return response; - }) - .build(); - } - - executeAndStream( - clientToUse, - request, - resultCallback, - new FramedSink(resultCallback) - ); - } - - @Override - public void post(TypeReference typeReference, ResultCallback resultCallback, InputStream body) { - Request request = requestBuilder - .post(toRequestBody(body, null)) - .build(); - - executeAndStream( - request, - resultCallback, - new JsonSink<>(objectMapper, typeReference, resultCallback) - ); - } - - @Override - public void postStream(InputStream body) { - Request request = requestBuilder - .post(toRequestBody(body, null)) - .build(); - - execute(request).close(); - } - - @Override - public InputStream get() { - Request request = requestBuilder - .get() - .build(); - - return execute(request).body().byteStream(); - } - - @Override - public void put(InputStream body, com.github.dockerjava.core.MediaType mediaType) { - Request request = requestBuilder - .put(toRequestBody(body, mediaType.toString())) - .build(); - - execute(request).close(); - } - - protected RequestBody toRequestBody(InputStream body, String mediaType) { - return new RequestBody() { - @Override - public MediaType contentType() { - if (mediaType == null) { - return null; - } - return MediaType.parse(mediaType); - } - - @Override - public void writeTo(BufferedSink sink) throws IOException { - try (Source source = Okio.source(body)) { - sink.writeAll(source); - } - } - }; - } - - protected Response execute(Request request) { - return execute(okHttpClient, request); - } - - protected Response execute(OkHttpClient okHttpClient, Request request) { - try { - Response response = okHttpClient.newCall(request).execute(); - if (!response.isSuccessful()) { - String body = response.body().string(); - switch (response.code()) { - case 304: - throw new NotModifiedException(body); - case 400: - throw new BadRequestException(body); - case 401: - throw new UnauthorizedException(body); - case 404: - throw new NotFoundException(body); - case 406: - throw new NotAcceptableException(body); - case 409: - throw new ConflictException(body); - case 500: - throw new InternalServerErrorException(body); - default: - throw new DockerException(body, response.code()); - } - } else { - return response; - } - } catch (IOException e) { - throw new RuntimeException(e); - } - } - - protected void executeAndStream(Request request, ResultCallback callback, Consumer sourceConsumer) { - executeAndStream(okHttpClient, request, callback, sourceConsumer); - } - - protected void executeAndStream( - OkHttpClient okHttpClient, - Request request, - ResultCallback callback, - Consumer sourceConsumer - ) { - Thread thread = new Thread(() -> { - try ( - Response response = execute(okHttpClient, request.newBuilder().tag("streaming").build()); - BufferedSource source = response.body().source(); - ) { - callback.onStart(() -> { - boolean previous = CLOSING.get(); - CLOSING.set(true); - try { - response.close(); - } finally { - CLOSING.set(previous); - } - }); - - sourceConsumer.accept(source); - callback.onComplete(); - } catch (Exception e) { - callback.onError(e); - } - }, "tc-okhttp-stream-" + Objects.hashCode(request)); - thread.setDaemon(true); - - thread.start(); - } - - private static class JsonSink implements Consumer { - - private final ObjectMapper objectMapper; - - private final TypeReference typeReference; - - private final ResultCallback resultCallback; - - JsonSink(ObjectMapper objectMapper, TypeReference typeReference, ResultCallback resultCallback) { - this.objectMapper = objectMapper; - this.typeReference = typeReference; - this.resultCallback = resultCallback; - } - - @Override - public void accept(BufferedSource source) { - try { - while (true) { - String line = source.readUtf8Line(); - if (line == null) { - break; - } - - resultCallback.onNext(objectMapper.readValue(line, typeReference)); - } - } catch (Exception e) { - resultCallback.onError(e); - } - } - } - -} diff --git a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpWebTarget.java b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpWebTarget.java deleted file mode 100644 index e90012f8d..000000000 --- a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpWebTarget.java +++ /dev/null @@ -1,148 +0,0 @@ -package com.github.dockerjava.okhttp; - -import com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.databind.ObjectMapper; -import com.github.dockerjava.core.InvocationBuilder; -import com.github.dockerjava.core.WebTarget; -import com.google.common.collect.HashMultimap; -import com.google.common.collect.ImmutableList; -import com.google.common.collect.SetMultimap; -import okhttp3.HttpUrl; -import okhttp3.OkHttpClient; -import org.apache.commons.lang.StringUtils; - -import java.util.Collection; -import java.util.Map; -import java.util.Objects; -import java.util.Set; -import java.util.stream.Collectors; - -class OkHttpWebTarget implements WebTarget { - - final ObjectMapper objectMapper; - - final OkHttpClient okHttpClient; - - final HttpUrl baseUrl; - - final ImmutableList path; - - final SetMultimap queryParams; - - OkHttpWebTarget( - ObjectMapper objectMapper, - OkHttpClient okHttpClient, - HttpUrl baseUrl, - ImmutableList path, - SetMultimap queryParams - ) { - this.objectMapper = objectMapper; - this.okHttpClient = okHttpClient; - this.baseUrl = baseUrl; - this.path = path; - this.queryParams = queryParams; - } - - @Override - public InvocationBuilder request() { - String resource = StringUtils.join(path, "/"); - - if (!resource.startsWith("/")) { - resource = "/" + resource; - } - - HttpUrl.Builder baseUrlBuilder = baseUrl.newBuilder() - .encodedPath(resource); - - for (Map.Entry> queryParamEntry : queryParams.asMap().entrySet()) { - String key = queryParamEntry.getKey(); - for (String paramValue : queryParamEntry.getValue()) { - baseUrlBuilder.addQueryParameter(key, paramValue); - } - } - - return new OkHttpInvocationBuilder( - objectMapper, - okHttpClient, - baseUrlBuilder.build() - ); - } - - @Override - public OkHttpWebTarget path(String... components) { - ImmutableList newPath = ImmutableList.builder() - .addAll(path) - .add(components) - .build(); - return new OkHttpWebTarget( - objectMapper, - okHttpClient, - baseUrl, - newPath, - queryParams - ); - } - - @Override - public OkHttpWebTarget resolveTemplate(String name, Object value) { - ImmutableList.Builder newPath = ImmutableList.builder(); - for (String component : path) { - component = component.replaceAll("\\{" + name + "\\}", value.toString()); - newPath.add(component); - } - - return new OkHttpWebTarget( - objectMapper, - okHttpClient, - baseUrl, - newPath.build(), - queryParams - ); - } - - @Override - public OkHttpWebTarget queryParam(String name, Object value) { - if (value == null) { - return this; - } - - SetMultimap newQueryParams = HashMultimap.create(queryParams); - newQueryParams.put(name, value.toString()); - - return new OkHttpWebTarget( - objectMapper, - okHttpClient, - baseUrl, - path, - newQueryParams - ); - } - - @Override - public OkHttpWebTarget queryParamsSet(String name, Set values) { - SetMultimap newQueryParams = HashMultimap.create(queryParams); - newQueryParams.replaceValues(name, values.stream().filter(Objects::nonNull).map(Object::toString).collect(Collectors.toSet())); - - return new OkHttpWebTarget( - objectMapper, - okHttpClient, - baseUrl, - path, - newQueryParams - ); - } - - @Override - public OkHttpWebTarget queryParamsJsonMap(String name, Map values) { - if (values == null || values.isEmpty()) { - return this; - } - - // when param value is JSON string - try { - return queryParam(name, objectMapper.writeValueAsString(values)); - } catch (JsonProcessingException e) { - throw new RuntimeException(e); - } - } -} diff --git a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/UnixSocketFactory.java b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/UnixSocketFactory.java index a32c28845..6c9dbe10b 100644 --- a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/UnixSocketFactory.java +++ b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/UnixSocketFactory.java @@ -1,5 +1,7 @@ package com.github.dockerjava.okhttp; +import com.github.dockerjava.okhttp.OkDockerHttpClient.OkResponse; + import javax.net.SocketFactory; import java.io.FilterInputStream; import java.io.FilterOutputStream; @@ -37,7 +39,7 @@ public void close() throws IOException { @Override public int read(byte[] b, int off, int len) throws IOException { - if (OkHttpInvocationBuilder.CLOSING.get()) { + if (OkResponse.CLOSING.get()) { return 0; } return super.read(b, off, len); diff --git a/docker-java/src/test/java/com/github/dockerjava/junit/DockerCmdExecFactoryDelegate.java b/docker-java/src/test/java/com/github/dockerjava/junit/DockerCmdExecFactoryDelegate.java index dba592a4d..807b92a4b 100644 --- a/docker-java/src/test/java/com/github/dockerjava/junit/DockerCmdExecFactoryDelegate.java +++ b/docker-java/src/test/java/com/github/dockerjava/junit/DockerCmdExecFactoryDelegate.java @@ -1,85 +1,11 @@ package com.github.dockerjava.junit; - -import com.github.dockerjava.api.command.AttachContainerCmd; -import com.github.dockerjava.api.command.AuthCmd; -import com.github.dockerjava.api.command.BuildImageCmd; -import com.github.dockerjava.api.command.CommitCmd; -import com.github.dockerjava.api.command.ConnectToNetworkCmd; -import com.github.dockerjava.api.command.ContainerDiffCmd; -import com.github.dockerjava.api.command.CopyArchiveFromContainerCmd; -import com.github.dockerjava.api.command.CopyArchiveToContainerCmd; -import com.github.dockerjava.api.command.CopyFileFromContainerCmd; -import com.github.dockerjava.api.command.CreateContainerCmd; -import com.github.dockerjava.api.command.CreateImageCmd; -import com.github.dockerjava.api.command.CreateNetworkCmd; -import com.github.dockerjava.api.command.CreateSecretCmd; -import com.github.dockerjava.api.command.CreateServiceCmd; -import com.github.dockerjava.api.command.CreateVolumeCmd; -import com.github.dockerjava.api.command.DisconnectFromNetworkCmd; +import com.github.dockerjava.api.command.DelegatingDockerCmdExecFactory; import com.github.dockerjava.api.command.DockerCmdExecFactory; -import com.github.dockerjava.api.command.EventsCmd; -import com.github.dockerjava.api.command.ExecCreateCmd; -import com.github.dockerjava.api.command.ExecStartCmd; -import com.github.dockerjava.api.command.InfoCmd; -import com.github.dockerjava.api.command.InitializeSwarmCmd; -import com.github.dockerjava.api.command.InspectContainerCmd; -import com.github.dockerjava.api.command.InspectExecCmd; -import com.github.dockerjava.api.command.InspectImageCmd; -import com.github.dockerjava.api.command.InspectNetworkCmd; -import com.github.dockerjava.api.command.InspectServiceCmd; -import com.github.dockerjava.api.command.InspectSwarmCmd; -import com.github.dockerjava.api.command.InspectSwarmNodeCmd; -import com.github.dockerjava.api.command.InspectVolumeCmd; -import com.github.dockerjava.api.command.JoinSwarmCmd; -import com.github.dockerjava.api.command.KillContainerCmd; -import com.github.dockerjava.api.command.LeaveSwarmCmd; -import com.github.dockerjava.api.command.ListContainersCmd; -import com.github.dockerjava.api.command.ListImagesCmd; -import com.github.dockerjava.api.command.ListNetworksCmd; -import com.github.dockerjava.api.command.ListSecretsCmd; -import com.github.dockerjava.api.command.ListServicesCmd; -import com.github.dockerjava.api.command.ListSwarmNodesCmd; -import com.github.dockerjava.api.command.ListTasksCmd; -import com.github.dockerjava.api.command.ListVolumesCmd; -import com.github.dockerjava.api.command.LoadImageCmd; -import com.github.dockerjava.api.command.LogContainerCmd; -import com.github.dockerjava.api.command.LogSwarmObjectCmd; -import com.github.dockerjava.api.command.PauseContainerCmd; -import com.github.dockerjava.api.command.PingCmd; -import com.github.dockerjava.api.command.PruneCmd; -import com.github.dockerjava.api.command.PullImageCmd; -import com.github.dockerjava.api.command.PushImageCmd; -import com.github.dockerjava.api.command.RemoveContainerCmd; -import com.github.dockerjava.api.command.RemoveImageCmd; -import com.github.dockerjava.api.command.RemoveNetworkCmd; -import com.github.dockerjava.api.command.RemoveSecretCmd; -import com.github.dockerjava.api.command.RemoveServiceCmd; -import com.github.dockerjava.api.command.RemoveSwarmNodeCmd; -import com.github.dockerjava.api.command.RemoveVolumeCmd; -import com.github.dockerjava.api.command.RenameContainerCmd; -import com.github.dockerjava.api.command.RestartContainerCmd; -import com.github.dockerjava.api.command.SaveImageCmd; -import com.github.dockerjava.api.command.SaveImagesCmd; -import com.github.dockerjava.api.command.SearchImagesCmd; -import com.github.dockerjava.api.command.StartContainerCmd; -import com.github.dockerjava.api.command.StatsCmd; -import com.github.dockerjava.api.command.StopContainerCmd; -import com.github.dockerjava.api.command.TagImageCmd; -import com.github.dockerjava.api.command.TopContainerCmd; -import com.github.dockerjava.api.command.UnpauseContainerCmd; -import com.github.dockerjava.api.command.UpdateContainerCmd; -import com.github.dockerjava.api.command.UpdateServiceCmd; -import com.github.dockerjava.api.command.UpdateSwarmCmd; -import com.github.dockerjava.api.command.UpdateSwarmNodeCmd; -import com.github.dockerjava.api.command.VersionCmd; -import com.github.dockerjava.api.command.WaitContainerCmd; import com.github.dockerjava.core.DockerClientConfig; import com.github.dockerjava.core.DockerClientConfigAware; -import java.io.IOException; - -class DockerCmdExecFactoryDelegate implements DockerCmdExecFactory, DockerClientConfigAware { +class DockerCmdExecFactoryDelegate extends DelegatingDockerCmdExecFactory implements DockerClientConfigAware { final DockerCmdExecFactory delegate; @@ -87,375 +13,15 @@ class DockerCmdExecFactoryDelegate implements DockerCmdExecFactory, DockerClient this.delegate = delegate; } + @Override + public final DockerCmdExecFactory getDockerCmdExecFactory() { + return delegate; + } + @Override public void init(DockerClientConfig dockerClientConfig) { if (delegate instanceof DockerClientConfigAware) { ((DockerClientConfigAware) delegate).init(dockerClientConfig); } } - - @Override - public AuthCmd.Exec createAuthCmdExec() { - return delegate.createAuthCmdExec(); - } - - @Override - public InfoCmd.Exec createInfoCmdExec() { - return delegate.createInfoCmdExec(); - } - - @Override - public PingCmd.Exec createPingCmdExec() { - return delegate.createPingCmdExec(); - } - - @Override - public ExecCreateCmd.Exec createExecCmdExec() { - return delegate.createExecCmdExec(); - } - - @Override - public VersionCmd.Exec createVersionCmdExec() { - return delegate.createVersionCmdExec(); - } - - @Override - public PullImageCmd.Exec createPullImageCmdExec() { - return delegate.createPullImageCmdExec(); - } - - @Override - public PushImageCmd.Exec createPushImageCmdExec() { - return delegate.createPushImageCmdExec(); - } - - @Override - public SaveImageCmd.Exec createSaveImageCmdExec() { - return delegate.createSaveImageCmdExec(); - } - - @Override - public SaveImagesCmd.Exec createSaveImagesCmdExec() { - return delegate.createSaveImagesCmdExec(); - } - - @Override - public CreateImageCmd.Exec createCreateImageCmdExec() { - return delegate.createCreateImageCmdExec(); - } - - @Override - public LoadImageCmd.Exec createLoadImageCmdExec() { - return delegate.createLoadImageCmdExec(); - } - - @Override - public SearchImagesCmd.Exec createSearchImagesCmdExec() { - return delegate.createSearchImagesCmdExec(); - } - - @Override - public RemoveImageCmd.Exec createRemoveImageCmdExec() { - return delegate.createRemoveImageCmdExec(); - } - - @Override - public ListImagesCmd.Exec createListImagesCmdExec() { - return delegate.createListImagesCmdExec(); - } - - @Override - public InspectImageCmd.Exec createInspectImageCmdExec() { - return delegate.createInspectImageCmdExec(); - } - - @Override - public ListContainersCmd.Exec createListContainersCmdExec() { - return delegate.createListContainersCmdExec(); - } - - @Override - public CreateContainerCmd.Exec createCreateContainerCmdExec() { - return delegate.createCreateContainerCmdExec(); - } - - @Override - public StartContainerCmd.Exec createStartContainerCmdExec() { - return delegate.createStartContainerCmdExec(); - } - - @Override - public InspectContainerCmd.Exec createInspectContainerCmdExec() { - return delegate.createInspectContainerCmdExec(); - } - - @Override - public RemoveContainerCmd.Exec createRemoveContainerCmdExec() { - return delegate.createRemoveContainerCmdExec(); - } - - @Override - public WaitContainerCmd.Exec createWaitContainerCmdExec() { - return delegate.createWaitContainerCmdExec(); - } - - @Override - public AttachContainerCmd.Exec createAttachContainerCmdExec() { - return delegate.createAttachContainerCmdExec(); - } - - @Override - public ExecStartCmd.Exec createExecStartCmdExec() { - return delegate.createExecStartCmdExec(); - } - - @Override - public InspectExecCmd.Exec createInspectExecCmdExec() { - return delegate.createInspectExecCmdExec(); - } - - @Override - public LogContainerCmd.Exec createLogContainerCmdExec() { - return delegate.createLogContainerCmdExec(); - } - - @Override - public CopyFileFromContainerCmd.Exec createCopyFileFromContainerCmdExec() { - return delegate.createCopyFileFromContainerCmdExec(); - } - - @Override - public CopyArchiveFromContainerCmd.Exec createCopyArchiveFromContainerCmdExec() { - return delegate.createCopyArchiveFromContainerCmdExec(); - } - - @Override - public CopyArchiveToContainerCmd.Exec createCopyArchiveToContainerCmdExec() { - return delegate.createCopyArchiveToContainerCmdExec(); - } - - @Override - public StopContainerCmd.Exec createStopContainerCmdExec() { - return delegate.createStopContainerCmdExec(); - } - - @Override - public ContainerDiffCmd.Exec createContainerDiffCmdExec() { - return delegate.createContainerDiffCmdExec(); - } - - @Override - public KillContainerCmd.Exec createKillContainerCmdExec() { - return delegate.createKillContainerCmdExec(); - } - - @Override - public UpdateContainerCmd.Exec createUpdateContainerCmdExec() { - return delegate.createUpdateContainerCmdExec(); - } - - @Override - public RenameContainerCmd.Exec createRenameContainerCmdExec() { - return delegate.createRenameContainerCmdExec(); - } - - @Override - public RestartContainerCmd.Exec createRestartContainerCmdExec() { - return delegate.createRestartContainerCmdExec(); - } - - @Override - public CommitCmd.Exec createCommitCmdExec() { - return delegate.createCommitCmdExec(); - } - - @Override - public BuildImageCmd.Exec createBuildImageCmdExec() { - return delegate.createBuildImageCmdExec(); - } - - @Override - public TopContainerCmd.Exec createTopContainerCmdExec() { - return delegate.createTopContainerCmdExec(); - } - - @Override - public TagImageCmd.Exec createTagImageCmdExec() { - return delegate.createTagImageCmdExec(); - } - - @Override - public PauseContainerCmd.Exec createPauseContainerCmdExec() { - return delegate.createPauseContainerCmdExec(); - } - - @Override - public UnpauseContainerCmd.Exec createUnpauseContainerCmdExec() { - return delegate.createUnpauseContainerCmdExec(); - } - - @Override - public EventsCmd.Exec createEventsCmdExec() { - return delegate.createEventsCmdExec(); - } - - @Override - public StatsCmd.Exec createStatsCmdExec() { - return delegate.createStatsCmdExec(); - } - - @Override - public CreateVolumeCmd.Exec createCreateVolumeCmdExec() { - return delegate.createCreateVolumeCmdExec(); - } - - @Override - public InspectVolumeCmd.Exec createInspectVolumeCmdExec() { - return delegate.createInspectVolumeCmdExec(); - } - - @Override - public RemoveVolumeCmd.Exec createRemoveVolumeCmdExec() { - return delegate.createRemoveVolumeCmdExec(); - } - - @Override - public ListVolumesCmd.Exec createListVolumesCmdExec() { - return delegate.createListVolumesCmdExec(); - } - - @Override - public ListNetworksCmd.Exec createListNetworksCmdExec() { - return delegate.createListNetworksCmdExec(); - } - - @Override - public InspectNetworkCmd.Exec createInspectNetworkCmdExec() { - return delegate.createInspectNetworkCmdExec(); - } - - @Override - public CreateNetworkCmd.Exec createCreateNetworkCmdExec() { - return delegate.createCreateNetworkCmdExec(); - } - - @Override - public RemoveNetworkCmd.Exec createRemoveNetworkCmdExec() { - return delegate.createRemoveNetworkCmdExec(); - } - - @Override - public ConnectToNetworkCmd.Exec createConnectToNetworkCmdExec() { - return delegate.createConnectToNetworkCmdExec(); - } - - @Override - public DisconnectFromNetworkCmd.Exec createDisconnectFromNetworkCmdExec() { - return delegate.createDisconnectFromNetworkCmdExec(); - } - - @Override - public InitializeSwarmCmd.Exec createInitializeSwarmCmdExec() { - return delegate.createInitializeSwarmCmdExec(); - } - - @Override - public InspectSwarmCmd.Exec createInspectSwarmCmdExec() { - return delegate.createInspectSwarmCmdExec(); - } - - @Override - public JoinSwarmCmd.Exec createJoinSwarmCmdExec() { - return delegate.createJoinSwarmCmdExec(); - } - - @Override - public LeaveSwarmCmd.Exec createLeaveSwarmCmdExec() { - return delegate.createLeaveSwarmCmdExec(); - } - - @Override - public UpdateSwarmCmd.Exec createUpdateSwarmCmdExec() { - return delegate.createUpdateSwarmCmdExec(); - } - - @Override - public ListServicesCmd.Exec createListServicesCmdExec() { - return delegate.createListServicesCmdExec(); - } - - @Override - public CreateServiceCmd.Exec createCreateServiceCmdExec() { - return delegate.createCreateServiceCmdExec(); - } - - @Override - public InspectServiceCmd.Exec createInspectServiceCmdExec() { - return delegate.createInspectServiceCmdExec(); - } - - @Override - public UpdateServiceCmd.Exec createUpdateServiceCmdExec() { - return delegate.createUpdateServiceCmdExec(); - } - - @Override - public RemoveServiceCmd.Exec createRemoveServiceCmdExec() { - return delegate.createRemoveServiceCmdExec(); - } - - @Override - public LogSwarmObjectCmd.Exec logSwarmObjectExec(String endpoint) { - return delegate.logSwarmObjectExec(endpoint); - } - - @Override - public ListSwarmNodesCmd.Exec listSwarmNodeCmdExec() { - return delegate.listSwarmNodeCmdExec(); - } - - @Override - public InspectSwarmNodeCmd.Exec inspectSwarmNodeCmdExec() { - return delegate.inspectSwarmNodeCmdExec(); - } - - @Override - public RemoveSwarmNodeCmd.Exec removeSwarmNodeCmdExec() { - return delegate.removeSwarmNodeCmdExec(); - } - - @Override - public UpdateSwarmNodeCmd.Exec updateSwarmNodeCmdExec() { - return delegate.updateSwarmNodeCmdExec(); - } - - @Override - public ListTasksCmd.Exec listTasksCmdExec() { - return delegate.listTasksCmdExec(); - } - - @Override - public PruneCmd.Exec pruneCmdExec() { - return delegate.pruneCmdExec(); - } - - @Override - public ListSecretsCmd.Exec createListSecretsCmdExec() { - return delegate.createListSecretsCmdExec(); - } - - @Override - public CreateSecretCmd.Exec createCreateSecretCmdExec() { - return delegate.createCreateSecretCmdExec(); - } - - @Override - public RemoveSecretCmd.Exec createRemoveSecretCmdExec() { - return delegate.createRemoveSecretCmdExec(); - } - - @Override - public void close() throws IOException { - delegate.close(); - } } From 03e846933474659f91dfed71c44f1a88f728a056 Mon Sep 17 00:00:00 2001 From: Sergei Egorov Date: Thu, 19 Mar 2020 20:43:31 +0100 Subject: [PATCH 02/13] fix checkstyle issues --- .../com/github/dockerjava/core/DefaultDockerCmdExecFactory.java | 2 +- .../com/github/dockerjava/core/DefaultInvocationBuilder.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/docker-java-core/src/main/java/com/github/dockerjava/core/DefaultDockerCmdExecFactory.java b/docker-java-core/src/main/java/com/github/dockerjava/core/DefaultDockerCmdExecFactory.java index 8f3dc3b04..5e13824ef 100644 --- a/docker-java-core/src/main/java/com/github/dockerjava/core/DefaultDockerCmdExecFactory.java +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/DefaultDockerCmdExecFactory.java @@ -50,7 +50,7 @@ private class DefaultWebTarget implements WebTarget { final SetMultimap queryParams; - public DefaultWebTarget() { + DefaultWebTarget() { this( ImmutableList.of(), MultimapBuilder.hashKeys().hashSetValues().build() 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 f94ac19c7..5b75a3fbe 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 @@ -29,7 +29,7 @@ class DefaultInvocationBuilder implements InvocationBuilder { private final DockerHttpClient dockerHttpClient; private final ObjectMapper objectMapper; - public DefaultInvocationBuilder(DockerHttpClient dockerHttpClient, ObjectMapper objectMapper, String path) { + DefaultInvocationBuilder(DockerHttpClient dockerHttpClient, ObjectMapper objectMapper, String path) { this.requestBuilder = DockerHttpClient.Request.builder().path(path); this.dockerHttpClient = dockerHttpClient; this.objectMapper = objectMapper; From d018fe422b7aad0d0653fcf16ed19b685887a98d Mon Sep 17 00:00:00 2001 From: Sergei Egorov Date: Thu, 19 Mar 2020 20:48:42 +0100 Subject: [PATCH 03/13] fix checkstyle issues (2) --- .../github/dockerjava/okhttp/OkDockerHttpClient.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) 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 cbe849d01..79de193b4 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 @@ -159,10 +159,10 @@ public Response execute(Request request) { @Override public void close() throws IOException { - for (OkHttpClient client : new OkHttpClient[]{client, streamingClient}) { - client.dispatcher().cancelAll(); - client.dispatcher().executorService().shutdown(); - client.connectionPool().evictAll(); + for (OkHttpClient clientToClose : new OkHttpClient[]{client, streamingClient}) { + clientToClose.dispatcher().cancelAll(); + clientToClose.dispatcher().executorService().shutdown(); + clientToClose.connectionPool().evictAll(); } } @@ -172,7 +172,7 @@ static class OkResponse implements Response { private final okhttp3.Response response; - public OkResponse(okhttp3.Response response) { + OkResponse(okhttp3.Response response) { this.response = response; } From c3ab8e0f8cc1a377f033107615287802d9a7d52a Mon Sep 17 00:00:00 2001 From: Sergei Egorov Date: Thu, 19 Mar 2020 20:52:43 +0100 Subject: [PATCH 04/13] Fix CmdIT --- docker-java/src/test/java/com/github/dockerjava/cmd/CmdIT.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 215a552b6..8981d1de2 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 @@ -35,7 +35,7 @@ public DockerCmdExecFactory createExecFactory() { OKHTTP(true) { @Override public DockerCmdExecFactory createExecFactory() { - return new OkHttpDockerCmdExecFactory().withConnectTimeout(30 * 1000); + return new OkHttpDockerCmdExecFactory(); } }; From e562fec97f2769b1528f0a2a771b50b169eef7ca Mon Sep 17 00:00:00 2001 From: Sergei Egorov Date: Fri, 20 Mar 2020 13:49:18 +0100 Subject: [PATCH 05/13] Introduce `OkDockerHttpClient.Factory` and restore backward compatibility --- .../dockerjava/okhttp/OkDockerHttpClient.java | 61 ++++++++++++++++++- .../okhttp/OkHttpDockerCmdExecFactory.java | 34 ++++++++++- 2 files changed, 92 insertions(+), 3 deletions(-) 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 79de193b4..8853dd815 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 @@ -26,7 +26,47 @@ import java.util.Map; import java.util.concurrent.TimeUnit; -public class OkDockerHttpClient implements DockerHttpClient { +public final class OkDockerHttpClient implements DockerHttpClient { + + public static final class Factory { + + private DockerClientConfig dockerClientConfig = null; + + private Integer readTimeout = null; + + private Integer connectTimeout = null; + + private Boolean retryOnConnectionFailure = null; + + public Factory dockerClientConfig(DockerClientConfig value) { + this.dockerClientConfig = value; + return this; + } + + public Factory readTimeout(Integer value) { + this.readTimeout = value; + return this; + } + + public Factory connectTimeout(Integer value) { + this.connectTimeout = value; + return this; + } + + Factory retryOnConnectionFailure(Boolean value) { + this.retryOnConnectionFailure = value; + return this; + } + + public OkDockerHttpClient build() { + return new OkDockerHttpClient( + dockerClientConfig, + readTimeout, + connectTimeout, + retryOnConnectionFailure + ); + } + } private static final String SOCKET_SUFFIX = ".socket"; @@ -36,12 +76,29 @@ public class OkDockerHttpClient implements DockerHttpClient { private final HttpUrl baseUrl; - public OkDockerHttpClient(DockerClientConfig dockerClientConfig) { + private OkDockerHttpClient( + DockerClientConfig dockerClientConfig, + Integer readTimeout, + Integer connectTimeout, + Boolean retryOnConnectionFailure + ) { okhttp3.OkHttpClient.Builder clientBuilder = new okhttp3.OkHttpClient.Builder() .addNetworkInterceptor(new HijackingInterceptor()) .readTimeout(0, TimeUnit.MILLISECONDS) .retryOnConnectionFailure(true); + if (readTimeout != null) { + clientBuilder.readTimeout(readTimeout, TimeUnit.MILLISECONDS); + } + + if (connectTimeout != null) { + clientBuilder.connectTimeout(connectTimeout, TimeUnit.MILLISECONDS); + } + + if (retryOnConnectionFailure != null) { + clientBuilder.retryOnConnectionFailure(retryOnConnectionFailure); + } + URI dockerHost = dockerClientConfig.getDockerHost(); switch (dockerHost.getScheme()) { case "unix": diff --git a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpDockerCmdExecFactory.java b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpDockerCmdExecFactory.java index e586fb145..90f37b7d8 100644 --- a/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpDockerCmdExecFactory.java +++ b/docker-java-transport-okhttp/src/main/java/com/github/dockerjava/okhttp/OkHttpDockerCmdExecFactory.java @@ -14,8 +14,39 @@ @Deprecated public class OkHttpDockerCmdExecFactory extends DelegatingDockerCmdExecFactory implements DockerClientConfigAware { + private OkDockerHttpClient.Factory clientFactory = new OkDockerHttpClient.Factory(); + + @Deprecated + protected Integer connectTimeout; + + @Deprecated + protected Integer readTimeout; + private DefaultDockerCmdExecFactory dockerCmdExecFactory; + /** + * Configure connection timeout in milliseconds + */ + public OkHttpDockerCmdExecFactory withConnectTimeout(Integer connectTimeout) { + clientFactory = clientFactory.connectTimeout(connectTimeout); + this.connectTimeout = connectTimeout; + return this; + } + + /** + * Configure read timeout in milliseconds + */ + public OkHttpDockerCmdExecFactory withReadTimeout(Integer readTimeout) { + clientFactory = clientFactory.readTimeout(readTimeout); + this.readTimeout = readTimeout; + return this; + } + + public OkHttpDockerCmdExecFactory setRetryOnConnectionFailure(Boolean retryOnConnectionFailure) { + this.clientFactory = clientFactory.retryOnConnectionFailure(retryOnConnectionFailure); + return this; + } + @Override public final DockerCmdExecFactory getDockerCmdExecFactory() { return dockerCmdExecFactory; @@ -23,8 +54,9 @@ public final DockerCmdExecFactory getDockerCmdExecFactory() { @Override public void init(DockerClientConfig dockerClientConfig) { + clientFactory = clientFactory.dockerClientConfig(dockerClientConfig); dockerCmdExecFactory = new DefaultDockerCmdExecFactory( - new OkDockerHttpClient(dockerClientConfig), + clientFactory.build(), dockerClientConfig.getObjectMapper() ); dockerCmdExecFactory.init(dockerClientConfig); From 3223d54889265190082b0261e0a5b57991b50922 Mon Sep 17 00:00:00 2001 From: Sergei Egorov Date: Fri, 20 Mar 2020 15:54:03 +0100 Subject: [PATCH 06/13] fix FramedInputStreamConsumer --- .../core/DefaultInvocationBuilder.java | 2 +- .../core/FramedInputStreamConsumer.java | 44 +++++++----- .../core/async/FrameStreamProcessor.java | 2 + .../core/async/JsonStreamProcessor.java | 2 + .../dockerjava/core/command/FrameReader.java | 1 + .../dockerjava/cmd/CreateContainerCmdIT.java | 4 +- .../core/async/JsonStreamProcessorTest.java | 67 ------------------- 7 files changed, 36 insertions(+), 86 deletions(-) delete mode 100644 docker-java/src/test/java/com/github/dockerjava/core/async/JsonStreamProcessorTest.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 5b75a3fbe..50bd4612b 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 @@ -251,7 +251,7 @@ protected void executeAndStream( } catch (Exception e) { callback.onError(e); } - }, "docker-java-okhttp-stream-" + Objects.hashCode(request)); + }, "docker-java-stream-" + Objects.hashCode(request)); thread.setDaemon(true); thread.start(); diff --git a/docker-java-core/src/main/java/com/github/dockerjava/core/FramedInputStreamConsumer.java b/docker-java-core/src/main/java/com/github/dockerjava/core/FramedInputStreamConsumer.java index ce707eb58..62390c0c7 100644 --- a/docker-java-core/src/main/java/com/github/dockerjava/core/FramedInputStreamConsumer.java +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/FramedInputStreamConsumer.java @@ -5,14 +5,11 @@ import com.github.dockerjava.api.model.StreamType; import java.io.InputStream; -import java.nio.ByteBuffer; import java.util.Arrays; import java.util.function.Consumer; class FramedInputStreamConsumer implements Consumer { - private static final int HEADER_SIZE = 8; - private final ResultCallback resultCallback; FramedInputStreamConsumer(ResultCallback resultCallback) { @@ -23,37 +20,52 @@ class FramedInputStreamConsumer implements Consumer { public void accept(DockerHttpClient.Response response) { try { InputStream body = response.getBody(); + + byte[] buffer = new byte[1024]; while (true) { // See https://docs.docker.com/engine/api/v1.37/#operation/ContainerAttach // [8]byte{STREAM_TYPE, 0, 0, 0, SIZE1, SIZE2, SIZE3, SIZE4}[]byte{OUTPUT} - byte[] header = new byte[HEADER_SIZE]; - if (body.read(header) < 0) { - // TODO log? + int streamTypeByte = body.read(); + if (streamTypeByte < 0) { return; } - int streamTypeByte = header[0]; - StreamType streamType = streamType(streamTypeByte); - byte[] buffer = new byte[1024]; - if (streamType == StreamType.RAW) { - // FIXME should check the header instead - resultCallback.onNext(new Frame(StreamType.RAW, header)); + resultCallback.onNext(new Frame(StreamType.RAW, new byte[] { (byte) streamTypeByte })); int readBytes; while ((readBytes = body.read(buffer)) >= 0) { - resultCallback.onNext(new Frame(StreamType.RAW, Arrays.copyOf(buffer, readBytes))); + if (readBytes == buffer.length) { + resultCallback.onNext(new Frame(StreamType.RAW, buffer)); + } else { + resultCallback.onNext(new Frame(StreamType.RAW, Arrays.copyOf(buffer, readBytes))); + } } return; } - int bytesToRead = ByteBuffer.wrap(header, 4, 4).getInt(); + // Skip 3 bytes + for (int i = 0; i < 3; i++) { + if (body.read() < 0) { + return; + } + } + + // uint32 encoded as big endian. + int bytesToRead = 0; + for (int i = 0; i < 4; i++) { + int readByte = body.read(); + if (readByte < 0) { + return; + } + bytesToRead |= (readByte & 0xff) << (8 * (3 - i)); + } do { - int readBytes = body.read(buffer); + int readBytes = body.read(buffer, 0, Math.min(buffer.length, bytesToRead)); if (readBytes < 0) { // TODO log? return; @@ -64,7 +76,7 @@ public void accept(DockerHttpClient.Response response) { } else { resultCallback.onNext(new Frame(streamType, Arrays.copyOf(buffer, readBytes))); } - bytesToRead -= buffer.length; + bytesToRead -= readBytes; } while (bytesToRead > 0); } } catch (Exception e) { diff --git a/docker-java-core/src/main/java/com/github/dockerjava/core/async/FrameStreamProcessor.java b/docker-java-core/src/main/java/com/github/dockerjava/core/async/FrameStreamProcessor.java index 3070930d6..4c3136417 100644 --- a/docker-java-core/src/main/java/com/github/dockerjava/core/async/FrameStreamProcessor.java +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/async/FrameStreamProcessor.java @@ -16,6 +16,8 @@ * @author Marcus Linke * */ +@SuppressWarnings("unused") +@Deprecated public class FrameStreamProcessor implements ResponseStreamProcessor { @Override diff --git a/docker-java-core/src/main/java/com/github/dockerjava/core/async/JsonStreamProcessor.java b/docker-java-core/src/main/java/com/github/dockerjava/core/async/JsonStreamProcessor.java index 8829decac..64aabd999 100644 --- a/docker-java-core/src/main/java/com/github/dockerjava/core/async/JsonStreamProcessor.java +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/async/JsonStreamProcessor.java @@ -20,6 +20,8 @@ * @author Marcus Linke * */ +@SuppressWarnings("unused") +@Deprecated public class JsonStreamProcessor implements ResponseStreamProcessor { private static final JsonFactory JSON_FACTORY = new JsonFactory(); diff --git a/docker-java-core/src/main/java/com/github/dockerjava/core/command/FrameReader.java b/docker-java-core/src/main/java/com/github/dockerjava/core/command/FrameReader.java index 9c9c31e74..1dc5d3503 100644 --- a/docker-java-core/src/main/java/com/github/dockerjava/core/command/FrameReader.java +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/command/FrameReader.java @@ -14,6 +14,7 @@ *

* See: {@link }http://docs.docker.com/v1.6/reference/api/docker_remote_api_v1.13/#attach-to-a-container} */ +@Deprecated public class FrameReader implements AutoCloseable { private static final int HEADER_SIZE = 8; diff --git a/docker-java/src/test/java/com/github/dockerjava/cmd/CreateContainerCmdIT.java b/docker-java/src/test/java/com/github/dockerjava/cmd/CreateContainerCmdIT.java index 1d15e447e..227c0acc9 100644 --- a/docker-java/src/test/java/com/github/dockerjava/cmd/CreateContainerCmdIT.java +++ b/docker-java/src/test/java/com/github/dockerjava/cmd/CreateContainerCmdIT.java @@ -935,7 +935,7 @@ public void testWithStopSignal() throws Exception { .awaitCompletion(); String log = callback.builder.toString(); - assertThat(log, is("exit trapped 10")); + assertThat(log.trim(), is("exit trapped 10")); } private static class StringBuilderLogReader extends ResultCallback.Adapter { @@ -947,7 +947,7 @@ public StringBuilderLogReader(StringBuilder builder) { @Override public void onNext(Frame item) { - builder.append(new String(item.getPayload()).trim()); + builder.append(new String(item.getPayload())); super.onNext(item); } } diff --git a/docker-java/src/test/java/com/github/dockerjava/core/async/JsonStreamProcessorTest.java b/docker-java/src/test/java/com/github/dockerjava/core/async/JsonStreamProcessorTest.java deleted file mode 100644 index f98d7d281..000000000 --- a/docker-java/src/test/java/com/github/dockerjava/core/async/JsonStreamProcessorTest.java +++ /dev/null @@ -1,67 +0,0 @@ -/* - * Created on 16.02.2016 - */ -package com.github.dockerjava.core.async; - -import com.fasterxml.jackson.core.type.TypeReference; -import com.github.dockerjava.api.async.ResultCallback; -import com.github.dockerjava.api.model.PullResponseItem; -import com.github.dockerjava.test.serdes.JSONTestHelper; -import org.junit.Test; - -import java.io.ByteArrayInputStream; -import java.io.Closeable; -import java.io.IOException; -import java.io.InputStream; -import java.util.ArrayList; -import java.util.List; - -import static org.junit.Assert.assertFalse; - - -/** - * - * @author Marcus Linke - * - */ -public class JsonStreamProcessorTest { - - @Test - public void processEmptyJson() throws Exception { - - InputStream response = new ByteArrayInputStream("{}".getBytes()); - - JsonStreamProcessor jsonStreamProcessor = new JsonStreamProcessor<>( - JSONTestHelper.getMapper(), new TypeReference() {}); - - final List completed = new ArrayList<>(); - - jsonStreamProcessor.processResponseStream(response, new ResultCallback() { - - @Override - public void close() throws IOException { - } - - @Override - public void onStart(Closeable closeable) { - } - - @Override - public void onNext(PullResponseItem object) { - assertFalse("onNext called for empty json", true); - } - - @Override - public void onError(Throwable throwable) { - } - - @Override - public void onComplete() { - completed.add(true); - } - }); - - assertFalse("Stream processing not completed", completed.isEmpty()); - } - -} From c584ac207ec3a66404ffff98cfe3ea19d85bf4fe Mon Sep 17 00:00:00 2001 From: Sergei Egorov Date: Fri, 20 Mar 2020 15:54:40 +0100 Subject: [PATCH 07/13] Add `Response#getHeader(String)` --- .../github/dockerjava/core/DockerHttpClient.java | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/docker-java-core/src/main/java/com/github/dockerjava/core/DockerHttpClient.java b/docker-java-core/src/main/java/com/github/dockerjava/core/DockerHttpClient.java index 7c412fd78..6b56fd07f 100644 --- a/docker-java-core/src/main/java/com/github/dockerjava/core/DockerHttpClient.java +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/DockerHttpClient.java @@ -2,6 +2,7 @@ import org.immutables.value.Value; +import javax.annotation.Nonnull; import javax.annotation.Nullable; import java.io.Closeable; import java.io.InputStream; @@ -22,6 +23,18 @@ interface Response extends Closeable { @Override void close(); + + @Nullable + default String getHeader(@Nonnull String name) { + for (Map.Entry> entry : getHeaders().entrySet()) { + if (name.equalsIgnoreCase(entry.getKey())) { + List values = entry.getValue(); + return values.isEmpty() ? null : values.get(0); + } + } + + return null; + } } @Value.Immutable From a8c38da89a87b04fe537008c2e20a7f3be2dbc13 Mon Sep 17 00:00:00 2001 From: Sergei Egorov Date: Fri, 20 Mar 2020 15:56:48 +0100 Subject: [PATCH 08/13] test improvements --- .../dockerjava/cmd/AttachContainerCmdIT.java | 18 ++++++++++-------- .../dockerjava/cmd/LogContainerCmdIT.java | 4 ++-- .../dockerjava/cmd/StartContainerCmdIT.java | 8 ++++---- 3 files changed, 16 insertions(+), 14 deletions(-) diff --git a/docker-java/src/test/java/com/github/dockerjava/cmd/AttachContainerCmdIT.java b/docker-java/src/test/java/com/github/dockerjava/cmd/AttachContainerCmdIT.java index 41eac8d27..299213216 100644 --- a/docker-java/src/test/java/com/github/dockerjava/cmd/AttachContainerCmdIT.java +++ b/docker-java/src/test/java/com/github/dockerjava/cmd/AttachContainerCmdIT.java @@ -73,21 +73,23 @@ public void onNext(Frame frame) { } }; - PipedOutputStream out = new PipedOutputStream(); - PipedInputStream in = new PipedInputStream(out); - - dockerClient.attachContainerCmd(container.getId()) + try ( + PipedOutputStream out = new PipedOutputStream(); + PipedInputStream in = new PipedInputStream(out); + ) { + dockerClient.attachContainerCmd(container.getId()) .withStdErr(true) .withStdOut(true) .withFollowStream(true) .withStdIn(in) .exec(callback); - out.write((snippet + "\n").getBytes()); - out.flush(); + out.write((snippet + "\n").getBytes()); + out.flush(); - callback.awaitCompletion(15, SECONDS); - callback.close(); + callback.awaitCompletion(15, SECONDS); + callback.close(); + } assertThat(callback.toString(), containsString(snippet)); } 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 6be307d88..37bf5f393 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 @@ -82,7 +82,7 @@ public void asyncLogContainerWithTtyDisabled() throws Exception { assertTrue(loggingCallback.toString().contains("hello")); - assertEquals(loggingCallback.getCollectedFrames().get(0).getStreamType(), StreamType.STDOUT); + assertEquals(StreamType.STDOUT, loggingCallback.getCollectedFrames().get(0).getStreamType()); } @Test @@ -92,7 +92,7 @@ public void asyncLogNonExistingContainer() throws Exception { @Override public void onError(Throwable throwable) { - assertEquals(throwable.getClass().getName(), NotFoundException.class.getName()); + assertEquals(NotFoundException.class.getName(), throwable.getClass().getName()); try { // close the callback to prevent the call to onComplete diff --git a/docker-java/src/test/java/com/github/dockerjava/cmd/StartContainerCmdIT.java b/docker-java/src/test/java/com/github/dockerjava/cmd/StartContainerCmdIT.java index f4180e81a..76e4fe329 100644 --- a/docker-java/src/test/java/com/github/dockerjava/cmd/StartContainerCmdIT.java +++ b/docker-java/src/test/java/com/github/dockerjava/cmd/StartContainerCmdIT.java @@ -54,7 +54,7 @@ public void startContainerWithVolumes() throws Exception { CreateContainerResponse container = dockerRule.getClient().createContainerCmd("busybox").withVolumes(volume1, volume2) .withCmd("true") .withHostConfig(newHostConfig() - .withBinds(new Bind("/src/webapp1", volume1, ro), new Bind("/src/webapp2", volume2))) + .withBinds(new Bind("/tmp/webapp1", volume1, ro), new Bind("/tmp/webapp2", volume2))) .exec(); LOG.info("Created container {}", container.toString()); @@ -78,9 +78,9 @@ public void startContainerWithVolumes() throws Exception { assertThat(mounts, hasSize(2)); final InspectContainerResponse.Mount mount1 = new InspectContainerResponse.Mount() - .withRw(false).withMode("ro").withDestination(volume1).withSource("/src/webapp1"); + .withRw(false).withMode("ro").withDestination(volume1).withSource("/tmp/webapp1"); final InspectContainerResponse.Mount mount2 = new InspectContainerResponse.Mount() - .withRw(true).withMode("rw").withDestination(volume2).withSource("/src/webapp2"); + .withRw(true).withMode("rw").withDestination(volume2).withSource("/tmp/webapp2"); assertThat(mounts, containsInAnyOrder(mount1, mount2)); } @@ -96,7 +96,7 @@ public void startContainerWithVolumesFrom() throws DockerException { CreateContainerResponse container1 = dockerRule.getClient().createContainerCmd("busybox").withCmd("sleep", "9999") .withName(container1Name) .withHostConfig(newHostConfig() - .withBinds(new Bind("/src/webapp1", volume1), new Bind("/src/webapp2", volume2))) + .withBinds(new Bind("/tmp/webapp1", volume1), new Bind("/tmp/webapp2", volume2))) .exec(); LOG.info("Created container1 {}", container1.toString()); From 7dfab442d98bfdd453f30ec918dd397476d4790d Mon Sep 17 00:00:00 2001 From: Sergei Egorov Date: Fri, 20 Mar 2020 16:13:48 +0100 Subject: [PATCH 09/13] restore `withConnectTimeout` --- docker-java/src/test/java/com/github/dockerjava/cmd/CmdIT.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 8981d1de2..215a552b6 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 @@ -35,7 +35,7 @@ public DockerCmdExecFactory createExecFactory() { OKHTTP(true) { @Override public DockerCmdExecFactory createExecFactory() { - return new OkHttpDockerCmdExecFactory(); + return new OkHttpDockerCmdExecFactory().withConnectTimeout(30 * 1000); } }; From b08ba85fbef55b72231e8f11670ef18c01ad11b6 Mon Sep 17 00:00:00 2001 From: Sergei Egorov Date: Fri, 20 Mar 2020 17:56:10 +0100 Subject: [PATCH 10/13] fix code style in `FramedInputStreamConsumer` --- .../com/github/dockerjava/core/FramedInputStreamConsumer.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docker-java-core/src/main/java/com/github/dockerjava/core/FramedInputStreamConsumer.java b/docker-java-core/src/main/java/com/github/dockerjava/core/FramedInputStreamConsumer.java index 62390c0c7..fde82647e 100644 --- a/docker-java-core/src/main/java/com/github/dockerjava/core/FramedInputStreamConsumer.java +++ b/docker-java-core/src/main/java/com/github/dockerjava/core/FramedInputStreamConsumer.java @@ -34,7 +34,7 @@ public void accept(DockerHttpClient.Response response) { StreamType streamType = streamType(streamTypeByte); if (streamType == StreamType.RAW) { - resultCallback.onNext(new Frame(StreamType.RAW, new byte[] { (byte) streamTypeByte })); + resultCallback.onNext(new Frame(StreamType.RAW, new byte[]{(byte) streamTypeByte})); int readBytes; while ((readBytes = body.read(buffer)) >= 0) { From e9fdc65d00a3b2a29a843dbf95e183e7fbf8f412 Mon Sep 17 00:00:00 2001 From: Sergei Egorov Date: Fri, 20 Mar 2020 22:20:51 +0100 Subject: [PATCH 11/13] Add JerseyDockerHttpClient --- .../dockerjava/jaxrs/ApacheUnixSocket.java | 5 +- .../jaxrs/JerseyDockerCmdExecFactory.java | 297 +++----------- .../jaxrs/JerseyDockerHttpClient.java | 381 ++++++++++++++++++ .../jaxrs/JerseyInvocationBuilder.java | 193 --------- .../dockerjava/jaxrs/JerseyWebTarget.java | 79 ---- .../jaxrs/UnixConnectionSocketFactory.java | 5 +- .../jaxrs/async/AbstractCallbackNotifier.java | 96 ----- .../jaxrs/async/GETCallbackNotifier.java | 29 -- .../jaxrs/async/POSTCallbackNotifier.java | 33 -- 9 files changed, 431 insertions(+), 687 deletions(-) create mode 100644 docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyDockerHttpClient.java delete mode 100644 docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyInvocationBuilder.java delete mode 100644 docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyWebTarget.java delete mode 100644 docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/async/AbstractCallbackNotifier.java delete mode 100644 docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/async/GETCallbackNotifier.java delete mode 100644 docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/async/POSTCallbackNotifier.java diff --git a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/ApacheUnixSocket.java b/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/ApacheUnixSocket.java index e64e252f8..72a851560 100644 --- a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/ApacheUnixSocket.java +++ b/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/ApacheUnixSocket.java @@ -41,14 +41,13 @@ * * This class also noop's any calls to setReuseAddress, which is called by the Apache client but isn't supported by AFUnixSocket. */ -@Deprecated -public class ApacheUnixSocket extends Socket { +class ApacheUnixSocket extends Socket { private final AFUNIXSocket inner; private final Queue optionsToSet = new ArrayDeque<>(); - public ApacheUnixSocket() throws IOException { + ApacheUnixSocket() throws IOException { this.inner = AFUNIXSocket.newInstance(); } diff --git a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyDockerCmdExecFactory.java b/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyDockerCmdExecFactory.java index 97ec72384..82d7b8324 100644 --- a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyDockerCmdExecFactory.java +++ b/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyDockerCmdExecFactory.java @@ -1,300 +1,95 @@ package com.github.dockerjava.jaxrs; -import com.fasterxml.jackson.jaxrs.json.JacksonJsonProvider; -import com.github.dockerjava.api.exception.DockerClientException; -import com.github.dockerjava.core.AbstractDockerCmdExecFactory; +import com.github.dockerjava.api.command.DelegatingDockerCmdExecFactory; +import com.github.dockerjava.api.command.DockerCmdExecFactory; +import com.github.dockerjava.core.DefaultDockerCmdExecFactory; import com.github.dockerjava.core.DockerClientConfig; -import com.github.dockerjava.core.SSLConfig; -import com.github.dockerjava.jaxrs.filter.JsonClientFilter; -import com.github.dockerjava.jaxrs.filter.ResponseStatusExceptionFilter; -import com.github.dockerjava.jaxrs.filter.SelectiveLoggingFilter; -import org.apache.http.client.config.RequestConfig; -import org.apache.http.config.RegistryBuilder; -import org.apache.http.conn.socket.ConnectionSocketFactory; -import org.apache.http.conn.socket.PlainConnectionSocketFactory; -import org.apache.http.conn.ssl.SSLConnectionSocketFactory; -import org.apache.http.impl.conn.PoolingHttpClientConnectionManager; -import org.glassfish.jersey.CommonProperties; -import org.glassfish.jersey.apache.connector.ApacheClientProperties; -import org.glassfish.jersey.apache.connector.ApacheConnectorProvider; -import org.glassfish.jersey.client.ClientConfig; -import org.glassfish.jersey.client.ClientProperties; +import com.github.dockerjava.core.DockerClientConfigAware; +import com.github.dockerjava.core.DockerClientImpl; +import com.github.dockerjava.core.DockerHttpClient; import org.glassfish.jersey.client.RequestEntityProcessing; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; -import javax.net.ssl.SSLContext; -import javax.ws.rs.client.Client; -import javax.ws.rs.client.ClientBuilder; import javax.ws.rs.client.ClientRequestFilter; import javax.ws.rs.client.ClientResponseFilter; -import java.io.IOException; -import java.net.InetSocketAddress; -import java.net.Proxy; -import java.net.ProxySelector; -import java.net.URI; -import java.net.URISyntaxException; -import java.util.List; -import java.util.concurrent.TimeUnit; - -import static com.google.common.base.Preconditions.checkNotNull; //import org.glassfish.jersey.apache.connector.ApacheConnectorProvider; // see https://github.com/docker-java/docker-java/issues/196 +/** + * @deprecated use {@link JerseyDockerHttpClient} with {@link DockerClientImpl#withHttpClient(DockerHttpClient)} + */ +@Deprecated +public class JerseyDockerCmdExecFactory extends DelegatingDockerCmdExecFactory implements DockerClientConfigAware { -public class JerseyDockerCmdExecFactory extends AbstractDockerCmdExecFactory { - - private static final Logger LOGGER = LoggerFactory.getLogger(JerseyDockerCmdExecFactory.class.getName()); - - private Client client; - - private JerseyWebTarget baseResource; - - private Integer maxTotalConnections = null; - - private Integer maxPerRouteConnections = null; - - private Integer connectionRequestTimeout = null; - - private ClientRequestFilter[] clientRequestFilters = null; + private JerseyDockerHttpClient.Factory clientFactory = new JerseyDockerHttpClient.Factory(); - private ClientResponseFilter[] clientResponseFilters = null; + @Deprecated + protected Integer connectTimeout; - private DockerClientConfig dockerClientConfig; + @Deprecated + protected Integer readTimeout; - private PoolingHttpClientConnectionManager connManager = null; + private DefaultDockerCmdExecFactory dockerCmdExecFactory; - private RequestEntityProcessing requestEntityProcessing; + @Override + public final DockerCmdExecFactory getDockerCmdExecFactory() { + return dockerCmdExecFactory; + } @Override public void init(DockerClientConfig dockerClientConfig) { - checkNotNull(dockerClientConfig, "config was not specified"); - this.dockerClientConfig = dockerClientConfig; - - ClientConfig clientConfig = new ClientConfig(); - clientConfig.connectorProvider(new ApacheConnectorProvider()); - clientConfig.property(CommonProperties.FEATURE_AUTO_DISCOVERY_DISABLE, true); - - if (requestEntityProcessing != null) { - clientConfig.property(ClientProperties.REQUEST_ENTITY_PROCESSING, requestEntityProcessing); - } - - clientConfig.register(new ResponseStatusExceptionFilter(dockerClientConfig.getObjectMapper())); - clientConfig.register(JsonClientFilter.class); - RequestConfig.Builder requestConfigBuilder = RequestConfig.custom(); - - clientConfig.register(new JacksonJsonProvider(dockerClientConfig.getObjectMapper())); - - // logging may disabled via log level - clientConfig.register(new SelectiveLoggingFilter(LOGGER, true)); - - if (readTimeout != null) { - requestConfigBuilder.setSocketTimeout(readTimeout); - clientConfig.property(ClientProperties.READ_TIMEOUT, readTimeout); - } - - if (connectTimeout != null) { - requestConfigBuilder.setConnectTimeout(connectTimeout); - clientConfig.property(ClientProperties.CONNECT_TIMEOUT, connectTimeout); - } - - if (clientResponseFilters != null) { - for (ClientResponseFilter clientResponseFilter : clientResponseFilters) { - if (clientResponseFilter != null) { - clientConfig.register(clientResponseFilter); - } - } - } - - if (clientRequestFilters != null) { - for (ClientRequestFilter clientRequestFilter : clientRequestFilters) { - if (clientRequestFilter != null) { - clientConfig.register(clientRequestFilter); - } - } - } - - URI originalUri = dockerClientConfig.getDockerHost(); - - String protocol = null; - - SSLContext sslContext = null; - - try { - final SSLConfig sslConfig = dockerClientConfig.getSSLConfig(); - if (sslConfig != null) { - sslContext = sslConfig.getSSLContext(); - } - } catch (Exception ex) { - throw new DockerClientException("Error in SSL Configuration", ex); - } - - if (sslContext != null) { - protocol = "https"; - } else { - protocol = "http"; - } - - switch (originalUri.getScheme()) { - case "unix": - break; - case "tcp": - try { - originalUri = new URI(originalUri.toString().replaceFirst("tcp", protocol)); - } catch (URISyntaxException e) { - throw new RuntimeException(e); - } - - configureProxy(clientConfig, originalUri, protocol); - break; - default: - throw new IllegalArgumentException("Unsupported protocol scheme: " + originalUri); - } - - connManager = new PoolingHttpClientConnectionManager(getSchemeRegistry( - originalUri, sslContext)) { - - @Override - public void close() { - super.shutdown(); - } - - @Override - public void shutdown() { - // Disable shutdown of the pool. This will be done later, when this factory is closed - // This is a workaround for finalize method on jerseys ClientRuntime which - // closes the client and shuts down the connection pool when it is garbage collected - } - }; - - if (maxTotalConnections != null) { - connManager.setMaxTotal(maxTotalConnections); - } - if (maxPerRouteConnections != null) { - connManager.setDefaultMaxPerRoute(maxPerRouteConnections); - } - - clientConfig.property(ApacheClientProperties.CONNECTION_MANAGER, connManager); - - // Configure connection pool timeout - if (connectionRequestTimeout != null) { - requestConfigBuilder.setConnectionRequestTimeout(connectionRequestTimeout); - } - clientConfig.property(ApacheClientProperties.REQUEST_CONFIG, requestConfigBuilder.build()); - ClientBuilder clientBuilder = ClientBuilder.newBuilder().withConfig(clientConfig); - - if (sslContext != null) { - clientBuilder.sslContext(sslContext); - } - - client = clientBuilder.build(); - - baseResource = new JerseyWebTarget( - dockerClientConfig.getObjectMapper(), - client.target(sanitizeUrl(originalUri).toString()) - .path(dockerClientConfig.getApiVersion().asWebPathPart()) + clientFactory = clientFactory.dockerClientConfig(dockerClientConfig); + dockerCmdExecFactory = new DefaultDockerCmdExecFactory( + clientFactory.build(), + dockerClientConfig.getObjectMapper() ); - - super.init(dockerClientConfig); - } - - private URI sanitizeUrl(URI originalUri) { - if (originalUri.getScheme().equals("unix")) { - return UnixConnectionSocketFactory.sanitizeUri(originalUri); - } - return originalUri; + dockerCmdExecFactory.init(dockerClientConfig); } - private void configureProxy(ClientConfig clientConfig, URI originalUri, String protocol) { - - List proxies = ProxySelector.getDefault().select(originalUri); - - for (Proxy proxy : proxies) { - InetSocketAddress address = (InetSocketAddress) proxy.address(); - if (address != null) { - String hostname = address.getHostName(); - int port = address.getPort(); - - clientConfig.property(ClientProperties.PROXY_URI, "http://" + hostname + ":" + port); - - String httpProxyUser = System.getProperty(protocol + ".proxyUser"); - if (httpProxyUser != null) { - clientConfig.property(ClientProperties.PROXY_USERNAME, httpProxyUser); - String httpProxyPassword = System.getProperty(protocol + ".proxyPassword"); - if (httpProxyPassword != null) { - clientConfig.property(ClientProperties.PROXY_PASSWORD, httpProxyPassword); - } - } - } - } - } - - private org.apache.http.config.Registry getSchemeRegistry(final URI originalUri, - SSLContext sslContext) { - RegistryBuilder registryBuilder = RegistryBuilder.create(); - registryBuilder.register("http", PlainConnectionSocketFactory.getSocketFactory()); - if (sslContext != null) { - registryBuilder.register("https", new SSLConnectionSocketFactory(sslContext)); - } - registryBuilder.register("unix", new UnixConnectionSocketFactory(originalUri)); - return registryBuilder.build(); - } - - protected JerseyWebTarget getBaseResource() { - checkNotNull(baseResource, "Factory not initialized, baseResource not set. You probably forgot to call init()!"); - return baseResource; - } - - protected DockerClientConfig getDockerClientConfig() { - checkNotNull(dockerClientConfig, - "Factor not initialized, dockerClientConfig not set. You probably forgot to call init()!"); - return dockerClientConfig; + /** + * Configure connection timeout in milliseconds + */ + public JerseyDockerCmdExecFactory withConnectTimeout(Integer connectTimeout) { + clientFactory = clientFactory.connectTimeout(connectTimeout); + this.connectTimeout = connectTimeout; + return this; } - @Override - public void close() throws IOException { - checkNotNull(client, "Factory not initialized. You probably forgot to call init()!"); - client.close(); - connManager.close(); + /** + * Configure read timeout in milliseconds + */ + public JerseyDockerCmdExecFactory withReadTimeout(Integer readTimeout) { + clientFactory = clientFactory.readTimeout(readTimeout); + this.readTimeout = readTimeout; + return this; } public JerseyDockerCmdExecFactory withMaxTotalConnections(Integer maxTotalConnections) { - this.maxTotalConnections = maxTotalConnections; + clientFactory = clientFactory.maxTotalConnections(maxTotalConnections); return this; } public JerseyDockerCmdExecFactory withMaxPerRouteConnections(Integer maxPerRouteConnections) { - this.maxPerRouteConnections = maxPerRouteConnections; + clientFactory = clientFactory.maxPerRouteConnections(maxPerRouteConnections); return this; } public JerseyDockerCmdExecFactory withConnectionRequestTimeout(Integer connectionRequestTimeout) { - this.connectionRequestTimeout = connectionRequestTimeout; + clientFactory = clientFactory.connectionRequestTimeout(connectionRequestTimeout); return this; } public JerseyDockerCmdExecFactory withClientResponseFilters(ClientResponseFilter... clientResponseFilter) { - this.clientResponseFilters = clientResponseFilter; + clientFactory = clientFactory.clientResponseFilters(clientResponseFilter); return this; } public JerseyDockerCmdExecFactory withClientRequestFilters(ClientRequestFilter... clientRequestFilters) { - this.clientRequestFilters = clientRequestFilters; + clientFactory = clientFactory.clientRequestFilters(clientRequestFilters); return this; } public JerseyDockerCmdExecFactory withRequestEntityProcessing(RequestEntityProcessing requestEntityProcessing) { - this.requestEntityProcessing = requestEntityProcessing; + clientFactory = clientFactory.requestEntityProcessing(requestEntityProcessing); return this; } - - /** - * release connections from the pool - * - * @param idleSeconds idle seconds, longer than the configured value will be evicted - */ - public void releaseConnection(long idleSeconds) { - this.connManager.closeExpiredConnections(); - this.connManager.closeIdleConnections(idleSeconds, TimeUnit.SECONDS); - } } 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 new file mode 100644 index 000000000..241dfc9ff --- /dev/null +++ b/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyDockerHttpClient.java @@ -0,0 +1,381 @@ +package com.github.dockerjava.jaxrs; + +import com.fasterxml.jackson.jaxrs.json.JacksonJsonProvider; +import com.github.dockerjava.api.exception.DockerClientException; +import com.github.dockerjava.api.exception.DockerException; +import com.github.dockerjava.core.DockerClientConfig; +import com.github.dockerjava.core.DockerHttpClient; +import com.github.dockerjava.core.SSLConfig; +import com.github.dockerjava.jaxrs.filter.ResponseStatusExceptionFilter; +import com.github.dockerjava.jaxrs.filter.SelectiveLoggingFilter; +import org.apache.http.client.config.RequestConfig; +import org.apache.http.config.Registry; +import org.apache.http.config.RegistryBuilder; +import org.apache.http.conn.socket.ConnectionSocketFactory; +import org.apache.http.conn.socket.PlainConnectionSocketFactory; +import org.apache.http.conn.ssl.SSLConnectionSocketFactory; +import org.apache.http.impl.conn.PoolingHttpClientConnectionManager; +import org.apache.http.impl.io.EmptyInputStream; +import org.glassfish.jersey.CommonProperties; +import org.glassfish.jersey.apache.connector.ApacheClientProperties; +import org.glassfish.jersey.apache.connector.ApacheConnectorProvider; +import org.glassfish.jersey.client.ClientConfig; +import org.glassfish.jersey.client.ClientProperties; +import org.glassfish.jersey.client.RequestEntityProcessing; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import javax.net.ssl.SSLContext; +import javax.ws.rs.ProcessingException; +import javax.ws.rs.client.Client; +import javax.ws.rs.client.ClientBuilder; +import javax.ws.rs.client.ClientRequestFilter; +import javax.ws.rs.client.ClientResponseFilter; +import javax.ws.rs.client.Entity; +import javax.ws.rs.client.Invocation; +import javax.ws.rs.core.MediaType; +import java.io.InputStream; +import java.net.InetSocketAddress; +import java.net.Proxy; +import java.net.ProxySelector; +import java.net.URI; +import java.net.URISyntaxException; +import java.util.List; +import java.util.Map; + +public final class JerseyDockerHttpClient implements DockerHttpClient { + + public static final class Factory { + + private DockerClientConfig dockerClientConfig = null; + + private Integer readTimeout = null; + + private Integer connectTimeout = null; + + private Integer maxTotalConnections = null; + + private Integer maxPerRouteConnections = null; + + private Integer connectionRequestTimeout = null; + + private ClientRequestFilter[] clientRequestFilters = null; + + private ClientResponseFilter[] clientResponseFilters = null; + + private RequestEntityProcessing requestEntityProcessing; + + public Factory dockerClientConfig(DockerClientConfig value) { + this.dockerClientConfig = value; + return this; + } + + public Factory readTimeout(Integer value) { + this.readTimeout = value; + return this; + } + + public Factory connectTimeout(Integer value) { + this.connectTimeout = value; + return this; + } + + public Factory maxTotalConnections(Integer value) { + this.maxTotalConnections = value; + return this; + } + + public Factory maxPerRouteConnections(Integer value) { + this.maxPerRouteConnections = value; + return this; + } + + public Factory connectionRequestTimeout(Integer value) { + this.connectionRequestTimeout = value; + return this; + } + + public Factory clientResponseFilters(ClientResponseFilter[] value) { + this.clientResponseFilters = value; + return this; + } + + public Factory clientRequestFilters(ClientRequestFilter[] value) { + this.clientRequestFilters = value; + return this; + } + + public Factory requestEntityProcessing(RequestEntityProcessing value) { + this.requestEntityProcessing = value; + return this; + } + + public JerseyDockerHttpClient build() { + return new JerseyDockerHttpClient( + dockerClientConfig, + maxTotalConnections, + maxPerRouteConnections, + connectionRequestTimeout, + readTimeout, + connectTimeout, + clientRequestFilters, + clientResponseFilters, + requestEntityProcessing + ); + } + } + + private static final Logger LOGGER = LoggerFactory.getLogger(JerseyDockerHttpClient.class.getName()); + + private final Client client; + + private final PoolingHttpClientConnectionManager connManager; + + private final URI originalUri; + + private JerseyDockerHttpClient( + DockerClientConfig dockerClientConfig, + Integer maxTotalConnections, + Integer maxPerRouteConnections, + Integer connectionRequestTimeout, + Integer readTimeout, + Integer connectTimeout, + ClientRequestFilter[] clientRequestFilters, + ClientResponseFilter[] clientResponseFilters, + RequestEntityProcessing requestEntityProcessing + ) { + ClientConfig clientConfig = new ClientConfig(); + clientConfig.connectorProvider(new ApacheConnectorProvider()); + clientConfig.property(CommonProperties.FEATURE_AUTO_DISCOVERY_DISABLE, true); + + if (requestEntityProcessing != null) { + clientConfig.property(ClientProperties.REQUEST_ENTITY_PROCESSING, requestEntityProcessing); + } + + clientConfig.register(new ResponseStatusExceptionFilter(dockerClientConfig.getObjectMapper())); + // clientConfig.register(JsonClientFilter.class); + RequestConfig.Builder requestConfigBuilder = RequestConfig.custom(); + + clientConfig.register(new JacksonJsonProvider(dockerClientConfig.getObjectMapper())); + + // logging may disabled via log level + clientConfig.register(new SelectiveLoggingFilter(LOGGER, true)); + + if (readTimeout != null) { + requestConfigBuilder.setSocketTimeout(readTimeout); + clientConfig.property(ClientProperties.READ_TIMEOUT, readTimeout); + } + + if (connectTimeout != null) { + requestConfigBuilder.setConnectTimeout(connectTimeout); + clientConfig.property(ClientProperties.CONNECT_TIMEOUT, connectTimeout); + } + + if (clientResponseFilters != null) { + for (ClientResponseFilter clientResponseFilter : clientResponseFilters) { + if (clientResponseFilter != null) { + clientConfig.register(clientResponseFilter); + } + } + } + + if (clientRequestFilters != null) { + for (ClientRequestFilter clientRequestFilter : clientRequestFilters) { + if (clientRequestFilter != null) { + clientConfig.register(clientRequestFilter); + } + } + } + + URI originalUri = dockerClientConfig.getDockerHost(); + + SSLContext sslContext = null; + + try { + final SSLConfig sslConfig = dockerClientConfig.getSSLConfig(); + if (sslConfig != null) { + sslContext = sslConfig.getSSLContext(); + } + } catch (Exception ex) { + throw new DockerClientException("Error in SSL Configuration", ex); + } + + final String protocol = sslContext != null ? "https" : "http"; + + switch (originalUri.getScheme()) { + case "unix": + break; + case "tcp": + try { + originalUri = new URI(originalUri.toString().replaceFirst("tcp", protocol)); + } catch (URISyntaxException e) { + throw new RuntimeException(e); + } + + configureProxy(clientConfig, originalUri, protocol); + break; + default: + throw new IllegalArgumentException("Unsupported protocol scheme: " + originalUri); + } + + connManager = new PoolingHttpClientConnectionManager(getSchemeRegistry(originalUri, sslContext)) { + + @Override + public void close() { + super.shutdown(); + } + + @Override + public void shutdown() { + // Disable shutdown of the pool. This will be done later, when this factory is closed + // This is a workaround for finalize method on jerseys ClientRuntime which + // closes the client and shuts down the connection pool when it is garbage collected + } + }; + + if (maxTotalConnections != null) { + connManager.setMaxTotal(maxTotalConnections); + } + if (maxPerRouteConnections != null) { + connManager.setDefaultMaxPerRoute(maxPerRouteConnections); + } + + clientConfig.property(ApacheClientProperties.CONNECTION_MANAGER, connManager); + + // Configure connection pool timeout + if (connectionRequestTimeout != null) { + requestConfigBuilder.setConnectionRequestTimeout(connectionRequestTimeout); + } + clientConfig.property(ApacheClientProperties.REQUEST_CONFIG, requestConfigBuilder.build()); + ClientBuilder clientBuilder = ClientBuilder.newBuilder().withConfig(clientConfig); + + if (sslContext != null) { + clientBuilder.sslContext(sslContext); + } + + client = clientBuilder.build(); + + this.originalUri = originalUri; + } + + private URI sanitizeUrl(URI originalUri) { + if (originalUri.getScheme().equals("unix")) { + return UnixConnectionSocketFactory.sanitizeUri(originalUri); + } + return originalUri; + } + + private Registry getSchemeRegistry(URI originalUri, SSLContext sslContext) { + RegistryBuilder registryBuilder = RegistryBuilder.create(); + registryBuilder.register("http", PlainConnectionSocketFactory.getSocketFactory()); + if (sslContext != null) { + registryBuilder.register("https", new SSLConnectionSocketFactory(sslContext)); + } + registryBuilder.register("unix", new UnixConnectionSocketFactory(originalUri)); + return registryBuilder.build(); + } + + @Override + public Response execute(Request request) { + if (request.hijackedInput() != null) { + throw new UnsupportedOperationException("Does not support hijacking"); + } + String url = sanitizeUrl(originalUri).toString(); + if (url.endsWith("/") && request.path().startsWith("/")) { + url = url.substring(0, url.length() - 1); + } + + Invocation.Builder builder = client.target(url + request.path()).request(); + + request.headers().forEach(builder::header); + + try { + return new JerseyResponse( + builder.build(request.method(), toEntity(request)).invoke() + ); + } catch (ProcessingException e) { + if (e.getCause() instanceof DockerException) { + throw (DockerException) e.getCause(); + } + throw e; + } + } + + private Entity toEntity(Request request) { + InputStream body = request.body(); + if (body != null) { + return Entity.entity(body, MediaType.APPLICATION_JSON_TYPE); + } + switch (request.method()) { + case "POST": + return Entity.json(null); + default: + return null; + } + } + + private void configureProxy(ClientConfig clientConfig, URI originalUri, String protocol) { + List proxies = ProxySelector.getDefault().select(originalUri); + + for (Proxy proxy : proxies) { + InetSocketAddress address = (InetSocketAddress) proxy.address(); + if (address != null) { + String hostname = address.getHostName(); + int port = address.getPort(); + + clientConfig.property(ClientProperties.PROXY_URI, "http://" + hostname + ":" + port); + + String httpProxyUser = System.getProperty(protocol + ".proxyUser"); + if (httpProxyUser != null) { + clientConfig.property(ClientProperties.PROXY_USERNAME, httpProxyUser); + String httpProxyPassword = System.getProperty(protocol + ".proxyPassword"); + if (httpProxyPassword != null) { + clientConfig.property(ClientProperties.PROXY_PASSWORD, httpProxyPassword); + } + } + } + } + } + + @Override + public void close() { + if (client != null) { + client.close(); + } + + if (connManager != null) { + connManager.close(); + } + } + + private static class JerseyResponse implements Response { + + private final javax.ws.rs.core.Response response; + + public JerseyResponse(javax.ws.rs.core.Response response) { + this.response = response; + } + + @Override + public int getStatusCode() { + return response.getStatus(); + } + + @Override + public Map> getHeaders() { + return response.getStringHeaders(); + } + + @Override + public InputStream getBody() { + return response.hasEntity() + ? response.readEntity(InputStream.class) + : EmptyInputStream.INSTANCE; + } + + @Override + public void close() { + response.close(); + } + } +} diff --git a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyInvocationBuilder.java b/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyInvocationBuilder.java deleted file mode 100644 index 761129867..000000000 --- a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyInvocationBuilder.java +++ /dev/null @@ -1,193 +0,0 @@ -package com.github.dockerjava.jaxrs; - -import com.fasterxml.jackson.core.type.TypeReference; -import com.fasterxml.jackson.databind.ObjectMapper; -import com.github.dockerjava.api.async.ResultCallback; -import com.github.dockerjava.api.exception.UnauthorizedException; -import com.github.dockerjava.api.model.Frame; -import com.github.dockerjava.core.InvocationBuilder; -import com.github.dockerjava.core.MediaType; -import com.github.dockerjava.core.async.FrameStreamProcessor; -import com.github.dockerjava.core.async.JsonStreamProcessor; -import com.github.dockerjava.jaxrs.async.GETCallbackNotifier; -import com.github.dockerjava.jaxrs.async.POSTCallbackNotifier; -import com.github.dockerjava.jaxrs.util.WrappedResponseInputStream; - -import javax.ws.rs.client.Entity; -import javax.ws.rs.client.Invocation; -import javax.ws.rs.core.Response; -import java.io.IOException; -import java.io.InputStream; - -class JerseyInvocationBuilder implements InvocationBuilder { - - private final ObjectMapper objectMapper; - - private final Invocation.Builder resource; - - JerseyInvocationBuilder(ObjectMapper objectMapper, Invocation.Builder resource) { - this.objectMapper = objectMapper; - this.resource = resource; - } - - @Override - public InvocationBuilder accept(MediaType mediaType) { - resource.accept(mediaType.getMediaType()); - return this; - } - - @Override - public InvocationBuilder header(String name, String value) { - resource.header(name, value); - return this; - } - - @Override - public void delete() { - resource.delete().close(); - } - - @Override - public void get(ResultCallback resultCallback) { - try { - GETCallbackNotifier getCallbackNotifier = new GETCallbackNotifier<>( - new FrameStreamProcessor(), - resultCallback, - resource - ); - getCallbackNotifier.start(); - } catch (Exception e) { - throw new RuntimeException(e); - } - } - - @Override - public T get(TypeReference typeReference) { - try (Response response = resource.get()) { - return objectMapper.readValue(response.readEntity(InputStream.class), typeReference); - } catch (IOException e) { - throw new RuntimeException(e); - } - } - - @Override - public void get(TypeReference typeReference, ResultCallback resultCallback) { - try { - GETCallbackNotifier getCallbackNotifier = new GETCallbackNotifier( - new JsonStreamProcessor<>(objectMapper, typeReference), - resultCallback, - resource - ); - getCallbackNotifier.start(); - } catch (Exception e) { - throw new RuntimeException(e); - } - } - - @Override - public InputStream post(Object entity) { - return new WrappedResponseInputStream(resource.post( - toEntity(entity, javax.ws.rs.core.MediaType.APPLICATION_JSON) - )); - } - - @Override - public void post(Object entity, InputStream stdin, ResultCallback resultCallback) { - if (stdin != null) { - throw new UnsupportedOperationException("Passing stdin to the container is currently not supported."); - } - - POSTCallbackNotifier postCallbackNotifier = new POSTCallbackNotifier<>( - new FrameStreamProcessor(), - resultCallback, - resource, - toEntity(entity, javax.ws.rs.core.MediaType.APPLICATION_JSON) - ); - - postCallbackNotifier.start(); - } - - @Override - public T post(Object entity, TypeReference typeReference) { - Response response = resource.post( - toEntity(entity, javax.ws.rs.core.MediaType.APPLICATION_JSON) - ); - - if (response.getStatus() == 401) { - throw new UnauthorizedException("Unauthorized"); - } - - try (InputStream inputStream = response.readEntity(InputStream.class)) { - return objectMapper.readValue(inputStream, typeReference); - } catch (IOException e) { - throw new RuntimeException(e); - } - } - - @Override - public void post(Object entity, TypeReference typeReference, ResultCallback resultCallback) { - try { - POSTCallbackNotifier postCallbackNotifier = new POSTCallbackNotifier<>( - new JsonStreamProcessor<>(objectMapper, typeReference), - resultCallback, - resource, - toEntity(entity, javax.ws.rs.core.MediaType.APPLICATION_JSON) - ); - postCallbackNotifier.start(); - } catch (Exception e) { - throw new RuntimeException(e); - } - } - - @Override - public T post(TypeReference typeReference, InputStream body) { - try ( - Response response = resource.post( - toEntity(body, javax.ws.rs.core.MediaType.APPLICATION_OCTET_STREAM) - ) - ) { - InputStream inputStream = response.readEntity(InputStream.class); - return objectMapper.readValue(inputStream, typeReference); - } catch (IOException e) { - throw new RuntimeException(e); - } - } - - @Override - public void post(TypeReference typeReference, ResultCallback resultCallback, InputStream body) { - try { - POSTCallbackNotifier postCallbackNotifier = new POSTCallbackNotifier( - new JsonStreamProcessor<>(objectMapper, typeReference), - resultCallback, - resource, - toEntity(body, "application/tar") - ); - postCallbackNotifier.start(); - } catch (Exception e) { - throw new RuntimeException(e); - } - } - - @Override - public void postStream(InputStream body) { - resource.post(toEntity(body, javax.ws.rs.core.MediaType.APPLICATION_OCTET_STREAM)).close(); - } - - @Override - public InputStream get() { - return new WrappedResponseInputStream(resource.get()); - } - - @Override - public void put(InputStream body, MediaType mediaType) { - resource.put(toEntity(body, mediaType.getMediaType())).close(); - } - - private static Entity toEntity(T entity, String mediaType) { - if (entity == null) { - return null; - } - - return Entity.entity(entity, mediaType); - } -} diff --git a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyWebTarget.java b/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyWebTarget.java deleted file mode 100644 index 123d2f7ea..000000000 --- a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/JerseyWebTarget.java +++ /dev/null @@ -1,79 +0,0 @@ -package com.github.dockerjava.jaxrs; - -import com.fasterxml.jackson.databind.ObjectMapper; -import com.github.dockerjava.core.InvocationBuilder; -import com.github.dockerjava.core.WebTarget; - -import java.io.IOException; -import java.util.Map; -import java.util.Set; - -import static com.google.common.net.UrlEscapers.urlPathSegmentEscaper; - -class JerseyWebTarget implements WebTarget { - - private static final String PATH_SEPARATOR = "/"; - - private final ObjectMapper objectMapper; - - private final javax.ws.rs.client.WebTarget webTarget; - - JerseyWebTarget(ObjectMapper objectMapper, javax.ws.rs.client.WebTarget webTarget) { - this.objectMapper = objectMapper; - this.webTarget = webTarget; - } - - @Override - public InvocationBuilder request() { - return new JerseyInvocationBuilder(objectMapper, webTarget.request()); - } - - @Override - public JerseyWebTarget path(String... components) { - return new JerseyWebTarget( - objectMapper, - webTarget.path(String.join(PATH_SEPARATOR, components)) - ); - } - - @Override - public JerseyWebTarget resolveTemplate(String name, Object value) { - return new JerseyWebTarget( - objectMapper, - webTarget.resolveTemplate(name, value) - ); - } - - @Override - public JerseyWebTarget queryParam(String name, Object value) { - if (value instanceof String) { - value = urlPathSegmentEscaper().escape((String) value); - } - return new JerseyWebTarget( - objectMapper, - webTarget.queryParam(name, value) - ); - } - - @Override - public JerseyWebTarget queryParamsSet(String name, Set values) { - return new JerseyWebTarget( - objectMapper, - webTarget.queryParam(name, values.toArray()) - ); - } - - @Override - public JerseyWebTarget queryParamsJsonMap(String name, Map values) { - if (values != null && !values.isEmpty()) { - try { - // when param value is JSON string - return queryParam(name, objectMapper.writeValueAsString(values)); - } catch (IOException e) { - throw new RuntimeException(e); - } - } else { - return this; - } - } -} diff --git a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/UnixConnectionSocketFactory.java b/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/UnixConnectionSocketFactory.java index a45561de9..84a72f077 100644 --- a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/UnixConnectionSocketFactory.java +++ b/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/UnixConnectionSocketFactory.java @@ -40,12 +40,11 @@ * Provides a ConnectionSocketFactory for connecting Apache HTTP clients to Unix sockets. */ @Contract(threading = ThreadingBehavior.IMMUTABLE_CONDITIONAL) -@Deprecated -public class UnixConnectionSocketFactory implements ConnectionSocketFactory { +class UnixConnectionSocketFactory implements ConnectionSocketFactory { private File socketFile; - public UnixConnectionSocketFactory(final URI socketUri) { + UnixConnectionSocketFactory(final URI socketUri) { super(); final String filename = socketUri.toString().replaceAll("^unix:///", "unix://localhost/") diff --git a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/async/AbstractCallbackNotifier.java b/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/async/AbstractCallbackNotifier.java deleted file mode 100644 index b8db73b04..000000000 --- a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/async/AbstractCallbackNotifier.java +++ /dev/null @@ -1,96 +0,0 @@ -/* - * Created on 17.06.2015 - */ -package com.github.dockerjava.jaxrs.async; - -import static com.google.common.base.Preconditions.checkNotNull; - -import java.io.InputStream; -import java.util.concurrent.Callable; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.Future; -import java.util.concurrent.ThreadFactory; - -import javax.ws.rs.ProcessingException; -import javax.ws.rs.client.Invocation.Builder; -import javax.ws.rs.core.Response; - -import com.github.dockerjava.api.async.ResultCallback; -import com.github.dockerjava.core.async.ResponseStreamProcessor; -import com.github.dockerjava.jaxrs.util.WrappedResponseInputStream; -import com.google.common.util.concurrent.ThreadFactoryBuilder; - -@Deprecated -public abstract class AbstractCallbackNotifier implements Callable { - - private final ResponseStreamProcessor responseStreamProcessor; - - private final ResultCallback resultCallback; - - private static final ThreadFactory FACTORY = - new ThreadFactoryBuilder().setDaemon(true).setNameFormat("dockerjava-jaxrs-async-%d").build(); - - protected final Builder requestBuilder; - - protected AbstractCallbackNotifier(ResponseStreamProcessor responseStreamProcessor, - ResultCallback resultCallback, Builder requestBuilder) { - checkNotNull(requestBuilder, "An WebTarget must be provided"); - checkNotNull(responseStreamProcessor, "A ResponseStreamProcessor must be provided"); - this.responseStreamProcessor = responseStreamProcessor; - this.resultCallback = resultCallback; - this.requestBuilder = requestBuilder; - } - - @Override - public Void call() { - - Response response; - - try { - response = response(); - } catch (ProcessingException e) { - if (resultCallback != null) { - resultCallback.onError(e.getCause()); - } - return null; - } catch (Exception e) { - if (resultCallback != null) { - resultCallback.onError(e); - } - return null; - } - if (resultCallback != null) { - resultCallback.onStart(response::close); - } - - try (InputStream inputStream = new WrappedResponseInputStream(response)) { - - if (resultCallback != null) { - responseStreamProcessor.processResponseStream(inputStream, resultCallback); - } - - return null; - } catch (Exception e) { - if (resultCallback != null) { - resultCallback.onError(e); - } - - return null; - } - } - - protected abstract Response response(); - - public static Future startAsyncProcessing(AbstractCallbackNotifier callbackNotifier) { - - ExecutorService executorService = Executors.newSingleThreadExecutor(FACTORY); - Future response = executorService.submit(callbackNotifier); - executorService.shutdown(); - return response; - } - - public void start() { - FACTORY.newThread(this::call).start(); - } -} diff --git a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/async/GETCallbackNotifier.java b/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/async/GETCallbackNotifier.java deleted file mode 100644 index 9297c2551..000000000 --- a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/async/GETCallbackNotifier.java +++ /dev/null @@ -1,29 +0,0 @@ -/* - * Created on 23.06.2015 - */ -package com.github.dockerjava.jaxrs.async; - -import javax.ws.rs.client.Invocation.Builder; -import javax.ws.rs.core.Response; - -import com.github.dockerjava.api.async.ResultCallback; -import com.github.dockerjava.core.async.ResponseStreamProcessor; - -/** - * - * @author Marcus Linke - * - */ -@Deprecated -public class GETCallbackNotifier extends AbstractCallbackNotifier { - - public GETCallbackNotifier(ResponseStreamProcessor responseStreamProcessor, ResultCallback resultCallback, - Builder requestBuilder) { - super(responseStreamProcessor, resultCallback, requestBuilder); - } - - protected Response response() { - return requestBuilder.get(); - } - -} diff --git a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/async/POSTCallbackNotifier.java b/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/async/POSTCallbackNotifier.java deleted file mode 100644 index 76fd540fe..000000000 --- a/docker-java-transport-jersey/src/main/java/com/github/dockerjava/jaxrs/async/POSTCallbackNotifier.java +++ /dev/null @@ -1,33 +0,0 @@ -/* - * Created on 23.06.2015 - */ -package com.github.dockerjava.jaxrs.async; - -import javax.ws.rs.client.Entity; -import javax.ws.rs.client.Invocation.Builder; -import javax.ws.rs.core.Response; - -import com.github.dockerjava.api.async.ResultCallback; -import com.github.dockerjava.core.async.ResponseStreamProcessor; - -/** - * - * @author Marcus Linke - * - */ -@Deprecated -public class POSTCallbackNotifier extends AbstractCallbackNotifier { - - Entity entity = null; - - public POSTCallbackNotifier(ResponseStreamProcessor responseStreamProcessor, ResultCallback resultCallback, - Builder requestBuilder, Entity entity) { - super(responseStreamProcessor, resultCallback, requestBuilder); - this.entity = entity; - } - - protected Response response() { - return requestBuilder.post(entity, Response.class); - } - -} From 30cfcac5aaf965fcbe114d89d3361bf397af3662 Mon Sep 17 00:00:00 2001 From: Sergei Egorov Date: Fri, 20 Mar 2020 22:24:41 +0100 Subject: [PATCH 12/13] checkstyle... --- .../dockerjava/jaxrs/JerseyDockerHttpClient.java | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) 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 241dfc9ff..c2f9f78da 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 @@ -187,7 +187,7 @@ private JerseyDockerHttpClient( } } - URI originalUri = dockerClientConfig.getDockerHost(); + URI dockerHost = dockerClientConfig.getDockerHost(); SSLContext sslContext = null; @@ -202,23 +202,23 @@ private JerseyDockerHttpClient( final String protocol = sslContext != null ? "https" : "http"; - switch (originalUri.getScheme()) { + switch (dockerHost.getScheme()) { case "unix": break; case "tcp": try { - originalUri = new URI(originalUri.toString().replaceFirst("tcp", protocol)); + dockerHost = new URI(dockerHost.toString().replaceFirst("tcp", protocol)); } catch (URISyntaxException e) { throw new RuntimeException(e); } - configureProxy(clientConfig, originalUri, protocol); + configureProxy(clientConfig, dockerHost, protocol); break; default: - throw new IllegalArgumentException("Unsupported protocol scheme: " + originalUri); + throw new IllegalArgumentException("Unsupported protocol scheme: " + dockerHost); } - connManager = new PoolingHttpClientConnectionManager(getSchemeRegistry(originalUri, sslContext)) { + connManager = new PoolingHttpClientConnectionManager(getSchemeRegistry(dockerHost, sslContext)) { @Override public void close() { @@ -255,7 +255,7 @@ public void shutdown() { client = clientBuilder.build(); - this.originalUri = originalUri; + this.originalUri = dockerHost; } private URI sanitizeUrl(URI originalUri) { @@ -352,7 +352,7 @@ private static class JerseyResponse implements Response { private final javax.ws.rs.core.Response response; - public JerseyResponse(javax.ws.rs.core.Response response) { + JerseyResponse(javax.ws.rs.core.Response response) { this.response = response; } From ae527ef2903b2f9e66c40e2e572e9a73938d73e4 Mon Sep 17 00:00:00 2001 From: Sergei Egorov Date: Fri, 3 Apr 2020 16:51:39 +0200 Subject: [PATCH 13/13] Add per-module japicmp --- docker-java-api/pom.xml | 58 ---------------------------- docker-java-transport-jersey/pom.xml | 27 +++++++++++++ docker-java-transport-okhttp/pom.xml | 16 ++++++++ docker-java/pom.xml | 26 ------------- pom.xml | 48 +++++++++++++++++++++++ 5 files changed, 91 insertions(+), 84 deletions(-) diff --git a/docker-java-api/pom.xml b/docker-java-api/pom.xml index f84228cac..f23df928c 100644 --- a/docker-java-api/pom.xml +++ b/docker-java-api/pom.xml @@ -55,64 +55,6 @@ - - - - com.github.siom79.japicmp - japicmp-maven-plugin - 0.14.3 - - - - com.github.docker-java - docker-java - 3.1.5 - jar - - - - - ${project.build.directory}/${project.artifactId}-${project.version}.jar - - - - true - public - true - - com.github.dockerjava.api - - - - com.github.dockerjava.api.model.*$Serializer - com.github.dockerjava.api.model.*$Deserializer - - com.github.dockerjava.api.command.DockerCmdExecFactory#init(com.github.dockerjava.core.DockerClientConfig) - com.github.dockerjava.api.model.Identifier#tag - - - - METHOD_NEW_DEFAULT - true - true - - - METHOD_ABSTRACT_NOW_DEFAULT - true - true - - - - - - - verify - - cmp - - - - diff --git a/docker-java-transport-jersey/pom.xml b/docker-java-transport-jersey/pom.xml index 7165cdf39..37de78f58 100644 --- a/docker-java-transport-jersey/pom.xml +++ b/docker-java-transport-jersey/pom.xml @@ -72,6 +72,33 @@ + + + com.github.siom79.japicmp + japicmp-maven-plugin + + + + + com.github.dockerjava.jaxrs.ApacheUnixSocket + com.github.dockerjava.jaxrs.async.AbstractCallbackNotifier + com.github.dockerjava.jaxrs.async.GETCallbackNotifier + com.github.dockerjava.jaxrs.async.GETCallbackNotifier + com.github.dockerjava.jaxrs.async.POSTCallbackNotifier + com.github.dockerjava.jaxrs.UnixConnectionSocketFactory + com.github.dockerjava.jaxrs.JerseyDockerCmdExecFactory#releaseConnection(long) + + + + SUPERCLASS_REMOVED + true + true + + + + + + org.apache.felix maven-bundle-plugin diff --git a/docker-java-transport-okhttp/pom.xml b/docker-java-transport-okhttp/pom.xml index f7303d894..3e6de70c2 100644 --- a/docker-java-transport-okhttp/pom.xml +++ b/docker-java-transport-okhttp/pom.xml @@ -37,6 +37,22 @@ + + com.github.siom79.japicmp + japicmp-maven-plugin + + + + + SUPERCLASS_REMOVED + true + true + + + + + + org.apache.felix maven-bundle-plugin diff --git a/docker-java/pom.xml b/docker-java/pom.xml index b52fcdcd6..beaf5d4c9 100644 --- a/docker-java/pom.xml +++ b/docker-java/pom.xml @@ -151,32 +151,6 @@ - - - - com.github.siom79.japicmp - japicmp-maven-plugin - 0.14.1 - - - - com.github.docker-java - docker-java - 3.1.0 - jar - - - - - ${project.build.directory}/${project.artifactId}-${project.version}.jar - - - - public - true - - - diff --git a/pom.xml b/pom.xml index 38880e712..2d3fb78b2 100644 --- a/pom.xml +++ b/pom.xml @@ -229,6 +229,54 @@ maven-bundle-plugin 4.2.1 + + + + + com.github.siom79.japicmp + japicmp-maven-plugin + 0.14.3 + + + + com.github.docker-java + ${project.artifactId} + 3.2.0 + jar + + + + + ${project.build.directory}/${project.artifactId}-${project.version}.jar + + + + true + public + true + + + METHOD_NEW_DEFAULT + true + true + + + METHOD_ABSTRACT_NOW_DEFAULT + true + true + + + + + + + verify + + cmp + + + +