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."
}
]
}