From 005b696bfb999634e26e42d3ad6d16be448ea055 Mon Sep 17 00:00:00 2001 From: Matthias Kurz Date: Mon, 5 Oct 2026 14:03:42 +0200 Subject: [PATCH] Stop inflating a suspended HTTP/1.1 response Motivation: HTTP/1.1 automatic decompression ran Netty's HttpContentDecompressor, which inflates every chunk completely as soon as it is read. Suspending a response through ResponseBodyControl stops further socket reads, but not that. One 64 KiB read of a highly compressible gzip body (about 1032:1) still became some 67 MB of body parts for a handler that had asked for none; only maxDecompressedResponseSize (256 MiB) bounded it. AHC 2 behaved the same way, so this is hardening, not a regression fix. Modification: Http1ContentDecompressor is now a ChannelDuplexHandler built on Netty's pull-based Decompressor (Netty 4.2.17+; Netty still marks this API unstable). It keeps the codings (gzip, including bodies of several gzip members, deflate, br, snappy, zstd), the header rewrites and the maxDecompressedResponseSize accounting. It hands the body on in parts of at most 64 KiB, splitting larger decompressor output, and checks before each part whether the response is suspended. If it is, the compressed input stays in the handler, and so does every message behind it. The read request that resume() issues continues decompression before it reaches the socket. Without suspension, each part goes straight through as before. The pipeline is torn down with its channel, so if the connection closes while input is held back, the rest is inflated and delivered before channelInactive. Once the exchange has ended (cancel, ABORT, timeout), held input is dropped instead of inflated. A framing failure on the final chunk now reaches the handler; Netty's decoder used to replace it with success. Revapi reports the changed superclass of this public handler class and the members inherited from Netty's decoder; they are justified in the pom. HTTP/2 is unchanged. After END_STREAM, Netty's stream channel hands every queued frame over on the next read and then closes. A handler in the stream pipeline therefore cannot hold compressed input back. Result: A response whose handler suspends it after each body part receives one part of at most 64 KiB per resume, and decompression runs ahead of the delivered parts by at most one decompressor output buffer. For 64 MiB of gzipped zeros, main delivered 1.96 MB in 30 more parts past the suspending part, then up to 33.7 MB per resume; now nothing past it, and 65,536 bytes per resume. Http1DecompressionSuspensionTest covers auto-read on and off, chunked bodies with trailers, close-delimited bodies, cancel, ABORT, bodiless responses and connection reuse; its two bound tests fail on main. Http1ContentDecompressorTest covers the handler on its own, including gzip bodies of several members, corrupt and trailing data, and every coding across repeated suspension. Http1ContentDecompressorLargeOutputTest covers decompressor output larger than a part, in a JVM with Netty's io.netty.compression.defaultMaxForwardBytes raised. Loopback throughput without suspension: unchanged for gzipped zeros, about 5% more elapsed time for gzipped text. Claude Code on behalf of Matthias Kurz Co-Authored-By: Claude Opus 5.5 --- .../DefaultAsyncHttpClientConfig.java | 2 +- .../netty/channel/ChannelManager.java | 3 +- .../handler/Http1ContentDecompressor.java | 430 +++++++++++-- .../Http1DecompressionSuspensionTest.java | 448 ++++++++++++++ ...tp1ContentDecompressorLargeOutputTest.java | 238 ++++++++ .../handler/Http1ContentDecompressorTest.java | 572 ++++++++++++++++++ pom.xml | 63 ++ 7 files changed, 1698 insertions(+), 58 deletions(-) create mode 100644 client/src/test/java/org/asynchttpclient/Http1DecompressionSuspensionTest.java create mode 100644 client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorLargeOutputTest.java create mode 100644 client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java diff --git a/client/src/main/java/org/asynchttpclient/DefaultAsyncHttpClientConfig.java b/client/src/main/java/org/asynchttpclient/DefaultAsyncHttpClientConfig.java index a316cc2e9..419e29ec3 100644 --- a/client/src/main/java/org/asynchttpclient/DefaultAsyncHttpClientConfig.java +++ b/client/src/main/java/org/asynchttpclient/DefaultAsyncHttpClientConfig.java @@ -1249,7 +1249,7 @@ public Builder setCompressionEnforced(boolean compressionEnforced) { } /* - * If true (default), AHC will add a Netty HttpContentDecompressor, so compressed + * If true (default), AHC will add a content decompressor, so compressed * content will automatically get decompressed. * * If set to false, response will be delivered as is received. Decompression must diff --git a/client/src/main/java/org/asynchttpclient/netty/channel/ChannelManager.java b/client/src/main/java/org/asynchttpclient/netty/channel/ChannelManager.java index 18dce043d..ebeb63377 100755 --- a/client/src/main/java/org/asynchttpclient/netty/channel/ChannelManager.java +++ b/client/src/main/java/org/asynchttpclient/netty/channel/ChannelManager.java @@ -470,8 +470,7 @@ protected void initChannel(Channel ch) { * The HTTP/1.1 decompressor, bounded so a decompression bomb — a small, highly compressible body that * inflates without limit — fails the exchange instead of exhausting the heap. The ceiling comes from * {@link AsyncHttpClientConfig#getMaxDecompressedResponseSize()} (256 MiB by default, {@code 0} to - * disable); see {@link Http1ContentDecompressor} for why Netty's own {@code maxAllocation} argument - * does not provide this. + * disable). It also stops inflating while the response is suspended; see {@link Http1ContentDecompressor}. */ private Http1ContentDecompressor newHttpContentDecompressor() { return new Http1ContentDecompressor(config.isKeepEncodingHeader(), config.getMaxDecompressedResponseSize()); diff --git a/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java b/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java index c990f50eb..537de6a5a 100644 --- a/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java +++ b/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java @@ -16,100 +16,420 @@ package org.asynchttpclient.netty.handler; import io.netty.buffer.ByteBuf; +import io.netty.buffer.ByteBufAllocator; +import io.netty.buffer.Unpooled; +import io.netty.channel.Channel; +import io.netty.channel.ChannelDuplexHandler; import io.netty.channel.ChannelHandlerContext; -import io.netty.channel.ChannelInboundHandlerAdapter; -import io.netty.channel.embedded.EmbeddedChannel; +import io.netty.handler.codec.compression.Brotli; +import io.netty.handler.codec.compression.BrotliDecompressor; import io.netty.handler.codec.compression.DecompressionException; -import io.netty.handler.codec.http.HttpContentDecompressor; +import io.netty.handler.codec.compression.Decompressor; +import io.netty.handler.codec.compression.JdkZlibDecompressor; +import io.netty.handler.codec.compression.SnappyFrameDecompressor; +import io.netty.handler.codec.compression.ZlibWrapper; +import io.netty.handler.codec.compression.Zstd; +import io.netty.handler.codec.compression.ZstdDecompressor; +import io.netty.handler.codec.http.DefaultHttpContent; +import io.netty.handler.codec.http.DefaultHttpResponse; +import io.netty.handler.codec.http.DefaultLastHttpContent; +import io.netty.handler.codec.http.HttpContent; +import io.netty.handler.codec.http.HttpHeaderNames; +import io.netty.handler.codec.http.HttpHeaderValues; +import io.netty.handler.codec.http.HttpHeaders; +import io.netty.handler.codec.http.HttpResponse; +import io.netty.handler.codec.http.LastHttpContent; import io.netty.util.ReferenceCountUtil; +import org.asynchttpclient.netty.DiscardEvent; +import org.asynchttpclient.netty.NettyResponseBodyControl; +import org.asynchttpclient.netty.NettyResponseFuture; +import org.asynchttpclient.netty.channel.Channels; +import org.jetbrains.annotations.Nullable; + +import java.util.ArrayDeque; /** - * HTTP/1.1 content decompressor that bounds how far a response body may inflate, mirroring what - * {@link Http2ContentDecompressor} enforces on HTTP/2 streams. + * HTTP/1.1 content decompressor that inflates a response body only as fast as the response consumes it, and + * bounds how far the body may inflate in total. + *

+ * Netty's {@code HttpContentDecompressor} inflates each compressed chunk completely as soon as it is read. + * Suspending a response through {@link org.asynchttpclient.ResponseBodyControl} stops further socket reads but + * cannot stop that, so one socket read of a highly compressible body still became tens of megabytes of body parts + * nobody had asked for. This handler uses Netty's pull-based {@link Decompressor} instead and hands the body on + * one part of at most {@value #MAX_PART_SIZE} bytes at a time. Before each part it checks whether the response is + * suspended. If it is, the compressed input stays here, together with every message that arrived after it, and + * the next read request, such as the one {@code resume()} issues, inflates that input before it is passed on to + * the transport. A suspension that takes effect on the event loop, as one made from a body callback does, + * therefore stops decompression before the next part. + *

+ * Closing the channel tears its pipeline down, so if the connection closes while input is held back, the rest of + * it is inflated and handed on at once, as it was before. *

- * Netty's own {@code maxAllocation} argument is not such a bound. It caps the capacity of the one output - * buffer produced by a single {@code ZlibDecoder.decode()} call, and nothing accumulates across calls. - * AHC hands the decompressor one {@code HttpContent} of at most {@code httpClientCodecMaxChunkSize} - * (8 KiB by default) at a time, so even DEFLATE's ~1032:1 ceiling keeps a single call's output around - * 8 MiB — under any sane limit — while the response as a whole inflates without bound. The cap therefore - * never fires on an ordinary decompression bomb. + * The cumulative decompressed size of each response is bounded by {@code maxDecompressedBytes}; exceeding it + * fails the exchange with a {@link DecompressionException}, the same exception type the HTTP/2 path uses. Only + * what decompression produced is counted, so an unencoded response, which is passed through untouched, is never + * failed for being large; its size is the caller's own choice, not a bomb. *

- * The counting is done inside the decoder's own {@link EmbeddedChannel} rather than around - * {@code HttpContentDecoder#decode}: that class forwards decompressed output straight down the outer - * pipeline through an internal forwarder installed at the end of the embedded pipeline, so the output - * never passes through the {@code decode} out-list where it could be measured. Sitting between the - * {@code ZlibDecoder} and that forwarder gives an exact count of what decompression produced — and only - * of that, so an unencoded response, which never gets a decoder, is passed through untouched and is - * never failed for being large; its size is the caller's own choice, not a bomb. + * Content codings, and the {@code Content-Encoding}, {@code Content-Length} and {@code Transfer-Encoding} header + * rewrites, follow Netty's {@code HttpContentDecompressor}. *

- * A decoder is created per response, so each response gets a fresh counter even though this handler is - * shared by every response on a keep-alive connection. + * Netty still marks its {@link Decompressor} API as unstable while it migrates its codecs to it, so a Netty upgrade + * may require changes here. */ -public class Http1ContentDecompressor extends HttpContentDecompressor { +public class Http1ContentDecompressor extends ChannelDuplexHandler { + + // Largest body part handed on at once. Larger decompressor output, such as zstd's when Netty's + // io.netty.compression.defaultMaxForwardBytes is raised, is split. + static final int MAX_PART_SIZE = 64 * 1024; private final boolean keepEncodingHeader; // Maximum cumulative decompressed bytes for one response; 0 disables the limit. private final long maxDecompressedBytes; + // Inbound messages not yet handled, in arrival order. They queue up only behind a suspended response body. + private final ArrayDeque pending = new ArrayDeque<>(); + // Null when the current response body is not compressed, or once decompressing it failed or became pointless. + private @Nullable Decompressor decompressor; + // Whether the current response body is compressed. Stays set after a failure so the rest of the body is dropped. + private boolean decoding; + // Output already taken from the decompressor that did not fit into the previous part. + private @Nullable ByteBuf carriedOutput; + // The terminating chunk of the current response, handed on once the rest of its body has been. + private @Nullable LastHttpContent lastContent; + private long decompressedBytes; + private boolean continueResponse; + private boolean draining; + // Whether decompressed output is held back for a suspended response. + private boolean paused; + private boolean readPending; + private boolean inputClosed; + public Http1ContentDecompressor(boolean keepEncodingHeader, long maxDecompressedBytes) { - // maxAllocation=0 is what the no-arg HttpContentDecompressor() constructor passes, so a single - // decode call behaves exactly as it did before any bound was attempted. The counter below is what - // enforces the limit. - super(0); this.keepEncodingHeader = keepEncodingHeader; this.maxDecompressedBytes = maxDecompressedBytes; } @Override - protected EmbeddedChannel newContentDecoder(String contentEncoding) throws Exception { - EmbeddedChannel decoder = super.newContentDecoder(contentEncoding); - if (decoder != null && maxDecompressedBytes > 0) { - // Appended before HttpContentDecoder appends its forwarder, so every decompressed buffer is - // counted before it leaves for the outer pipeline. - decoder.pipeline().addLast(new DecompressedSizeLimiter()); + public void channelRead(ChannelHandlerContext ctx, Object msg) { + pending.add(msg); + if (!draining) { + drain(ctx); + } + } + + @Override + public void read(ChannelHandlerContext ctx) { + if (draining || paused) { + // Input that was already read is used up before the transport is asked for more. + readPending = true; + if (!draining) { + drain(ctx); + } + } else { + ctx.read(); + } + } + + @Override + public void channelInactive(ChannelHandlerContext ctx) { + // The pipeline is about to be torn down, so anything held back has to go now, even to a suspended response. + inputClosed = true; + try { + if (paused && !draining) { + drain(ctx); + } + } catch (Throwable cause) { + ctx.fireExceptionCaught(cause); + } finally { + release(); } - return decoder; + ctx.fireChannelInactive(); } @Override - protected String getTargetContentEncoding(String contentEncoding) throws Exception { - // Leaving Content-Encoding in place lets the caller see what the server sent, at the cost of a - // header that no longer describes the body handed up. Opt-in via keepEncodingHeader. - return keepEncodingHeader ? contentEncoding : super.getTargetContentEncoding(contentEncoding); + public void handlerRemoved(ChannelHandlerContext ctx) { + release(); + } + + private void drain(ChannelHandlerContext ctx) { + draining = true; + try { + paused = !process(ctx); + } finally { + draining = false; + } + if (!paused && readPending) { + readPending = false; + ctx.read(); + } } /** - * Fails the response once its decompressed body passes {@link #maxDecompressedBytes}. The thrown - * {@link DecompressionException} is recorded by the {@link EmbeddedChannel} and rethrown out of the - * {@code writeInbound} call driving the decoder, so it surfaces exactly like a corrupt-body error and - * fails the exchange — the same exception type the HTTP/2 path uses. + * Hands on everything that can be handed on. Returns {@code false} if it stopped because the response is + * suspended while it has decompressed output to deliver. */ - private final class DecompressedSizeLimiter extends ChannelInboundHandlerAdapter { + private boolean process(ChannelHandlerContext ctx) { + try { + for (;;) { + Decompressor decompressor = this.decompressor; + if (decompressor != null && hasOutput(decompressor)) { + Channel channel = ctx.channel(); + if (isAbandoned(channel)) { + // The exchange is over; inflating the rest of its body would be wasted work. + closeDecompressor(); + continue; + } + if (!inputClosed && isSuspended(channel)) { + return false; + } + ByteBuf part = nextPart(decompressor, ctx.alloc()); + if (part != null) { + ctx.fireChannelRead(new DefaultHttpContent(part)); + } + } else if (lastContent != null) { + LastHttpContent last = lastContent; + lastContent = null; + endResponse(); + ctx.fireChannelRead(last); + } else { + Object msg = pending.poll(); + if (msg == null) { + return true; + } + handle(ctx, msg); + } + } + } catch (RuntimeException e) { + // A corrupt or oversized body fails the exchange; the rest of it is dropped. + closeDecompressor(); + throw e; + } + } - private long totalDecompressedBytes; - private boolean exceeded; + private void handle(ChannelHandlerContext ctx, Object msg) { + // A 100 Continue, and everything up to the LastHttpContent that ends it, passes through untouched, as with + // Netty's HttpContentDecoder. + if (continueResponse || msg instanceof HttpResponse && ((HttpResponse) msg).status().code() == 100) { + continueResponse = !(msg instanceof LastHttpContent); + ctx.fireChannelRead(msg); + return; + } - @Override - public void channelRead(ChannelHandlerContext ctx, Object msg) { - if (exceeded) { - // The decoder is torn down with finishAndReleaseAll(), which re-runs it over anything left - // cumulated. Drop that quietly: throwing a second time out of cleanup would replace the - // error the caller sees and can leak the decoder itself. + if (msg instanceof HttpResponse) { + HttpResponse response = (HttpResponse) msg; + try { + startResponse(response, ctx.alloc()); + } catch (RuntimeException e) { ReferenceCountUtil.release(msg); + throw e; + } + if (!decoding || !(msg instanceof HttpContent)) { + ctx.fireChannelRead(msg); return; } + // Only the content of a response that carries its own content is decompressed. + DefaultHttpResponse head = new DefaultHttpResponse(response.protocolVersion(), response.status(), + response.headers()); + head.setDecoderResult(response.decoderResult()); + ctx.fireChannelRead(head); + } + + if (decoding && msg instanceof HttpContent) { + decode((HttpContent) msg); + } else { + ctx.fireChannelRead(msg); + } + } + + private void startResponse(HttpResponse response, ByteBufAllocator alloc) { + endResponse(); + HttpHeaders headers = response.headers(); + String contentEncoding = contentEncoding(headers); + decompressor = newDecompressor(contentEncoding, alloc); + if (contentEncoding == null || decompressor == null) { + return; + } + + decoding = true; + decompressedBytes = 0; + if (headers.contains(HttpHeaderNames.CONTENT_LENGTH)) { + headers.remove(HttpHeaderNames.CONTENT_LENGTH); + headers.set(HttpHeaderNames.TRANSFER_ENCODING, HttpHeaderValues.CHUNKED); + } + // Leaving Content-Encoding in place lets the caller see what the server sent, at the cost of a header that + // no longer describes the body handed up. Opt-in via keepEncodingHeader. + if (keepEncodingHeader) { + headers.set(HttpHeaderNames.CONTENT_ENCODING, contentEncoding); + } else { + headers.remove(HttpHeaderNames.CONTENT_ENCODING); + } + } - if (msg instanceof ByteBuf) { - totalDecompressedBytes += ((ByteBuf) msg).readableBytes(); - if (totalDecompressedBytes > maxDecompressedBytes) { - exceeded = true; - ReferenceCountUtil.release(msg); + private void decode(HttpContent content) { + try { + Decompressor decompressor = this.decompressor; + ByteBuf input = content.content(); + // Once the compressed stream is complete, further bytes are dropped, as Netty's decoders drop them. A gzip + // stream is only complete at the end of the input, because another member may follow. + if (decompressor != null && input.isReadable() + && decompressor.status() == Decompressor.Status.NEED_INPUT) { + decompressor.addInput(input.retain()); + } + if (content instanceof LastHttpContent) { + lastContent = trailer((LastHttpContent) content); + } + } finally { + content.release(); + } + } + + private boolean hasOutput(Decompressor decompressor) { + return carriedOutput != null || decompressor.status() == Decompressor.Status.NEED_OUTPUT; + } + + /** + * Takes up to {@link #MAX_PART_SIZE} bytes of output. Small output buffers are copied together so a highly + * compressible body is not handed on in hundreds of tiny parts. + */ + private @Nullable ByteBuf nextPart(Decompressor decompressor, ByteBufAllocator alloc) { + ByteBuf part = carriedOutput; + carriedOutput = null; + if (part != null && part.readableBytes() > MAX_PART_SIZE) { + return split(part); + } + boolean merged = false; + try { + while (decompressor.status() == Decompressor.Status.NEED_OUTPUT) { + ByteBuf output = decompressor.takeOutput(); + int length = output.readableBytes(); + if (maxDecompressedBytes > 0 && (decompressedBytes += length) > maxDecompressedBytes) { + output.release(); throw new DecompressionException("HTTP/1.1 response body exceeds the maximum decompressed size of " + maxDecompressedBytes + " bytes"); } + + if (length == 0) { + output.release(); + } else if (part == null) { + if (length > MAX_PART_SIZE) { + return split(output); + } + part = output; + } else if (part.readableBytes() + length > MAX_PART_SIZE) { + carriedOutput = output; + break; + } else { + if (!merged) { + ByteBuf first = part; + part = alloc.heapBuffer(MAX_PART_SIZE); + part.writeBytes(first); + first.release(); + merged = true; + } + part.writeBytes(output); + output.release(); + } } + } catch (RuntimeException e) { + ReferenceCountUtil.release(part); + throw e; + } + return part; + } - ctx.fireChannelRead(msg); + /** Hands on the first {@link #MAX_PART_SIZE} bytes of {@code output} and carries the rest over. */ + private ByteBuf split(ByteBuf output) { + ByteBuf part = output.readRetainedSlice(MAX_PART_SIZE); + carriedOutput = output; + return part; + } + + private void endResponse() { + closeDecompressor(); + decoding = false; + } + + private void closeDecompressor() { + if (carriedOutput != null) { + carriedOutput.release(); + carriedOutput = null; + } + Decompressor decompressor = this.decompressor; + if (decompressor != null) { + this.decompressor = null; + decompressor.close(); + } + } + + private void release() { + Object msg; + while ((msg = pending.poll()) != null) { + ReferenceCountUtil.release(msg); + } + lastContent = null; + paused = false; + endResponse(); + } + + private static boolean isSuspended(Channel channel) { + Object attribute = Channels.getAttribute(channel); + return attribute instanceof NettyResponseFuture + && NettyResponseBodyControl.isSuspended((NettyResponseFuture) attribute, channel); + } + + private static boolean isAbandoned(Channel channel) { + Object attribute = Channels.getAttribute(channel); + return attribute == DiscardEvent.DISCARD + || attribute instanceof NettyResponseFuture && ((NettyResponseFuture) attribute).isDone(); + } + + private static @Nullable String contentEncoding(HttpHeaders headers) { + String contentEncoding = headers.get(HttpHeaderNames.CONTENT_ENCODING); + if (contentEncoding != null) { + return contentEncoding.trim(); + } + String transferEncoding = headers.get(HttpHeaderNames.TRANSFER_ENCODING); + if (transferEncoding == null) { + return null; + } + int comma = transferEncoding.indexOf(','); + return (comma != -1 ? transferEncoding.substring(0, comma) : transferEncoding).trim(); + } + + private static @Nullable Decompressor newDecompressor(@Nullable String contentEncoding, ByteBufAllocator alloc) { + if (contentEncoding == null) { + return null; + } + if (HttpHeaderValues.GZIP.contentEqualsIgnoreCase(contentEncoding) + || HttpHeaderValues.X_GZIP.contentEqualsIgnoreCase(contentEncoding)) { + // A gzip body may consist of several members (RFC 1952, section 2.2). Netty's HttpContentDecompressor + // decodes all of them, and so does this one. + return JdkZlibDecompressor.builder().wrapper(ZlibWrapper.GZIP).decompressConcatenated(true) + .maxAllocation(MAX_PART_SIZE).build(alloc); + } + if (HttpHeaderValues.DEFLATE.contentEqualsIgnoreCase(contentEncoding) + || HttpHeaderValues.X_DEFLATE.contentEqualsIgnoreCase(contentEncoding)) { + return JdkZlibDecompressor.builder().wrapper(ZlibWrapper.ZLIB_OR_NONE).maxAllocation(MAX_PART_SIZE) + .build(alloc); + } + if (Brotli.isAvailable() && HttpHeaderValues.BR.contentEqualsIgnoreCase(contentEncoding)) { + return BrotliDecompressor.builder().maxOutputChunkSize(MAX_PART_SIZE).build(alloc); + } + if (HttpHeaderValues.SNAPPY.contentEqualsIgnoreCase(contentEncoding)) { + return SnappyFrameDecompressor.builder().build(alloc); + } + if (Zstd.isAvailable() && HttpHeaderValues.ZSTD.contentEqualsIgnoreCase(contentEncoding)) { + return ZstdDecompressor.builder().build(alloc); + } + return null; + } + + private static LastHttpContent trailer(LastHttpContent last) { + if (last.trailingHeaders().isEmpty() && last.decoderResult().isSuccess()) { + return LastHttpContent.EMPTY_LAST_CONTENT; } + LastHttpContent trailer = new DefaultLastHttpContent(Unpooled.EMPTY_BUFFER, last.trailingHeaders()); + trailer.setDecoderResult(last.decoderResult()); + return trailer; } } diff --git a/client/src/test/java/org/asynchttpclient/Http1DecompressionSuspensionTest.java b/client/src/test/java/org/asynchttpclient/Http1DecompressionSuspensionTest.java new file mode 100644 index 000000000..402e2970c --- /dev/null +++ b/client/src/test/java/org/asynchttpclient/Http1DecompressionSuspensionTest.java @@ -0,0 +1,448 @@ +/* + * Copyright (c) 2026 AsyncHttpClient Project. 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. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License 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 org.asynchttpclient; + +import io.github.nettyplus.leakdetector.junit.NettyLeakDetectorExtension; +import io.netty.bootstrap.ServerBootstrap; +import io.netty.buffer.Unpooled; +import io.netty.channel.Channel; +import io.netty.channel.ChannelFutureListener; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelInitializer; +import io.netty.channel.ChannelOption; +import io.netty.channel.SimpleChannelInboundHandler; +import io.netty.channel.group.ChannelGroup; +import io.netty.channel.group.DefaultChannelGroup; +import io.netty.channel.nio.NioEventLoopGroup; +import io.netty.channel.socket.nio.NioServerSocketChannel; +import io.netty.handler.codec.http.DefaultFullHttpResponse; +import io.netty.handler.codec.http.DefaultHttpContent; +import io.netty.handler.codec.http.DefaultHttpResponse; +import io.netty.handler.codec.http.DefaultLastHttpContent; +import io.netty.handler.codec.http.FullHttpRequest; +import io.netty.handler.codec.http.HttpHeaderNames; +import io.netty.handler.codec.http.HttpHeaderValues; +import io.netty.handler.codec.http.HttpHeaders; +import io.netty.handler.codec.http.HttpMethod; +import io.netty.handler.codec.http.HttpObjectAggregator; +import io.netty.handler.codec.http.HttpResponse; +import io.netty.handler.codec.http.HttpServerCodec; +import io.netty.handler.codec.http.HttpUtil; +import io.netty.handler.codec.http.LastHttpContent; +import io.netty.util.concurrent.GlobalEventExecutor; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; + +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.net.InetSocketAddress; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.Random; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; +import java.util.zip.GZIPOutputStream; + +import static io.netty.handler.codec.http.HttpResponseStatus.OK; +import static io.netty.handler.codec.http.HttpVersion.HTTP_1_1; +import static java.util.concurrent.TimeUnit.MILLISECONDS; +import static java.util.concurrent.TimeUnit.SECONDS; +import static org.asynchttpclient.Dsl.asyncHttpClient; +import static org.asynchttpclient.Dsl.config; +import static org.junit.jupiter.api.Assertions.assertArrayEquals; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * A suspended HTTP/1.1 response must not keep inflating the compressed bytes it has already read. A highly + * compressible body inflates ~1000:1, so a single 64 KiB socket read can otherwise produce tens of megabytes + * of body parts that no one asked for. + */ +@ExtendWith(NettyLeakDetectorExtension.class) +public class Http1DecompressionSuspensionTest { + + private static final int MAX_PART_SIZE = 64 * 1024; + + /** 64 MiB of zeros: gzip inflates it ~1000:1, so the whole compressed body is about one socket read. */ + private static final int ZEROS_SIZE = 64 * 1024 * 1024; + private static final byte[] ZEROS_GZIP = gzip(new byte[ZEROS_SIZE]); + + private static final byte[] TEXT = text(4 * 1024 * 1024); + private static final byte[] TEXT_GZIP = gzip(TEXT); + + private final AtomicInteger connectionCount = new AtomicInteger(); + + private NioEventLoopGroup serverGroup; + private Channel serverChannel; + private ChannelGroup serverChildChannels; + private int serverPort; + + private static byte[] gzip(byte[] payload) { + ByteArrayOutputStream compressed = new ByteArrayOutputStream(); + try (GZIPOutputStream gzip = new GZIPOutputStream(compressed)) { + gzip.write(payload); + } catch (IOException e) { + throw new IllegalStateException(e); + } + return compressed.toByteArray(); + } + + private static byte[] text(int length) { + String[] words = "the quick brown fox jumps over the lazy dog while streams wait for demand".split(" "); + Random random = new Random(42); + StringBuilder text = new StringBuilder(length + 16); + while (text.length() < length) { + text.append(words[random.nextInt(words.length)]).append(random.nextInt(9) == 0 ? '\n' : ' '); + } + return text.substring(0, length).getBytes(StandardCharsets.US_ASCII); + } + + @BeforeEach + public void startServer() throws InterruptedException { + serverGroup = new NioEventLoopGroup(1); + serverChildChannels = new DefaultChannelGroup("http1-decompression-suspension", GlobalEventExecutor.INSTANCE); + serverChannel = new ServerBootstrap() + .group(serverGroup) + .channel(NioServerSocketChannel.class) + .childHandler(new ChannelInitializer() { + @Override + protected void initChannel(Channel channel) { + serverChildChannels.add(channel); + connectionCount.incrementAndGet(); + channel.pipeline() + .addLast(new HttpServerCodec()) + .addLast(new HttpObjectAggregator(1024)) + .addLast(new CompressedResponseHandler()); + } + }) + .bind(0) + .sync() + .channel(); + serverPort = ((InetSocketAddress) serverChannel.localAddress()).getPort(); + } + + @AfterEach + public void stopServer() throws InterruptedException { + if (serverChildChannels != null) { + serverChildChannels.close().sync(); + } + if (serverChannel != null) { + serverChannel.close().sync(); + } + if (serverGroup != null) { + serverGroup.shutdownGracefully(0, 100, MILLISECONDS).sync(); + } + } + + @Test + public void suspensionStopsInflationUntilResumed() throws Exception { + assertDemandBoundedInflation(true); + } + + @Test + public void suspensionStopsInflationWithAutoReadDisabled() throws Exception { + assertDemandBoundedInflation(false); + } + + private void assertDemandBoundedInflation(boolean autoRead) throws Exception { + AtomicReference clientChannel = new AtomicReference<>(); + DefaultAsyncHttpClientConfig.Builder config = config() + .setMaxConnectionsPerHost(1) + .setRequestTimeout(Duration.ofSeconds(60)) + .setHttpAdditionalChannelInitializer(clientChannel::set); + if (!autoRead) { + config.addChannelOption(ChannelOption.AUTO_READ, false); + } + try (AsyncHttpClient client = asyncHttpClient(config)) { + SuspendingHandler handler = new SuspendingHandler(); + ListenableFuture response = client.prepareGet(url("/zeros")).execute(handler); + + int first = handler.parts.poll(10, SECONDS); + awaitEventLoop(clientChannel.get()); + assertTrue(first > 0 && first <= MAX_PART_SIZE, "unexpected first part size " + first); + assertEquals(first, handler.bytes.get(), "a suspended response inflated " + (handler.bytes.get() - first) + + " more bytes in " + (handler.partCount.get() - 1) + " parts from input it had already read"); + + int resumes = 0; + while (!handler.last.get()) { + long before = handler.bytes.get(); + handler.control.resume(); + resumes++; + int part = handler.parts.poll(10, SECONDS); + awaitEventLoop(clientChannel.get()); + Integer terminal = handler.parts.poll(); + if (terminal != null) { + assertEquals(0, terminal, "only the empty terminal part may follow a data part"); + assertNull(handler.parts.poll()); + } + assertTrue(part <= MAX_PART_SIZE, "a resume produced a part of " + part + " bytes"); + assertEquals(before + part, handler.bytes.get(), "one resume must produce at most one part"); + } + + response.get(10, SECONDS); + assertNull(handler.failure.get()); + assertEquals(ZEROS_SIZE, handler.bytes.get()); + assertFalse(handler.corrupt.get(), "the inflated body must be all zeros"); + assertTrue(resumes > ZEROS_SIZE / MAX_PART_SIZE / 2, "suspension must hold back most parts: " + resumes); + + Response pooled = client.prepareGet(url("/empty")).execute().get(10, SECONDS); + assertEquals(200, pooled.getStatusCode()); + assertEquals(1, connectionCount.get(), "a fully drained connection must be reused"); + } + } + + @Test + public void chunkedBodyAndTrailersAreDeliveredIntactAcrossSuspensions() throws Exception { + try (AsyncHttpClient client = asyncHttpClient(config() + .setMaxConnectionsPerHost(1) + .setRequestTimeout(Duration.ofSeconds(60)))) { + SuspendingHandler handler = new SuspendingHandler(); + handler.keepBody = true; + ListenableFuture response = client.prepareGet(url("/text-chunked")).execute(handler); + + while (!handler.last.get()) { + assertNotNull(handler.parts.poll(10, SECONDS)); + handler.control.resume(); + } + response.get(10, SECONDS); + + assertNull(handler.failure.get()); + assertArrayEquals(TEXT, handler.body.toByteArray()); + assertEquals("done", handler.trailer.get(), "trailers must arrive after the whole body"); + + Response pooled = client.prepareGet(url("/empty")).execute().get(10, SECONDS); + assertEquals(200, pooled.getStatusCode()); + assertEquals(1, connectionCount.get(), + "a fully drained chunked response must leave the connection reusable"); + } + } + + @Test + public void bodyEndedByConnectionCloseIsDeliveredIntactAcrossSuspensions() throws Exception { + try (AsyncHttpClient client = asyncHttpClient(config().setRequestTimeout(Duration.ofSeconds(60)))) { + SuspendingHandler handler = new SuspendingHandler(); + handler.keepBody = true; + ListenableFuture response = client.prepareGet(url("/text-close")).execute(handler); + + while (!handler.last.get()) { + assertNotNull(handler.parts.poll(10, SECONDS)); + handler.control.resume(); + } + response.get(10, SECONDS); + + assertNull(handler.failure.get()); + assertArrayEquals(TEXT, handler.body.toByteArray()); + } + } + + @Test + public void cancellationWhileCompressedInputIsBufferedClosesTheConnection() throws Exception { + AtomicReference clientChannel = new AtomicReference<>(); + try (AsyncHttpClient client = asyncHttpClient(config() + .setMaxConnectionsPerHost(1) + .setRequestTimeout(Duration.ofSeconds(60)) + .setHttpAdditionalChannelInitializer(clientChannel::set))) { + SuspendingHandler handler = new SuspendingHandler(); + ListenableFuture response = client.prepareGet(url("/zeros")).execute(handler); + + assertNotNull(handler.parts.poll(10, SECONDS)); + handler.control.cancel(); + response.get(10, SECONDS); + clientChannel.get().closeFuture().await(10, SECONDS); + assertFalse(clientChannel.get().isOpen(), "an unfinished response cannot leave its connection reusable"); + assertTrue(handler.bytes.get() < ZEROS_SIZE); + + Response next = client.prepareGet(url("/empty")).execute().get(10, SECONDS); + assertEquals(200, next.getStatusCode()); + assertEquals(2, connectionCount.get()); + } + } + + @Test + public void abortWhileCompressedInputIsBufferedClosesTheConnection() throws Exception { + AtomicReference clientChannel = new AtomicReference<>(); + try (AsyncHttpClient client = asyncHttpClient(config() + .setMaxConnectionsPerHost(1) + .setRequestTimeout(Duration.ofSeconds(60)) + .setHttpAdditionalChannelInitializer(clientChannel::set))) { + SuspendingHandler handler = new SuspendingHandler(); + handler.abortAfterFirstPart = true; + client.prepareGet(url("/zeros")).execute(handler).get(10, SECONDS); + + clientChannel.get().closeFuture().await(10, SECONDS); + assertFalse(clientChannel.get().isOpen()); + assertEquals(1, handler.partCount.get(), "no part may be delivered after ABORT"); + } + } + + @Test + public void compressedResponsesWithoutBodyComplete() throws Exception { + try (AsyncHttpClient client = asyncHttpClient(config().setMaxConnectionsPerHost(1))) { + Response head = client.prepareHead(url("/zeros")).execute().get(10, SECONDS); + assertEquals(200, head.getStatusCode()); + assertEquals(0, head.getResponseBodyAsBytes().length); + + Response empty = client.prepareGet(url("/empty")).execute().get(10, SECONDS); + assertEquals(200, empty.getStatusCode()); + assertEquals(0, empty.getResponseBodyAsBytes().length); + assertNull(empty.getHeader(HttpHeaderNames.CONTENT_ENCODING)); + + assertEquals(1, connectionCount.get()); + } + } + + private String url(String path) { + return "http://localhost:" + serverPort + path; + } + + private static void awaitEventLoop(Channel channel) throws InterruptedException { + channel.eventLoop().submit(() -> { + }).sync(); + } + + /** + * Suspends inline from every body callback, the way a demand-driven consumer does, and records what was + * delivered. + */ + private static final class SuspendingHandler implements AsyncHandler { + + final LinkedBlockingQueue parts = new LinkedBlockingQueue<>(); + final AtomicLong bytes = new AtomicLong(); + final AtomicInteger partCount = new AtomicInteger(); + final AtomicBoolean corrupt = new AtomicBoolean(); + final AtomicBoolean last = new AtomicBoolean(); + final AtomicReference trailer = new AtomicReference<>(); + final AtomicReference failure = new AtomicReference<>(); + final ByteArrayOutputStream body = new ByteArrayOutputStream(); + volatile ResponseBodyControl control; + volatile boolean keepBody; + volatile boolean abortAfterFirstPart; + + @Override + public State onStatusReceived(HttpResponseStatus responseStatus) { + return State.CONTINUE; + } + + @Override + public State onHeadersReceived(HttpHeaders headers) { + return State.CONTINUE; + } + + @Override + public State onResponseBodyStart(ResponseBodyControl control) { + this.control = control; + return State.CONTINUE; + } + + @Override + public State onBodyPartReceived(HttpResponseBodyPart bodyPart) { + control.suspend(); + byte[] bytes = bodyPart.getBodyPartBytes(); + if (keepBody) { + body.write(bytes, 0, bytes.length); + } else { + for (byte b : bytes) { + if (b != 0) { + corrupt.set(true); + break; + } + } + } + this.bytes.addAndGet(bytes.length); + partCount.incrementAndGet(); + if (bodyPart.isLast()) { + last.set(true); + } + parts.add(bytes.length); + return abortAfterFirstPart ? State.ABORT : State.CONTINUE; + } + + @Override + public State onTrailingHeadersReceived(HttpHeaders headers) { + trailer.set(last.get() || bytes.get() != TEXT.length ? "early" : headers.get("x-trailer")); + return State.CONTINUE; + } + + @Override + public void onThrowable(Throwable t) { + failure.set(t); + } + + @Override + public Void onCompleted() { + return null; + } + } + + private static final class CompressedResponseHandler extends SimpleChannelInboundHandler { + + @Override + protected void channelRead0(ChannelHandlerContext ctx, FullHttpRequest request) { + switch (request.uri()) { + case "/zeros": { + boolean head = HttpMethod.HEAD.equals(request.method()); + DefaultFullHttpResponse response = new DefaultFullHttpResponse(HTTP_1_1, OK, + head ? Unpooled.EMPTY_BUFFER : Unpooled.wrappedBuffer(ZEROS_GZIP)); + response.headers().set(HttpHeaderNames.CONTENT_ENCODING, HttpHeaderValues.GZIP); + HttpUtil.setContentLength(response, ZEROS_GZIP.length); + ctx.writeAndFlush(response); + break; + } + case "/empty": { + DefaultFullHttpResponse response = new DefaultFullHttpResponse(HTTP_1_1, OK, Unpooled.EMPTY_BUFFER); + response.headers().set(HttpHeaderNames.CONTENT_ENCODING, HttpHeaderValues.GZIP); + HttpUtil.setContentLength(response, 0); + ctx.writeAndFlush(response); + break; + } + case "/text-chunked": { + HttpResponse response = new DefaultHttpResponse(HTTP_1_1, OK); + response.headers().set(HttpHeaderNames.CONTENT_ENCODING, HttpHeaderValues.GZIP); + HttpUtil.setTransferEncodingChunked(response, true); + ctx.write(response); + for (int offset = 0; offset < TEXT_GZIP.length; offset += 10_000) { + int length = Math.min(10_000, TEXT_GZIP.length - offset); + ctx.write(new DefaultHttpContent(Unpooled.wrappedBuffer(TEXT_GZIP, offset, length))); + } + LastHttpContent last = new DefaultLastHttpContent(); + last.trailingHeaders().set("x-trailer", "done"); + ctx.writeAndFlush(last); + break; + } + case "/text-close": { + HttpResponse response = new DefaultHttpResponse(HTTP_1_1, OK); + response.headers().set(HttpHeaderNames.CONTENT_ENCODING, HttpHeaderValues.GZIP); + response.headers().set(HttpHeaderNames.CONNECTION, HttpHeaderValues.CLOSE); + ctx.write(response); + ctx.writeAndFlush(new DefaultHttpContent(Unpooled.wrappedBuffer(TEXT_GZIP))) + .addListener(ChannelFutureListener.CLOSE); + break; + } + default: + throw new IllegalArgumentException(request.uri()); + } + } + } +} diff --git a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorLargeOutputTest.java b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorLargeOutputTest.java new file mode 100644 index 000000000..1161808a7 --- /dev/null +++ b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorLargeOutputTest.java @@ -0,0 +1,238 @@ +/* + * Copyright (c) 2026 AsyncHttpClient Project. 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. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License 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 org.asynchttpclient.netty.handler; + +import com.github.luben.zstd.Zstd; +import io.netty.buffer.ByteBuf; +import io.netty.buffer.Unpooled; +import io.netty.buffer.AbstractByteBufAllocator; +import io.netty.buffer.UnpooledDirectByteBuf; +import io.netty.buffer.UnpooledHeapByteBuf; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelInboundHandlerAdapter; +import io.netty.channel.embedded.EmbeddedChannel; +import io.netty.handler.codec.compression.Decompressor; +import io.netty.handler.codec.compression.ZstdDecompressor; +import io.netty.handler.codec.http.DefaultHttpResponse; +import io.netty.handler.codec.http.DefaultLastHttpContent; +import io.netty.handler.codec.http.HttpContent; +import io.netty.handler.codec.http.HttpHeaderNames; +import io.netty.handler.codec.http.HttpResponseStatus; +import io.netty.handler.codec.http.LastHttpContent; +import io.netty.util.ReferenceCountUtil; +import org.asynchttpclient.AsyncCompletionHandlerBase; +import org.asynchttpclient.RequestBuilder; +import org.asynchttpclient.netty.NettyResponseBodyControl; +import org.asynchttpclient.netty.NettyResponseFuture; +import org.asynchttpclient.netty.channel.Channels; +import org.junit.jupiter.api.Test; + +import java.io.ByteArrayOutputStream; +import java.nio.charset.StandardCharsets; +import java.nio.file.Files; +import java.nio.file.Path; +import java.nio.file.Paths; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; + +import static io.netty.handler.codec.http.HttpVersion.HTTP_1_1; +import static java.util.concurrent.TimeUnit.SECONDS; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Netty's zstd decompressor hands out buffers of {@code io.netty.compression.defaultMaxForwardBytes}, which Netty + * reads once per JVM. This test raises it in a fresh JVM and checks that {@link Http1ContentDecompressor} still hands + * on parts of at most {@link Http1ContentDecompressor#MAX_PART_SIZE} bytes, one per resume, and releases what it + * holds. + */ +public class Http1ContentDecompressorLargeOutputTest { + + private static final int FORWARD_BYTES = 4 * Http1ContentDecompressor.MAX_PART_SIZE; + + @Test + public void splitsDecompressorOutputLargerThanOnePart() throws Exception { + String java = Paths.get(System.getProperty("java.home"), "bin", "java").toString(); + String classPath = System.getProperty("surefire.test.class.path", System.getProperty("java.class.path")); + // The child writes to a file rather than a pipe, so that a child that hangs cannot block the timed wait. + Path log = Files.createTempFile("http1-content-decompressor-large-output", ".log"); + Process process = null; + try { + process = new ProcessBuilder(java, "-Dio.netty.compression.defaultMaxForwardBytes=" + FORWARD_BYTES, + "-cp", classPath, Check.class.getName()) + .redirectErrorStream(true) + .redirectOutput(log.toFile()) + .start(); + boolean exited = process.waitFor(60, SECONDS); + String output = Files.readString(log, StandardCharsets.UTF_8); + assertTrue(exited, "the check did not finish within 60 seconds:\n" + output); + assertEquals(0, process.exitValue(), output); + assertTrue(output.contains("split output of "), output); + } finally { + if (process != null) { + process.destroyForcibly(); + process.waitFor(10, SECONDS); + } + Files.deleteIfExists(log); + } + } + + /** Runs in the fresh JVM; any failure exits with a non-zero status. */ + public static final class Check { + + public static void main(String[] args) { + byte[] payload = "zstd hands out large buffers\n".repeat(40_000).getBytes(StandardCharsets.US_ASCII); + byte[] compressed = Zstd.compress(payload); + TrackingAllocator allocator = new TrackingAllocator(); + + // Without this, the check below would pass trivially. + int rawOutput = largestRawOutput(compressed, allocator); + check(rawOutput > Http1ContentDecompressor.MAX_PART_SIZE, "the zstd output is not larger than one part: " + + rawOutput); + + Exchange exchange = new Exchange(allocator, compressed); + for (int resumes = 0; !exchange.last; resumes++) { + check(resumes < 1_000, "the body must end"); + int before = exchange.parts; + exchange.control.resume(); + check(exchange.parts == before + 1 || exchange.last && exchange.parts == before, + "each resume must deliver exactly one more part"); + } + check(exchange.maxPart <= Http1ContentDecompressor.MAX_PART_SIZE, "a part of " + exchange.maxPart + + " bytes was handed on"); + check(Arrays.equals(payload, exchange.bytes.toByteArray()), "the body differs"); + check(!exchange.channel.finishAndReleaseAll(), "messages were left over"); + allocator.checkAllReleased(); + + // Removing the handler while it carries over the rest of a split output releases that rest. + Exchange removed = new Exchange(allocator, compressed); + removed.control.resume(); + check(removed.parts == 1, "one part before removal, got " + removed.parts); + removed.channel.pipeline().remove(Http1ContentDecompressor.class); + check(!removed.channel.finishAndReleaseAll(), "messages were left over"); + allocator.checkAllReleased(); + + System.out.println("split output of " + rawOutput + " bytes into parts of at most " + exchange.maxPart); + } + + private static int largestRawOutput(byte[] compressed, TrackingAllocator allocator) { + Decompressor decompressor = ZstdDecompressor.builder().build(allocator); + int largest = 0; + try { + check(decompressor.status() == Decompressor.Status.NEED_INPUT, "a new decompressor needs input"); + decompressor.addInput(Unpooled.wrappedBuffer(compressed)); + while (decompressor.status() == Decompressor.Status.NEED_OUTPUT) { + ByteBuf output = decompressor.takeOutput(); + largest = Math.max(largest, output.readableBytes()); + output.release(); + } + } finally { + decompressor.close(); + } + return largest; + } + + private static void check(boolean condition, String message) { + if (!condition) { + throw new AssertionError(message); + } + } + } + + /** A zstd response, suspended from the start and again after every body part. */ + private static final class Exchange { + final EmbeddedChannel channel = new EmbeddedChannel(); + final NettyResponseBodyControl control; + final ByteArrayOutputStream bytes = new ByteArrayOutputStream(); + int parts; + int maxPart; + boolean last; + + Exchange(TrackingAllocator allocator, byte[] compressed) { + channel.config().setAllocator(allocator); + NettyResponseFuture future = new NettyResponseFuture<>( + new RequestBuilder().setUrl("http://localhost/").build(), new AsyncCompletionHandlerBase(), + null, 0, null, null, null); + Channels.setAttribute(channel, future); + control = NettyResponseBodyControl.create(future, channel, () -> { + }, ignored -> { + }); + channel.pipeline().addLast(new Http1ContentDecompressor(false, 0), new ChannelInboundHandlerAdapter() { + @Override + public void channelRead(ChannelHandlerContext ctx, Object msg) { + try { + if (msg instanceof HttpContent) { + ByteBuf content = ((HttpContent) msg).content(); + if (content.isReadable()) { + parts++; + maxPart = Math.max(maxPart, content.readableBytes()); + byte[] data = new byte[content.readableBytes()]; + content.readBytes(data); + bytes.writeBytes(data); + } + last = msg instanceof LastHttpContent; + control.suspend(); + } + } finally { + ReferenceCountUtil.release(msg); + } + } + }); + DefaultHttpResponse response = new DefaultHttpResponse(HTTP_1_1, HttpResponseStatus.OK); + response.headers().set(HttpHeaderNames.CONTENT_ENCODING, "zstd"); + response.headers().set(HttpHeaderNames.CONTENT_LENGTH, compressed.length); + control.suspend(); + channel.writeInbound(response); + channel.writeInbound(new DefaultLastHttpContent(Unpooled.wrappedBuffer(compressed))); + } + } + + /** Records every buffer it allocates, so that leaks of the decompressor's output can be detected. */ + private static final class TrackingAllocator extends AbstractByteBufAllocator { + private final List allocated = new ArrayList<>(); + + TrackingAllocator() { + super(false); + } + + @Override + protected ByteBuf newHeapBuffer(int initialCapacity, int maxCapacity) { + ByteBuf buf = new UnpooledHeapByteBuf(this, initialCapacity, maxCapacity); + allocated.add(buf); + return buf; + } + + @Override + protected ByteBuf newDirectBuffer(int initialCapacity, int maxCapacity) { + ByteBuf buf = new UnpooledDirectByteBuf(this, initialCapacity, maxCapacity); + allocated.add(buf); + return buf; + } + + @Override + public boolean isDirectBufferPooled() { + return false; + } + + void checkAllReleased() { + for (ByteBuf buf : allocated) { + Check.check(buf.refCnt() == 0, "a buffer of " + buf.capacity() + " bytes was not released"); + } + allocated.clear(); + } + } +} diff --git a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java new file mode 100644 index 000000000..5f2c57345 --- /dev/null +++ b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java @@ -0,0 +1,572 @@ +/* + * Copyright (c) 2026 AsyncHttpClient Project. 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. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License 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 org.asynchttpclient.netty.handler; + +import com.aayushatharva.brotli4j.encoder.Encoder; +import com.github.luben.zstd.Zstd; +import io.netty.buffer.ByteBuf; +import io.netty.buffer.Unpooled; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelInboundHandlerAdapter; +import io.netty.channel.embedded.EmbeddedChannel; +import io.netty.handler.codec.DecoderResult; +import io.netty.handler.codec.compression.Brotli; +import io.netty.handler.codec.compression.DecompressionException; +import io.netty.handler.codec.compression.SnappyFrameEncoder; +import io.netty.handler.codec.http.DefaultHttpContent; +import io.netty.handler.codec.http.DefaultHttpResponse; +import io.netty.handler.codec.http.DefaultLastHttpContent; +import io.netty.handler.codec.http.HttpContent; +import io.netty.handler.codec.http.HttpHeaderNames; +import io.netty.handler.codec.http.HttpResponse; +import io.netty.handler.codec.http.HttpResponseStatus; +import io.netty.handler.codec.http.LastHttpContent; +import io.netty.util.ReferenceCountUtil; +import org.asynchttpclient.AsyncCompletionHandlerBase; +import org.asynchttpclient.RequestBuilder; +import org.asynchttpclient.netty.NettyResponseBodyControl; +import org.asynchttpclient.netty.NettyResponseFuture; +import org.asynchttpclient.netty.channel.Channels; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import java.io.ByteArrayOutputStream; +import java.io.IOException; +import java.io.OutputStream; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.List; +import java.util.Random; +import java.util.zip.Deflater; +import java.util.zip.DeflaterOutputStream; +import java.util.zip.GZIPOutputStream; + +import static io.netty.handler.codec.http.HttpVersion.HTTP_1_1; +import static org.junit.jupiter.api.Assertions.assertArrayEquals; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * Unit tests for {@link Http1ContentDecompressor} on an {@link EmbeddedChannel}, so suspension, channel closure + * and handler removal can be driven deterministically and reference counts checked exactly. + */ +public class Http1ContentDecompressorTest { + + private static final byte[] TEXT = "the quick brown fox jumps over the lazy dog\n".repeat(2_000) + .getBytes(StandardCharsets.US_ASCII); + private static final byte[] OTHER_TEXT = "pack my box with five dozen liquor jugs\n".repeat(2_000) + .getBytes(StandardCharsets.US_ASCII); + private static final int ZEROS_SIZE = 8 * 1024 * 1024; + private static final byte[] ZEROS_GZIP = gzip(new byte[ZEROS_SIZE]); + + private static byte[] gzip(byte[] payload) { + return compress(payload, GZIPOutputStream::new); + } + + private interface Compressor { + OutputStream wrap(OutputStream out) throws IOException; + } + + private static byte[] compress(byte[] payload, Compressor compressor) { + ByteArrayOutputStream compressed = new ByteArrayOutputStream(); + try (OutputStream out = compressor.wrap(compressed)) { + out.write(payload); + } catch (IOException e) { + throw new IllegalStateException(e); + } + return compressed.toByteArray(); + } + + private static byte[] concat(byte[]... arrays) { + ByteArrayOutputStream out = new ByteArrayOutputStream(); + for (byte[] array : arrays) { + out.writeBytes(array); + } + return out.toByteArray(); + } + + private static byte[] compress(String contentEncoding, byte[] payload) { + switch (contentEncoding) { + case "gzip": + return gzip(payload); + case "deflate": + return compress(payload, DeflaterOutputStream::new); + case "x-deflate": + return compress(payload, + out -> new DeflaterOutputStream(out, new Deflater(Deflater.DEFAULT_COMPRESSION, true))); + case "snappy": + return snappy(payload); + case "zstd": + return Zstd.compress(payload); + case "br": + try { + Brotli.ensureAvailability(); + return Encoder.compress(payload); + } catch (Throwable e) { + throw new IllegalStateException(e); + } + default: + throw new IllegalArgumentException(contentEncoding); + } + } + + private static byte[] snappy(byte[] payload) { + EmbeddedChannel encoder = new EmbeddedChannel(new SnappyFrameEncoder()); + encoder.writeOutbound(Unpooled.wrappedBuffer(payload)); + encoder.finish(); + ByteArrayOutputStream compressed = new ByteArrayOutputStream(); + ByteBuf buf; + while ((buf = encoder.readOutbound()) != null) { + compressed.writeBytes(toBytes(buf)); + buf.release(); + } + return compressed.toByteArray(); + } + + private static byte[] toBytes(ByteBuf buf) { + byte[] bytes = new byte[buf.readableBytes()]; + buf.getBytes(buf.readerIndex(), bytes); + return bytes; + } + + private static HttpResponse response(String contentEncoding, int contentLength) { + HttpResponse response = new DefaultHttpResponse(HTTP_1_1, HttpResponseStatus.OK); + if (contentEncoding != null) { + response.headers().set(HttpHeaderNames.CONTENT_ENCODING, contentEncoding); + } + response.headers().set(HttpHeaderNames.CONTENT_LENGTH, contentLength); + return response; + } + + /** Collects what the decompressor hands on, releasing each message. */ + private static final class Body { + final ByteArrayOutputStream bytes = new ByteArrayOutputStream(); + final List messages = new ArrayList<>(); + int parts; + int maxPart; + boolean last; + + Body drain(EmbeddedChannel channel) { + Object msg; + while ((msg = channel.readInbound()) != null) { + assertFalse(last, "nothing may follow the LastHttpContent"); + if (msg instanceof HttpContent) { + ByteBuf content = ((HttpContent) msg).content(); + if (content.isReadable()) { + parts++; + maxPart = Math.max(maxPart, content.readableBytes()); + bytes.writeBytes(toBytes(content)); + } + last = msg instanceof LastHttpContent; + } + messages.add(msg instanceof HttpContent ? ((HttpContent) msg).getClass() : msg); + ReferenceCountUtil.release(msg); + } + return this; + } + } + + /** Writes the body in fragments of {@code fragmentSize} bytes, as separate reads would deliver it. */ + private static Body decompressFragmented(String contentEncoding, byte[] compressed, int fragmentSize) { + EmbeddedChannel channel = new EmbeddedChannel(new Http1ContentDecompressor(false, 0)); + channel.writeInbound(response(contentEncoding, compressed.length)); + for (int offset = 0; offset < compressed.length; offset += fragmentSize) { + int length = Math.min(fragmentSize, compressed.length - offset); + channel.writeInbound(new DefaultHttpContent(Unpooled.wrappedBuffer(compressed, offset, length))); + } + channel.writeInbound(LastHttpContent.EMPTY_LAST_CONTENT); + Body body = new Body().drain(channel); + assertFalse(channel.finishAndReleaseAll()); + return body; + } + + /** Asserts that decompressing {@code compressed} fails, and returns the failure. */ + private static DecompressionException decompressionFailure(String contentEncoding, byte[] compressed, + long maxDecompressedBytes) { + EmbeddedChannel channel = new EmbeddedChannel(new Http1ContentDecompressor(false, maxDecompressedBytes)); + channel.writeInbound(response(contentEncoding, compressed.length)); + ByteBuf input = Unpooled.wrappedBuffer(compressed); + DecompressionException failure = assertThrows(DecompressionException.class, + () -> channel.writeInbound(new DefaultLastHttpContent(input))); + assertEquals(0, input.refCnt()); + channel.finishAndReleaseAll(); + return failure; + } + + private static Body decompress(String contentEncoding, byte[] compressed, boolean keepEncodingHeader) { + EmbeddedChannel channel = new EmbeddedChannel(new Http1ContentDecompressor(keepEncodingHeader, 0)); + channel.writeInbound(response(contentEncoding, compressed.length)); + channel.writeInbound(new DefaultLastHttpContent(Unpooled.wrappedBuffer(compressed))); + Body body = new Body().drain(channel); + assertFalse(channel.finishAndReleaseAll()); + return body; + } + + @Test + public void decompressesEveryContentCodingOfNettysDecompressor() throws Throwable { + assertArrayEquals(TEXT, decompress("gzip", gzip(TEXT), false).bytes.toByteArray()); + assertArrayEquals(TEXT, decompress("X-GZIP", gzip(TEXT), false).bytes.toByteArray()); + assertArrayEquals(TEXT, decompress("deflate", compress(TEXT, DeflaterOutputStream::new), false) + .bytes.toByteArray()); + assertArrayEquals(TEXT, decompress("x-deflate", compress(TEXT, + out -> new DeflaterOutputStream(out, new Deflater(Deflater.DEFAULT_COMPRESSION, true))), false) + .bytes.toByteArray(), "deflate without the zlib wrapper is accepted as Netty accepts it"); + assertArrayEquals(TEXT, decompress("snappy", snappy(TEXT), false).bytes.toByteArray()); + assertArrayEquals(TEXT, decompress("zstd", Zstd.compress(TEXT), false).bytes.toByteArray()); + Brotli.ensureAvailability(); + assertArrayEquals(TEXT, decompress("br", Encoder.compress(TEXT), false).bytes.toByteArray()); + } + + @Test + public void rewritesHeadersLikeNettysDecompressor() { + EmbeddedChannel channel = new EmbeddedChannel(new Http1ContentDecompressor(false, 0)); + channel.writeInbound(response("GZip", 10)); + HttpResponse decoded = channel.readInbound(); + assertNull(decoded.headers().get(HttpHeaderNames.CONTENT_ENCODING)); + assertNull(decoded.headers().get(HttpHeaderNames.CONTENT_LENGTH)); + assertEquals("chunked", decoded.headers().get(HttpHeaderNames.TRANSFER_ENCODING)); + assertFalse(channel.finishAndReleaseAll()); + + channel = new EmbeddedChannel(new Http1ContentDecompressor(true, 0)); + channel.writeInbound(response("GZip", 10)); + decoded = channel.readInbound(); + assertEquals("GZip", decoded.headers().get(HttpHeaderNames.CONTENT_ENCODING)); + assertNull(decoded.headers().get(HttpHeaderNames.CONTENT_LENGTH)); + assertFalse(channel.finishAndReleaseAll()); + } + + @Test + public void passesUnsupportedCodingsThroughUntouched() { + EmbeddedChannel channel = new EmbeddedChannel(new Http1ContentDecompressor(false, 0)); + HttpResponse response = response("compress", 3); + channel.writeInbound(response); + ByteBuf content = Unpooled.wrappedBuffer(new byte[]{1, 2, 3}); + LastHttpContent last = new DefaultLastHttpContent(content); + channel.writeInbound(last); + assertSame(response, channel.readInbound()); + assertEquals("compress", response.headers().get(HttpHeaderNames.CONTENT_ENCODING)); + assertEquals("3", response.headers().get(HttpHeaderNames.CONTENT_LENGTH)); + assertSame(last, channel.readInbound()); + last.release(); + assertFalse(channel.finishAndReleaseAll()); + } + + @Test + public void bodilessCompressedResponsesEndNormally() { + // HEAD, 204 and 304 responses routinely carry Content-Encoding without a body. + Body body = decompress("gzip", new byte[0], false); + assertEquals(0, body.parts); + assertTrue(body.last); + } + + @Test + public void trailersAndDecoderFailuresFollowTheBody() { + EmbeddedChannel channel = new EmbeddedChannel(new Http1ContentDecompressor(false, 0)); + channel.writeInbound(response("gzip", 0)); + byte[] gzip = gzip(TEXT); + channel.writeInbound(new DefaultHttpContent(Unpooled.wrappedBuffer(gzip, 0, gzip.length / 2))); + LastHttpContent last = new DefaultLastHttpContent( + Unpooled.wrappedBuffer(gzip, gzip.length / 2, gzip.length - gzip.length / 2)); + last.trailingHeaders().set("x-trailer", "done"); + IllegalStateException failure = new IllegalStateException("bad chunk"); + last.setDecoderResult(DecoderResult.failure(failure)); + channel.writeInbound(last); + + assertInstanceOf(HttpResponse.class, channel.readInbound()); + ByteArrayOutputStream bytes = new ByteArrayOutputStream(); + Object msg; + LastHttpContent trailer = null; + while ((msg = channel.readInbound()) != null) { + assertNull(trailer, "the trailer must come last"); + HttpContent content = (HttpContent) msg; + bytes.writeBytes(toBytes(content.content())); + if (msg instanceof LastHttpContent) { + trailer = (LastHttpContent) msg; + } + content.release(); + } + assertArrayEquals(TEXT, bytes.toByteArray()); + assertEquals("done", trailer.trailingHeaders().get("x-trailer")); + assertSame(failure, trailer.decoderResult().cause(), "a framing failure must not be turned into success"); + assertFalse(channel.finishAndReleaseAll()); + } + + @Test + public void continueResponsesPassThroughUntouched() { + EmbeddedChannel channel = new EmbeddedChannel(new Http1ContentDecompressor(false, 0)); + HttpResponse interim = new DefaultHttpResponse(HTTP_1_1, HttpResponseStatus.CONTINUE); + interim.headers().set(HttpHeaderNames.CONTENT_ENCODING, "gzip"); + channel.writeInbound(interim); + channel.writeInbound(LastHttpContent.EMPTY_LAST_CONTENT); + assertSame(interim, channel.readInbound()); + assertEquals("gzip", interim.headers().get(HttpHeaderNames.CONTENT_ENCODING)); + assertSame(LastHttpContent.EMPTY_LAST_CONTENT, channel.readInbound()); + + channel.writeInbound(response("gzip", 0)); + channel.writeInbound(new DefaultLastHttpContent(Unpooled.wrappedBuffer(gzip(TEXT)))); + assertInstanceOf(HttpResponse.class, channel.readInbound()); + assertArrayEquals(TEXT, new Body().drain(channel).bytes.toByteArray()); + assertFalse(channel.finishAndReleaseAll()); + } + + /** A channel whose response is suspended by the handler after every body part, as a demand-driven consumer does. */ + private static final class SuspendedExchange { + final EmbeddedChannel channel; + final NettyResponseBodyControl control; + final List input = new ArrayList<>(); + final ByteBuf compressed; + final boolean suspendedBeforeBody; + + SuspendedExchange(long maxDecompressedBytes) { + this(maxDecompressedBytes, false); + } + + SuspendedExchange(long maxDecompressedBytes, boolean suspendBeforeBody) { + this("gzip", ZEROS_GZIP, ZEROS_GZIP.length, maxDecompressedBytes, suspendBeforeBody, true); + } + + /** Writes the body in fragments of {@code fragmentSize} bytes, the last one as the {@link LastHttpContent}. */ + SuspendedExchange(String contentEncoding, byte[] body, int fragmentSize, long maxDecompressedBytes, + boolean suspendBeforeBody, boolean autoRead) { + suspendedBeforeBody = suspendBeforeBody; + channel = new EmbeddedChannel(); + channel.config().setAutoRead(autoRead); + NettyResponseFuture future = new NettyResponseFuture<>( + new RequestBuilder().setUrl("http://localhost/").build(), new AsyncCompletionHandlerBase(), + null, 0, null, null, null); + Channels.setAttribute(channel, future); + control = NettyResponseBodyControl.create(future, channel, () -> { + }, ignored -> { + }); + channel.pipeline().addLast(new Http1ContentDecompressor(false, maxDecompressedBytes), + new ChannelInboundHandlerAdapter() { + @Override + public void channelRead(ChannelHandlerContext ctx, Object msg) { + if (msg instanceof HttpContent) { + control.suspend(); + } + ctx.fireChannelRead(msg); + } + }); + channel.writeInbound(response(contentEncoding, body.length)); + ReferenceCountUtil.release(channel.readInbound()); + if (suspendBeforeBody) { + control.suspend(); + } + for (int offset = 0; offset < body.length; offset += fragmentSize) { + int length = Math.min(fragmentSize, body.length - offset); + ByteBuf fragment = Unpooled.wrappedBuffer(body, offset, length); + input.add(fragment); + channel.writeInbound(offset + length < body.length ? new DefaultHttpContent(fragment) + : new DefaultLastHttpContent(fragment)); + } + compressed = input.get(0); + } + + /** Resumes until the body has ended, checking that each resume delivers exactly one more part. */ + Body resumeToTheEnd() { + Body body = new Body().drain(channel); + assertEquals(suspendedBeforeBody ? 0 : 1, body.parts, "nothing more may arrive while suspended"); + for (int resumes = 0; !body.last; resumes++) { + assertTrue(resumes < 10_000, "the body must end"); + int before = body.parts; + control.resume(); + body.drain(channel); + assertTrue(body.parts == before + 1 || body.last && body.parts == before, + "each resume must deliver exactly one more part"); + } + assertTrue(body.maxPart <= Http1ContentDecompressor.MAX_PART_SIZE); + for (ByteBuf fragment : input) { + assertEquals(0, fragment.refCnt()); + } + assertFalse(channel.finishAndReleaseAll()); + return body; + } + } + + @Test + public void suspensionBeforeTheBodyHoldsEveryPartBack() { + SuspendedExchange exchange = new SuspendedExchange(0, true); + Body body = new Body().drain(exchange.channel); + assertEquals(0, body.parts, "a response suspended from onResponseBodyStart must not receive any part"); + + exchange.control.resume(); + body.drain(exchange.channel); + assertEquals(1, body.parts); + exchange.channel.pipeline().remove(Http1ContentDecompressor.class); + assertEquals(0, exchange.compressed.refCnt()); + assertFalse(exchange.channel.finishAndReleaseAll()); + } + + @Test + public void suspensionHoldsInputUntilReadIsRequested() { + SuspendedExchange exchange = new SuspendedExchange(0); + Body body = new Body().drain(exchange.channel); + assertEquals(1, body.parts, "decompression must stop at the part that suspended the response"); + + int resumes = 0; + while (!body.last) { + int before = body.parts; + exchange.control.resume(); + resumes++; + body.drain(exchange.channel); + assertEquals(before + 1, body.parts, "each read request must inflate exactly one more part"); + } + assertEquals(ZEROS_SIZE, body.bytes.size()); + assertTrue(body.maxPart <= Http1ContentDecompressor.MAX_PART_SIZE); + assertEquals(resumes + 1, body.parts); + assertEquals(0, exchange.compressed.refCnt()); + assertFalse(exchange.channel.finishAndReleaseAll()); + } + + @Test + public void channelClosureHandsOnHeldInputBeforeChannelInactive() { + SuspendedExchange exchange = new SuspendedExchange(0); + List events = new ArrayList<>(); + exchange.channel.pipeline().addLast(new ChannelInboundHandlerAdapter() { + @Override + public void channelInactive(ChannelHandlerContext ctx) { + events.add("inactive"); + } + }); + Body body = new Body().drain(exchange.channel); + assertEquals(1, body.parts); + + exchange.channel.pipeline().fireChannelInactive(); + body.drain(exchange.channel); + assertTrue(body.last); + assertEquals(ZEROS_SIZE, body.bytes.size()); + assertEquals(List.of("inactive"), events); + assertEquals(0, exchange.compressed.refCnt()); + exchange.channel.finishAndReleaseAll(); + } + + @Test + public void removalReleasesHeldInput() { + SuspendedExchange exchange = new SuspendedExchange(0); + Body body = new Body().drain(exchange.channel); + assertEquals(1, body.parts); + assertEquals(1, exchange.compressed.refCnt(), "the compressed input is held while suspended"); + + exchange.channel.pipeline().remove(Http1ContentDecompressor.class); + assertEquals(0, exchange.compressed.refCnt()); + assertFalse(exchange.channel.finishAndReleaseAll()); + } + + @Test + public void limitExceededWhileResumingFailsTheExchangeAndReleasesInput() { + SuspendedExchange exchange = new SuspendedExchange(1024 * 1024); + Body body = new Body().drain(exchange.channel); + DecompressionException failure = null; + while (failure == null) { + assertFalse(body.last); + int before = body.parts; + exchange.control.resume(); + body.drain(exchange.channel); + try { + exchange.channel.checkException(); + } catch (DecompressionException e) { + failure = e; + } + assertTrue(failure != null || body.parts > before, "each read request must inflate more or fail"); + } + assertTrue(failure.getMessage().contains("maximum decompressed size of 1048576 bytes")); + assertTrue(body.bytes.size() <= 1024 * 1024); + assertEquals(0, exchange.compressed.refCnt()); + + // The rest of the failed body is dropped; only its end is still passed on. + long delivered = body.bytes.size(); + exchange.channel.read(); + body.drain(exchange.channel); + assertTrue(body.last); + assertEquals(delivered, body.bytes.size()); + assertFalse(exchange.channel.finishAndReleaseAll()); + } + + @Test + public void decompressesEveryMemberOfAConcatenatedGzipBody() { + byte[] compressed = concat(gzip(TEXT), gzip(OTHER_TEXT)); + byte[] expected = concat(TEXT, OTHER_TEXT); + assertArrayEquals(expected, decompress("gzip", compressed, false).bytes.toByteArray()); + assertArrayEquals(expected, decompressFragmented("x-gzip", compressed, 1).bytes.toByteArray()); + assertArrayEquals(expected, decompressFragmented("gzip", compressed, gzip(TEXT).length).bytes.toByteArray()); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + public void suspensionCarriesOverIntoTheNextGzipMember(boolean autoRead) { + byte[] first = new byte[300 * 1024]; + byte[] second = "x".repeat(300 * 1024).getBytes(StandardCharsets.US_ASCII); + SuspendedExchange exchange = new SuspendedExchange("gzip", concat(gzip(first), gzip(second)), 512, 0, true, + autoRead); + assertArrayEquals(concat(first, second), exchange.resumeToTheEnd().bytes.toByteArray()); + } + + @Test + public void corruptLaterGzipMemberFailsTheBody() { + // A member ends with the CRC-32 of its data and the data's length. Corrupt compressed data can instead look + // like a stream that ends early, which is accepted like any truncated body, as Netty's decoder accepts it. + byte[] second = gzip(OTHER_TEXT); + byte[] corruptChecksum = second.clone(); + corruptChecksum[second.length - 8] ^= (byte) 0xFF; + byte[] corruptLength = second.clone(); + corruptLength[second.length - 4] ^= (byte) 0xFF; + assertTrue(decompressionFailure("gzip", concat(gzip(TEXT), corruptChecksum), 0).getMessage() + .contains("CRC value mismatch")); + assertTrue(decompressionFailure("gzip", concat(gzip(TEXT), corruptLength), 0).getMessage() + .contains("Number of bytes mismatch")); + } + + @Test + public void bytesAfterTheCompressedStreamAreHandledAsNettysDecompressorHandlesThem() { + byte[] junk = "not compressed at all".getBytes(StandardCharsets.US_ASCII); + // After a gzip member, they would have to start another member. + decompressionFailure("gzip", concat(gzip(TEXT), junk), 0); + // Too few of them to make up a gzip member header, so the body ends before they are looked at. + assertArrayEquals(TEXT, decompress("gzip", concat(gzip(TEXT), new byte[]{1, 2, 3}), false) + .bytes.toByteArray()); + // A deflate stream cannot be continued, so they are dropped. + assertArrayEquals(TEXT, decompress("deflate", concat(compress("deflate", TEXT), junk), false) + .bytes.toByteArray()); + } + + @Test + public void limitCountsEveryGzipMember() { + DecompressionException failure = decompressionFailure("gzip", concat(gzip(TEXT), gzip(TEXT)), + TEXT.length + 1); + assertTrue(failure.getMessage().contains("maximum decompressed size of " + (TEXT.length + 1) + " bytes")); + } + + @ParameterizedTest + @ValueSource(booleans = {true, false}) + public void everyCodingDeliversBoundedPartsAcrossRepeatedSuspension(boolean autoRead) { + // Incompressible, highly compressible and text data, so that both large and small decompressor output occurs. + byte[] random = new byte[130_000]; + new Random(29).nextBytes(random); + byte[] payload = concat(random, new byte[200_000], TEXT); + for (String contentEncoding : List.of("gzip", "deflate", "x-deflate", "snappy", "zstd", "br")) { + SuspendedExchange exchange = new SuspendedExchange(contentEncoding, compress(contentEncoding, payload), + 4096, 0, true, autoRead); + assertArrayEquals(payload, exchange.resumeToTheEnd().bytes.toByteArray(), contentEncoding); + } + } +} diff --git a/pom.xml b/pom.xml index f34f516d9..c019bb814 100644 --- a/pom.xml +++ b/pom.xml @@ -526,6 +526,69 @@ "new": "method void org.asynchttpclient.netty.handler.Http2Handler::userEventTriggered(io.netty.channel.ChannelHandlerContext, java.lang.Object) throws java.lang.Exception", "annotationType": "io.netty.channel.ChannelHandlerMask.Skip", "justification": "Http2Handler overrides userEventTriggered to handle RST_STREAM (Http2ResetFrame); Netty's internal @ChannelHandlerMask.Skip no-op marker is correctly dropped because the method now does real work. Http2Handler is public final and internal-package, so there is no binary, source, or behavioral change for consumers. Scoped to this exact method so any real future change to it is still flagged." + }, + { + "code": "java.field.removed", + "old": "field io.netty.handler.codec.http.HttpContentDecoder.ctx @ org.asynchttpclient.netty.handler.Http1ContentDecompressor", + "justification": "Http1ContentDecompressor no longer extends Netty's HttpContentDecompressor, which inflates every chunk as soon as it is read and cannot pause; it is now a ChannelDuplexHandler on Netty's pull-based Decompressor that stops while its response is suspended. The class is in the internal handler package and is only instantiated by ChannelManager; its public constructor is unchanged. This member was inherited from Netty's decoder or was one of its hooks." + }, + { + "code": "java.method.removed", + "old": "method boolean io.netty.handler.codec.MessageToMessageDecoder::acceptInboundMessage(java.lang.Object) throws java.lang.Exception @ org.asynchttpclient.netty.handler.Http1ContentDecompressor", + "justification": "Http1ContentDecompressor no longer extends Netty's HttpContentDecompressor, which inflates every chunk as soon as it is read and cannot pause; it is now a ChannelDuplexHandler on Netty's pull-based Decompressor that stops while its response is suspended. The class is in the internal handler package and is only instantiated by ChannelManager; its public constructor is unchanged. This member was inherited from Netty's decoder or was one of its hooks." + }, + { + "code": "java.method.exception.checkedRemoved", + "old": "method void io.netty.handler.codec.http.HttpContentDecoder::channelInactive(io.netty.channel.ChannelHandlerContext) throws java.lang.Exception @ org.asynchttpclient.netty.handler.Http1ContentDecompressor", + "new": "method void org.asynchttpclient.netty.handler.Http1ContentDecompressor::channelInactive(io.netty.channel.ChannelHandlerContext)", + "justification": "Http1ContentDecompressor no longer extends Netty's HttpContentDecompressor, which inflates every chunk as soon as it is read and cannot pause; it is now a ChannelDuplexHandler on Netty's pull-based Decompressor that stops while its response is suspended. The class is in the internal handler package and is only instantiated by ChannelManager; its public constructor is unchanged. This handler method no longer declares checked exceptions; callers are unaffected." + }, + { + "code": "java.method.exception.checkedRemoved", + "old": "method void io.netty.handler.codec.MessageToMessageDecoder::channelRead(io.netty.channel.ChannelHandlerContext, java.lang.Object) throws java.lang.Exception @ org.asynchttpclient.netty.handler.Http1ContentDecompressor", + "new": "method void org.asynchttpclient.netty.handler.Http1ContentDecompressor::channelRead(io.netty.channel.ChannelHandlerContext, java.lang.Object)", + "justification": "Http1ContentDecompressor no longer extends Netty's HttpContentDecompressor, which inflates every chunk as soon as it is read and cannot pause; it is now a ChannelDuplexHandler on Netty's pull-based Decompressor that stops while its response is suspended. The class is in the internal handler package and is only instantiated by ChannelManager; its public constructor is unchanged. This handler method no longer declares checked exceptions; callers are unaffected." + }, + { + "code": "java.annotation.added", + "old": "method void io.netty.handler.codec.http.HttpContentDecoder::channelReadComplete(io.netty.channel.ChannelHandlerContext) throws java.lang.Exception @ org.asynchttpclient.netty.handler.Http1ContentDecompressor", + "new": "method void io.netty.channel.ChannelInboundHandlerAdapter::channelReadComplete(io.netty.channel.ChannelHandlerContext) throws java.lang.Exception @ org.asynchttpclient.netty.handler.Http1ContentDecompressor", + "annotationType": "io.netty.channel.ChannelHandlerMask.Skip", + "justification": "Http1ContentDecompressor no longer extends Netty's HttpContentDecompressor, which inflates every chunk as soon as it is read and cannot pause; it is now a ChannelDuplexHandler on Netty's pull-based Decompressor that stops while its response is suspended. The class is in the internal handler package and is only instantiated by ChannelManager; its public constructor is unchanged. Netty's @Skip marker now comes from ChannelInboundHandlerAdapter instead of Netty's decoder; channelReadComplete is passed through as before." + }, + { + "code": "java.method.visibilityReduced", + "old": "method void io.netty.handler.codec.http.HttpContentDecoder::decode(io.netty.channel.ChannelHandlerContext, io.netty.handler.codec.http.HttpObject, java.util.List) throws java.lang.Exception @ org.asynchttpclient.netty.handler.Http1ContentDecompressor", + "new": "method void org.asynchttpclient.netty.handler.Http1ContentDecompressor::decode(io.netty.handler.codec.http.HttpContent)", + "justification": "Http1ContentDecompressor no longer extends Netty's HttpContentDecompressor, which inflates every chunk as soon as it is read and cannot pause; it is now a ChannelDuplexHandler on Netty's pull-based Decompressor that stops while its response is suspended. The class is in the internal handler package and is only instantiated by ChannelManager; its public constructor is unchanged. This member was inherited from Netty's decoder or was one of its hooks." + }, + { + "code": "java.method.removed", + "old": "method java.lang.String org.asynchttpclient.netty.handler.Http1ContentDecompressor::getTargetContentEncoding(java.lang.String) throws java.lang.Exception", + "justification": "Http1ContentDecompressor no longer extends Netty's HttpContentDecompressor, which inflates every chunk as soon as it is read and cannot pause; it is now a ChannelDuplexHandler on Netty's pull-based Decompressor that stops while its response is suspended. The class is in the internal handler package and is only instantiated by ChannelManager; its public constructor is unchanged. This member was inherited from Netty's decoder or was one of its hooks." + }, + { + "code": "java.method.exception.checkedRemoved", + "old": "method void io.netty.handler.codec.http.HttpContentDecoder::handlerRemoved(io.netty.channel.ChannelHandlerContext) throws java.lang.Exception @ org.asynchttpclient.netty.handler.Http1ContentDecompressor", + "new": "method void org.asynchttpclient.netty.handler.Http1ContentDecompressor::handlerRemoved(io.netty.channel.ChannelHandlerContext)", + "justification": "Http1ContentDecompressor no longer extends Netty's HttpContentDecompressor, which inflates every chunk as soon as it is read and cannot pause; it is now a ChannelDuplexHandler on Netty's pull-based Decompressor that stops while its response is suspended. The class is in the internal handler package and is only instantiated by ChannelManager; its public constructor is unchanged. This handler method no longer declares checked exceptions; callers are unaffected." + }, + { + "code": "java.method.removed", + "old": "method io.netty.channel.embedded.EmbeddedChannel org.asynchttpclient.netty.handler.Http1ContentDecompressor::newContentDecoder(java.lang.String) throws java.lang.Exception", + "justification": "Http1ContentDecompressor no longer extends Netty's HttpContentDecompressor, which inflates every chunk as soon as it is read and cannot pause; it is now a ChannelDuplexHandler on Netty's pull-based Decompressor that stops while its response is suspended. The class is in the internal handler package and is only instantiated by ChannelManager; its public constructor is unchanged. This member was inherited from Netty's decoder or was one of its hooks." + }, + { + "code": "java.class.noLongerInheritsFromClass", + "old": "class org.asynchttpclient.netty.handler.Http1ContentDecompressor", + "new": "class org.asynchttpclient.netty.handler.Http1ContentDecompressor", + "justification": "Http1ContentDecompressor no longer extends Netty's HttpContentDecompressor, which inflates every chunk as soon as it is read and cannot pause; it is now a ChannelDuplexHandler on Netty's pull-based Decompressor that stops while its response is suspended. The class is in the internal handler package and is only instantiated by ChannelManager; its public constructor is unchanged." + }, + { + "code": "java.class.nonFinalClassInheritsFromNewClass", + "old": "class org.asynchttpclient.netty.handler.Http1ContentDecompressor", + "new": "class org.asynchttpclient.netty.handler.Http1ContentDecompressor", + "justification": "Http1ContentDecompressor no longer extends Netty's HttpContentDecompressor, which inflates every chunk as soon as it is read and cannot pause; it is now a ChannelDuplexHandler on Netty's pull-based Decompressor that stops while its response is suspended. The class is in the internal handler package and is only instantiated by ChannelManager; its public constructor is unchanged." } ] }