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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions .changes/next-release/bugfix-AWSSDKforJavav2-c8ce5ba.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
{
"type": "bugfix",
"category": "AWS SDK for Java v2",
"contributor": "",
"description": "Added opt-in GZIPInputStream compatibility for blocking response streams to prevent concatenated (multi-member) gzip responses from being truncated when available() temporarily returns 0 at a member boundary. Enable it with ResponseTransformer.toInputStream(true) or AsyncResponseTransformer.toBlockingInputStream(true)."
}
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,6 @@ public ResponseInputStream(ResponseT resp, AbortableInputStream in, Duration tim
super(in);
this.response = Validate.paramNotNull(resp, "response");
this.abortable = Validate.paramNotNull(in, "abortableInputStream");

Duration resolvedTimeout = timeout != null ? timeout : DEFAULT_TIMEOUT;
scheduleTimeoutTask(resolvedTimeout);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -386,7 +386,27 @@ ResponsePublisher<ResponseT>> toPublisher(Duration timeout) {
*/
static <ResponseT extends SdkResponse>
AsyncResponseTransformer<ResponseT, ResponseInputStream<ResponseT>> toBlockingInputStream() {
return new InputStreamResponseTransformer<>();
return new InputStreamResponseTransformer<>(false);
}

/**
* Creates an {@link AsyncResponseTransformer} that allows reading the response body content as an {@link InputStream}.
* You are responsible for performing blocking reads from this input stream and closing the stream when you are finished.
*
* <p>When enabled, gzip response streams are adapted so that {@link InputStream#available()} does not temporarily return
* {@code 0} while the stream is still open. This works around {@link java.util.zip.GZIPInputStream} treating a temporary
* {@code 0} at a concatenated gzip member boundary as the end of the complete stream. Because this can cause a read after
* {@code available()} to block, it should only be enabled when the response will be read with {@code GZIPInputStream}.
*
* @param gzipInputStreamCompatibilityEnabled Whether to enable {@code GZIPInputStream} compatibility for concatenated gzip.
* @param <ResponseT> Type of unmarshalled response POJO.
* @return AsyncResponseTransformer instance.
* @see #toBlockingInputStream()
*/
static <ResponseT extends SdkResponse>
AsyncResponseTransformer<ResponseT, ResponseInputStream<ResponseT>> toBlockingInputStream(
boolean gzipInputStreamCompatibilityEnabled) {
return new InputStreamResponseTransformer<>(gzipInputStreamCompatibilityEnabled);
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@

package software.amazon.awssdk.core.internal.async;

import java.io.InputStream;
import java.nio.ByteBuffer;
import java.util.concurrent.CompletableFuture;
import org.reactivestreams.Subscriber;
Expand All @@ -24,6 +25,7 @@
import software.amazon.awssdk.core.SdkResponse;
import software.amazon.awssdk.core.async.AsyncResponseTransformer;
import software.amazon.awssdk.core.async.SdkPublisher;
import software.amazon.awssdk.core.internal.io.GzipAvailabilityInputStream;
import software.amazon.awssdk.http.async.AbortableInputStreamSubscriber;

/**
Expand All @@ -38,6 +40,15 @@ public class InputStreamResponseTransformer<ResponseT extends SdkResponse>
private volatile CompletableFuture<ResponseInputStream<ResponseT>> future;
private volatile ResponseT response;
private volatile WaitForSubscribeOnErrorWrapper subscriber;
private final boolean gzipInputStreamCompatibilityEnabled;

public InputStreamResponseTransformer() {
this(false);
}

public InputStreamResponseTransformer(boolean gzipInputStreamCompatibilityEnabled) {
this.gzipInputStreamCompatibilityEnabled = gzipInputStreamCompatibilityEnabled;
}

@Override
public CompletableFuture<ResponseInputStream<ResponseT>> prepare() {
Expand All @@ -59,7 +70,10 @@ public void onStream(SdkPublisher<ByteBuffer> publisher) {
this.subscriber = waitForSubscribeSubscriber;

publisher.subscribe(waitForSubscribeSubscriber);
future.complete(new ResponseInputStream<>(response, inputStreamSubscriber));
InputStream content = gzipInputStreamCompatibilityEnabled
? GzipAvailabilityInputStream.wrap(inputStreamSubscriber, inputStreamSubscriber)
: inputStreamSubscriber;
future.complete(new ResponseInputStream<>(response, content));
}

@Override
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,176 @@
/*
* Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
*
* Licensed under the Apache License, Version 2.0 (the "License").
* You may not use this file except in compliance with the License.
* A copy of the License is located at
*
* http://aws.amazon.com/apache2.0
*
* or in the "license" file accompanying this file. This file is distributed
* on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either
* express or implied. See the License for the specific language governing
* permissions and limitations under the License.
*/

package software.amazon.awssdk.core.internal.io;

import java.io.FilterInputStream;
import java.io.IOException;
import java.io.InputStream;
import software.amazon.awssdk.annotations.SdkInternalApi;
import software.amazon.awssdk.http.Abortable;
import software.amazon.awssdk.http.AbortableInputStream;
import software.amazon.awssdk.utils.IoUtils;

/**
* Wraps a response body so {@code available()} never returns {@code 0} for gzip content while the stream is open.
* {@link java.util.zip.GZIPInputStream} treats a transient {@code 0} from {@code available()} at a member boundary
* as end of stream and stops, truncating concatenated gzip. Gzip is detected passively from the leading bytes
* ({@code 1f 8b 08}); non-gzip streams keep honest {@code available()}.
*/
@SdkInternalApi
public final class GzipAvailabilityInputStream extends FilterInputStream implements Releasable {

private static final int GZIP_MAGIC_1 = 0x1f;
private static final int GZIP_MAGIC_2 = 0x8b;
private static final int GZIP_METHOD_DEFLATE = 0x08;
private static final int HEADER_LENGTH = 3;

private final byte[] header = new byte[HEADER_LENGTH];
private volatile int headerLen;
private volatile boolean classified;
private volatile boolean gzipDetected;
private volatile boolean eof;
private volatile boolean closed;

private int markHeaderLen;
private boolean markClassified;
private boolean markGzipDetected;
private boolean markEof;
private boolean marked;

public GzipAvailabilityInputStream(InputStream in) {
super(in);
}

/**
* Wraps a response body's content so {@code available()} is gzip-safe, while preserving {@code abort()} on the
* original stream. Applied by the SDK at the sync and async blocking-stream boundaries (where the caller's
* {@link java.util.zip.GZIPInputStream} reads), so that {@link software.amazon.awssdk.core.ResponseInputStream}
* itself stays content-type agnostic.
*/
public static AbortableInputStream wrap(InputStream content, Abortable abortable) {
return AbortableInputStream.create(new GzipAvailabilityInputStream(content), abortable);
}

@Override
public int read() throws IOException {
int b = in.read();
if (b == -1) {
eof = true;
} else {
if (eof) {
eof = false;
}
observe((byte) b);
}
return b;
}

@Override
public int read(byte[] b, int off, int len) throws IOException {
int n = in.read(b, off, len);
if (n == -1) {
eof = true;
} else if (n > 0) {
if (eof) {
eof = false;
}
observe(b, off, n);
}
return n;
}

@Override
public int available() throws IOException {
if (closed) {
return 0;
}
int available = in.available();
return available == 0 && gzipDetected && !eof ? 1 : available;
}

@Override
public long skip(long n) throws IOException {
long skipped = in.skip(n);
if (skipped > 0) {
classified = true;
}
return skipped;
}

@Override
public synchronized void mark(int readlimit) {
markHeaderLen = headerLen;
markClassified = classified;
markGzipDetected = gzipDetected;
markEof = eof;
marked = true;
in.mark(readlimit);
}

@Override
public synchronized void reset() throws IOException {
in.reset();
if (marked) {
headerLen = markHeaderLen;
classified = markClassified;
gzipDetected = markGzipDetected;
eof = markEof;
} else {
headerLen = 0;
classified = false;
gzipDetected = false;
eof = false;
}
}

@Override
public void close() throws IOException {
closed = true;
in.close();
}

@Override
public void release() {
IoUtils.closeQuietly(this, null);
if (in instanceof Releasable) {
((Releasable) in).release();
}
}

private void observe(byte[] b, int off, int len) {
if (classified) {
return;
}
for (int i = 0; i < len && !classified; i++) {
observe(b[off + i]);
}
}

private void observe(byte b) {
if (classified) {
return;
}
int len = headerLen;
header[len] = b;
headerLen = len + 1;
if (headerLen == HEADER_LENGTH) {
classified = true;
gzipDetected = (header[0] & 0xff) == GZIP_MAGIC_1
&& (header[1] & 0xff) == GZIP_MAGIC_2
&& (header[2] & 0xff) == GZIP_METHOD_DEFLATE;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
import software.amazon.awssdk.core.exception.SdkClientException;
import software.amazon.awssdk.core.exception.SdkException;
import software.amazon.awssdk.core.internal.http.InterruptMonitor;
import software.amazon.awssdk.core.internal.io.GzipAvailabilityInputStream;
import software.amazon.awssdk.core.retry.RetryPolicy;
import software.amazon.awssdk.http.AbortableInputStream;
import software.amazon.awssdk.utils.IoUtils;
Expand Down Expand Up @@ -261,17 +262,30 @@ public String name() {
* @see #toInputStream(Duration)
*/
static <ResponseT> ResponseTransformer<ResponseT, ResponseInputStream<ResponseT>> toInputStream() {
return unmanaged(new ResponseTransformer<ResponseT, ResponseInputStream<ResponseT>>() {
@Override
public ResponseInputStream<ResponseT> transform(ResponseT response, AbortableInputStream inputStream) {
return new ResponseInputStream<>(response, inputStream);
}
return toInputStream(null, false);
}

@Override
public String name() {
return TransformerType.STREAM.getName();
}
});
/**
* Creates a response transformer that returns an unmanaged input stream with the response content. This input stream must
* be explicitly closed to release the connection.
*
* <p>The stream has the default first-read timeout of 60 seconds. Use {@link #toInputStream(Duration, boolean)} to specify a
* custom timeout.
*
* <p>When enabled, gzip response streams are adapted so that {@link InputStream#available()} does not temporarily return
* {@code 0} while the stream is still open. This works around {@link java.util.zip.GZIPInputStream} treating a temporary
* {@code 0} at a concatenated gzip member boundary as the end of the complete stream. Because this can cause a read after
* {@code available()} to block, it should only be enabled when the response will be read with {@code GZIPInputStream}.
*
* @param gzipInputStreamCompatibilityEnabled Whether to enable {@code GZIPInputStream} compatibility for concatenated gzip.
* @param <ResponseT> Type of unmarshalled response POJO.
* @return ResponseTransformer instance.
* @see #toInputStream()
* @see #toInputStream(Duration, boolean)
*/
static <ResponseT> ResponseTransformer<ResponseT, ResponseInputStream<ResponseT>> toInputStream(
boolean gzipInputStreamCompatibilityEnabled) {
return toInputStream(null, gzipInputStreamCompatibilityEnabled);
}

/**
Expand All @@ -289,10 +303,33 @@ public String name() {
* @see #toInputStream()
*/
static <ResponseT> ResponseTransformer<ResponseT, ResponseInputStream<ResponseT>> toInputStream(Duration timeout) {
return toInputStream(timeout, false);
}

/**
* Creates a response transformer that returns an unmanaged input stream with the response content and a custom timeout.
* This input stream must be explicitly closed to release the connection.
*
* <p>When enabled, gzip response streams are adapted so that {@link InputStream#available()} does not temporarily return
* {@code 0} while the stream is still open. This works around {@link java.util.zip.GZIPInputStream} treating a temporary
* {@code 0} at a concatenated gzip member boundary as the end of the complete stream. Because this can cause a read after
* {@code available()} to block, it should only be enabled when the response will be read with {@code GZIPInputStream}.
*
* @param timeout Maximum time to wait for first read operation before aborting. Use {@link Duration#ZERO} or a negative
* {@link Duration} to disable timeout.
* @param gzipInputStreamCompatibilityEnabled Whether to enable {@code GZIPInputStream} compatibility for concatenated gzip.
* @param <ResponseT> Type of unmarshalled response POJO.
* @return ResponseTransformer instance.
*/
static <ResponseT> ResponseTransformer<ResponseT, ResponseInputStream<ResponseT>> toInputStream(
Duration timeout, boolean gzipInputStreamCompatibilityEnabled) {
return unmanaged(new ResponseTransformer<ResponseT, ResponseInputStream<ResponseT>>() {
@Override
public ResponseInputStream<ResponseT> transform(ResponseT response, AbortableInputStream inputStream) {
return new ResponseInputStream<>(response, inputStream, timeout);
AbortableInputStream content = gzipInputStreamCompatibilityEnabled
? GzipAvailabilityInputStream.wrap(inputStream, inputStream)
: inputStream;
return new ResponseInputStream<>(response, content, timeout);
}

@Override
Expand Down
Loading
Loading