From 6d1308318367073bd6a77e8fdbf6fa9d759c5200 Mon Sep 17 00:00:00 2001 From: Pavel Ptashyts <49400901+pavel-ptashyts@users.noreply.github.com> Date: Thu, 1 Oct 2026 14:57:57 +0200 Subject: [PATCH 1/8] Inflate HTTP/1.1 gzip and deflate in place HttpContentDecompressor builds an EmbeddedChannel and a new Inflater for every compressed response, and its JdkZlibDecoder copies each direct input buffer onto the heap before inflating it. On a client that reads many small gzip responses this shows up as steady allocation and native zlib init/free on the event loop. Http1ContentDecompressor already lives in the pipeline of a single connection, so it now inflates gzip, x-gzip, deflate and x-deflate itself: one Inflater per wrapper for the life of the connection, reset between responses, fed the input buffer as it is. Header and trailer checks, concatenated members, zlib-or-raw detection and the quiet end of a truncated stream follow JdkZlibDecoder, and a differential test holds the output and errors to Netty's across many chunk splits. Other encodings still go through Netty. Claude Code on behalf of Pavel Ptashyts Co-Authored-By: Claude Opus 5.5 --- .../Http1ContentDecompressorBenchmark.java | 127 ++++ .../handler/Http1ContentDecompressor.java | 558 +++++++++++++++- .../handler/Http1ContentDecompressorTest.java | 606 ++++++++++++++++++ 3 files changed, 1282 insertions(+), 9 deletions(-) create mode 100644 client/src/jmh/java/org/asynchttpclient/bench/Http1ContentDecompressorBenchmark.java create mode 100644 client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java diff --git a/client/src/jmh/java/org/asynchttpclient/bench/Http1ContentDecompressorBenchmark.java b/client/src/jmh/java/org/asynchttpclient/bench/Http1ContentDecompressorBenchmark.java new file mode 100644 index 000000000..48c260cf6 --- /dev/null +++ b/client/src/jmh/java/org/asynchttpclient/bench/Http1ContentDecompressorBenchmark.java @@ -0,0 +1,127 @@ +/* + * 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.bench; + +import io.netty.buffer.ByteBuf; +import io.netty.buffer.PooledByteBufAllocator; +import io.netty.channel.ChannelHandler; +import io.netty.channel.embedded.EmbeddedChannel; +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.HttpContentDecompressor; +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.HttpVersion; +import io.netty.util.ReferenceCountUtil; +import org.asynchttpclient.netty.handler.Http1ContentDecompressor; +import org.openjdk.jmh.annotations.Benchmark; +import org.openjdk.jmh.annotations.BenchmarkMode; +import org.openjdk.jmh.annotations.Level; +import org.openjdk.jmh.annotations.Mode; +import org.openjdk.jmh.annotations.OutputTimeUnit; +import org.openjdk.jmh.annotations.Param; +import org.openjdk.jmh.annotations.Scope; +import org.openjdk.jmh.annotations.Setup; +import org.openjdk.jmh.annotations.State; +import org.openjdk.jmh.annotations.TearDown; + +import java.io.ByteArrayOutputStream; +import java.nio.charset.StandardCharsets; +import java.util.concurrent.TimeUnit; +import java.util.zip.GZIPOutputStream; + +/** + * Decodes one gzip response per operation on a keep-alive connection, through Netty's + * {@link HttpContentDecompressor} and through {@link Http1ContentDecompressor}, which inflates gzip without an + * {@code EmbeddedChannel} and with one {@code Inflater} per connection. + *

+ * The compressed body arrives in direct pooled buffers of at most {@code chunkSize} bytes, as + * {@code HttpClientCodec} hands it on under the default allocator. + *

+ * Run with: {@code /tmp/run-jmh.sh Http1ContentDecompressorBenchmark -prof gc -f 1 -wi 5 -i 5} + */ +@State(Scope.Thread) +@BenchmarkMode(Mode.AverageTime) +@OutputTimeUnit(TimeUnit.NANOSECONDS) +public class Http1ContentDecompressorBenchmark { + + @Param({"2048", "20480", "204800"}) + public int bodySize; + + @Param({"8192"}) + public int chunkSize; + + @Param({"netty", "ahc"}) + public String decoder; + + private EmbeddedChannel channel; + private byte[] compressed; + + @Setup(Level.Trial) + public void setup() throws Exception { + StringBuilder json = new StringBuilder(); + for (int i = 0; json.length() < bodySize; i++) { + json.append("{\"id\":\"").append(i).append("\",\"impid\":\"").append(i % 7) + .append("\",\"price\":").append(i * 0.013).append(",\"adm\":\"") + .append(Integer.toHexString(i * 7919)).append("\"},"); + } + byte[] body = json.substring(0, bodySize).getBytes(StandardCharsets.US_ASCII); + ByteArrayOutputStream bos = new ByteArrayOutputStream(); + try (GZIPOutputStream gzip = new GZIPOutputStream(bos)) { + gzip.write(body); + } + compressed = bos.toByteArray(); + + @SuppressWarnings("deprecation") + ChannelHandler handler = "netty".equals(decoder) + ? new HttpContentDecompressor() + : new Http1ContentDecompressor(false, 256L * 1024 * 1024); + channel = new EmbeddedChannel(); + channel.config().setAllocator(PooledByteBufAllocator.DEFAULT); + channel.pipeline().addLast(handler); + } + + @TearDown(Level.Trial) + public void tearDown() { + channel.finishAndReleaseAll(); + } + + @Benchmark + public int decodeResponse() { + HttpResponse response = new DefaultHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK); + response.headers().set(HttpHeaderNames.CONTENT_ENCODING, "gzip"); + response.headers().set(HttpHeaderNames.CONTENT_LENGTH, compressed.length); + channel.writeInbound(response); + for (int offset = 0; offset < compressed.length; offset += chunkSize) { + int length = Math.min(chunkSize, compressed.length - offset); + ByteBuf chunk = PooledByteBufAllocator.DEFAULT.directBuffer(length).writeBytes(compressed, offset, length); + channel.writeInbound(offset + length == compressed.length + ? new DefaultLastHttpContent(chunk) + : new DefaultHttpContent(chunk)); + } + int decoded = 0; + for (Object msg; (msg = channel.readInbound()) != null; ) { + if (msg instanceof HttpContent) { + decoded += ((HttpContent) msg).content().readableBytes(); + } + ReferenceCountUtil.release(msg); + } + return decoded; + } +} 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..aa1fab8c0 100644 --- a/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java +++ b/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java @@ -16,12 +16,30 @@ package org.asynchttpclient.netty.handler; 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.compression.DecompressionException; +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.HttpContentDecompressor; +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.HttpObject; +import io.netty.handler.codec.http.HttpResponse; +import io.netty.handler.codec.http.LastHttpContent; import io.netty.util.ReferenceCountUtil; +import org.jetbrains.annotations.Nullable; + +import java.util.List; +import java.util.zip.CRC32; +import java.util.zip.DataFormatException; +import java.util.zip.Deflater; +import java.util.zip.Inflater; /** * HTTP/1.1 content decompressor that bounds how far a response body may inflate, mirroring what @@ -34,23 +52,76 @@ * 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 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. + * {@code gzip}, {@code x-gzip}, {@code deflate} and {@code x-deflate} named by {@code Content-Encoding} are + * inflated here directly. {@link HttpContentDecompressor} builds an {@link EmbeddedChannel} and a fresh + * {@link Inflater} for every such response, and its {@code JdkZlibDecoder} copies each direct input buffer + * onto the heap before inflating it. This handler sits in the pipeline of one connection, whose responses + * arrive one after another, so it keeps one {@link Inflater} per wrapper for the life of the connection, + * resets it between responses, and feeds it the input buffer as it is. The format handling (header and + * trailer checks, concatenated gzip members, zlib-or-raw detection for {@code deflate}, a truncated stream + * ending quietly) follows {@code JdkZlibDecoder} as {@link HttpContentDecompressor} configures it, so the + * decoded body and the errors are unchanged. *

- * 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. + * Every other encoding still goes through {@link HttpContentDecompressor}. There 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 decoder 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. */ public class Http1ContentDecompressor extends HttpContentDecompressor { + private static final int FHCRC = 0x02; + private static final int FEXTRA = 0x04; + private static final int FNAME = 0x08; + private static final int FCOMMENT = 0x10; + private static final int FRESERVED = 0xE0; + + // Same floor and forwarding threshold as JdkZlibDecoder, so the body is cut into parts the same way. + private static final int MIN_OUTPUT_BUFFER_SIZE = 512; + private static final int MAX_FORWARD_BYTES = 64 * 1024; + + private enum GzipState { + HEADER_START, + FLG_READ, + XLEN_READ, + SKIP_FNAME, + SKIP_COMMENT, + PROCESS_FHCRC, + HEADER_END, + FOOTER_START, + } + private final boolean keepEncodingHeader; // Maximum cumulative decompressed bytes for one response; 0 disables the limit. private final long maxDecompressedBytes; + // Reused by every response on this connection and released when the handler leaves the pipeline. + private @Nullable Inflater nowrapInflater; + private @Nullable Inflater zlibInflater; + private final CRC32 crc = new CRC32(); + + // The response being inflated by this class rather than by HttpContentDecompressor. + private boolean inflating; + private boolean gzip; + // Null until the first two bytes of a deflate body have told zlib from raw deflate. + private @Nullable Inflater inflater; + private boolean failed; + private boolean finished; + private GzipState gzipState = GzipState.HEADER_START; + private int flags = -1; + private int xlen = -1; + private long decompressedBytes; + // Bytes of a gzip header or trailer split across two HttpContent messages. + private @Nullable ByteBuf pending; + + // HttpContentDecoder keeps its read-request flag private, so channelReadComplete defers to it only when + // the last message of the read was its own. + private boolean lastDecodedBySuper = true; + private boolean needRead = true; + 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 @@ -60,6 +131,68 @@ public Http1ContentDecompressor(boolean keepEncodingHeader, long maxDecompressed this.maxDecompressedBytes = maxDecompressedBytes; } + @Override + protected void decode(ChannelHandlerContext ctx, HttpObject msg, List out) throws Exception { + if (msg instanceof HttpResponse) { + endResponse(); + HttpResponse response = (HttpResponse) msg; + String contentEncoding = response.headers().get(HttpHeaderNames.CONTENT_ENCODING); + if (contentEncoding != null && response.status().code() != 100) { + contentEncoding = contentEncoding.trim(); + if (isGzip(contentEncoding) || isDeflate(contentEncoding)) { + lastDecodedBySuper = false; + needRead = true; + startResponse(ctx, response, contentEncoding); + if (response instanceof HttpContent) { + decodeContent(ctx, (HttpContent) response); + } + return; + } + } + } else if (inflating) { + lastDecodedBySuper = false; + needRead = true; + decodeContent(ctx, (HttpContent) msg); + return; + } + lastDecodedBySuper = true; + super.decode(ctx, msg, out); + } + + @Override + public void channelReadComplete(ChannelHandlerContext ctx) throws Exception { + if (lastDecodedBySuper) { + super.channelReadComplete(ctx); + return; + } + boolean needRead = this.needRead; + this.needRead = true; + try { + ctx.fireChannelReadComplete(); + } finally { + // A message that was swallowed whole, such as a chunk holding only part of a gzip header, gives + // the handlers after this one nothing to answer with a read of their own. + if (needRead && !ctx.channel().config().isAutoRead()) { + ctx.read(); + } + } + } + + @Override + public void handlerRemoved(ChannelHandlerContext ctx) throws Exception { + try { + releaseInflaters(); + } finally { + super.handlerRemoved(ctx); + } + } + + @Override + public void channelInactive(ChannelHandlerContext ctx) throws Exception { + releaseInflaters(); + super.channelInactive(ctx); + } + @Override protected EmbeddedChannel newContentDecoder(String contentEncoding) throws Exception { EmbeddedChannel decoder = super.newContentDecoder(contentEncoding); @@ -78,6 +211,413 @@ protected String getTargetContentEncoding(String contentEncoding) throws Excepti return keepEncodingHeader ? contentEncoding : super.getTargetContentEncoding(contentEncoding); } + private static boolean isGzip(String contentEncoding) { + return HttpHeaderValues.GZIP.contentEqualsIgnoreCase(contentEncoding) + || HttpHeaderValues.X_GZIP.contentEqualsIgnoreCase(contentEncoding); + } + + private static boolean isDeflate(String contentEncoding) { + return HttpHeaderValues.DEFLATE.contentEqualsIgnoreCase(contentEncoding) + || HttpHeaderValues.X_DEFLATE.contentEqualsIgnoreCase(contentEncoding); + } + + /** + * Rewrites the headers exactly as {@code HttpContentDecoder} does once it has installed a decoder. + */ + private void startResponse(ChannelHandlerContext ctx, HttpResponse response, String contentEncoding) throws Exception { + inflating = true; + gzip = isGzip(contentEncoding); + if (gzip) { + inflater = nowrapInflater(); + } + + HttpHeaders headers = response.headers(); + // The decoded length is known only once the whole body is through. + if (headers.contains(HttpHeaderNames.CONTENT_LENGTH)) { + headers.remove(HttpHeaderNames.CONTENT_LENGTH); + headers.set(HttpHeaderNames.TRANSFER_ENCODING, HttpHeaderValues.CHUNKED); + } + String targetContentEncoding = getTargetContentEncoding(contentEncoding); + if (HttpHeaderValues.IDENTITY.contentEquals(targetContentEncoding)) { + headers.remove(HttpHeaderNames.CONTENT_ENCODING); + } else { + headers.set(HttpHeaderNames.CONTENT_ENCODING, targetContentEncoding); + } + + HttpResponse head = response; + if (response instanceof HttpContent) { + // The body is decoded and sent separately, and a message that carries content of its own would + // read downstream as the end of the response. + head = new DefaultHttpResponse(response.protocolVersion(), response.status()); + head.headers().set(headers); + head.setDecoderResult(response.decoderResult()); + } + needRead = false; + ctx.fireChannelRead(head); + } + + private void decodeContent(ChannelHandlerContext ctx, HttpContent content) { + if (failed) { + // The exchange has already been failed; what is still in flight for it is dropped. + return; + } + ByteBuf in = content.content(); + if (in.isReadable()) { + try { + inflate(ctx, in); + } catch (Throwable t) { + endResponse(); + inflating = true; + failed = true; + throw t; + } + } + if (content instanceof LastHttpContent) { + endResponse(); + HttpHeaders trailers = ((LastHttpContent) content).trailingHeaders(); + needRead = false; + ctx.fireChannelRead(trailers.isEmpty() + ? LastHttpContent.EMPTY_LAST_CONTENT + : new DefaultLastHttpContent(Unpooled.EMPTY_BUFFER, trailers)); + } + } + + private void inflate(ChannelHandlerContext ctx, ByteBuf content) { + ByteBuf in = content; + ByteBuf pending = this.pending; + if (pending != null) { + pending.writeBytes(content); + in = pending; + } + + int readable; + do { + readable = in.readableBytes(); + inflateOnce(ctx, in); + } while (in.isReadable() && in.readableBytes() != readable); + + if (pending != null) { + if (pending.isReadable()) { + pending.discardReadBytes(); + } else { + pending.release(); + this.pending = null; + } + } else if (in.isReadable()) { + // Only a gzip header or trailer that has not fully arrived is left over; the inflater consumes + // all the deflate data it is given. + this.pending = ctx.alloc().heapBuffer(in.readableBytes()).writeBytes(in); + } + } + + /** + * One step of {@code JdkZlibDecoder#decode}, which {@code ByteToMessageDecoder} would call again for as + * long as it consumes input. + */ + private void inflateOnce(ChannelHandlerContext ctx, ByteBuf in) { + if (finished) { + in.skipBytes(in.readableBytes()); + return; + } + + Inflater inflater = this.inflater; + if (inflater == null) { + if (in.readableBytes() < 2) { + return; + } + // RFC 1950 says deflate is zlib-wrapped, but some servers send raw deflate. + inflater = looksLikeZlib(in.getShort(in.readerIndex())) ? zlibInflater() : nowrapInflater(); + this.inflater = inflater; + } + + if (gzip && gzipState != GzipState.HEADER_END) { + if (gzipState == GzipState.FOOTER_START && !readGzipFooter(in, inflater)) { + return; + } + if (!readGzipHeader(in) || !in.isReadable()) { + return; + } + } + + int readable = in.readableBytes(); + if (in.hasArray()) { + inflater.setInput(in.array(), in.arrayOffset() + in.readerIndex(), readable); + } else if (in.nioBufferCount() == 1) { + inflater.setInput(in.nioBuffer(in.readerIndex(), readable)); + } else { + byte[] array = new byte[readable]; + in.getBytes(in.readerIndex(), array); + inflater.setInput(array); + } + + ByteBuf decompressed = null; + boolean forward = false; + try { + boolean readFooter = false; + // A full output buffer means the inflater may still hold decoded data after it has taken all the + // input, so needsInput() alone would end the loop too early. + boolean pendingOutput = false; + while (pendingOutput || !inflater.needsInput()) { + int preferredSize = Math.max(inflater.getRemaining() << 1, MIN_OUTPUT_BUFFER_SIZE); + if (decompressed == null) { + decompressed = ctx.alloc().heapBuffer(preferredSize); + } else { + decompressed.ensureWritable(preferredSize); + } + byte[] outArray = decompressed.array(); + int writerIndex = decompressed.writerIndex(); + int outIndex = decompressed.arrayOffset() + writerIndex; + int writable = decompressed.writableBytes(); + int outputLength = inflater.inflate(outArray, outIndex, writable); + pendingOutput = outputLength == writable; + if (outputLength > 0) { + decompressedBytes += outputLength; + if (maxDecompressedBytes > 0 && decompressedBytes > maxDecompressedBytes) { + throw new DecompressionException("HTTP/1.1 response body exceeds the maximum decompressed size of " + + maxDecompressedBytes + " bytes"); + } + decompressed.writerIndex(writerIndex + outputLength); + if (gzip) { + crc.update(outArray, outIndex, outputLength); + } + if (decompressed.readableBytes() >= MAX_FORWARD_BYTES) { + ByteBuf buffer = decompressed; + decompressed = null; + fireContent(ctx, buffer); + } + } else if (inflater.needsDictionary()) { + throw new DecompressionException("decompression failure, unable to set dictionary as non was specified"); + } + + if (inflater.finished()) { + if (gzip) { + readFooter = true; + } else { + finished = true; + } + break; + } + } + + in.skipBytes(readable - inflater.getRemaining()); + + if (readFooter) { + gzipState = GzipState.FOOTER_START; + readGzipFooter(in, inflater); + } + forward = true; + } catch (DataFormatException e) { + throw new DecompressionException("decompression failure", e); + } finally { + if (decompressed != null) { + if (forward && decompressed.isReadable()) { + fireContent(ctx, decompressed); + } else { + decompressed.release(); + } + } + } + } + + private void fireContent(ChannelHandlerContext ctx, ByteBuf buffer) { + needRead = false; + ctx.fireChannelRead(new DefaultHttpContent(buffer)); + } + + private boolean readGzipHeader(ByteBuf in) { + switch (gzipState) { + case HEADER_START: + if (in.readableBytes() < 10) { + return false; + } + int magic0 = in.readByte(); + int magic1 = in.readByte(); + // Like JdkZlibDecoder, only the first magic byte is checked. + if (magic0 != 31) { + throw new DecompressionException("Input is not in the GZIP format"); + } + crc.update(magic0); + crc.update(magic1); + + int method = in.readUnsignedByte(); + if (method != Deflater.DEFLATED) { + throw new DecompressionException("Unsupported compression method " + method + " in the GZIP header"); + } + crc.update(method); + + flags = in.readUnsignedByte(); + crc.update(flags); + if ((flags & FRESERVED) != 0) { + throw new DecompressionException("Reserved flags are set in the GZIP header"); + } + + // MTIME, XFL and OS. + updateCrc(in, in.readerIndex(), 6); + in.skipBytes(6); + gzipState = GzipState.FLG_READ; + // fall through + case FLG_READ: + if ((flags & FEXTRA) != 0) { + if (in.readableBytes() < 2) { + return false; + } + int xlen1 = in.readUnsignedByte(); + int xlen2 = in.readUnsignedByte(); + crc.update(xlen1); + crc.update(xlen2); + xlen = xlen2 << 8 | xlen1; + } + gzipState = GzipState.XLEN_READ; + // fall through + case XLEN_READ: + if (xlen != -1) { + if (in.readableBytes() < xlen) { + return false; + } + updateCrc(in, in.readerIndex(), xlen); + in.skipBytes(xlen); + } + gzipState = GzipState.SKIP_FNAME; + // fall through + case SKIP_FNAME: + if (!skipIfNeeded(in, FNAME)) { + return false; + } + gzipState = GzipState.SKIP_COMMENT; + // fall through + case SKIP_COMMENT: + if (!skipIfNeeded(in, FCOMMENT)) { + return false; + } + gzipState = GzipState.PROCESS_FHCRC; + // fall through + case PROCESS_FHCRC: + if ((flags & FHCRC) != 0) { + if (in.readableBytes() < 2) { + return false; + } + int crc16Value = in.readUnsignedShortLE(); + int readCrc16 = (int) (crc.getValue() & 0xFFFF); + if (crc16Value != readCrc16) { + throw new DecompressionException("CRC16 value mismatch. Expected: " + crc16Value + ", Got: " + readCrc16); + } + } + crc.reset(); + gzipState = GzipState.HEADER_END; + // fall through + case HEADER_END: + return true; + default: + throw new IllegalStateException(); + } + } + + private boolean skipIfNeeded(ByteBuf in, int flagMask) { + if ((flags & flagMask) != 0) { + for (;;) { + if (!in.isReadable()) { + return false; + } + int b = in.readUnsignedByte(); + crc.update(b); + if (b == 0x00) { + break; + } + } + } + return true; + } + + /** + * Checks CRC32 and ISIZE, then readies the inflater for a further gzip member, which + * {@link HttpContentDecompressor} decodes too. + */ + private boolean readGzipFooter(ByteBuf in, Inflater inflater) { + if (in.readableBytes() < 8) { + return false; + } + long crcValue = in.readUnsignedIntLE(); + long readCrc = crc.getValue(); + if (crcValue != readCrc) { + throw new DecompressionException("CRC value mismatch. Expected: " + crcValue + ", Got: " + readCrc); + } + int dataLength = in.readIntLE(); + int readLength = inflater.getTotalOut(); + if (dataLength != readLength) { + throw new DecompressionException("Number of bytes mismatch. Expected: " + dataLength + ", Got: " + readLength); + } + inflater.reset(); + crc.reset(); + xlen = -1; + gzipState = GzipState.HEADER_START; + return true; + } + + private void updateCrc(ByteBuf in, int index, int length) { + if (in.hasArray()) { + crc.update(in.array(), in.arrayOffset() + index, length); + } else { + crc.update(in.nioBuffer(index, length)); + } + } + + // RFC 1950 section 2.2, written as JdkZlibDecoder writes it so the same bodies are taken for zlib. + private static boolean looksLikeZlib(short cmfFlg) { + return (cmfFlg & 0x7800) == 0x7800 && cmfFlg % 31 == 0; + } + + private Inflater nowrapInflater() { + Inflater inflater = nowrapInflater; + if (inflater == null) { + inflater = new Inflater(true); + nowrapInflater = inflater; + } + return inflater; + } + + private Inflater zlibInflater() { + Inflater inflater = zlibInflater; + if (inflater == null) { + inflater = new Inflater(); + zlibInflater = inflater; + } + return inflater; + } + + private void endResponse() { + Inflater inflater = this.inflater; + if (inflater != null) { + // Also drops the inflater's reference to the last input buffer, which is released after decode. + inflater.reset(); + this.inflater = null; + } + ByteBuf pending = this.pending; + if (pending != null) { + pending.release(); + this.pending = null; + } + inflating = false; + failed = false; + finished = false; + gzipState = GzipState.HEADER_START; + flags = -1; + xlen = -1; + decompressedBytes = 0; + crc.reset(); + } + + private void releaseInflaters() { + endResponse(); + if (nowrapInflater != null) { + nowrapInflater.end(); + nowrapInflater = null; + } + if (zlibInflater != null) { + zlibInflater.end(); + zlibInflater = null; + } + } + /** * 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 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..59fb0c4f3 --- /dev/null +++ b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java @@ -0,0 +1,606 @@ +/* + * 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 io.netty.buffer.ByteBuf; +import io.netty.buffer.Unpooled; +import io.netty.buffer.UnpooledByteBufAllocator; +import io.netty.channel.ChannelHandler; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelOutboundHandlerAdapter; +import io.netty.channel.embedded.EmbeddedChannel; +import io.netty.handler.codec.compression.DecompressionException; +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.HttpContentDecompressor; +import io.netty.handler.codec.http.HttpHeaderNames; +import io.netty.handler.codec.http.HttpHeaderValues; +import io.netty.handler.codec.http.HttpResponse; +import io.netty.handler.codec.http.HttpResponseStatus; +import io.netty.handler.codec.http.HttpVersion; +import io.netty.handler.codec.http.LastHttpContent; +import io.netty.util.ReferenceCountUtil; +import org.junit.jupiter.api.Test; + +import java.io.ByteArrayOutputStream; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.Random; +import java.util.function.Supplier; +import java.util.zip.CRC32; +import java.util.zip.Deflater; +import java.util.zip.GZIPOutputStream; + +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.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * {@link Http1ContentDecompressor} inflates gzip and deflate itself instead of through Netty's + * {@link HttpContentDecompressor}. Most tests here feed both handlers the same body, cut into chunks in + * several ways, and require the same decoded bytes and the same error, so the two cannot drift apart. + */ +public class Http1ContentDecompressorTest { + + private static final int FTEXT = 0x01; + private static final int FHCRC = 0x02; + private static final int FEXTRA = 0x04; + private static final int FNAME = 0x08; + private static final int FCOMMENT = 0x10; + + private static final byte[] EMPTY = new byte[0]; + private static final byte[] TEXT = "{\"id\":\"1\",\"seatbid\":[{\"bid\":[{\"price\":1.5}]}]}".repeat(200) + .getBytes(StandardCharsets.US_ASCII); + // Large and incompressible, so one chunk inflates past the 64 KiB forwarding threshold only when it is + // big itself, and the response spans many chunks. + private static final byte[] RANDOM = randomBytes(300 * 1024, 42); + // Compressible enough that a single chunk inflates well past the forwarding threshold. + private static final byte[] ZEROS = new byte[2 * 1024 * 1024]; + + private static byte[] randomBytes(int length, long seed) { + byte[] bytes = new byte[length]; + new Random(seed).nextBytes(bytes); + return bytes; + } + + private static byte[] gzip(byte[] payload) { + ByteArrayOutputStream bos = new ByteArrayOutputStream(); + try (GZIPOutputStream gz = new GZIPOutputStream(bos)) { + gz.write(payload); + } catch (Exception e) { + throw new IllegalStateException(e); + } + return bos.toByteArray(); + } + + private static byte[] deflate(byte[] payload, boolean nowrap) { + Deflater deflater = new Deflater(Deflater.DEFAULT_COMPRESSION, nowrap); + try { + deflater.setInput(payload); + deflater.finish(); + ByteArrayOutputStream bos = new ByteArrayOutputStream(); + byte[] buffer = new byte[8192]; + while (!deflater.finished()) { + bos.write(buffer, 0, deflater.deflate(buffer)); + } + return bos.toByteArray(); + } finally { + deflater.end(); + } + } + + /** + * A gzip member with every optional header field the flags ask for, built by hand because + * {@link GZIPOutputStream} writes none of them. + */ + private static byte[] gzipWithHeaderFields(byte[] payload, int flags) { + ByteArrayOutputStream header = new ByteArrayOutputStream(); + header.write(0x1f); + header.write(0x8b); + header.write(8); + header.write(flags); + header.writeBytes(new byte[]{1, 2, 3, 4, 0, (byte) 255}); + if ((flags & FEXTRA) != 0) { + byte[] extra = "AB\u0003\u0000xyz-extra-field".getBytes(StandardCharsets.US_ASCII); + header.write(extra.length & 0xff); + header.write(extra.length >>> 8); + header.writeBytes(extra); + } + if ((flags & FNAME) != 0) { + header.writeBytes("response.json".getBytes(StandardCharsets.US_ASCII)); + header.write(0); + } + if ((flags & FCOMMENT) != 0) { + header.writeBytes("a comment".getBytes(StandardCharsets.US_ASCII)); + header.write(0); + } + if ((flags & FHCRC) != 0) { + CRC32 crc = new CRC32(); + crc.update(header.toByteArray()); + int crc16 = (int) crc.getValue() & 0xffff; + header.write(crc16 & 0xff); + header.write(crc16 >>> 8); + } + ByteArrayOutputStream member = new ByteArrayOutputStream(); + member.writeBytes(header.toByteArray()); + member.writeBytes(deflate(payload, true)); + CRC32 crc = new CRC32(); + crc.update(payload); + writeIntLE(member, (int) crc.getValue()); + writeIntLE(member, payload.length); + return member.toByteArray(); + } + + private static void writeIntLE(ByteArrayOutputStream out, int value) { + out.write(value); + out.write(value >>> 8); + out.write(value >>> 16); + out.write(value >>> 24); + } + + private static byte[] concat(byte[]... parts) { + ByteArrayOutputStream bos = new ByteArrayOutputStream(); + for (byte[] part : parts) { + bos.writeBytes(part); + } + return bos.toByteArray(); + } + + private static byte[] withByte(byte[] body, int index, int value) { + byte[] copy = body.clone(); + copy[index < 0 ? copy.length + index : index] = (byte) value; + return copy; + } + + private static byte[] flipped(byte[] body, int index) { + byte[] copy = body.clone(); + int i = index < 0 ? copy.length + index : index; + copy[i] = (byte) ~copy[i]; + return copy; + } + + private static final class Body { + final String name; + final String contentEncoding; + final byte[] encoded; + + Body(String name, String contentEncoding, byte[] encoded) { + this.name = name; + this.contentEncoding = contentEncoding; + this.encoded = encoded; + } + } + + private static List bodies() { + List bodies = new ArrayList<>(); + byte[] gzipText = gzip(TEXT); + byte[] gzipRandom = gzip(RANDOM); + bodies.add(new Body("gzip empty", "gzip", gzip(EMPTY))); + bodies.add(new Body("gzip text", "gzip", gzipText)); + bodies.add(new Body("gzip random", "gzip", gzipRandom)); + bodies.add(new Body("gzip zeros", "gzip", gzip(ZEROS))); + bodies.add(new Body("x-gzip mixed case", "X-GZip", gzipText)); + bodies.add(new Body("gzip all header fields", "gzip", + gzipWithHeaderFields(TEXT, FTEXT | FHCRC | FEXTRA | FNAME | FCOMMENT))); + bodies.add(new Body("gzip name and comment", "gzip", gzipWithHeaderFields(TEXT, FNAME | FCOMMENT))); + bodies.add(new Body("gzip extra only", "gzip", gzipWithHeaderFields(RANDOM, FEXTRA))); + bodies.add(new Body("gzip concatenated", "gzip", concat(gzipText, gzip(RANDOM), gzip(EMPTY), gzipText))); + bodies.add(new Body("deflate zlib", "deflate", deflate(TEXT, false))); + bodies.add(new Body("deflate zlib random", "deflate", deflate(RANDOM, false))); + bodies.add(new Body("deflate raw", "deflate", deflate(TEXT, true))); + bodies.add(new Body("x-deflate raw zeros", "X-Deflate", deflate(ZEROS, true))); + bodies.add(new Body("deflate empty raw", "deflate", deflate(EMPTY, true))); + bodies.add(new Body("deflate trailing bytes", "deflate", concat(deflate(TEXT, false), TEXT))); + + bodies.add(new Body("gzip bad crc", "gzip", flipped(gzipText, -8))); + bodies.add(new Body("gzip bad isize", "gzip", flipped(gzipText, -1))); + bodies.add(new Body("gzip bad header crc", "gzip", + flipped(gzipWithHeaderFields(TEXT, FHCRC | FNAME), 24))); + bodies.add(new Body("gzip bad magic", "gzip", withByte(gzipText, 0, 0x1e))); + bodies.add(new Body("gzip second magic byte ignored", "gzip", withByte(gzipText, 1, 0))); + bodies.add(new Body("gzip bad method", "gzip", withByte(gzipText, 2, 7))); + bodies.add(new Body("gzip reserved flag", "gzip", withByte(gzipText, 3, 0x20))); + bodies.add(new Body("gzip corrupt data", "gzip", flipped(gzipRandom, gzipRandom.length / 2))); + bodies.add(new Body("gzip truncated data", "gzip", Arrays.copyOf(gzipRandom, gzipRandom.length / 2))); + bodies.add(new Body("gzip truncated trailer", "gzip", Arrays.copyOf(gzipText, gzipText.length - 3))); + bodies.add(new Body("gzip truncated header", "gzip", Arrays.copyOf(gzipText, 6))); + bodies.add(new Body("gzip short garbage after member", "gzip", concat(gzipText, new byte[]{'\r', '\n'}))); + bodies.add(new Body("gzip long garbage after member", "gzip", concat(gzipText, TEXT))); + bodies.add(new Body("deflate corrupt data", "deflate", flipped(deflate(RANDOM, false), 1000))); + bodies.add(new Body("deflate bad adler", "deflate", flipped(deflate(TEXT, false), -1))); + bodies.add(new Body("deflate one byte", "deflate", new byte[]{0x78})); + bodies.add(new Body("deflate preset dictionary", "deflate", new byte[]{0x78, (byte) 0xbb, 0, 0, 0, 1, 1, 1})); + return bodies; + } + + private static List splits(int length) { + List splits = new ArrayList<>(); + splits.add(new int[]{length}); + splits.add(sizes(length, 8192)); + // Tiny chunks put every header and trailer boundary on a chunk edge; a large body gets them only + // from the random cuts. + boolean small = length <= 64 * 1024; + if (small) { + splits.add(sizes(length, 1)); + splits.add(sizes(length, 7)); + } + Random random = new Random(length); + for (int i = small ? 0 : 1; i < 3; i++) { + List sizes = new ArrayList<>(); + int remaining = length; + while (remaining > 0) { + int size = Math.min(remaining, 1 + random.nextInt(i == 0 ? 16 : 9000)); + sizes.add(size); + remaining -= size; + } + splits.add(sizes.stream().mapToInt(Integer::intValue).toArray()); + } + return splits; + } + + private static int[] sizes(int length, int chunk) { + int[] sizes = new int[Math.max(1, (length + chunk - 1) / chunk)]; + Arrays.fill(sizes, chunk); + if (length % chunk != 0) { + sizes[sizes.length - 1] = length % chunk; + } + if (length == 0) { + sizes[0] = 0; + } + return sizes; + } + + private static final class Outcome { + final ByteArrayOutputStream body = new ByteArrayOutputStream(); + HttpResponse head; + int parts; + boolean ended; + Throwable failure; + } + + private static HttpResponse response(String contentEncoding, long contentLength) { + HttpResponse response = new DefaultHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK); + if (contentEncoding != null) { + response.headers().set(HttpHeaderNames.CONTENT_ENCODING, contentEncoding); + } + if (contentLength >= 0) { + response.headers().set(HttpHeaderNames.CONTENT_LENGTH, contentLength); + } + return response; + } + + private static ByteBuf buffer(byte[] bytes, int offset, int length, boolean direct) { + ByteBuf buf = direct ? Unpooled.directBuffer(Math.max(length, 1)) : Unpooled.buffer(Math.max(length, 1)); + return buf.writeBytes(bytes, offset, length); + } + + /** + * Sends one response through the channel and collects what comes out, stopping at the first error. + */ + private static Outcome exchange(EmbeddedChannel channel, String contentEncoding, byte[] encoded, int[] sizes, + boolean direct) { + Outcome outcome = new Outcome(); + List messages = new ArrayList<>(); + messages.add(response(contentEncoding, encoded.length)); + int offset = 0; + for (int i = 0; i < sizes.length; i++) { + ByteBuf content = buffer(encoded, offset, sizes[i], direct); + offset += sizes[i]; + messages.add(i == sizes.length - 1 ? new DefaultLastHttpContent(content) : new DefaultHttpContent(content)); + } + for (int i = 0; i < messages.size(); i++) { + try { + channel.writeInbound(messages.get(i)); + } catch (Throwable t) { + outcome.failure = t; + for (int j = i + 1; j < messages.size(); j++) { + ReferenceCountUtil.release(messages.get(j)); + } + break; + } finally { + drain(channel, outcome); + } + } + return outcome; + } + + private static void drain(EmbeddedChannel channel, Outcome outcome) { + for (Object msg; (msg = channel.readInbound()) != null; ) { + try { + if (msg instanceof HttpResponse) { + assertNull(outcome.head, "a second response head"); + assertFalse(msg instanceof HttpContent, "the head must not carry content"); + outcome.head = (HttpResponse) msg; + } + if (msg instanceof HttpContent) { + ByteBuf content = ((HttpContent) msg).content(); + if (content.isReadable()) { + outcome.parts++; + content.readBytes(outcome.body, content.readableBytes()); + } + if (msg instanceof LastHttpContent) { + outcome.ended = true; + } + } + } catch (Exception e) { + throw new IllegalStateException(e); + } finally { + ReferenceCountUtil.release(msg); + } + } + } + + private static EmbeddedChannel channel(ChannelHandler decompressor, UnpooledByteBufAllocator allocator) { + EmbeddedChannel channel = new EmbeddedChannel(); + channel.config().setAllocator(allocator); + channel.pipeline().addLast(decompressor); + return channel; + } + + private static UnpooledByteBufAllocator allocator() { + return new UnpooledByteBufAllocator(false, true); + } + + @SuppressWarnings("deprecation") + private static HttpContentDecompressor nettyDecompressor() { + return new HttpContentDecompressor(); + } + + @Test + void decodesLikeNettyForEveryBodyAndSplit() { + int checked = 0; + for (Body body : bodies()) { + for (int[] sizes : splits(body.encoded.length)) { + for (boolean direct : new boolean[]{false, true}) { + String scenario = body.name + ", " + sizes.length + " chunks, " + (direct ? "direct" : "heap"); + UnpooledByteBufAllocator ahcAllocator = allocator(); + EmbeddedChannel ahc = channel(new Http1ContentDecompressor(false, 0), ahcAllocator); + EmbeddedChannel netty = channel(nettyDecompressor(), allocator()); + Outcome expected = exchange(netty, body.contentEncoding, body.encoded, sizes, direct); + Outcome actual = exchange(ahc, body.contentEncoding, body.encoded, sizes, direct); + + assertSameOutcome(expected, actual, scenario); + ahc.finishAndReleaseAll(); + try { + netty.finishAndReleaseAll(); + } catch (DecompressionException ignored) { + // Netty's decoder runs once more over a corrupt body when its channel is torn down. + } + assertEquals(0, ahcAllocator.metric().usedHeapMemory(), scenario + ": heap left allocated"); + checked++; + } + } + } + assertTrue(checked > 300, "only " + checked + " scenarios ran"); + } + + private static void assertSameOutcome(Outcome expected, Outcome actual, String scenario) { + assertNotNull(actual.head, scenario); + assertEquals(expected.head.headers().entries().toString(), actual.head.headers().entries().toString(), + scenario + ": headers"); + if (expected.failure == null) { + assertNull(actual.failure, scenario + ": unexpected failure"); + assertArrayEquals(expected.body.toByteArray(), actual.body.toByteArray(), scenario + ": body"); + assertEquals(expected.ended, actual.ended, scenario + ": end of response"); + } else { + assertNotNull(actual.failure, scenario + ": expected " + expected.failure); + assertEquals(expected.failure.getClass(), actual.failure.getClass(), scenario + ": failure type"); + assertEquals(expected.failure.getMessage(), actual.failure.getMessage(), scenario + ": failure message"); + assertFalse(actual.ended, scenario + ": failed response must not end normally"); + } + } + + @Test + void forwardsOnePartPerChunkUnlessItInflatesPastTheThreshold() { + byte[] encoded = gzip(TEXT); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator()); + Outcome outcome = exchange(channel, "gzip", encoded, new int[]{encoded.length}, true); + assertNull(outcome.failure); + assertEquals(1, outcome.parts); + assertArrayEquals(TEXT, outcome.body.toByteArray()); + channel.finishAndReleaseAll(); + } + + @Test + void reusesTheInflaterAcrossResponsesOnOneConnection() { + UnpooledByteBufAllocator allocator = allocator(); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator); + byte[] gzipText = gzip(TEXT); + byte[] zlib = deflate(RANDOM, false); + byte[] raw = deflate(TEXT, true); + + assertDecoded(TEXT, exchange(channel, "gzip", gzipText, new int[]{gzipText.length}, true)); + Outcome identity = exchange(channel, null, TEXT, new int[]{100, TEXT.length - 100}, false); + assertDecoded(TEXT, identity); + assertEquals(String.valueOf(TEXT.length), identity.head.headers().get(HttpHeaderNames.CONTENT_LENGTH)); + assertDecoded(RANDOM, exchange(channel, "deflate", zlib, sizes(zlib.length, 4000), true)); + assertDecoded(TEXT, exchange(channel, "deflate", raw, sizes(raw.length, 3), false)); + // A response cut short leaves state behind that the next one must not see. + Outcome cut = exchange(channel, "gzip", Arrays.copyOf(gzipText, 40), new int[]{40}, true); + assertNull(cut.failure); + assertDecoded(TEXT, exchange(channel, "gzip", gzipText, sizes(gzipText.length, 5), true)); + // So does a corrupt one. + Outcome corrupt = exchange(channel, "gzip", flipped(gzipText, -8), new int[]{gzipText.length}, true); + assertInstanceOf(DecompressionException.class, corrupt.failure); + assertDecoded(TEXT, exchange(channel, "gzip", gzipText, new int[]{gzipText.length}, true)); + + channel.finishAndReleaseAll(); + assertEquals(0, allocator.metric().usedHeapMemory()); + } + + @Test + void dropsWhatIsLeftOfAFailedResponse() { + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator()); + byte[] corrupt = withByte(gzip(TEXT), 0, 0); + channel.writeInbound(response("gzip", -1)); + assertInstanceOf(HttpResponse.class, channel.readInbound()); + try { + channel.writeInbound(new DefaultHttpContent(Unpooled.wrappedBuffer(corrupt))); + } catch (DecompressionException expected) { + // the exchange is failed by this + } + ByteBuf rest = Unpooled.wrappedBuffer(TEXT); + channel.writeInbound(new DefaultLastHttpContent(rest)); + assertEquals(0, rest.refCnt()); + assertNull(channel.readInbound()); + channel.finishAndReleaseAll(); + } + + private static void assertDecoded(byte[] expected, Outcome outcome) { + assertNull(outcome.failure); + assertTrue(outcome.ended); + assertArrayEquals(expected, outcome.body.toByteArray()); + } + + @Test + void rewritesHeadersLikeNetty() { + Outcome outcome = exchange(channel(new Http1ContentDecompressor(false, 0), allocator()), "gzip", gzip(TEXT), + new int[]{10, 20}, false); + assertNull(outcome.head.headers().get(HttpHeaderNames.CONTENT_ENCODING)); + assertNull(outcome.head.headers().get(HttpHeaderNames.CONTENT_LENGTH)); + assertEquals(HttpHeaderValues.CHUNKED.toString(), outcome.head.headers().get(HttpHeaderNames.TRANSFER_ENCODING)); + } + + @Test + void keepEncodingHeaderLeavesContentEncoding() { + Outcome outcome = exchange(channel(new Http1ContentDecompressor(true, 0), allocator()), "X-Gzip", gzip(RANDOM), + sizes(gzip(RANDOM).length, 2000), false); + assertEquals("X-Gzip", outcome.head.headers().get(HttpHeaderNames.CONTENT_ENCODING)); + assertNull(outcome.head.headers().get(HttpHeaderNames.CONTENT_LENGTH)); + assertArrayEquals(RANDOM, outcome.body.toByteArray()); + } + + @Test + void trailersSurviveDecompression() { + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator()); + byte[] encoded = gzip(TEXT); + channel.writeInbound(response("gzip", -1)); + LastHttpContent last = new DefaultLastHttpContent(Unpooled.wrappedBuffer(encoded)); + last.trailingHeaders().set("x-checksum", "abc"); + channel.writeInbound(last); + assertInstanceOf(HttpResponse.class, channel.readInbound()); + HttpContent body = channel.readInbound(); + assertEquals(TEXT.length, body.content().readableBytes()); + body.release(); + LastHttpContent end = channel.readInbound(); + assertEquals("abc", end.trailingHeaders().get("x-checksum")); + assertFalse(end.content().isReadable()); + end.release(); + channel.finishAndReleaseAll(); + } + + @Test + void continueResponseIsPassedThrough() { + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator()); + HttpResponse interim = new DefaultHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.CONTINUE); + interim.headers().set(HttpHeaderNames.CONTENT_ENCODING, "gzip"); + channel.writeInbound(interim, LastHttpContent.EMPTY_LAST_CONTENT); + HttpResponse passed = channel.readInbound(); + assertEquals("gzip", passed.headers().get(HttpHeaderNames.CONTENT_ENCODING)); + assertInstanceOf(LastHttpContent.class, channel.readInbound()); + byte[] encoded = gzip(TEXT); + assertDecoded(TEXT, exchange(channel, "gzip", encoded, new int[]{encoded.length}, false)); + channel.finishAndReleaseAll(); + } + + @Test + void limitFailsTheResponseOnceTheBodyPassesIt() { + byte[] encoded = gzip(ZEROS); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, ZEROS.length - 1), allocator()); + Outcome outcome = exchange(channel, "gzip", encoded, sizes(encoded.length, 100), true); + assertInstanceOf(DecompressionException.class, outcome.failure); + assertEquals("HTTP/1.1 response body exceeds the maximum decompressed size of " + (ZEROS.length - 1) + " bytes", + outcome.failure.getMessage()); + assertTrue(outcome.body.size() < ZEROS.length); + // The limit counts one response, not the connection. + assertDecoded(TEXT, exchange(channel, "gzip", gzip(TEXT), new int[]{50, 56}, true)); + channel.finishAndReleaseAll(); + } + + @Test + void limitAllowsABodyOfExactlyTheLimit() { + byte[] encoded = gzip(ZEROS); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, ZEROS.length), allocator()); + assertDecoded(ZEROS, exchange(channel, "gzip", encoded, new int[]{encoded.length}, false)); + channel.finishAndReleaseAll(); + } + + @Test + void removingTheHandlerMidResponseReleasesWhatItHolds() { + UnpooledByteBufAllocator allocator = allocator(); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator); + byte[] encoded = gzipWithHeaderFields(TEXT, FNAME); + channel.writeInbound(response("gzip", -1)); + // Five bytes of a header are held back until the rest arrives. + channel.writeInbound(new DefaultHttpContent(Unpooled.wrappedBuffer(encoded, 0, 5))); + assertTrue(allocator.metric().usedHeapMemory() > 0); + channel.pipeline().removeFirst(); + assertEquals(0, allocator.metric().usedHeapMemory()); + ReferenceCountUtil.release(channel.readInbound()); + channel.finishAndReleaseAll(); + } + + @Test + void requestsReadsLikeNettyWhenAutoReadIsOff() { + byte[] encoded = gzipWithHeaderFields(TEXT, FNAME | FCOMMENT); + int[] sizes = sizes(encoded.length, 3); + assertEquals(readsRequested(nettyDecompressor(), encoded, sizes), + readsRequested(new Http1ContentDecompressor(false, 0), encoded, sizes)); + } + + private static List readsRequested(ChannelHandler decompressor, byte[] encoded, int[] sizes) { + int[] reads = new int[1]; + EmbeddedChannel channel = new EmbeddedChannel(); + channel.config().setAutoRead(false); + channel.pipeline().addLast(new ChannelOutboundHandlerAdapter() { + @Override + public void read(ChannelHandlerContext ctx) { + reads[0]++; + } + }); + channel.pipeline().addLast(decompressor); + List perMessage = new ArrayList<>(); + Supplier step = () -> { + int before = reads[0]; + channel.pipeline().fireChannelReadComplete(); + for (Object msg; (msg = channel.readInbound()) != null; ) { + ReferenceCountUtil.release(msg); + } + return reads[0] - before; + }; + channel.pipeline().fireChannelRead(response("gzip", -1)); + perMessage.add(step.get()); + int offset = 0; + for (int i = 0; i < sizes.length; i++) { + ByteBuf content = Unpooled.wrappedBuffer(encoded, offset, sizes[i]); + offset += sizes[i]; + channel.pipeline().fireChannelRead(i == sizes.length - 1 ? new DefaultLastHttpContent(content) + : new DefaultHttpContent(content)); + perMessage.add(step.get()); + } + channel.finishAndReleaseAll(); + assertTrue(perMessage.contains(1) && perMessage.contains(0), perMessage.toString()); + return perMessage; + } +} From 14851921c3f7b3fac4100a3cf8dc298f33a79101 Mon Sep 17 00:00:00 2001 From: Pavel Ptashyts <49400901+pavel-ptashyts@users.noreply.github.com> Date: Thu, 1 Oct 2026 15:12:40 +0200 Subject: [PATCH 2/8] Pool inflaters per event loop, survive removal An Inflater kept per connection pins zlib's native window on every idle pooled connection, so a client with thousands of keep-alive connections would hold that memory for nothing. Inflaters now live in a small per-event-loop pool between responses. A handler further on may remove this one while a decoded chunk is being forwarded; decoding now stops there instead of carrying on with state that was already handed back. Tests cover removal mid-chunk, close mid-response, a full response, encodings left to Netty alternating with gzip on one connection, composite input, the deflate limit and more read-request shapes. Claude Code on behalf of Pavel Ptashyts Co-Authored-By: Claude Opus 5.5 --- .../handler/Http1ContentDecompressor.java | 128 ++++++---- .../handler/Http1ContentDecompressorTest.java | 228 +++++++++++++++--- 2 files changed, 268 insertions(+), 88 deletions(-) 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 aa1fab8c0..fee7c74a9 100644 --- a/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java +++ b/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java @@ -33,8 +33,10 @@ import io.netty.handler.codec.http.HttpResponse; import io.netty.handler.codec.http.LastHttpContent; import io.netty.util.ReferenceCountUtil; +import io.netty.util.concurrent.FastThreadLocal; import org.jetbrains.annotations.Nullable; +import java.util.ArrayDeque; import java.util.List; import java.util.zip.CRC32; import java.util.zip.DataFormatException; @@ -55,12 +57,11 @@ * {@code gzip}, {@code x-gzip}, {@code deflate} and {@code x-deflate} named by {@code Content-Encoding} are * inflated here directly. {@link HttpContentDecompressor} builds an {@link EmbeddedChannel} and a fresh * {@link Inflater} for every such response, and its {@code JdkZlibDecoder} copies each direct input buffer - * onto the heap before inflating it. This handler sits in the pipeline of one connection, whose responses - * arrive one after another, so it keeps one {@link Inflater} per wrapper for the life of the connection, - * resets it between responses, and feeds it the input buffer as it is. The format handling (header and - * trailer checks, concatenated gzip members, zlib-or-raw detection for {@code deflate}, a truncated stream - * ending quietly) follows {@code JdkZlibDecoder} as {@link HttpContentDecompressor} configures it, so the - * decoded body and the errors are unchanged. + * onto the heap before inflating it. This handler borrows a reset {@link Inflater} from a small pool kept + * per event loop for the length of one response, and feeds it the input buffer as it is. The format + * handling (header and trailer checks, concatenated gzip members, zlib-or-raw detection for + * {@code deflate}, a truncated stream ending quietly) follows {@code JdkZlibDecoder} as + * {@link HttpContentDecompressor} configures it, so the decoded body and the errors are unchanged. *

* Every other encoding still goes through {@link HttpContentDecompressor}. There the counting is done * inside the decoder's own {@link EmbeddedChannel} rather than around {@code HttpContentDecoder#decode}: @@ -98,9 +99,12 @@ private enum GzipState { // Maximum cumulative decompressed bytes for one response; 0 disables the limit. private final long maxDecompressedBytes; - // Reused by every response on this connection and released when the handler leaves the pipeline. - private @Nullable Inflater nowrapInflater; - private @Nullable Inflater zlibInflater; + // Inflaters between responses, per event loop. Holding them per connection instead would pin zlib's + // native window on every idle pooled connection. + private static final int MAX_IDLE_INFLATERS = 16; + private static final FastThreadLocal> IDLE_NOWRAP_INFLATERS = idleInflaters(); + private static final FastThreadLocal> IDLE_ZLIB_INFLATERS = idleInflaters(); + private final CRC32 crc = new CRC32(); // The response being inflated by this class rather than by HttpContentDecompressor. @@ -108,6 +112,7 @@ private enum GzipState { private boolean gzip; // Null until the first two bytes of a deflate body have told zlib from raw deflate. private @Nullable Inflater inflater; + private boolean nowrap; private boolean failed; private boolean finished; private GzipState gzipState = GzipState.HEADER_START; @@ -143,7 +148,7 @@ protected void decode(ChannelHandlerContext ctx, HttpObject msg, List ou lastDecodedBySuper = false; needRead = true; startResponse(ctx, response, contentEncoding); - if (response instanceof HttpContent) { + if (response instanceof HttpContent && !ctx.isRemoved()) { decodeContent(ctx, (HttpContent) response); } return; @@ -181,7 +186,7 @@ public void channelReadComplete(ChannelHandlerContext ctx) throws Exception { @Override public void handlerRemoved(ChannelHandlerContext ctx) throws Exception { try { - releaseInflaters(); + endResponse(); } finally { super.handlerRemoved(ctx); } @@ -189,8 +194,11 @@ public void handlerRemoved(ChannelHandlerContext ctx) throws Exception { @Override public void channelInactive(ChannelHandlerContext ctx) throws Exception { - releaseInflaters(); - super.channelInactive(ctx); + try { + endResponse(); + } finally { + super.channelInactive(ctx); + } } @Override @@ -224,11 +232,12 @@ private static boolean isDeflate(String contentEncoding) { /** * Rewrites the headers exactly as {@code HttpContentDecoder} does once it has installed a decoder. */ - private void startResponse(ChannelHandlerContext ctx, HttpResponse response, String contentEncoding) throws Exception { + private void startResponse(ChannelHandlerContext ctx, HttpResponse response, String contentEncoding) + throws Exception { inflating = true; gzip = isGzip(contentEncoding); if (gzip) { - inflater = nowrapInflater(); + acquireInflater(true); } HttpHeaders headers = response.headers(); @@ -271,6 +280,9 @@ private void decodeContent(ChannelHandlerContext ctx, HttpContent content) { failed = true; throw t; } + if (ctx.isRemoved()) { + return; + } } if (content instanceof LastHttpContent) { endResponse(); @@ -290,11 +302,16 @@ private void inflate(ChannelHandlerContext ctx, ByteBuf content) { in = pending; } - int readable; - do { - readable = in.readableBytes(); + for (;;) { + int readable = in.readableBytes(); inflateOnce(ctx, in); - } while (in.isReadable() && in.readableBytes() != readable); + if (ctx.isRemoved()) { + return; + } + if (!in.isReadable() || in.readableBytes() == readable) { + break; + } + } if (pending != null) { if (pending.isReadable()) { @@ -326,8 +343,7 @@ private void inflateOnce(ChannelHandlerContext ctx, ByteBuf in) { return; } // RFC 1950 says deflate is zlib-wrapped, but some servers send raw deflate. - inflater = looksLikeZlib(in.getShort(in.readerIndex())) ? zlibInflater() : nowrapInflater(); - this.inflater = inflater; + inflater = acquireInflater(!looksLikeZlib(in.getShort(in.readerIndex()))); } if (gzip && gzipState != GzipState.HEADER_END) { @@ -373,7 +389,8 @@ private void inflateOnce(ChannelHandlerContext ctx, ByteBuf in) { if (outputLength > 0) { decompressedBytes += outputLength; if (maxDecompressedBytes > 0 && decompressedBytes > maxDecompressedBytes) { - throw new DecompressionException("HTTP/1.1 response body exceeds the maximum decompressed size of " + throw new DecompressionException( + "HTTP/1.1 response body exceeds the maximum decompressed size of " + maxDecompressedBytes + " bytes"); } decompressed.writerIndex(writerIndex + outputLength); @@ -384,9 +401,15 @@ private void inflateOnce(ChannelHandlerContext ctx, ByteBuf in) { ByteBuf buffer = decompressed; decompressed = null; fireContent(ctx, buffer); + if (ctx.isRemoved()) { + // A handler further on took this one out of the pipeline, which already gave the + // inflater back. + return; + } } } else if (inflater.needsDictionary()) { - throw new DecompressionException("decompression failure, unable to set dictionary as non was specified"); + throw new DecompressionException( + "decompression failure, unable to set dictionary as non was specified"); } if (inflater.finished()) { @@ -441,7 +464,8 @@ private boolean readGzipHeader(ByteBuf in) { int method = in.readUnsignedByte(); if (method != Deflater.DEFLATED) { - throw new DecompressionException("Unsupported compression method " + method + " in the GZIP header"); + throw new DecompressionException( + "Unsupported compression method " + method + " in the GZIP header"); } crc.update(method); @@ -499,7 +523,8 @@ private boolean readGzipHeader(ByteBuf in) { int crc16Value = in.readUnsignedShortLE(); int readCrc16 = (int) (crc.getValue() & 0xFFFF); if (crc16Value != readCrc16) { - throw new DecompressionException("CRC16 value mismatch. Expected: " + crc16Value + ", Got: " + readCrc16); + throw new DecompressionException( + "CRC16 value mismatch. Expected: " + crc16Value + ", Got: " + readCrc16); } } crc.reset(); @@ -544,7 +569,8 @@ private boolean readGzipFooter(ByteBuf in, Inflater inflater) { int dataLength = in.readIntLE(); int readLength = inflater.getTotalOut(); if (dataLength != readLength) { - throw new DecompressionException("Number of bytes mismatch. Expected: " + dataLength + ", Got: " + readLength); + throw new DecompressionException( + "Number of bytes mismatch. Expected: " + dataLength + ", Got: " + readLength); } inflater.reset(); crc.reset(); @@ -566,30 +592,44 @@ private static boolean looksLikeZlib(short cmfFlg) { return (cmfFlg & 0x7800) == 0x7800 && cmfFlg % 31 == 0; } - private Inflater nowrapInflater() { - Inflater inflater = nowrapInflater; - if (inflater == null) { - inflater = new Inflater(true); - nowrapInflater = inflater; - } - return inflater; + private static FastThreadLocal> idleInflaters() { + return new FastThreadLocal>() { + @Override + protected ArrayDeque initialValue() { + return new ArrayDeque<>(); + } + + @Override + protected void onRemoval(ArrayDeque inflaters) { + for (Inflater inflater; (inflater = inflaters.pollFirst()) != null; ) { + inflater.end(); + } + } + }; } - private Inflater zlibInflater() { - Inflater inflater = zlibInflater; + private Inflater acquireInflater(boolean nowrap) { + Inflater inflater = (nowrap ? IDLE_NOWRAP_INFLATERS : IDLE_ZLIB_INFLATERS).get().pollLast(); if (inflater == null) { - inflater = new Inflater(); - zlibInflater = inflater; + inflater = new Inflater(nowrap); } + this.inflater = inflater; + this.nowrap = nowrap; return inflater; } private void endResponse() { Inflater inflater = this.inflater; if (inflater != null) { + this.inflater = null; // Also drops the inflater's reference to the last input buffer, which is released after decode. inflater.reset(); - this.inflater = null; + ArrayDeque idle = (nowrap ? IDLE_NOWRAP_INFLATERS : IDLE_ZLIB_INFLATERS).get(); + if (idle.size() < MAX_IDLE_INFLATERS) { + idle.offerLast(inflater); + } else { + inflater.end(); + } } ByteBuf pending = this.pending; if (pending != null) { @@ -606,18 +646,6 @@ private void endResponse() { crc.reset(); } - private void releaseInflaters() { - endResponse(); - if (nowrapInflater != null) { - nowrapInflater.end(); - nowrapInflater = null; - } - if (zlibInflater != null) { - zlibInflater.end(); - zlibInflater = null; - } - } - /** * 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 diff --git a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java index 59fb0c4f3..54dd83e33 100644 --- a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java +++ b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java @@ -20,13 +20,17 @@ import io.netty.buffer.UnpooledByteBufAllocator; import io.netty.channel.ChannelHandler; import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelInboundHandlerAdapter; import io.netty.channel.ChannelOutboundHandlerAdapter; import io.netty.channel.embedded.EmbeddedChannel; import io.netty.handler.codec.compression.DecompressionException; +import io.netty.handler.codec.compression.SnappyFrameEncoder; +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.HttpContent; +import io.netty.handler.codec.http.FullHttpResponse; import io.netty.handler.codec.http.HttpContentDecompressor; import io.netty.handler.codec.http.HttpHeaderNames; import io.netty.handler.codec.http.HttpHeaderValues; @@ -290,8 +294,19 @@ private static HttpResponse response(String contentEncoding, long contentLength) return response; } - private static ByteBuf buffer(byte[] bytes, int offset, int length, boolean direct) { - ByteBuf buf = direct ? Unpooled.directBuffer(Math.max(length, 1)) : Unpooled.buffer(Math.max(length, 1)); + private enum Kind { HEAP, DIRECT, COMPOSITE } + + private static ByteBuf buffer(byte[] bytes, int offset, int length, Kind kind) { + if (kind == Kind.COMPOSITE && length > 1) { + // Two direct components, so the inflater cannot be handed one NIO buffer. + int half = length / 2; + return Unpooled.compositeBuffer(2) + .addComponent(true, Unpooled.directBuffer(half).writeBytes(bytes, offset, half)) + .addComponent(true, Unpooled.directBuffer(length - half) + .writeBytes(bytes, offset + half, length - half)); + } + int capacity = Math.max(length, 1); + ByteBuf buf = kind == Kind.HEAP ? Unpooled.buffer(capacity) : Unpooled.directBuffer(capacity); return buf.writeBytes(bytes, offset, length); } @@ -299,13 +314,13 @@ private static ByteBuf buffer(byte[] bytes, int offset, int length, boolean dire * Sends one response through the channel and collects what comes out, stopping at the first error. */ private static Outcome exchange(EmbeddedChannel channel, String contentEncoding, byte[] encoded, int[] sizes, - boolean direct) { + Kind kind) { Outcome outcome = new Outcome(); List messages = new ArrayList<>(); messages.add(response(contentEncoding, encoded.length)); int offset = 0; for (int i = 0; i < sizes.length; i++) { - ByteBuf content = buffer(encoded, offset, sizes[i], direct); + ByteBuf content = buffer(encoded, offset, sizes[i], kind); offset += sizes[i]; messages.add(i == sizes.length - 1 ? new DefaultLastHttpContent(content) : new DefaultHttpContent(content)); } @@ -372,13 +387,13 @@ void decodesLikeNettyForEveryBodyAndSplit() { int checked = 0; for (Body body : bodies()) { for (int[] sizes : splits(body.encoded.length)) { - for (boolean direct : new boolean[]{false, true}) { - String scenario = body.name + ", " + sizes.length + " chunks, " + (direct ? "direct" : "heap"); + for (Kind kind : Kind.values()) { + String scenario = body.name + ", " + sizes.length + " chunks, " + kind; UnpooledByteBufAllocator ahcAllocator = allocator(); EmbeddedChannel ahc = channel(new Http1ContentDecompressor(false, 0), ahcAllocator); EmbeddedChannel netty = channel(nettyDecompressor(), allocator()); - Outcome expected = exchange(netty, body.contentEncoding, body.encoded, sizes, direct); - Outcome actual = exchange(ahc, body.contentEncoding, body.encoded, sizes, direct); + Outcome expected = exchange(netty, body.contentEncoding, body.encoded, sizes, kind); + Outcome actual = exchange(ahc, body.contentEncoding, body.encoded, sizes, kind); assertSameOutcome(expected, actual, scenario); ahc.finishAndReleaseAll(); @@ -413,12 +428,57 @@ private static void assertSameOutcome(Outcome expected, Outcome actual, String s @Test void forwardsOnePartPerChunkUnlessItInflatesPastTheThreshold() { - byte[] encoded = gzip(TEXT); EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator()); - Outcome outcome = exchange(channel, "gzip", encoded, new int[]{encoded.length}, true); - assertNull(outcome.failure); - assertEquals(1, outcome.parts); - assertArrayEquals(TEXT, outcome.body.toByteArray()); + byte[] text = gzip(TEXT); + Outcome small = exchange(channel, "gzip", text, new int[]{text.length}, Kind.DIRECT); + assertDecoded(TEXT, small); + assertEquals(1, small.parts); + // 2 MiB from one chunk leaves in 64 KiB parts rather than as one buffer grown to 2 MiB. + byte[] zeros = gzip(ZEROS); + Outcome large = exchange(channel, "gzip", zeros, new int[]{zeros.length}, Kind.DIRECT); + assertDecoded(ZEROS, large); + assertEquals(ZEROS.length / (64 * 1024), large.parts); + channel.finishAndReleaseAll(); + } + + @Test + void stopsWhenAHandlerFurtherOnRemovesItMidChunk() { + UnpooledByteBufAllocator allocator = allocator(); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator); + int[] parts = new int[1]; + List afterRemoval = new ArrayList<>(); + channel.pipeline().addLast(new ChannelInboundHandlerAdapter() { + @Override + public void channelRead(ChannelHandlerContext ctx, Object msg) { + if (parts[0] > 0) { + afterRemoval.add(msg); + } else if (msg instanceof HttpContent && ((HttpContent) msg).content().isReadable()) { + parts[0]++; + ctx.pipeline().remove(Http1ContentDecompressor.class); + } + ReferenceCountUtil.release(msg); + } + }); + byte[] encoded = gzip(ZEROS); + channel.writeInbound(response("gzip", encoded.length)); + channel.writeInbound(new DefaultLastHttpContent(Unpooled.wrappedBuffer(encoded))); + assertEquals(1, parts[0]); + assertEquals(List.of(), afterRemoval); + assertNull(channel.pipeline().get(Http1ContentDecompressor.class)); + channel.finishAndReleaseAll(); + assertEquals(0, allocator.metric().usedHeapMemory()); + } + + @Test + void closingTheChannelMidResponseReleasesWhatItHolds() { + UnpooledByteBufAllocator allocator = allocator(); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator); + byte[] encoded = gzipWithHeaderFields(TEXT, FCOMMENT); + channel.writeInbound(response("gzip", -1)); + channel.writeInbound(new DefaultHttpContent(Unpooled.wrappedBuffer(encoded, 0, 7))); + assertTrue(allocator.metric().usedHeapMemory() > 0); + channel.close(); + assertEquals(0, allocator.metric().usedHeapMemory()); channel.finishAndReleaseAll(); } @@ -430,20 +490,20 @@ void reusesTheInflaterAcrossResponsesOnOneConnection() { byte[] zlib = deflate(RANDOM, false); byte[] raw = deflate(TEXT, true); - assertDecoded(TEXT, exchange(channel, "gzip", gzipText, new int[]{gzipText.length}, true)); - Outcome identity = exchange(channel, null, TEXT, new int[]{100, TEXT.length - 100}, false); + assertDecoded(TEXT, exchange(channel, "gzip", gzipText, new int[]{gzipText.length}, Kind.DIRECT)); + Outcome identity = exchange(channel, null, TEXT, new int[]{100, TEXT.length - 100}, Kind.HEAP); assertDecoded(TEXT, identity); assertEquals(String.valueOf(TEXT.length), identity.head.headers().get(HttpHeaderNames.CONTENT_LENGTH)); - assertDecoded(RANDOM, exchange(channel, "deflate", zlib, sizes(zlib.length, 4000), true)); - assertDecoded(TEXT, exchange(channel, "deflate", raw, sizes(raw.length, 3), false)); + assertDecoded(RANDOM, exchange(channel, "deflate", zlib, sizes(zlib.length, 4000), Kind.DIRECT)); + assertDecoded(TEXT, exchange(channel, "deflate", raw, sizes(raw.length, 3), Kind.HEAP)); // A response cut short leaves state behind that the next one must not see. - Outcome cut = exchange(channel, "gzip", Arrays.copyOf(gzipText, 40), new int[]{40}, true); + Outcome cut = exchange(channel, "gzip", Arrays.copyOf(gzipText, 40), new int[]{40}, Kind.DIRECT); assertNull(cut.failure); - assertDecoded(TEXT, exchange(channel, "gzip", gzipText, sizes(gzipText.length, 5), true)); + assertDecoded(TEXT, exchange(channel, "gzip", gzipText, sizes(gzipText.length, 5), Kind.DIRECT)); // So does a corrupt one. - Outcome corrupt = exchange(channel, "gzip", flipped(gzipText, -8), new int[]{gzipText.length}, true); + Outcome corrupt = exchange(channel, "gzip", flipped(gzipText, -8), new int[]{gzipText.length}, Kind.DIRECT); assertInstanceOf(DecompressionException.class, corrupt.failure); - assertDecoded(TEXT, exchange(channel, "gzip", gzipText, new int[]{gzipText.length}, true)); + assertDecoded(TEXT, exchange(channel, "gzip", gzipText, new int[]{gzipText.length}, Kind.DIRECT)); channel.finishAndReleaseAll(); assertEquals(0, allocator.metric().usedHeapMemory()); @@ -475,20 +535,89 @@ private static void assertDecoded(byte[] expected, Outcome outcome) { @Test void rewritesHeadersLikeNetty() { - Outcome outcome = exchange(channel(new Http1ContentDecompressor(false, 0), allocator()), "gzip", gzip(TEXT), - new int[]{10, 20}, false); + byte[] encoded = gzip(TEXT); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator()); + Outcome outcome = exchange(channel, "gzip", encoded, sizes(encoded.length, 50), Kind.HEAP); + channel.finishAndReleaseAll(); + assertDecoded(TEXT, outcome); assertNull(outcome.head.headers().get(HttpHeaderNames.CONTENT_ENCODING)); assertNull(outcome.head.headers().get(HttpHeaderNames.CONTENT_LENGTH)); - assertEquals(HttpHeaderValues.CHUNKED.toString(), outcome.head.headers().get(HttpHeaderNames.TRANSFER_ENCODING)); + assertEquals(HttpHeaderValues.CHUNKED.toString(), + outcome.head.headers().get(HttpHeaderNames.TRANSFER_ENCODING)); } @Test void keepEncodingHeaderLeavesContentEncoding() { - Outcome outcome = exchange(channel(new Http1ContentDecompressor(true, 0), allocator()), "X-Gzip", gzip(RANDOM), - sizes(gzip(RANDOM).length, 2000), false); + byte[] encoded = gzip(RANDOM); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(true, 0), allocator()); + Outcome outcome = exchange(channel, "X-Gzip", encoded, sizes(encoded.length, 2000), Kind.HEAP); + channel.finishAndReleaseAll(); + assertDecoded(RANDOM, outcome); assertEquals("X-Gzip", outcome.head.headers().get(HttpHeaderNames.CONTENT_ENCODING)); assertNull(outcome.head.headers().get(HttpHeaderNames.CONTENT_LENGTH)); - assertArrayEquals(RANDOM, outcome.body.toByteArray()); + } + + @Test + void fullResponseIsSplitIntoHeadAndBodyLikeNetty() { + byte[] encoded = gzip(RANDOM); + List outcomes = new ArrayList<>(); + for (ChannelHandler decompressor : new ChannelHandler[]{nettyDecompressor(), + new Http1ContentDecompressor(false, 0)}) { + EmbeddedChannel channel = channel(decompressor, allocator()); + FullHttpResponse response = new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK, + Unpooled.wrappedBuffer(encoded)); + response.headers().set(HttpHeaderNames.CONTENT_ENCODING, "gzip"); + response.headers().set(HttpHeaderNames.CONTENT_LENGTH, encoded.length); + response.trailingHeaders().set("x-trailer", "1"); + channel.writeInbound(response); + Outcome outcome = new Outcome(); + drain(channel, outcome); + channel.finishAndReleaseAll(); + outcomes.add(outcome); + } + assertSameOutcome(outcomes.get(0), outcomes.get(1), "full response"); + assertDecoded(RANDOM, outcomes.get(1)); + } + + @Test + void alternatesWithEncodingsLeftToNettyOnOneConnection() { + byte[] snappyText = snappy(TEXT); + byte[] snappyRandom = snappy(RANDOM); + byte[] gzipText = gzip(TEXT); + byte[] zlib = deflate(RANDOM, false); + UnpooledByteBufAllocator allocator = allocator(); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator); + + assertDecoded(TEXT, exchange(channel, "snappy", snappyText, sizes(snappyText.length, 10), Kind.HEAP)); + assertDecoded(TEXT, exchange(channel, "gzip", gzipText, sizes(gzipText.length, 10), Kind.DIRECT)); + assertDecoded(RANDOM, exchange(channel, "snappy", snappyRandom, sizes(snappyRandom.length, 5000), + Kind.DIRECT)); + assertDecoded(RANDOM, exchange(channel, "deflate", zlib, sizes(zlib.length, 5000), Kind.DIRECT)); + // A snappy response that never ends leaves Netty's decoder behind; the next gzip response must not care. + channel.writeInbound(response("snappy", -1), new DefaultHttpContent(Unpooled.wrappedBuffer(snappyRandom, 0, 3000))); + drain(channel, new Outcome()); + assertDecoded(TEXT, exchange(channel, "gzip", gzipText, new int[]{gzipText.length}, Kind.HEAP)); + assertDecoded(TEXT, exchange(channel, "snappy", snappyText, new int[]{snappyText.length}, Kind.HEAP)); + + channel.finishAndReleaseAll(); + assertEquals(0, allocator.metric().usedHeapMemory()); + } + + private static byte[] snappy(byte[] payload) { + EmbeddedChannel encoder = new EmbeddedChannel(new SnappyFrameEncoder()); + encoder.writeOutbound(Unpooled.wrappedBuffer(payload)); + ByteArrayOutputStream bos = new ByteArrayOutputStream(); + for (ByteBuf buf; (buf = encoder.readOutbound()) != null; ) { + try { + buf.readBytes(bos, buf.readableBytes()); + } catch (Exception e) { + throw new IllegalStateException(e); + } finally { + buf.release(); + } + } + encoder.finishAndReleaseAll(); + return bos.toByteArray(); } @Test @@ -520,7 +649,7 @@ void continueResponseIsPassedThrough() { assertEquals("gzip", passed.headers().get(HttpHeaderNames.CONTENT_ENCODING)); assertInstanceOf(LastHttpContent.class, channel.readInbound()); byte[] encoded = gzip(TEXT); - assertDecoded(TEXT, exchange(channel, "gzip", encoded, new int[]{encoded.length}, false)); + assertDecoded(TEXT, exchange(channel, "gzip", encoded, new int[]{encoded.length}, Kind.HEAP)); channel.finishAndReleaseAll(); } @@ -528,13 +657,23 @@ void continueResponseIsPassedThrough() { void limitFailsTheResponseOnceTheBodyPassesIt() { byte[] encoded = gzip(ZEROS); EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, ZEROS.length - 1), allocator()); - Outcome outcome = exchange(channel, "gzip", encoded, sizes(encoded.length, 100), true); + Outcome outcome = exchange(channel, "gzip", encoded, sizes(encoded.length, 100), Kind.DIRECT); assertInstanceOf(DecompressionException.class, outcome.failure); assertEquals("HTTP/1.1 response body exceeds the maximum decompressed size of " + (ZEROS.length - 1) + " bytes", outcome.failure.getMessage()); assertTrue(outcome.body.size() < ZEROS.length); // The limit counts one response, not the connection. - assertDecoded(TEXT, exchange(channel, "gzip", gzip(TEXT), new int[]{50, 56}, true)); + assertDecoded(TEXT, exchange(channel, "gzip", gzip(TEXT), new int[]{50, 56}, Kind.DIRECT)); + channel.finishAndReleaseAll(); + } + + @Test + void limitAppliesToDeflateToo() { + byte[] encoded = deflate(ZEROS, false); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 100_000), allocator()); + Outcome outcome = exchange(channel, "deflate", encoded, new int[]{encoded.length}, Kind.HEAP); + assertInstanceOf(DecompressionException.class, outcome.failure); + assertTrue(outcome.body.size() <= 100_000); channel.finishAndReleaseAll(); } @@ -542,7 +681,7 @@ void limitFailsTheResponseOnceTheBodyPassesIt() { void limitAllowsABodyOfExactlyTheLimit() { byte[] encoded = gzip(ZEROS); EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, ZEROS.length), allocator()); - assertDecoded(ZEROS, exchange(channel, "gzip", encoded, new int[]{encoded.length}, false)); + assertDecoded(ZEROS, exchange(channel, "gzip", encoded, new int[]{encoded.length}, Kind.HEAP)); channel.finishAndReleaseAll(); } @@ -563,13 +702,27 @@ void removingTheHandlerMidResponseReleasesWhatItHolds() { @Test void requestsReadsLikeNettyWhenAutoReadIsOff() { - byte[] encoded = gzipWithHeaderFields(TEXT, FNAME | FCOMMENT); - int[] sizes = sizes(encoded.length, 3); - assertEquals(readsRequested(nettyDecompressor(), encoded, sizes), - readsRequested(new Http1ContentDecompressor(false, 0), encoded, sizes)); + byte[] header = gzipWithHeaderFields(TEXT, FNAME | FCOMMENT); + byte[] zeros = gzip(ZEROS); + byte[] raw = deflate(TEXT, true); + byte[][] bodies = {header, zeros, raw}; + String[] encodings = {"gzip", "gzip", "deflate"}; + int[] chunk = {3, 97, 1}; + for (int i = 0; i < bodies.length; i++) { + int[] sizes = sizes(bodies[i].length, chunk[i]); + List expected = readsRequested(nettyDecompressor(), encodings[i], bodies[i], sizes); + assertEquals(expected, + readsRequested(new Http1ContentDecompressor(false, 0), encodings[i], bodies[i], sizes), + encodings[i] + " in chunks of " + chunk[i]); + if (i == 0) { + // Chunks holding only header bytes ask for a read; the head and chunks with output do not. + assertTrue(expected.contains(1) && expected.contains(0), expected.toString()); + } + } } - private static List readsRequested(ChannelHandler decompressor, byte[] encoded, int[] sizes) { + private static List readsRequested(ChannelHandler decompressor, String contentEncoding, byte[] encoded, + int[] sizes) { int[] reads = new int[1]; EmbeddedChannel channel = new EmbeddedChannel(); channel.config().setAutoRead(false); @@ -589,7 +742,7 @@ public void read(ChannelHandlerContext ctx) { } return reads[0] - before; }; - channel.pipeline().fireChannelRead(response("gzip", -1)); + channel.pipeline().fireChannelRead(response(contentEncoding, -1)); perMessage.add(step.get()); int offset = 0; for (int i = 0; i < sizes.length; i++) { @@ -600,7 +753,6 @@ public void read(ChannelHandlerContext ctx) { perMessage.add(step.get()); } channel.finishAndReleaseAll(); - assertTrue(perMessage.contains(1) && perMessage.contains(0), perMessage.toString()); return perMessage; } } From 67c668c2571cb183eebd01bee58eda0dfe4b5976 Mon Sep 17 00:00:00 2001 From: Pavel Ptashyts <49400901+pavel-ptashyts@users.noreply.github.com> Date: Thu, 1 Oct 2026 15:22:39 +0200 Subject: [PATCH 3/8] Stop decoding once a response ends mid-forward Checking ctx.isRemoved() after forwarding decoded content only caught a handler further on removing this one. Any way of ending the response from that callback, such as closing the channel, should stop the decoding that made the call, so each dropped response now bumps a counter that the decoding loop checks. The forwarding threshold reads the same system property as Netty's JdkZlibDecoder, and the differential test now also holds the number of decoded parts to Netty's. New tests cover the per-thread inflater pool and its cap, reuse of an inflater handed back by a removal, a close from a handler further on, and read requests for a read that mixes Netty's messages with this handler's, including a failed body. Claude Code on behalf of Pavel Ptashyts Co-Authored-By: Claude Opus 5.5 --- .../Http1ContentDecompressorBenchmark.java | 2 +- .../handler/Http1ContentDecompressor.java | 59 ++++--- .../handler/Http1ContentDecompressorTest.java | 155 +++++++++++++++++- 3 files changed, 191 insertions(+), 25 deletions(-) diff --git a/client/src/jmh/java/org/asynchttpclient/bench/Http1ContentDecompressorBenchmark.java b/client/src/jmh/java/org/asynchttpclient/bench/Http1ContentDecompressorBenchmark.java index 48c260cf6..c46df670f 100644 --- a/client/src/jmh/java/org/asynchttpclient/bench/Http1ContentDecompressorBenchmark.java +++ b/client/src/jmh/java/org/asynchttpclient/bench/Http1ContentDecompressorBenchmark.java @@ -49,7 +49,7 @@ /** * Decodes one gzip response per operation on a keep-alive connection, through Netty's * {@link HttpContentDecompressor} and through {@link Http1ContentDecompressor}, which inflates gzip without an - * {@code EmbeddedChannel} and with one {@code Inflater} per connection. + * {@code EmbeddedChannel} and with an {@code Inflater} borrowed from a per-event-loop pool for each response. *

* The compressed body arrives in direct pooled buffers of at most {@code chunkSize} bytes, as * {@code HttpClientCodec} hands it on under the default allocator. 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 fee7c74a9..deb1c9f97 100644 --- a/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java +++ b/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java @@ -34,6 +34,7 @@ import io.netty.handler.codec.http.LastHttpContent; import io.netty.util.ReferenceCountUtil; import io.netty.util.concurrent.FastThreadLocal; +import io.netty.util.internal.SystemPropertyUtil; import org.jetbrains.annotations.Nullable; import java.util.ArrayDeque; @@ -82,7 +83,14 @@ public class Http1ContentDecompressor extends HttpContentDecompressor { // Same floor and forwarding threshold as JdkZlibDecoder, so the body is cut into parts the same way. private static final int MIN_OUTPUT_BUFFER_SIZE = 512; - private static final int MAX_FORWARD_BYTES = 64 * 1024; + private static final int MAX_FORWARD_BYTES = + SystemPropertyUtil.getInt("io.netty.compression.defaultMaxForwardBytes", 64 * 1024); + + // Inflaters between responses, per event loop. Holding them per connection instead would pin zlib's + // native window on every idle pooled connection. + private static final int MAX_IDLE_INFLATERS = 16; + private static final FastThreadLocal> IDLE_NOWRAP_INFLATERS = idleInflaters(); + private static final FastThreadLocal> IDLE_ZLIB_INFLATERS = idleInflaters(); private enum GzipState { HEADER_START, @@ -99,14 +107,12 @@ private enum GzipState { // Maximum cumulative decompressed bytes for one response; 0 disables the limit. private final long maxDecompressedBytes; - // Inflaters between responses, per event loop. Holding them per connection instead would pin zlib's - // native window on every idle pooled connection. - private static final int MAX_IDLE_INFLATERS = 16; - private static final FastThreadLocal> IDLE_NOWRAP_INFLATERS = idleInflaters(); - private static final FastThreadLocal> IDLE_ZLIB_INFLATERS = idleInflaters(); - private final CRC32 crc = new CRC32(); + // Bumped whenever the state of a response is dropped. A handler further on may end the response from + // inside a call that forwards decoded content, by removing this handler or closing the channel; the + // decoding that made the call must then stop. + private int generation; // The response being inflated by this class rather than by HttpContentDecompressor. private boolean inflating; private boolean gzip; @@ -147,8 +153,9 @@ protected void decode(ChannelHandlerContext ctx, HttpObject msg, List ou if (isGzip(contentEncoding) || isDeflate(contentEncoding)) { lastDecodedBySuper = false; needRead = true; + int generation = this.generation; startResponse(ctx, response, contentEncoding); - if (response instanceof HttpContent && !ctx.isRemoved()) { + if (response instanceof HttpContent && generation == this.generation) { decodeContent(ctx, (HttpContent) response); } return; @@ -272,15 +279,18 @@ private void decodeContent(ChannelHandlerContext ctx, HttpContent content) { } ByteBuf in = content.content(); if (in.isReadable()) { + int generation = this.generation; try { - inflate(ctx, in); + inflate(ctx, in, generation); } catch (Throwable t) { - endResponse(); - inflating = true; - failed = true; + if (generation == this.generation) { + endResponse(); + inflating = true; + failed = true; + } throw t; } - if (ctx.isRemoved()) { + if (generation != this.generation) { return; } } @@ -294,7 +304,7 @@ private void decodeContent(ChannelHandlerContext ctx, HttpContent content) { } } - private void inflate(ChannelHandlerContext ctx, ByteBuf content) { + private void inflate(ChannelHandlerContext ctx, ByteBuf content, int generation) { ByteBuf in = content; ByteBuf pending = this.pending; if (pending != null) { @@ -304,8 +314,8 @@ private void inflate(ChannelHandlerContext ctx, ByteBuf content) { for (;;) { int readable = in.readableBytes(); - inflateOnce(ctx, in); - if (ctx.isRemoved()) { + inflateOnce(ctx, in, generation); + if (generation != this.generation) { return; } if (!in.isReadable() || in.readableBytes() == readable) { @@ -321,8 +331,8 @@ private void inflate(ChannelHandlerContext ctx, ByteBuf content) { this.pending = null; } } else if (in.isReadable()) { - // Only a gzip header or trailer that has not fully arrived is left over; the inflater consumes - // all the deflate data it is given. + // Only a gzip header or trailer that has not fully arrived, or the first byte of a deflate body, is + // left over; the inflater consumes all the deflate data it is given. this.pending = ctx.alloc().heapBuffer(in.readableBytes()).writeBytes(in); } } @@ -331,7 +341,7 @@ private void inflate(ChannelHandlerContext ctx, ByteBuf content) { * One step of {@code JdkZlibDecoder#decode}, which {@code ByteToMessageDecoder} would call again for as * long as it consumes input. */ - private void inflateOnce(ChannelHandlerContext ctx, ByteBuf in) { + private void inflateOnce(ChannelHandlerContext ctx, ByteBuf in, int generation) { if (finished) { in.skipBytes(in.readableBytes()); return; @@ -401,9 +411,8 @@ private void inflateOnce(ChannelHandlerContext ctx, ByteBuf in) { ByteBuf buffer = decompressed; decompressed = null; fireContent(ctx, buffer); - if (ctx.isRemoved()) { - // A handler further on took this one out of the pipeline, which already gave the - // inflater back. + if (generation != this.generation) { + // The inflater has already been given back. return; } } @@ -618,7 +627,13 @@ private Inflater acquireInflater(boolean nowrap) { return inflater; } + // Visible for tests. + static int idleInflaterCount(boolean nowrap) { + return (nowrap ? IDLE_NOWRAP_INFLATERS : IDLE_ZLIB_INFLATERS).get().size(); + } + private void endResponse() { + generation++; Inflater inflater = this.inflater; if (inflater != null) { this.inflater = null; diff --git a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java index 54dd83e33..2db6a166c 100644 --- a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java +++ b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java @@ -29,8 +29,8 @@ 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.FullHttpResponse; +import io.netty.handler.codec.http.HttpContent; import io.netty.handler.codec.http.HttpContentDecompressor; import io.netty.handler.codec.http.HttpHeaderNames; import io.netty.handler.codec.http.HttpHeaderValues; @@ -47,6 +47,7 @@ import java.util.Arrays; import java.util.List; import java.util.Random; +import java.util.concurrent.atomic.AtomicReference; import java.util.function.Supplier; import java.util.zip.CRC32; import java.util.zip.Deflater; @@ -418,6 +419,7 @@ private static void assertSameOutcome(Outcome expected, Outcome actual, String s assertNull(actual.failure, scenario + ": unexpected failure"); assertArrayEquals(expected.body.toByteArray(), actual.body.toByteArray(), scenario + ": body"); assertEquals(expected.ended, actual.ended, scenario + ": end of response"); + assertEquals(expected.parts, actual.parts, scenario + ": parts"); } else { assertNotNull(actual.failure, scenario + ": expected " + expected.failure); assertEquals(expected.failure.getClass(), actual.failure.getClass(), scenario + ": failure type"); @@ -594,7 +596,8 @@ void alternatesWithEncodingsLeftToNettyOnOneConnection() { Kind.DIRECT)); assertDecoded(RANDOM, exchange(channel, "deflate", zlib, sizes(zlib.length, 5000), Kind.DIRECT)); // A snappy response that never ends leaves Netty's decoder behind; the next gzip response must not care. - channel.writeInbound(response("snappy", -1), new DefaultHttpContent(Unpooled.wrappedBuffer(snappyRandom, 0, 3000))); + channel.writeInbound(response("snappy", -1), + new DefaultHttpContent(Unpooled.wrappedBuffer(snappyRandom, 0, 3000))); drain(channel, new Outcome()); assertDecoded(TEXT, exchange(channel, "gzip", gzipText, new int[]{gzipText.length}, Kind.HEAP)); assertDecoded(TEXT, exchange(channel, "snappy", snappyText, new int[]{snappyText.length}, Kind.HEAP)); @@ -755,4 +758,152 @@ public void read(ChannelHandlerContext ctx) { channel.finishAndReleaseAll(); return perMessage; } + + @Test + void closingTheChannelFromAHandlerFurtherOnMidChunkIsSafe() { + UnpooledByteBufAllocator allocator = allocator(); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator); + int[] parts = new int[1]; + channel.pipeline().addLast(new ChannelInboundHandlerAdapter() { + @Override + public void channelRead(ChannelHandlerContext ctx, Object msg) { + if (msg instanceof HttpContent && ((HttpContent) msg).content().isReadable() && parts[0]++ == 0) { + ctx.channel().close(); + } + ReferenceCountUtil.release(msg); + } + }); + byte[] encoded = gzip(ZEROS); + channel.writeInbound(response("gzip", encoded.length)); + channel.writeInbound(new DefaultLastHttpContent(Unpooled.wrappedBuffer(encoded))); + assertTrue(parts[0] >= 1); + channel.finishAndReleaseAll(); + assertEquals(0, allocator.metric().usedHeapMemory()); + } + + @Test + void inflatersGoBackToABoundedPoolPerThread() throws Exception { + onFreshThread(() -> { + byte[] encoded = gzip(TEXT); + List channels = new ArrayList<>(); + for (int i = 0; i < 20; i++) { + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator()); + channel.writeInbound(response("gzip", -1), + new DefaultHttpContent(Unpooled.wrappedBuffer(encoded, 0, 20))); + channels.add(channel); + } + assertEquals(0, Http1ContentDecompressor.idleInflaterCount(true)); + for (EmbeddedChannel channel : channels) { + channel.writeInbound(new DefaultLastHttpContent( + Unpooled.wrappedBuffer(encoded, 20, encoded.length - 20))); + channel.finishAndReleaseAll(); + } + assertEquals(16, Http1ContentDecompressor.idleInflaterCount(true)); + + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator()); + assertDecoded(TEXT, exchange(channel, "gzip", encoded, new int[]{encoded.length}, Kind.HEAP)); + assertEquals(16, Http1ContentDecompressor.idleInflaterCount(true)); + byte[] zlib = deflate(TEXT, false); + assertDecoded(TEXT, exchange(channel, "deflate", zlib, new int[]{zlib.length}, Kind.HEAP)); + assertEquals(1, Http1ContentDecompressor.idleInflaterCount(false)); + channel.finishAndReleaseAll(); + }); + } + + @Test + void anInflaterGivenBackByARemovalMidChunkDecodesTheNextResponse() throws Exception { + onFreshThread(() -> { + byte[] zeros = gzip(ZEROS); + EmbeddedChannel removed = channel(new Http1ContentDecompressor(false, 0), allocator()); + removed.pipeline().addLast(new ChannelInboundHandlerAdapter() { + @Override + public void channelRead(ChannelHandlerContext ctx, Object msg) { + if (msg instanceof HttpContent && ((HttpContent) msg).content().isReadable()) { + ctx.pipeline().remove(Http1ContentDecompressor.class); + } + ReferenceCountUtil.release(msg); + } + }); + removed.writeInbound(response("gzip", -1), new DefaultLastHttpContent(Unpooled.wrappedBuffer(zeros))); + removed.finishAndReleaseAll(); + assertEquals(1, Http1ContentDecompressor.idleInflaterCount(true)); + + EmbeddedChannel next = channel(new Http1ContentDecompressor(false, 0), allocator()); + assertDecoded(ZEROS, exchange(next, "gzip", zeros, sizes(zeros.length, 500), Kind.DIRECT)); + byte[] raw = deflate(RANDOM, true); + assertDecoded(RANDOM, exchange(next, "deflate", raw, sizes(raw.length, 3000), Kind.DIRECT)); + assertEquals(1, Http1ContentDecompressor.idleInflaterCount(true)); + next.finishAndReleaseAll(); + }); + } + + private static void onFreshThread(Runnable test) throws Exception { + AtomicReference failure = new AtomicReference<>(); + Thread thread = new Thread(() -> { + try { + test.run(); + } catch (Throwable t) { + failure.set(t); + } + }); + thread.start(); + thread.join(); + if (failure.get() instanceof Error) { + throw (Error) failure.get(); + } + if (failure.get() != null) { + throw new IllegalStateException(failure.get()); + } + } + + @Test + void requestsReadsLikeNettyWhenOneReadMixesItsMessagesWithNettys() { + byte[] snappy = snappy(TEXT); + byte[] gzipText = gzip(TEXT); + byte[] corrupt = withByte(gzip(TEXT), 0, 0); + List>> batches = List.of( + () -> List.of(response("snappy", -1), new DefaultLastHttpContent(Unpooled.wrappedBuffer(snappy)), + response("gzip", -1), new DefaultHttpContent(Unpooled.wrappedBuffer(gzipText, 0, 3))), + () -> List.of(new DefaultLastHttpContent(Unpooled.wrappedBuffer(gzipText, 3, gzipText.length - 3)), + response("snappy", -1), new DefaultHttpContent(Unpooled.wrappedBuffer(snappy, 0, 5))), + () -> List.of(new DefaultLastHttpContent(Unpooled.wrappedBuffer(snappy, 5, snappy.length - 5)), + response("gzip", -1)), + () -> List.of(new DefaultLastHttpContent(Unpooled.wrappedBuffer(gzipText))), + // A body rejected before anything is decoded still asks for a read. + () -> List.of(response("gzip", -1), new DefaultHttpContent(Unpooled.wrappedBuffer(corrupt)))); + List expected = readsPerBatch(nettyDecompressor(), batches); + assertEquals(expected, readsPerBatch(new Http1ContentDecompressor(false, 0), batches)); + assertEquals(List.of(1, 1, 0, 0, 1), expected); + } + + private static List readsPerBatch(ChannelHandler decompressor, List>> batches) { + int[] reads = new int[1]; + EmbeddedChannel channel = new EmbeddedChannel(); + channel.config().setAutoRead(false); + channel.pipeline().addLast(new ChannelOutboundHandlerAdapter() { + @Override + public void read(ChannelHandlerContext ctx) { + reads[0]++; + } + }); + channel.pipeline().addLast(decompressor); + List perBatch = new ArrayList<>(); + for (Supplier> batch : batches) { + int before = reads[0]; + for (Object msg : batch.get()) { + channel.pipeline().fireChannelRead(msg); + } + channel.pipeline().fireChannelReadComplete(); + perBatch.add(reads[0] - before); + for (Object msg; (msg = channel.readInbound()) != null; ) { + ReferenceCountUtil.release(msg); + } + } + try { + channel.finishAndReleaseAll(); + } catch (DecompressionException ignored) { + // The corrupt body of the last batch. + } + return perBatch; + } } From 0c9eb221eda05e8d89e72319d5415f32eebaec27 Mon Sep 17 00:00:00 2001 From: Pavel Ptashyts <49400901+pavel-ptashyts@users.noreply.github.com> Date: Thu, 1 Oct 2026 15:32:25 +0200 Subject: [PATCH 4/8] Stop decoding when a forward closes the channel Netty queues channelInactive, so a handler further on that closes the channel while decoded content is forwarded did not end the response: the rest of the body was still inflated and forwarded. The response is now abandoned when the channel went inactive during the forward. A channel that was already inactive when the message arrived is left alone, because HttpClientCodec ends a body without a length from channelInactive and that last part must still decode. A buffer that cannot grow fails with Netty's DecompressionException rather than an IndexOutOfBoundsException. Tests add a body ending with the connection, Transfer-Encoding gzip, an empty body, read-only input, and the inflater going back to the pool on close and removal. Claude Code on behalf of Pavel Ptashyts Co-Authored-By: Claude Opus 5.5 --- .../handler/Http1ContentDecompressor.java | 39 ++++-- .../handler/Http1ContentDecompressorTest.java | 123 +++++++++++++++--- 2 files changed, 135 insertions(+), 27 deletions(-) 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 deb1c9f97..6df91ce6a 100644 --- a/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java +++ b/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java @@ -110,9 +110,11 @@ private enum GzipState { private final CRC32 crc = new CRC32(); // Bumped whenever the state of a response is dropped. A handler further on may end the response from - // inside a call that forwards decoded content, by removing this handler or closing the channel; the - // decoding that made the call must then stop. + // inside a call that forwards decoded content, and the decoding that made the call must then stop. private int generation; + // Whether the channel was active when the message being decoded arrived. The last part of a body that + // ends with the connection is decoded after the channel has gone inactive, and must still be decoded. + private boolean activeAtDecode; // The response being inflated by this class rather than by HttpContentDecompressor. private boolean inflating; private boolean gzip; @@ -152,10 +154,11 @@ protected void decode(ChannelHandlerContext ctx, HttpObject msg, List ou contentEncoding = contentEncoding.trim(); if (isGzip(contentEncoding) || isDeflate(contentEncoding)) { lastDecodedBySuper = false; + activeAtDecode = ctx.channel().isActive(); needRead = true; int generation = this.generation; startResponse(ctx, response, contentEncoding); - if (response instanceof HttpContent && generation == this.generation) { + if (response instanceof HttpContent && !abandoned(ctx, generation)) { decodeContent(ctx, (HttpContent) response); } return; @@ -163,6 +166,7 @@ protected void decode(ChannelHandlerContext ctx, HttpObject msg, List ou } } else if (inflating) { lastDecodedBySuper = false; + activeAtDecode = ctx.channel().isActive(); needRead = true; decodeContent(ctx, (HttpContent) msg); return; @@ -290,7 +294,7 @@ private void decodeContent(ChannelHandlerContext ctx, HttpContent content) { } throw t; } - if (generation != this.generation) { + if (abandoned(ctx, generation)) { return; } } @@ -315,7 +319,7 @@ private void inflate(ChannelHandlerContext ctx, ByteBuf content, int generation) for (;;) { int readable = in.readableBytes(); inflateOnce(ctx, in, generation); - if (generation != this.generation) { + if (abandoned(ctx, generation)) { return; } if (!in.isReadable() || in.readableBytes() == readable) { @@ -387,8 +391,9 @@ private void inflateOnce(ChannelHandlerContext ctx, ByteBuf in, int generation) int preferredSize = Math.max(inflater.getRemaining() << 1, MIN_OUTPUT_BUFFER_SIZE); if (decompressed == null) { decompressed = ctx.alloc().heapBuffer(preferredSize); - } else { - decompressed.ensureWritable(preferredSize); + } else if (decompressed.ensureWritable(preferredSize, true) == 1) { + throw new DecompressionException( + "Decompression buffer has reached maximum size: " + decompressed.maxCapacity()); } byte[] outArray = decompressed.array(); int writerIndex = decompressed.writerIndex(); @@ -411,7 +416,7 @@ private void inflateOnce(ChannelHandlerContext ctx, ByteBuf in, int generation) ByteBuf buffer = decompressed; decompressed = null; fireContent(ctx, buffer); - if (generation != this.generation) { + if (abandoned(ctx, generation)) { // The inflater has already been given back. return; } @@ -451,6 +456,24 @@ private void inflateOnce(ChannelHandlerContext ctx, ByteBuf in, int generation) } } + /** + * Whether the response being decoded when content was last forwarded is over. Removing this handler + * ends it at once. Closing the channel only queues the inactive event, so a channel that went inactive + * during the forward ends the response here, and what is left of its body is not inflated. + */ + private boolean abandoned(ChannelHandlerContext ctx, int generation) { + if (generation != this.generation) { + return true; + } + if (activeAtDecode && !ctx.channel().isActive()) { + endResponse(); + inflating = true; + failed = true; + return true; + } + return false; + } + private void fireContent(ChannelHandlerContext ctx, ByteBuf buffer) { needRead = false; ctx.fireChannelRead(new DefaultHttpContent(buffer)); diff --git a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java index 2db6a166c..ed4587c44 100644 --- a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java +++ b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java @@ -202,6 +202,7 @@ private static List bodies() { byte[] gzipText = gzip(TEXT); byte[] gzipRandom = gzip(RANDOM); bodies.add(new Body("gzip empty", "gzip", gzip(EMPTY))); + bodies.add(new Body("gzip no body", "gzip", EMPTY)); bodies.add(new Body("gzip text", "gzip", gzipText)); bodies.add(new Body("gzip random", "gzip", gzipRandom)); bodies.add(new Body("gzip zeros", "gzip", gzip(ZEROS))); @@ -295,7 +296,7 @@ private static HttpResponse response(String contentEncoding, long contentLength) return response; } - private enum Kind { HEAP, DIRECT, COMPOSITE } + private enum Kind { HEAP, DIRECT, COMPOSITE, READ_ONLY } private static ByteBuf buffer(byte[] bytes, int offset, int length, Kind kind) { if (kind == Kind.COMPOSITE && length > 1) { @@ -307,7 +308,11 @@ private static ByteBuf buffer(byte[] bytes, int offset, int length, Kind kind) { .writeBytes(bytes, offset + half, length - half)); } int capacity = Math.max(length, 1); - ByteBuf buf = kind == Kind.HEAP ? Unpooled.buffer(capacity) : Unpooled.directBuffer(capacity); + ByteBuf buf = kind == Kind.DIRECT ? Unpooled.directBuffer(capacity) : Unpooled.buffer(capacity); + if (kind == Kind.READ_ONLY) { + // Heap, but without an accessible array, so its NIO view goes to the inflater. + return buf.writeBytes(bytes, offset, length).asReadOnly(); + } return buf.writeBytes(bytes, offset, length); } @@ -760,25 +765,105 @@ public void read(ChannelHandlerContext ctx) { } @Test - void closingTheChannelFromAHandlerFurtherOnMidChunkIsSafe() { - UnpooledByteBufAllocator allocator = allocator(); - EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator); - int[] parts = new int[1]; - channel.pipeline().addLast(new ChannelInboundHandlerAdapter() { - @Override - public void channelRead(ChannelHandlerContext ctx, Object msg) { - if (msg instanceof HttpContent && ((HttpContent) msg).content().isReadable() && parts[0]++ == 0) { - ctx.channel().close(); + void closingTheChannelFromAHandlerFurtherOnMidChunkStopsDecoding() throws Exception { + onFreshThread(() -> { + UnpooledByteBufAllocator allocator = allocator(); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator); + int[] parts = new int[1]; + List afterClose = new ArrayList<>(); + channel.pipeline().addLast(new ChannelInboundHandlerAdapter() { + @Override + public void channelRead(ChannelHandlerContext ctx, Object msg) { + if (parts[0] > 0) { + afterClose.add(msg); + } else if (msg instanceof HttpContent && ((HttpContent) msg).content().isReadable()) { + parts[0]++; + // Netty runs channelInactive later, once this read is over. + ctx.channel().close(); + } + ReferenceCountUtil.release(msg); } - ReferenceCountUtil.release(msg); - } + }); + byte[] encoded = gzip(ZEROS); + channel.writeInbound(response("gzip", encoded.length)); + channel.writeInbound(new DefaultLastHttpContent(Unpooled.wrappedBuffer(encoded))); + assertEquals(1, parts[0]); + assertEquals(List.of(), afterClose); + assertEquals(1, Http1ContentDecompressor.idleInflaterCount(true)); + channel.finishAndReleaseAll(); + assertEquals(0, allocator.metric().usedHeapMemory()); }); - byte[] encoded = gzip(ZEROS); - channel.writeInbound(response("gzip", encoded.length)); - channel.writeInbound(new DefaultLastHttpContent(Unpooled.wrappedBuffer(encoded))); - assertTrue(parts[0] >= 1); - channel.finishAndReleaseAll(); - assertEquals(0, allocator.metric().usedHeapMemory()); + } + + @Test + void decodesTheEndOfABodyThatEndsWithTheConnectionLikeNetty() { + byte[] encoded = gzip(RANDOM); + int split = encoded.length / 2; + List outcomes = new ArrayList<>(); + for (ChannelHandler decompressor : new ChannelHandler[]{nettyDecompressor(), + new Http1ContentDecompressor(false, 0)}) { + EmbeddedChannel channel = new EmbeddedChannel(); + // Stands in for HttpClientCodec, which ends a body without a length when the channel goes inactive. + channel.pipeline().addLast(new ChannelInboundHandlerAdapter() { + @Override + public void channelInactive(ChannelHandlerContext ctx) { + ctx.fireChannelRead(new DefaultLastHttpContent( + Unpooled.wrappedBuffer(encoded, split, encoded.length - split))); + ctx.fireChannelInactive(); + } + }); + channel.pipeline().addLast(decompressor); + channel.writeInbound(response("gzip", -1), + new DefaultHttpContent(Unpooled.wrappedBuffer(encoded, 0, split))); + channel.close(); + Outcome outcome = new Outcome(); + drain(channel, outcome); + channel.finishAndReleaseAll(); + outcomes.add(outcome); + } + assertSameOutcome(outcomes.get(0), outcomes.get(1), "body ending with the connection"); + assertDecoded(RANDOM, outcomes.get(1)); + } + + @Test + void removingOrClosingMidResponseGivesTheInflaterBack() throws Exception { + onFreshThread(() -> { + byte[] encoded = gzip(RANDOM); + EmbeddedChannel removed = channel(new Http1ContentDecompressor(false, 0), allocator()); + removed.writeInbound(response("gzip", -1), new DefaultHttpContent(Unpooled.wrappedBuffer(encoded, 0, 100))); + ReferenceCountUtil.release(removed.readInbound()); + removed.pipeline().removeFirst(); + assertEquals(1, Http1ContentDecompressor.idleInflaterCount(true)); + + EmbeddedChannel closed = channel(new Http1ContentDecompressor(false, 0), allocator()); + closed.writeInbound(response("gzip", -1), new DefaultHttpContent(Unpooled.wrappedBuffer(encoded, 0, 100))); + // Took the idle inflater. + assertEquals(0, Http1ContentDecompressor.idleInflaterCount(true)); + closed.close(); + assertEquals(1, Http1ContentDecompressor.idleInflaterCount(true)); + + removed.finishAndReleaseAll(); + closed.finishAndReleaseAll(); + }); + } + + @Test + void transferEncodingGzipIsLeftToNettyAndDecodedTheSame() { + byte[] encoded = gzip(TEXT); + List outcomes = new ArrayList<>(); + for (ChannelHandler decompressor : new ChannelHandler[]{nettyDecompressor(), + new Http1ContentDecompressor(false, 0)}) { + EmbeddedChannel channel = channel(decompressor, allocator()); + HttpResponse response = response(null, -1); + response.headers().set(HttpHeaderNames.TRANSFER_ENCODING, "gzip, chunked"); + channel.writeInbound(response, new DefaultLastHttpContent(Unpooled.wrappedBuffer(encoded))); + Outcome outcome = new Outcome(); + drain(channel, outcome); + channel.finishAndReleaseAll(); + outcomes.add(outcome); + } + assertSameOutcome(outcomes.get(0), outcomes.get(1), "transfer-encoding gzip"); + assertDecoded(TEXT, outcomes.get(1)); } @Test From 9c8a119375284d96ec45ed3c9c3430613eee8cd3 Mon Sep 17 00:00:00 2001 From: Pavel Ptashyts <49400901+pavel-ptashyts@users.noreply.github.com> Date: Thu, 1 Oct 2026 16:11:58 +0200 Subject: [PATCH 5/8] Grow the output buffer where Netty does ZlibDecoder.prepareDecompressBuffer grows the output buffer after every unfinished inflate step, so an allocator that caps heap buffers fails there. Growing only at the start of the next step let a capped buffer pass where Netty throws; the growth now happens at the same point. A channel closed while the response head is forwarded now drops the body already queued behind the head in the same read. The class notes that gzip and deflate no longer reach newContentDecoder and do not follow io.netty.noJdkZlibDecoder. Tests check that every input buffer is released, add pooled direct input and a capped allocator. Claude Code on behalf of Pavel Ptashyts Co-Authored-By: Claude Opus 5.5 --- .../handler/Http1ContentDecompressor.java | 23 ++-- .../handler/Http1ContentDecompressorTest.java | 103 +++++++++++++++++- 2 files changed, 116 insertions(+), 10 deletions(-) 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 6df91ce6a..c24d4c868 100644 --- a/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java +++ b/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java @@ -62,7 +62,9 @@ * per event loop for the length of one response, and feeds it the input buffer as it is. The format * handling (header and trailer checks, concatenated gzip members, zlib-or-raw detection for * {@code deflate}, a truncated stream ending quietly) follows {@code JdkZlibDecoder} as - * {@link HttpContentDecompressor} configures it, so the decoded body and the errors are unchanged. + * {@link HttpContentDecompressor} configures it, so the decoded body and the errors are unchanged. These + * encodings no longer go through {@link #newContentDecoder(String)}, and {@code java.util.zip} is used even when + * {@code io.netty.noJdkZlibDecoder} asks Netty for JZlib. *

* Every other encoding still goes through {@link HttpContentDecompressor}. There the counting is done * inside the decoder's own {@link EmbeddedChannel} rather than around {@code HttpContentDecoder#decode}: @@ -158,7 +160,7 @@ protected void decode(ChannelHandlerContext ctx, HttpObject msg, List ou needRead = true; int generation = this.generation; startResponse(ctx, response, contentEncoding); - if (response instanceof HttpContent && !abandoned(ctx, generation)) { + if (!abandoned(ctx, generation) && response instanceof HttpContent) { decodeContent(ctx, (HttpContent) response); } return; @@ -388,12 +390,8 @@ private void inflateOnce(ChannelHandlerContext ctx, ByteBuf in, int generation) // input, so needsInput() alone would end the loop too early. boolean pendingOutput = false; while (pendingOutput || !inflater.needsInput()) { - int preferredSize = Math.max(inflater.getRemaining() << 1, MIN_OUTPUT_BUFFER_SIZE); if (decompressed == null) { - decompressed = ctx.alloc().heapBuffer(preferredSize); - } else if (decompressed.ensureWritable(preferredSize, true) == 1) { - throw new DecompressionException( - "Decompression buffer has reached maximum size: " + decompressed.maxCapacity()); + decompressed = ctx.alloc().heapBuffer(preferredOutputBufferSize(inflater)); } byte[] outArray = decompressed.array(); int writerIndex = decompressed.writerIndex(); @@ -434,6 +432,13 @@ private void inflateOnce(ChannelHandlerContext ctx, ByteBuf in, int generation) } break; } + // Grown after every unfinished step, as ZlibDecoder.prepareDecompressBuffer does, so a buffer + // that cannot grow fails at the same point. + if (decompressed != null + && decompressed.ensureWritable(preferredOutputBufferSize(inflater), true) == 1) { + throw new DecompressionException( + "Decompression buffer has reached maximum size: " + decompressed.maxCapacity()); + } } in.skipBytes(readable - inflater.getRemaining()); @@ -474,6 +479,10 @@ private boolean abandoned(ChannelHandlerContext ctx, int generation) { return false; } + private static int preferredOutputBufferSize(Inflater inflater) { + return Math.max(inflater.getRemaining() << 1, MIN_OUTPUT_BUFFER_SIZE); + } + private void fireContent(ChannelHandlerContext ctx, ByteBuf buffer) { needRead = false; ctx.fireChannelRead(new DefaultHttpContent(buffer)); diff --git a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java index ed4587c44..a11a8d53a 100644 --- a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java +++ b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java @@ -15,9 +15,14 @@ */ package org.asynchttpclient.netty.handler; +import io.netty.buffer.AbstractByteBufAllocator; import io.netty.buffer.ByteBuf; +import io.netty.buffer.ByteBufAllocator; +import io.netty.buffer.PooledByteBufAllocator; import io.netty.buffer.Unpooled; import io.netty.buffer.UnpooledByteBufAllocator; +import io.netty.buffer.UnpooledDirectByteBuf; +import io.netty.buffer.UnpooledHeapByteBuf; import io.netty.channel.ChannelHandler; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelInboundHandlerAdapter; @@ -279,6 +284,7 @@ private static int[] sizes(int length, int chunk) { private static final class Outcome { final ByteArrayOutputStream body = new ByteArrayOutputStream(); + final List inputs = new ArrayList<>(); HttpResponse head; int parts; boolean ended; @@ -296,7 +302,7 @@ private static HttpResponse response(String contentEncoding, long contentLength) return response; } - private enum Kind { HEAP, DIRECT, COMPOSITE, READ_ONLY } + private enum Kind { HEAP, DIRECT, POOLED_DIRECT, COMPOSITE, READ_ONLY } private static ByteBuf buffer(byte[] bytes, int offset, int length, Kind kind) { if (kind == Kind.COMPOSITE && length > 1) { @@ -308,7 +314,9 @@ private static ByteBuf buffer(byte[] bytes, int offset, int length, Kind kind) { .writeBytes(bytes, offset + half, length - half)); } int capacity = Math.max(length, 1); - ByteBuf buf = kind == Kind.DIRECT ? Unpooled.directBuffer(capacity) : Unpooled.buffer(capacity); + ByteBuf buf = kind == Kind.DIRECT ? Unpooled.directBuffer(capacity) + : kind == Kind.POOLED_DIRECT ? PooledByteBufAllocator.DEFAULT.directBuffer(capacity) + : Unpooled.buffer(capacity); if (kind == Kind.READ_ONLY) { // Heap, but without an accessible array, so its NIO view goes to the inflater. return buf.writeBytes(bytes, offset, length).asReadOnly(); @@ -327,6 +335,7 @@ private static Outcome exchange(EmbeddedChannel channel, String contentEncoding, int offset = 0; for (int i = 0; i < sizes.length; i++) { ByteBuf content = buffer(encoded, offset, sizes[i], kind); + outcome.inputs.add(content); offset += sizes[i]; messages.add(i == sizes.length - 1 ? new DefaultLastHttpContent(content) : new DefaultHttpContent(content)); } @@ -346,6 +355,13 @@ private static Outcome exchange(EmbeddedChannel channel, String contentEncoding, return outcome; } + // Netty keeps an unread input as its cumulation until the channel is finished, so this is checked after. + private static void assertInputsReleased(Outcome outcome, String scenario) { + for (ByteBuf input : outcome.inputs) { + assertEquals(0, input.refCnt(), scenario + ": an input buffer was not released"); + } + } + private static void drain(EmbeddedChannel channel, Outcome outcome) { for (Object msg; (msg = channel.readInbound()) != null; ) { try { @@ -372,7 +388,7 @@ private static void drain(EmbeddedChannel channel, Outcome outcome) { } } - private static EmbeddedChannel channel(ChannelHandler decompressor, UnpooledByteBufAllocator allocator) { + private static EmbeddedChannel channel(ChannelHandler decompressor, ByteBufAllocator allocator) { EmbeddedChannel channel = new EmbeddedChannel(); channel.config().setAllocator(allocator); channel.pipeline().addLast(decompressor); @@ -409,6 +425,8 @@ void decodesLikeNettyForEveryBodyAndSplit() { // Netty's decoder runs once more over a corrupt body when its channel is torn down. } assertEquals(0, ahcAllocator.metric().usedHeapMemory(), scenario + ": heap left allocated"); + assertInputsReleased(actual, scenario); + assertInputsReleased(expected, scenario + " (Netty)"); checked++; } } @@ -991,4 +1009,83 @@ public void read(ChannelHandlerContext ctx) { } return perBatch; } + + @Test + void closingTheChannelWhileTheHeadIsForwardedDropsTheBody() throws Exception { + onFreshThread(() -> { + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator()); + List afterClose = new ArrayList<>(); + channel.pipeline().addLast(new ChannelInboundHandlerAdapter() { + @Override + public void channelRead(ChannelHandlerContext ctx, Object msg) { + if (msg instanceof HttpResponse) { + ctx.channel().close(); + } else { + afterClose.add(msg); + } + ReferenceCountUtil.release(msg); + } + }); + byte[] encoded = gzip(TEXT); + // One read: the body is already queued behind the head when the head is forwarded. + channel.writeInbound(response("gzip", encoded.length), + new DefaultLastHttpContent(Unpooled.wrappedBuffer(encoded))); + assertEquals(List.of(), afterClose); + assertEquals(1, Http1ContentDecompressor.idleInflaterCount(true)); + channel.finishAndReleaseAll(); + }); + } + + @Test + void aBufferThatCannotGrowFailsLikeNetty() { + // A raw deflate stream flushed but not finished, so the inflater wants more room after its output. + Deflater deflater = new Deflater(Deflater.DEFAULT_COMPRESSION, true); + byte[] flushed = new byte[4096]; + deflater.setInput(randomBytes(900, 7)); + int length = deflater.deflate(flushed, 0, flushed.length, Deflater.SYNC_FLUSH); + deflater.end(); + byte[] encoded = Arrays.copyOf(flushed, length); + + List outcomes = new ArrayList<>(); + for (ChannelHandler decompressor : new ChannelHandler[]{nettyDecompressor(), + new Http1ContentDecompressor(false, 0)}) { + EmbeddedChannel channel = channel(decompressor, new CappedHeapAllocator(1024)); + outcomes.add(exchange(channel, "deflate", encoded, new int[]{encoded.length}, Kind.DIRECT)); + try { + channel.finishAndReleaseAll(); + } catch (DecompressionException ignored) { + // Netty's decoder runs once more when its channel is torn down. + } + } + assertSameOutcome(outcomes.get(0), outcomes.get(1), "capped buffer"); + assertEquals("Decompression buffer has reached maximum size: 1024", outcomes.get(1).failure.getMessage()); + assertInputsReleased(outcomes.get(1), "capped buffer"); + } + + /** + * Never hands out a heap buffer that can grow past {@code cap} bytes. + */ + private static final class CappedHeapAllocator extends AbstractByteBufAllocator { + + private final int cap; + + CappedHeapAllocator(int cap) { + this.cap = cap; + } + + @Override + protected ByteBuf newHeapBuffer(int initialCapacity, int maxCapacity) { + return new UnpooledHeapByteBuf(this, Math.min(initialCapacity, cap), Math.min(maxCapacity, cap)); + } + + @Override + protected ByteBuf newDirectBuffer(int initialCapacity, int maxCapacity) { + return new UnpooledDirectByteBuf(this, initialCapacity, maxCapacity); + } + + @Override + public boolean isDirectBufferPooled() { + return false; + } + } } From bdd383de1c37372d460f4b6d4c20965fa6b60292 Mon Sep 17 00:00:00 2001 From: Pavel Ptashyts <49400901+pavel-ptashyts@users.noreply.github.com> Date: Thu, 1 Oct 2026 16:19:49 +0200 Subject: [PATCH 6/8] Cover the pool on event-loop style threads The pool tests now run on FastThreadLocalThreads, which is what Netty event loops are, so they take the same thread-local path as production, and a test checks that FastThreadLocal.removeAll() drops the idle inflaters. Another checks that the size limit counts every member of a concatenated gzip body together. The two places that fail a response share one helper, a guard that could never be false is gone, and a comment says why the inflater may keep pointing at an input buffer after it is released. Claude Code on behalf of Pavel Ptashyts Co-Authored-By: Claude Opus 5.5 --- .../handler/Http1ContentDecompressor.java | 20 +++++++----- .../handler/Http1ContentDecompressorTest.java | 31 ++++++++++++++++++- 2 files changed, 42 insertions(+), 9 deletions(-) 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 c24d4c868..ad964f4ca 100644 --- a/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java +++ b/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java @@ -289,11 +289,7 @@ private void decodeContent(ChannelHandlerContext ctx, HttpContent content) { try { inflate(ctx, in, generation); } catch (Throwable t) { - if (generation == this.generation) { - endResponse(); - inflating = true; - failed = true; - } + failResponse(); throw t; } if (abandoned(ctx, generation)) { @@ -372,6 +368,9 @@ private void inflateOnce(ChannelHandlerContext ctx, ByteBuf in, int generation) } int readable = in.readableBytes(); + // The inflater keeps pointing at this input after the chunk is released. That is safe only because + // nothing reads from it again before the next setInput() or reset(): the loop below runs until the + // input is used up or the stream has finished, and a finished gzip member is reset by its trailer. if (in.hasArray()) { inflater.setInput(in.array(), in.arrayOffset() + in.readerIndex(), readable); } else if (in.nioBufferCount() == 1) { @@ -471,14 +470,19 @@ private boolean abandoned(ChannelHandlerContext ctx, int generation) { return true; } if (activeAtDecode && !ctx.channel().isActive()) { - endResponse(); - inflating = true; - failed = true; + failResponse(); return true; } return false; } + // Drops the response's state; whatever is still in flight for it is then discarded. + private void failResponse() { + endResponse(); + inflating = true; + failed = true; + } + private static int preferredOutputBufferSize(Inflater inflater) { return Math.max(inflater.getRemaining() << 1, MIN_OUTPUT_BUFFER_SIZE); } diff --git a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java index a11a8d53a..e348d31da 100644 --- a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java +++ b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java @@ -44,6 +44,8 @@ import io.netty.handler.codec.http.HttpVersion; import io.netty.handler.codec.http.LastHttpContent; import io.netty.util.ReferenceCountUtil; +import io.netty.util.concurrent.FastThreadLocal; +import io.netty.util.concurrent.FastThreadLocalThread; import org.junit.jupiter.api.Test; import java.io.ByteArrayOutputStream; @@ -703,6 +705,32 @@ void limitAppliesToDeflateToo() { channel.finishAndReleaseAll(); } + @Test + void idleInflatersAreDroppedWithTheThreadLocals() throws Exception { + onFreshThread(() -> { + byte[] encoded = gzip(TEXT); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator()); + assertDecoded(TEXT, exchange(channel, "gzip", encoded, new int[]{encoded.length}, Kind.HEAP)); + channel.finishAndReleaseAll(); + assertEquals(1, Http1ContentDecompressor.idleInflaterCount(true)); + // What an event loop does on its way out. + FastThreadLocal.removeAll(); + assertEquals(0, Http1ContentDecompressor.idleInflaterCount(true)); + }); + } + + @Test + void limitCountsEveryMemberOfAConcatenatedBody() { + // Each member is under the limit on its own; together they are over it. + byte[] member = gzip(Arrays.copyOf(ZEROS, 600_000)); + byte[] encoded = concat(member, member); + EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 1_000_000), allocator()); + Outcome outcome = exchange(channel, "gzip", encoded, new int[]{encoded.length}, Kind.HEAP); + assertInstanceOf(DecompressionException.class, outcome.failure); + assertTrue(outcome.body.size() <= 1_000_000); + channel.finishAndReleaseAll(); + } + @Test void limitAllowsABodyOfExactlyTheLimit() { byte[] encoded = gzip(ZEROS); @@ -942,7 +970,8 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) { private static void onFreshThread(Runnable test) throws Exception { AtomicReference failure = new AtomicReference<>(); - Thread thread = new Thread(() -> { + // Event loops run on FastThreadLocalThreads, so the pool takes the same path as in production. + Thread thread = new FastThreadLocalThread(() -> { try { test.run(); } catch (Throwable t) { From 39f64e6977beee760c80e1721f0546c6230e3956 Mon Sep 17 00:00:00 2001 From: Pavel Ptashyts <49400901+pavel-ptashyts@users.noreply.github.com> Date: Thu, 1 Oct 2026 16:26:55 +0200 Subject: [PATCH 7/8] Check idle inflaters are ended on thread exit The removeAll() test only showed that the pool was dropped: asking for the pool again creates an empty one, so it passed even without the end() calls in onRemoval. Tests now reach the pool itself and check that the pooled inflater was ended. A message that is not HttpContent while a response is being inflated now goes to HttpContentDecoder instead of failing a cast. Closes #2356 Claude Code on behalf of Pavel Ptashyts Co-Authored-By: Claude Opus 5.5 --- .../handler/Http1ContentDecompressor.java | 12 +++---- .../handler/Http1ContentDecompressorTest.java | 32 +++++++++++-------- 2 files changed, 24 insertions(+), 20 deletions(-) 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 ad964f4ca..dd4c9be95 100644 --- a/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java +++ b/client/src/main/java/org/asynchttpclient/netty/handler/Http1ContentDecompressor.java @@ -91,8 +91,8 @@ public class Http1ContentDecompressor extends HttpContentDecompressor { // Inflaters between responses, per event loop. Holding them per connection instead would pin zlib's // native window on every idle pooled connection. private static final int MAX_IDLE_INFLATERS = 16; - private static final FastThreadLocal> IDLE_NOWRAP_INFLATERS = idleInflaters(); - private static final FastThreadLocal> IDLE_ZLIB_INFLATERS = idleInflaters(); + private static final FastThreadLocal> IDLE_NOWRAP_INFLATERS = newIdleInflaterPool(); + private static final FastThreadLocal> IDLE_ZLIB_INFLATERS = newIdleInflaterPool(); private enum GzipState { HEADER_START, @@ -166,7 +166,7 @@ protected void decode(ChannelHandlerContext ctx, HttpObject msg, List ou return; } } - } else if (inflating) { + } else if (inflating && msg instanceof HttpContent) { lastDecodedBySuper = false; activeAtDecode = ctx.channel().isActive(); needRead = true; @@ -637,7 +637,7 @@ private static boolean looksLikeZlib(short cmfFlg) { return (cmfFlg & 0x7800) == 0x7800 && cmfFlg % 31 == 0; } - private static FastThreadLocal> idleInflaters() { + private static FastThreadLocal> newIdleInflaterPool() { return new FastThreadLocal>() { @Override protected ArrayDeque initialValue() { @@ -664,8 +664,8 @@ private Inflater acquireInflater(boolean nowrap) { } // Visible for tests. - static int idleInflaterCount(boolean nowrap) { - return (nowrap ? IDLE_NOWRAP_INFLATERS : IDLE_ZLIB_INFLATERS).get().size(); + static ArrayDeque idleInflaters(boolean nowrap) { + return (nowrap ? IDLE_NOWRAP_INFLATERS : IDLE_ZLIB_INFLATERS).get(); } private void endResponse() { diff --git a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java index e348d31da..46ed5eb60 100644 --- a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java +++ b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java @@ -59,6 +59,7 @@ import java.util.zip.CRC32; import java.util.zip.Deflater; import java.util.zip.GZIPOutputStream; +import java.util.zip.Inflater; import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; @@ -66,6 +67,7 @@ import static org.junit.jupiter.api.Assertions.assertInstanceOf; import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; /** @@ -706,16 +708,18 @@ void limitAppliesToDeflateToo() { } @Test - void idleInflatersAreDroppedWithTheThreadLocals() throws Exception { + void idleInflatersAreEndedWithTheThreadLocals() throws Exception { onFreshThread(() -> { byte[] encoded = gzip(TEXT); EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator()); assertDecoded(TEXT, exchange(channel, "gzip", encoded, new int[]{encoded.length}, Kind.HEAP)); channel.finishAndReleaseAll(); - assertEquals(1, Http1ContentDecompressor.idleInflaterCount(true)); + Inflater pooled = Http1ContentDecompressor.idleInflaters(true).peekLast(); + assertNotNull(pooled); // What an event loop does on its way out. FastThreadLocal.removeAll(); - assertEquals(0, Http1ContentDecompressor.idleInflaterCount(true)); + // An open inflater with no input returns 0 here; an ended one throws. + assertThrows(NullPointerException.class, () -> pooled.inflate(new byte[1])); }); } @@ -835,7 +839,7 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) { channel.writeInbound(new DefaultLastHttpContent(Unpooled.wrappedBuffer(encoded))); assertEquals(1, parts[0]); assertEquals(List.of(), afterClose); - assertEquals(1, Http1ContentDecompressor.idleInflaterCount(true)); + assertEquals(1, Http1ContentDecompressor.idleInflaters(true).size()); channel.finishAndReleaseAll(); assertEquals(0, allocator.metric().usedHeapMemory()); }); @@ -879,14 +883,14 @@ void removingOrClosingMidResponseGivesTheInflaterBack() throws Exception { removed.writeInbound(response("gzip", -1), new DefaultHttpContent(Unpooled.wrappedBuffer(encoded, 0, 100))); ReferenceCountUtil.release(removed.readInbound()); removed.pipeline().removeFirst(); - assertEquals(1, Http1ContentDecompressor.idleInflaterCount(true)); + assertEquals(1, Http1ContentDecompressor.idleInflaters(true).size()); EmbeddedChannel closed = channel(new Http1ContentDecompressor(false, 0), allocator()); closed.writeInbound(response("gzip", -1), new DefaultHttpContent(Unpooled.wrappedBuffer(encoded, 0, 100))); // Took the idle inflater. - assertEquals(0, Http1ContentDecompressor.idleInflaterCount(true)); + assertEquals(0, Http1ContentDecompressor.idleInflaters(true).size()); closed.close(); - assertEquals(1, Http1ContentDecompressor.idleInflaterCount(true)); + assertEquals(1, Http1ContentDecompressor.idleInflaters(true).size()); removed.finishAndReleaseAll(); closed.finishAndReleaseAll(); @@ -923,20 +927,20 @@ void inflatersGoBackToABoundedPoolPerThread() throws Exception { new DefaultHttpContent(Unpooled.wrappedBuffer(encoded, 0, 20))); channels.add(channel); } - assertEquals(0, Http1ContentDecompressor.idleInflaterCount(true)); + assertEquals(0, Http1ContentDecompressor.idleInflaters(true).size()); for (EmbeddedChannel channel : channels) { channel.writeInbound(new DefaultLastHttpContent( Unpooled.wrappedBuffer(encoded, 20, encoded.length - 20))); channel.finishAndReleaseAll(); } - assertEquals(16, Http1ContentDecompressor.idleInflaterCount(true)); + assertEquals(16, Http1ContentDecompressor.idleInflaters(true).size()); EmbeddedChannel channel = channel(new Http1ContentDecompressor(false, 0), allocator()); assertDecoded(TEXT, exchange(channel, "gzip", encoded, new int[]{encoded.length}, Kind.HEAP)); - assertEquals(16, Http1ContentDecompressor.idleInflaterCount(true)); + assertEquals(16, Http1ContentDecompressor.idleInflaters(true).size()); byte[] zlib = deflate(TEXT, false); assertDecoded(TEXT, exchange(channel, "deflate", zlib, new int[]{zlib.length}, Kind.HEAP)); - assertEquals(1, Http1ContentDecompressor.idleInflaterCount(false)); + assertEquals(1, Http1ContentDecompressor.idleInflaters(false).size()); channel.finishAndReleaseAll(); }); } @@ -957,13 +961,13 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) { }); removed.writeInbound(response("gzip", -1), new DefaultLastHttpContent(Unpooled.wrappedBuffer(zeros))); removed.finishAndReleaseAll(); - assertEquals(1, Http1ContentDecompressor.idleInflaterCount(true)); + assertEquals(1, Http1ContentDecompressor.idleInflaters(true).size()); EmbeddedChannel next = channel(new Http1ContentDecompressor(false, 0), allocator()); assertDecoded(ZEROS, exchange(next, "gzip", zeros, sizes(zeros.length, 500), Kind.DIRECT)); byte[] raw = deflate(RANDOM, true); assertDecoded(RANDOM, exchange(next, "deflate", raw, sizes(raw.length, 3000), Kind.DIRECT)); - assertEquals(1, Http1ContentDecompressor.idleInflaterCount(true)); + assertEquals(1, Http1ContentDecompressor.idleInflaters(true).size()); next.finishAndReleaseAll(); }); } @@ -1060,7 +1064,7 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) { channel.writeInbound(response("gzip", encoded.length), new DefaultLastHttpContent(Unpooled.wrappedBuffer(encoded))); assertEquals(List.of(), afterClose); - assertEquals(1, Http1ContentDecompressor.idleInflaterCount(true)); + assertEquals(1, Http1ContentDecompressor.idleInflaters(true).size()); channel.finishAndReleaseAll(); }); } From 645506f089fa81ce26ef6ce0d3b19b43cbf77c1d Mon Sep 17 00:00:00 2001 From: Pavel Ptashyts <49400901+pavel-ptashyts@users.noreply.github.com> Date: Thu, 1 Oct 2026 17:28:44 +0200 Subject: [PATCH 8/8] Accept JDK 25's exception from an ended Inflater An ended Inflater throws NullPointerException on JDK 11 to 21 and IllegalStateException on JDK 25, which failed the pool cleanup test on every JDK 25 CI job. An open inflater with no input throws nothing, so any RuntimeException still shows that the pooled inflater was ended. Claude Code on behalf of Pavel Ptashyts Co-Authored-By: Claude Opus 5.5 --- .../netty/handler/Http1ContentDecompressorTest.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java index 46ed5eb60..7327f66d9 100644 --- a/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java +++ b/client/src/test/java/org/asynchttpclient/netty/handler/Http1ContentDecompressorTest.java @@ -718,8 +718,9 @@ void idleInflatersAreEndedWithTheThreadLocals() throws Exception { assertNotNull(pooled); // What an event loop does on its way out. FastThreadLocal.removeAll(); - // An open inflater with no input returns 0 here; an ended one throws. - assertThrows(NullPointerException.class, () -> pooled.inflate(new byte[1])); + // An open inflater with no input returns 0 here; an ended one throws, a NullPointerException on + // JDK 11 to 21 and an IllegalStateException on JDK 25. + assertThrows(RuntimeException.class, () -> pooled.inflate(new byte[1])); }); }