Skip to content

Commit 912da36

Browse files
feat: add opt-in GZIPInputStream compatibility for blocking response streams (#7331)
* fix: prevent concatenated gzip response truncation when read via GZIPInputStream * Fix truncation of concatenated GZIP response streams * Add opt-out for concatenated GZIP stream support * Replace GZIP client option with transformer opt-in APIs This reverts commit f902dc2. * Use named factories for GZIP-compatible response streams
1 parent 881f20b commit 912da36

9 files changed

Lines changed: 1247 additions & 11 deletions

File tree

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
{
2+
"type": "feature",
3+
"category": "AWS SDK for Java v2",
4+
"contributor": "",
5+
"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.toGzipCompatibleInputStream() or AsyncResponseTransformer.toGzipCompatibleBlockingInputStream()."
6+
}

‎core/sdk-core/src/main/java/software/amazon/awssdk/core/async/AsyncResponseTransformer.java‎

Lines changed: 25 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -386,7 +386,31 @@ ResponsePublisher<ResponseT>> toPublisher(Duration timeout) {
386386
*/
387387
static <ResponseT extends SdkResponse>
388388
AsyncResponseTransformer<ResponseT, ResponseInputStream<ResponseT>> toBlockingInputStream() {
389-
return new InputStreamResponseTransformer<>();
389+
return new InputStreamResponseTransformer<>(false);
390+
}
391+
392+
/**
393+
* Creates an {@link AsyncResponseTransformer} that allows reading the response body content as an {@link InputStream}
394+
* adapted for reading concatenated GZIP content with {@link java.util.zip.GZIPInputStream}. You are responsible for
395+
* performing blocking reads from this input stream and closing the stream when you are finished.
396+
* <p>
397+
* When this transformer is used with an async client, the {@link CompletableFuture} that the client returns will be completed
398+
* once the {@link SdkResponse} is available and the response body <i>begins</i> streaming. This behavior differs from some
399+
* other transformers, like {@link #toFile(Path)} and {@link #toBytes()}, which only have their {@link CompletableFuture}
400+
* completed after the entire response body has finished streaming.
401+
* <p>
402+
* GZIP response streams are adapted so that {@link InputStream#available()} does not temporarily return {@code 0} while the
403+
* stream is still open. This works around {@code GZIPInputStream} treating a temporary {@code 0} at a concatenated GZIP
404+
* member boundary as the end of the complete stream. Because this can cause a read after {@code available()} to block, this
405+
* transformer should only be used when the response will be read with {@code GZIPInputStream}.
406+
*
407+
* @param <ResponseT> Type of unmarshalled response POJO.
408+
* @return AsyncResponseTransformer instance.
409+
* @see #toBlockingInputStream()
410+
*/
411+
static <ResponseT extends SdkResponse>
412+
AsyncResponseTransformer<ResponseT, ResponseInputStream<ResponseT>> toGzipCompatibleBlockingInputStream() {
413+
return new InputStreamResponseTransformer<>(true);
390414
}
391415

392416
/**

‎core/sdk-core/src/main/java/software/amazon/awssdk/core/internal/async/InputStreamResponseTransformer.java‎

Lines changed: 17 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515

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

18+
import java.io.InputStream;
1819
import java.nio.ByteBuffer;
1920
import java.util.concurrent.CompletableFuture;
2021
import org.reactivestreams.Subscriber;
@@ -24,12 +25,14 @@
2425
import software.amazon.awssdk.core.SdkResponse;
2526
import software.amazon.awssdk.core.async.AsyncResponseTransformer;
2627
import software.amazon.awssdk.core.async.SdkPublisher;
28+
import software.amazon.awssdk.core.internal.io.GzipAvailabilityInputStream;
2729
import software.amazon.awssdk.http.async.AbortableInputStreamSubscriber;
2830

2931
/**
3032
* A {@link AsyncResponseTransformer} that allows performing blocking reads on the response data.
3133
* <p>
32-
* Created with {@link AsyncResponseTransformer#toBlockingInputStream()}.
34+
* Created with {@link AsyncResponseTransformer#toBlockingInputStream()} or
35+
* {@link AsyncResponseTransformer#toGzipCompatibleBlockingInputStream()}.
3336
*/
3437
@SdkInternalApi
3538
public class InputStreamResponseTransformer<ResponseT extends SdkResponse>
@@ -38,6 +41,15 @@ public class InputStreamResponseTransformer<ResponseT extends SdkResponse>
3841
private volatile CompletableFuture<ResponseInputStream<ResponseT>> future;
3942
private volatile ResponseT response;
4043
private volatile WaitForSubscribeOnErrorWrapper subscriber;
44+
private final boolean gzipInputStreamCompatibilityEnabled;
45+
46+
public InputStreamResponseTransformer() {
47+
this(false);
48+
}
49+
50+
public InputStreamResponseTransformer(boolean gzipInputStreamCompatibilityEnabled) {
51+
this.gzipInputStreamCompatibilityEnabled = gzipInputStreamCompatibilityEnabled;
52+
}
4153

4254
@Override
4355
public CompletableFuture<ResponseInputStream<ResponseT>> prepare() {
@@ -59,7 +71,10 @@ public void onStream(SdkPublisher<ByteBuffer> publisher) {
5971
this.subscriber = waitForSubscribeSubscriber;
6072

6173
publisher.subscribe(waitForSubscribeSubscriber);
62-
future.complete(new ResponseInputStream<>(response, inputStreamSubscriber));
74+
InputStream content = gzipInputStreamCompatibilityEnabled
75+
? GzipAvailabilityInputStream.wrap(inputStreamSubscriber, inputStreamSubscriber)
76+
: inputStreamSubscriber;
77+
future.complete(new ResponseInputStream<>(response, content));
6378
}
6479

6580
@Override
Lines changed: 176 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,176 @@
1+
/*
2+
* Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
3+
*
4+
* Licensed under the Apache License, Version 2.0 (the "License").
5+
* You may not use this file except in compliance with the License.
6+
* A copy of the License is located at
7+
*
8+
* http://aws.amazon.com/apache2.0
9+
*
10+
* or in the "license" file accompanying this file. This file is distributed
11+
* on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either
12+
* express or implied. See the License for the specific language governing
13+
* permissions and limitations under the License.
14+
*/
15+
16+
package software.amazon.awssdk.core.internal.io;
17+
18+
import java.io.FilterInputStream;
19+
import java.io.IOException;
20+
import java.io.InputStream;
21+
import software.amazon.awssdk.annotations.SdkInternalApi;
22+
import software.amazon.awssdk.http.Abortable;
23+
import software.amazon.awssdk.http.AbortableInputStream;
24+
import software.amazon.awssdk.utils.IoUtils;
25+
26+
/**
27+
* Wraps a response body so {@code available()} never returns {@code 0} for gzip content while the stream is open.
28+
* {@link java.util.zip.GZIPInputStream} treats a transient {@code 0} from {@code available()} at a member boundary
29+
* as end of stream and stops, truncating concatenated gzip. Gzip is detected passively from the leading bytes
30+
* ({@code 1f 8b 08}); non-gzip streams keep honest {@code available()}.
31+
*/
32+
@SdkInternalApi
33+
public final class GzipAvailabilityInputStream extends FilterInputStream implements Releasable {
34+
35+
private static final int GZIP_MAGIC_1 = 0x1f;
36+
private static final int GZIP_MAGIC_2 = 0x8b;
37+
private static final int GZIP_METHOD_DEFLATE = 0x08;
38+
private static final int HEADER_LENGTH = 3;
39+
40+
private final byte[] header = new byte[HEADER_LENGTH];
41+
private volatile int headerLen;
42+
private volatile boolean classified;
43+
private volatile boolean gzipDetected;
44+
private volatile boolean eof;
45+
private volatile boolean closed;
46+
47+
private int markHeaderLen;
48+
private boolean markClassified;
49+
private boolean markGzipDetected;
50+
private boolean markEof;
51+
private boolean marked;
52+
53+
public GzipAvailabilityInputStream(InputStream in) {
54+
super(in);
55+
}
56+
57+
/**
58+
* Wraps a response body's content so {@code available()} is gzip-safe, while preserving {@code abort()} on the
59+
* original stream. Applied by the SDK at the sync and async blocking-stream boundaries (where the caller's
60+
* {@link java.util.zip.GZIPInputStream} reads), so that {@link software.amazon.awssdk.core.ResponseInputStream}
61+
* itself stays content-type agnostic.
62+
*/
63+
public static AbortableInputStream wrap(InputStream content, Abortable abortable) {
64+
return AbortableInputStream.create(new GzipAvailabilityInputStream(content), abortable);
65+
}
66+
67+
@Override
68+
public int read() throws IOException {
69+
int b = in.read();
70+
if (b == -1) {
71+
eof = true;
72+
} else {
73+
if (eof) {
74+
eof = false;
75+
}
76+
observe((byte) b);
77+
}
78+
return b;
79+
}
80+
81+
@Override
82+
public int read(byte[] b, int off, int len) throws IOException {
83+
int n = in.read(b, off, len);
84+
if (n == -1) {
85+
eof = true;
86+
} else if (n > 0) {
87+
if (eof) {
88+
eof = false;
89+
}
90+
observe(b, off, n);
91+
}
92+
return n;
93+
}
94+
95+
@Override
96+
public int available() throws IOException {
97+
if (closed) {
98+
return 0;
99+
}
100+
int available = in.available();
101+
return available == 0 && gzipDetected && !eof ? 1 : available;
102+
}
103+
104+
@Override
105+
public long skip(long n) throws IOException {
106+
long skipped = in.skip(n);
107+
if (skipped > 0) {
108+
classified = true;
109+
}
110+
return skipped;
111+
}
112+
113+
@Override
114+
public synchronized void mark(int readlimit) {
115+
markHeaderLen = headerLen;
116+
markClassified = classified;
117+
markGzipDetected = gzipDetected;
118+
markEof = eof;
119+
marked = true;
120+
in.mark(readlimit);
121+
}
122+
123+
@Override
124+
public synchronized void reset() throws IOException {
125+
in.reset();
126+
if (marked) {
127+
headerLen = markHeaderLen;
128+
classified = markClassified;
129+
gzipDetected = markGzipDetected;
130+
eof = markEof;
131+
} else {
132+
headerLen = 0;
133+
classified = false;
134+
gzipDetected = false;
135+
eof = false;
136+
}
137+
}
138+
139+
@Override
140+
public void close() throws IOException {
141+
closed = true;
142+
in.close();
143+
}
144+
145+
@Override
146+
public void release() {
147+
IoUtils.closeQuietly(this, null);
148+
if (in instanceof Releasable) {
149+
((Releasable) in).release();
150+
}
151+
}
152+
153+
private void observe(byte[] b, int off, int len) {
154+
if (classified) {
155+
return;
156+
}
157+
for (int i = 0; i < len && !classified; i++) {
158+
observe(b[off + i]);
159+
}
160+
}
161+
162+
private void observe(byte b) {
163+
if (classified) {
164+
return;
165+
}
166+
int len = headerLen;
167+
header[len] = b;
168+
headerLen = len + 1;
169+
if (headerLen == HEADER_LENGTH) {
170+
classified = true;
171+
gzipDetected = (header[0] & 0xff) == GZIP_MAGIC_1
172+
&& (header[1] & 0xff) == GZIP_MAGIC_2
173+
&& (header[2] & 0xff) == GZIP_METHOD_DEFLATE;
174+
}
175+
}
176+
}

‎core/sdk-core/src/main/java/software/amazon/awssdk/core/sync/ResponseTransformer.java‎

Lines changed: 56 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,7 @@
3636
import software.amazon.awssdk.core.exception.SdkClientException;
3737
import software.amazon.awssdk.core.exception.SdkException;
3838
import software.amazon.awssdk.core.internal.http.InterruptMonitor;
39+
import software.amazon.awssdk.core.internal.io.GzipAvailabilityInputStream;
3940
import software.amazon.awssdk.core.retry.RetryPolicy;
4041
import software.amazon.awssdk.http.AbortableInputStream;
4142
import software.amazon.awssdk.utils.IoUtils;
@@ -261,10 +262,49 @@ public String name() {
261262
* @see #toInputStream(Duration)
262263
*/
263264
static <ResponseT> ResponseTransformer<ResponseT, ResponseInputStream<ResponseT>> toInputStream() {
265+
return toInputStream(null);
266+
}
267+
268+
/**
269+
* Creates a response transformer that returns an unmanaged input stream adapted for reading concatenated GZIP content with
270+
* {@link java.util.zip.GZIPInputStream}. This input stream must be explicitly closed to release the connection.
271+
* <p>
272+
* GZIP response streams are adapted so that {@link InputStream#available()} does not temporarily return {@code 0} while the
273+
* stream is still open. This works around {@code GZIPInputStream} treating a temporary {@code 0} at a concatenated GZIP
274+
* member boundary as the end of the complete stream. Because this can cause a read after {@code available()} to block, this
275+
* transformer should only be used when the response will be read with {@code GZIPInputStream}.
276+
* <p>
277+
* The stream has the default first-read timeout of 60 seconds. Use
278+
* {@link #toGzipCompatibleInputStream(Duration)} to specify a custom timeout.
279+
*
280+
* @param <ResponseT> Type of unmarshalled response POJO.
281+
* @return ResponseTransformer instance.
282+
* @see #toInputStream()
283+
* @see #toGzipCompatibleInputStream(Duration)
284+
*/
285+
static <ResponseT> ResponseTransformer<ResponseT, ResponseInputStream<ResponseT>> toGzipCompatibleInputStream() {
286+
return toGzipCompatibleInputStream(null);
287+
}
288+
289+
/**
290+
* Creates a response transformer that returns an unmanaged input stream with the response content and a custom timeout.
291+
* This input stream must be explicitly closed to release the connection.
292+
* <p>
293+
* The timeout starts when the response stream is ready. If no read operation occurs within the specified timeout, the
294+
* connection will be automatically aborted. To disable the timeout, pass {@link Duration#ZERO} or a negative
295+
* {@link Duration}.
296+
*
297+
* @param timeout Maximum time to wait for first read operation before aborting. Use {@link Duration#ZERO} or a negative
298+
* {@link Duration} to disable timeout.
299+
* @param <ResponseT> Type of unmarshalled response POJO.
300+
* @return ResponseTransformer instance.
301+
* @see #toInputStream()
302+
*/
303+
static <ResponseT> ResponseTransformer<ResponseT, ResponseInputStream<ResponseT>> toInputStream(Duration timeout) {
264304
return unmanaged(new ResponseTransformer<ResponseT, ResponseInputStream<ResponseT>>() {
265305
@Override
266306
public ResponseInputStream<ResponseT> transform(ResponseT response, AbortableInputStream inputStream) {
267-
return new ResponseInputStream<>(response, inputStream);
307+
return new ResponseInputStream<>(response, inputStream, timeout);
268308
}
269309

270310
@Override
@@ -275,24 +315,33 @@ public String name() {
275315
}
276316

277317
/**
278-
* Creates a response transformer that returns an unmanaged input stream with the response content and a custom timeout.
279-
* This input stream must be explicitly closed to release the connection.
318+
* Creates a response transformer that returns an unmanaged input stream adapted for reading concatenated GZIP content with
319+
* {@link java.util.zip.GZIPInputStream} and a custom timeout. This input stream must be explicitly closed to release the
320+
* connection.
321+
* <p>
322+
* GZIP response streams are adapted so that {@link InputStream#available()} does not temporarily return {@code 0} while the
323+
* stream is still open. This works around {@code GZIPInputStream} treating a temporary {@code 0} at a concatenated GZIP
324+
* member boundary as the end of the complete stream. Because this can cause a read after {@code available()} to block, this
325+
* transformer should only be used when the response will be read with {@code GZIPInputStream}.
280326
* <p>
281327
* The timeout starts when the response stream is ready. If no read operation occurs within the specified timeout, the
282328
* connection will be automatically aborted. To disable the timeout, pass {@link Duration#ZERO} or a negative
283329
* {@link Duration}.
284330
*
285331
* @param timeout Maximum time to wait for first read operation before aborting. Use {@link Duration#ZERO} or a negative
286-
* {@link Duration} to disable timeout.
332+
* {@link Duration} to disable timeout.
287333
* @param <ResponseT> Type of unmarshalled response POJO.
288334
* @return ResponseTransformer instance.
289-
* @see #toInputStream()
335+
* @see #toGzipCompatibleInputStream()
336+
* @see #toInputStream(Duration)
290337
*/
291-
static <ResponseT> ResponseTransformer<ResponseT, ResponseInputStream<ResponseT>> toInputStream(Duration timeout) {
338+
static <ResponseT> ResponseTransformer<ResponseT, ResponseInputStream<ResponseT>> toGzipCompatibleInputStream(
339+
Duration timeout) {
292340
return unmanaged(new ResponseTransformer<ResponseT, ResponseInputStream<ResponseT>>() {
293341
@Override
294342
public ResponseInputStream<ResponseT> transform(ResponseT response, AbortableInputStream inputStream) {
295-
return new ResponseInputStream<>(response, inputStream, timeout);
343+
AbortableInputStream content = GzipAvailabilityInputStream.wrap(inputStream, inputStream);
344+
return new ResponseInputStream<>(response, content, timeout);
296345
}
297346

298347
@Override

0 commit comments

Comments
 (0)