From e8b6eff4cd2ea97ab37f99e7969c4b80026df31c Mon Sep 17 00:00:00 2001 From: Qther Date: Fri, 25 Sep 2026 19:37:46 +0800 Subject: [PATCH 1/8] feat: mod compression (only gzip for now) --- .../client/MultiThreadedDownloader.java | 124 +++++++++++++++++- .../server/RequestHandler.java | 62 +++++++-- .../utils/ByteBufferUtils.java | 43 ++++++ .../utils/CompressionUtils.java | 58 ++++++++ 4 files changed, 273 insertions(+), 14 deletions(-) create mode 100644 src/main/java/net/forgecraft/serverpacklocator/utils/ByteBufferUtils.java create mode 100644 src/main/java/net/forgecraft/serverpacklocator/utils/CompressionUtils.java diff --git a/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java b/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java index 9108788..e536816 100644 --- a/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java +++ b/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java @@ -1,31 +1,39 @@ package net.forgecraft.serverpacklocator.client; import com.google.common.hash.HashCode; +import io.netty.buffer.ByteBuf; +import io.netty.buffer.ByteBufInputStream; +import io.netty.buffer.Unpooled; import net.forgecraft.serverpacklocator.FileChecksumValidator; import net.forgecraft.serverpacklocator.ServerManifest; import net.forgecraft.serverpacklocator.secure.IConnectionSecurityManager; +import net.forgecraft.serverpacklocator.utils.CompressionUtils; import net.neoforged.fml.loading.progress.StartupNotificationManager; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import java.io.IOException; +import java.io.InputStream; import java.net.URI; import java.net.URLEncoder; import java.net.http.HttpClient; import java.net.http.HttpRequest; import java.net.http.HttpResponse; import java.nio.ByteBuffer; +import java.nio.channels.Channels; +import java.nio.channels.FileChannel; import java.nio.charset.StandardCharsets; import java.nio.file.Files; import java.nio.file.Path; import java.nio.file.StandardOpenOption; import java.util.ArrayList; -import java.util.Base64; import java.util.List; import java.util.Objects; import java.util.OptionalLong; +import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionStage; import java.util.concurrent.Flow; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Function; import java.util.regex.Pattern; import java.util.stream.Collectors; @@ -96,6 +104,7 @@ private HttpResponse makeRequest(String path, boolean authenticated, Http LOGGER.info("ServerPackLocator is requesting {}...", requestUri); var request = HttpRequest.newBuilder(requestUri); + request.header("Accept-Encoding", "gzip"); this.connectionSecurityManager.decorateClientRequest(request, authenticated); var response = httpClient.send(request.build(), bodyHandler); this.connectionSecurityManager.handleClientResponse(response); @@ -265,9 +274,29 @@ private void downloadFile(final FileToDownload fileToDownload, final ProgressLis makeRequest( "files/" + URLEncoder.encode(nextFile, StandardCharsets.UTF_8).replace("+", "%20"), true, - progressListener.trackBodyHandler( - HttpResponse.BodyHandlers.ofFile(destinationPath, StandardOpenOption.WRITE, StandardOpenOption.CREATE, StandardOpenOption.TRUNCATE_EXISTING) - ) + progressListener.trackBodyHandler((HttpResponse.BodyHandler) responseInfo -> { + Function decompressor = s -> s; + var encoding = responseInfo.headers().firstValue("Content-Encoding").orElse(null); + if (encoding != null) { + var methodNames = encoding.split(","); + for (int i = methodNames.length - 1; i >= 0; i--) { + var method = CompressionUtils.methodFromName(methodNames[i].trim()); + if (method != null) { + decompressor = decompressor.andThen(s -> { + try { + return method.decompress(s); + } catch (IOException e) { + throw new RuntimeException(e); + } + }); + } + } + } + + int initialCapacity = (int) fileToDownload.size(); + + return new CompressedFileSubscriber(fileToDownload, decompressor, initialCapacity); + }) ); // Validate that the downloaded file actually matches the expected checksum @@ -326,6 +355,93 @@ public void onComplete() { void onProgress(long downloaded, long expectedSize); } + static class CompressedFileSubscriber implements HttpResponse.BodySubscriber { + /** + * Adapted from {@link jdk.internal.net.http.ResponseSubscribers.PathSubscriber} + */ + + private final FileToDownload file; + private final Function decompressor; + + private final ByteBuf fullBuf; + private final CompletableFuture result = new CompletableFuture<>(); + + private final AtomicBoolean subscribed = new AtomicBoolean(); + private volatile Flow.Subscription subscription; + private volatile FileChannel out; + + CompressedFileSubscriber(FileToDownload file, Function decompressor, int initialCapacity) { + this.file = file; + this.decompressor = decompressor; + this.fullBuf = Unpooled.directBuffer(initialCapacity); + } + + @Override + public void onSubscribe(Flow.Subscription subscription) { + Objects.requireNonNull(subscription); + if (!subscribed.compareAndSet(false, true)) { + subscription.cancel(); + return; + } + + this.subscription = subscription; + try { + out = FileChannel.open(file.localFile(), StandardOpenOption.WRITE, StandardOpenOption.CREATE, StandardOpenOption.TRUNCATE_EXISTING); + } catch (IOException ioe) { + result.completeExceptionally(ioe); + subscription.cancel(); + return; + } + subscription.request(1); + } + + @Override + public void onNext(List items) { + for (var buf : items) { + fullBuf.writeBytes(buf); + buf.clear(); + } + + subscription.request(1); + } + + @Override + public void onError(Throwable e) { + result.completeExceptionally(e); + close(); + } + + @Override + public void onComplete() { + try ( + var inStream = decompressor.apply(new ByteBufInputStream(fullBuf)); + var outStream = Channels.newOutputStream(out) + ) { + var transferred = inStream.transferTo(outStream); + LOGGER.debug("Received {}/{} bytes for {}", transferred, file.size(), file.relativeDownloadPath()); + } catch (IOException ex) { + close(); + subscription.cancel(); + result.completeExceptionally(ex); + } + + close(); + + result.complete(file.localFile()); + } + + @Override + public CompletionStage getBody() { + return result; + } + + private void close() { + try { + out.close(); + } catch (IOException ignored) {} + } + } + record FileToDownload(String relativeUrl, String relativeDownloadPath, Path localFile, long size, HashCode checksum) { } diff --git a/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java b/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java index d3a0a1c..ff424c2 100644 --- a/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java +++ b/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java @@ -1,19 +1,30 @@ package net.forgecraft.serverpacklocator.server; -import net.forgecraft.serverpacklocator.ModAccessor; -import net.forgecraft.serverpacklocator.secure.IConnectionSecurityManager; import io.netty.buffer.ByteBuf; +import io.netty.buffer.ByteBufOutputStream; import io.netty.buffer.Unpooled; import io.netty.channel.ChannelHandlerContext; -import io.netty.channel.DefaultFileRegion; import io.netty.channel.SimpleChannelInboundHandler; -import io.netty.handler.codec.http.*; +import io.netty.handler.codec.http.DefaultFullHttpResponse; +import io.netty.handler.codec.http.FullHttpRequest; +import io.netty.handler.codec.http.FullHttpResponse; +import io.netty.handler.codec.http.HttpHeaderNames; +import io.netty.handler.codec.http.HttpMethod; +import io.netty.handler.codec.http.HttpResponseStatus; +import io.netty.handler.codec.http.HttpUtil; +import io.netty.handler.codec.http.HttpVersion; +import net.forgecraft.serverpacklocator.ModAccessor; +import net.forgecraft.serverpacklocator.secure.IConnectionSecurityManager; +import net.forgecraft.serverpacklocator.utils.CompressionUtils; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import javax.net.ssl.SSLException; +import java.io.IOException; +import java.io.OutputStream; import java.net.URLDecoder; import java.nio.charset.StandardCharsets; +import java.nio.file.Files; import java.util.Objects; class RequestHandler extends SimpleChannelInboundHandler { @@ -106,15 +117,46 @@ private void buildReply(final ChannelHandlerContext ctx, final FullHttpRequest m } private void buildFileReply(final ChannelHandlerContext ctx, final FullHttpRequest msg, final ServerFileManager.ExposedFile file) { - final HttpResponse response = new DefaultHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK); + ByteBuf content = Unpooled.buffer(); + + var usedEncodingsStr = ""; + try (OutputStream contentStream = new ByteBufOutputStream(content)) { + var out = contentStream; + var acceptedEncodings = msg.headers().get("Accept-Encoding", "").split(","); + var usedEncodings = new StringBuilder(); + for (var accepted : acceptedEncodings) { + var method = CompressionUtils.methodFromName(accepted.trim()); + if (method != null) { + out = method.compress(out); + if (!usedEncodings.isEmpty()) { + usedEncodings.append(','); + } + usedEncodings.append(method.name()); + } + } + + try (var fileStream = Files.newInputStream(file.path())) { + fileStream.transferTo(out); + } + + usedEncodingsStr = usedEncodings.toString(); + out.flush(); + out.close(); + } catch (IOException e) { + throw new RuntimeException(e); + } + + LOGGER.debug("Sending {} with {} compression ({} -> {} bytes)", file.name(), usedEncodingsStr.isEmpty() ? "no" : usedEncodingsStr, file.size(), content.writerIndex()); + + final FullHttpResponse response = new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK, content); HttpUtil.setKeepAlive(response, HttpUtil.isKeepAlive(msg)); response.headers().set(HttpHeaderNames.CONTENT_TYPE, "application/octet-stream"); response.headers().set("filename", file.name()); - HttpUtil.setContentLength(response, file.size()); - this.connectionSecurityManager.decorateServerResponse(ctx, msg, response); + if (!usedEncodingsStr.isEmpty()) { + response.headers().set("Content-Encoding", usedEncodingsStr); + } + HttpUtil.setContentLength(response, content.writerIndex()); - ctx.write(response); - ctx.write(new DefaultFileRegion(file.path().toFile(), 0, file.size())); - ctx.writeAndFlush(LastHttpContent.EMPTY_LAST_CONTENT); + ctx.writeAndFlush(response); } } diff --git a/src/main/java/net/forgecraft/serverpacklocator/utils/ByteBufferUtils.java b/src/main/java/net/forgecraft/serverpacklocator/utils/ByteBufferUtils.java new file mode 100644 index 0000000..86a6a5f --- /dev/null +++ b/src/main/java/net/forgecraft/serverpacklocator/utils/ByteBufferUtils.java @@ -0,0 +1,43 @@ +package net.forgecraft.serverpacklocator.utils; + +import org.jetbrains.annotations.NotNull; + +import java.io.ByteArrayInputStream; +import java.io.InputStream; +import java.nio.ByteBuffer; + +public class ByteBufferUtils { + // Adapted from https://stackoverflow.com/a/6603018 (https://creativecommons.org/licenses/by-sa/3.0/) + public static class ByteBufferBackedInputStream extends InputStream { + @NotNull + final ByteBuffer buf; + + public ByteBufferBackedInputStream(@NotNull ByteBuffer buf) { + this.buf = buf; + } + + @Override + public int read() { + if (!buf.hasRemaining()) { + return -1; + } + return buf.get() & 0xFF; + } + + @Override + public int read(byte @NotNull [] bytes, int off, int len) { + if (!buf.hasRemaining()) { + return -1; + } + + len = Math.min(len, buf.remaining()); + buf.get(bytes, off, len); + return len; + } + } + + @NotNull + public static InputStream inputStreamFromByteBuffer(@NotNull ByteBuffer buf) { + return buf.hasArray() ? new ByteArrayInputStream(buf.array()) : new ByteBufferBackedInputStream(buf); + } +} diff --git a/src/main/java/net/forgecraft/serverpacklocator/utils/CompressionUtils.java b/src/main/java/net/forgecraft/serverpacklocator/utils/CompressionUtils.java new file mode 100644 index 0000000..c6c492b --- /dev/null +++ b/src/main/java/net/forgecraft/serverpacklocator/utils/CompressionUtils.java @@ -0,0 +1,58 @@ +package net.forgecraft.serverpacklocator.utils; + +import org.jetbrains.annotations.NotNull; + +import java.io.IOException; +import java.io.InputStream; +import java.io.OutputStream; +import java.nio.channels.Channels; +import java.nio.channels.SeekableByteChannel; +import java.util.zip.GZIPInputStream; +import java.util.zip.GZIPOutputStream; + +public class CompressionUtils { + public abstract static class CompressionMethod { + CompressionMethod() {} + + public abstract @NotNull String name(); + + public abstract @NotNull InputStream decompress(InputStream input) throws IOException; + + public @NotNull InputStream decompress(SeekableByteChannel input) throws IOException { + return this.decompress(Channels.newInputStream(input)); + } + + public abstract @NotNull OutputStream compress(OutputStream input) throws IOException; + + public @NotNull OutputStream compress(SeekableByteChannel input) throws IOException { + return this.compress(Channels.newOutputStream(input)); + } + } + + public static class Gzip extends CompressionMethod { + public static final CompressionMethod INSTANCE = new Gzip(); + public static final String NAME = "gzip"; + + @Override + public @NotNull String name() { + return NAME; + } + + @Override + public @NotNull InputStream decompress(InputStream input) throws IOException { + return new GZIPInputStream(input); + } + + @Override + public @NotNull OutputStream compress(OutputStream input) throws IOException { + return new GZIPOutputStream(input); + } + } + + public static CompressionMethod methodFromName(String name) { + return switch (name) { + case Gzip.NAME -> Gzip.INSTANCE; + case null, default -> null; + }; + } +} From 8ceeda88905308b9a90895b010bed251c691874d Mon Sep 17 00:00:00 2001 From: Qther Date: Fri, 25 Sep 2026 22:13:02 +0800 Subject: [PATCH 2/8] remove unused overload --- .../serverpacklocator/utils/CompressionUtils.java | 10 ---------- 1 file changed, 10 deletions(-) diff --git a/src/main/java/net/forgecraft/serverpacklocator/utils/CompressionUtils.java b/src/main/java/net/forgecraft/serverpacklocator/utils/CompressionUtils.java index c6c492b..3292918 100644 --- a/src/main/java/net/forgecraft/serverpacklocator/utils/CompressionUtils.java +++ b/src/main/java/net/forgecraft/serverpacklocator/utils/CompressionUtils.java @@ -5,8 +5,6 @@ import java.io.IOException; import java.io.InputStream; import java.io.OutputStream; -import java.nio.channels.Channels; -import java.nio.channels.SeekableByteChannel; import java.util.zip.GZIPInputStream; import java.util.zip.GZIPOutputStream; @@ -18,15 +16,7 @@ public abstract static class CompressionMethod { public abstract @NotNull InputStream decompress(InputStream input) throws IOException; - public @NotNull InputStream decompress(SeekableByteChannel input) throws IOException { - return this.decompress(Channels.newInputStream(input)); - } - public abstract @NotNull OutputStream compress(OutputStream input) throws IOException; - - public @NotNull OutputStream compress(SeekableByteChannel input) throws IOException { - return this.compress(Channels.newOutputStream(input)); - } } public static class Gzip extends CompressionMethod { From 3a48c6b42acaf928c8283acfa8a5852efbc777b0 Mon Sep 17 00:00:00 2001 From: Qther Date: Fri, 25 Sep 2026 22:13:53 +0800 Subject: [PATCH 3/8] remove unused utils class --- .../utils/ByteBufferUtils.java | 43 ------------------- 1 file changed, 43 deletions(-) delete mode 100644 src/main/java/net/forgecraft/serverpacklocator/utils/ByteBufferUtils.java diff --git a/src/main/java/net/forgecraft/serverpacklocator/utils/ByteBufferUtils.java b/src/main/java/net/forgecraft/serverpacklocator/utils/ByteBufferUtils.java deleted file mode 100644 index 86a6a5f..0000000 --- a/src/main/java/net/forgecraft/serverpacklocator/utils/ByteBufferUtils.java +++ /dev/null @@ -1,43 +0,0 @@ -package net.forgecraft.serverpacklocator.utils; - -import org.jetbrains.annotations.NotNull; - -import java.io.ByteArrayInputStream; -import java.io.InputStream; -import java.nio.ByteBuffer; - -public class ByteBufferUtils { - // Adapted from https://stackoverflow.com/a/6603018 (https://creativecommons.org/licenses/by-sa/3.0/) - public static class ByteBufferBackedInputStream extends InputStream { - @NotNull - final ByteBuffer buf; - - public ByteBufferBackedInputStream(@NotNull ByteBuffer buf) { - this.buf = buf; - } - - @Override - public int read() { - if (!buf.hasRemaining()) { - return -1; - } - return buf.get() & 0xFF; - } - - @Override - public int read(byte @NotNull [] bytes, int off, int len) { - if (!buf.hasRemaining()) { - return -1; - } - - len = Math.min(len, buf.remaining()); - buf.get(bytes, off, len); - return len; - } - } - - @NotNull - public static InputStream inputStreamFromByteBuffer(@NotNull ByteBuffer buf) { - return buf.hasArray() ? new ByteArrayInputStream(buf.array()) : new ByteBufferBackedInputStream(buf); - } -} From 90e85f5fc1322e1daf480e086ec46964b3a27781 Mon Sep 17 00:00:00 2001 From: Qther Date: Tue, 29 Sep 2026 12:17:17 +0800 Subject: [PATCH 4/8] streamed decompression --- .../client/MultiThreadedDownloader.java | 107 ++---------------- 1 file changed, 12 insertions(+), 95 deletions(-) diff --git a/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java b/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java index e536816..8a7d3ac 100644 --- a/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java +++ b/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java @@ -1,9 +1,6 @@ package net.forgecraft.serverpacklocator.client; import com.google.common.hash.HashCode; -import io.netty.buffer.ByteBuf; -import io.netty.buffer.ByteBufInputStream; -import io.netty.buffer.Unpooled; import net.forgecraft.serverpacklocator.FileChecksumValidator; import net.forgecraft.serverpacklocator.ServerManifest; import net.forgecraft.serverpacklocator.secure.IConnectionSecurityManager; @@ -30,10 +27,8 @@ import java.util.List; import java.util.Objects; import java.util.OptionalLong; -import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionStage; import java.util.concurrent.Flow; -import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Function; import java.util.regex.Pattern; import java.util.stream.Collectors; @@ -286,6 +281,7 @@ private void downloadFile(final FileToDownload fileToDownload, final ProgressLis try { return method.decompress(s); } catch (IOException e) { + LOGGER.error(e); throw new RuntimeException(e); } }); @@ -293,9 +289,16 @@ private void downloadFile(final FileToDownload fileToDownload, final ProgressLis } } - int initialCapacity = (int) fileToDownload.size(); + var downstream = HttpResponse.BodySubscribers.mapping(HttpResponse.BodySubscribers.ofInputStream(), decompressor); + return HttpResponse.BodySubscribers.mapping(downstream, s -> { + try (var out = Channels.newOutputStream(FileChannel.open(fileToDownload.localFile(), StandardOpenOption.WRITE, StandardOpenOption.CREATE, StandardOpenOption.TRUNCATE_EXISTING))){ + s.transferTo(out); + } catch (IOException e) { + throw new RuntimeException(e); + } - return new CompressedFileSubscriber(fileToDownload, decompressor, initialCapacity); + return fileToDownload.localFile(); + }); }) ); @@ -355,94 +358,8 @@ public void onComplete() { void onProgress(long downloaded, long expectedSize); } - static class CompressedFileSubscriber implements HttpResponse.BodySubscriber { - /** - * Adapted from {@link jdk.internal.net.http.ResponseSubscribers.PathSubscriber} - */ - - private final FileToDownload file; - private final Function decompressor; - - private final ByteBuf fullBuf; - private final CompletableFuture result = new CompletableFuture<>(); - - private final AtomicBoolean subscribed = new AtomicBoolean(); - private volatile Flow.Subscription subscription; - private volatile FileChannel out; - - CompressedFileSubscriber(FileToDownload file, Function decompressor, int initialCapacity) { - this.file = file; - this.decompressor = decompressor; - this.fullBuf = Unpooled.directBuffer(initialCapacity); - } - - @Override - public void onSubscribe(Flow.Subscription subscription) { - Objects.requireNonNull(subscription); - if (!subscribed.compareAndSet(false, true)) { - subscription.cancel(); - return; - } - - this.subscription = subscription; - try { - out = FileChannel.open(file.localFile(), StandardOpenOption.WRITE, StandardOpenOption.CREATE, StandardOpenOption.TRUNCATE_EXISTING); - } catch (IOException ioe) { - result.completeExceptionally(ioe); - subscription.cancel(); - return; - } - subscription.request(1); - } - - @Override - public void onNext(List items) { - for (var buf : items) { - fullBuf.writeBytes(buf); - buf.clear(); - } - - subscription.request(1); - } - - @Override - public void onError(Throwable e) { - result.completeExceptionally(e); - close(); - } - - @Override - public void onComplete() { - try ( - var inStream = decompressor.apply(new ByteBufInputStream(fullBuf)); - var outStream = Channels.newOutputStream(out) - ) { - var transferred = inStream.transferTo(outStream); - LOGGER.debug("Received {}/{} bytes for {}", transferred, file.size(), file.relativeDownloadPath()); - } catch (IOException ex) { - close(); - subscription.cancel(); - result.completeExceptionally(ex); - } - - close(); - - result.complete(file.localFile()); - } - - @Override - public CompletionStage getBody() { - return result; - } - - private void close() { - try { - out.close(); - } catch (IOException ignored) {} - } - } - - record FileToDownload(String relativeUrl, String relativeDownloadPath, Path localFile, long size, HashCode checksum) { + record FileToDownload(String relativeUrl, String relativeDownloadPath, Path localFile, long size, + HashCode checksum) { } public record PreparedServerDownloadData(ServerManifest manifest, From 20ef6748ea4d271d6a8dd939260a460dd2985034 Mon Sep 17 00:00:00 2001 From: Qther Date: Tue, 29 Sep 2026 12:18:31 +0800 Subject: [PATCH 5/8] unformat FileToDownload --- .../serverpacklocator/client/MultiThreadedDownloader.java | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java b/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java index 8a7d3ac..7719d42 100644 --- a/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java +++ b/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java @@ -358,8 +358,7 @@ public void onComplete() { void onProgress(long downloaded, long expectedSize); } - record FileToDownload(String relativeUrl, String relativeDownloadPath, Path localFile, long size, - HashCode checksum) { + record FileToDownload(String relativeUrl, String relativeDownloadPath, Path localFile, long size, HashCode checksum) { } public record PreparedServerDownloadData(ServerManifest manifest, From 63629120466863222ea854cfe7964ebb1db7b4dd Mon Sep 17 00:00:00 2001 From: Qther Date: Tue, 29 Sep 2026 12:21:22 +0800 Subject: [PATCH 6/8] use direct buffer for upload --- .../net/forgecraft/serverpacklocator/server/RequestHandler.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java b/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java index ff424c2..4c311fb 100644 --- a/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java +++ b/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java @@ -117,7 +117,7 @@ private void buildReply(final ChannelHandlerContext ctx, final FullHttpRequest m } private void buildFileReply(final ChannelHandlerContext ctx, final FullHttpRequest msg, final ServerFileManager.ExposedFile file) { - ByteBuf content = Unpooled.buffer(); + ByteBuf content = Unpooled.directBuffer(); var usedEncodingsStr = ""; try (OutputStream contentStream = new ByteBufOutputStream(content)) { From ef5dfe999bdc2e8b1b39fba47545f859d5fef626 Mon Sep 17 00:00:00 2001 From: Qther Date: Tue, 29 Sep 2026 12:27:57 +0800 Subject: [PATCH 7/8] use netty ctx to alloc instead of Unpooled --- .../net/forgecraft/serverpacklocator/server/RequestHandler.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java b/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java index 4c311fb..3ec6d86 100644 --- a/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java +++ b/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java @@ -117,7 +117,7 @@ private void buildReply(final ChannelHandlerContext ctx, final FullHttpRequest m } private void buildFileReply(final ChannelHandlerContext ctx, final FullHttpRequest msg, final ServerFileManager.ExposedFile file) { - ByteBuf content = Unpooled.directBuffer(); + ByteBuf content = ctx.alloc().ioBuffer(); var usedEncodingsStr = ""; try (OutputStream contentStream = new ByteBufOutputStream(content)) { From a5b584bbff683c074b3366f9b7870da1f05cc567 Mon Sep 17 00:00:00 2001 From: Qther Date: Tue, 29 Sep 2026 13:20:32 +0800 Subject: [PATCH 8/8] chunked uploads --- .../client/MultiThreadedDownloader.java | 29 +++++++++- .../server/RequestHandler.java | 53 +++++-------------- .../server/SimpleHttpServer.java | 11 +++- 3 files changed, 50 insertions(+), 43 deletions(-) diff --git a/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java b/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java index 7719d42..7904bc8 100644 --- a/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java +++ b/src/main/java/net/forgecraft/serverpacklocator/client/MultiThreadedDownloader.java @@ -67,10 +67,35 @@ private PreparedServerDownloadData downloadManifest() throws IOException, Interr authenticate(); var progressBar = StartupNotificationManager.addProgressBar("Requesting server manifest...", 1); try { - var response = makeRequest("servermanifest.json", true, HttpResponse.BodyHandlers.ofString(StandardCharsets.UTF_8)); + var response = makeRequest("servermanifest.json", true, responseInfo -> HttpResponse.BodySubscribers.mapping(HttpResponse.BodySubscribers.ofInputStream(), in -> { + Function decompressor = s -> s; + var encoding = responseInfo.headers().firstValue("Content-Encoding").orElse(null); + if (encoding != null) { + var methodNames = encoding.split(","); + for (int i = methodNames.length - 1; i >= 0; i--) { + var method = CompressionUtils.methodFromName(methodNames[i].trim()); + if (method != null) { + decompressor = decompressor.andThen(s -> { + try { + return method.decompress(s); + } catch (IOException e) { + LOGGER.error(e); + throw new RuntimeException(e); + } + }); + } + } + } + + return decompressor.apply(in); + })); + this.connectionSecurityManager.handleClientResponse(response); - var serverManifest = ServerManifest.fromString(response.body()); + ServerManifest serverManifest; + try (var stream = response.body()) { + serverManifest = ServerManifest.fromString(String.valueOf(StandardCharsets.UTF_8.decode(ByteBuffer.wrap(stream.readAllBytes())))); + } // Write the file to the client system for debugging if (serverManifest != null) { diff --git a/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java b/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java index 3ec6d86..22e6285 100644 --- a/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java +++ b/src/main/java/net/forgecraft/serverpacklocator/server/RequestHandler.java @@ -1,27 +1,29 @@ package net.forgecraft.serverpacklocator.server; import io.netty.buffer.ByteBuf; -import io.netty.buffer.ByteBufOutputStream; import io.netty.buffer.Unpooled; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.SimpleChannelInboundHandler; import io.netty.handler.codec.http.DefaultFullHttpResponse; +import io.netty.handler.codec.http.DefaultHttpResponse; import io.netty.handler.codec.http.FullHttpRequest; import io.netty.handler.codec.http.FullHttpResponse; +import io.netty.handler.codec.http.HttpChunkedInput; import io.netty.handler.codec.http.HttpHeaderNames; +import io.netty.handler.codec.http.HttpHeaderValues; import io.netty.handler.codec.http.HttpMethod; +import io.netty.handler.codec.http.HttpResponse; import io.netty.handler.codec.http.HttpResponseStatus; import io.netty.handler.codec.http.HttpUtil; import io.netty.handler.codec.http.HttpVersion; +import io.netty.handler.stream.ChunkedStream; import net.forgecraft.serverpacklocator.ModAccessor; import net.forgecraft.serverpacklocator.secure.IConnectionSecurityManager; -import net.forgecraft.serverpacklocator.utils.CompressionUtils; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import javax.net.ssl.SSLException; import java.io.IOException; -import java.io.OutputStream; import java.net.URLDecoder; import java.nio.charset.StandardCharsets; import java.nio.file.Files; @@ -117,46 +119,19 @@ private void buildReply(final ChannelHandlerContext ctx, final FullHttpRequest m } private void buildFileReply(final ChannelHandlerContext ctx, final FullHttpRequest msg, final ServerFileManager.ExposedFile file) { - ByteBuf content = ctx.alloc().ioBuffer(); - - var usedEncodingsStr = ""; - try (OutputStream contentStream = new ByteBufOutputStream(content)) { - var out = contentStream; - var acceptedEncodings = msg.headers().get("Accept-Encoding", "").split(","); - var usedEncodings = new StringBuilder(); - for (var accepted : acceptedEncodings) { - var method = CompressionUtils.methodFromName(accepted.trim()); - if (method != null) { - out = method.compress(out); - if (!usedEncodings.isEmpty()) { - usedEncodings.append(','); - } - usedEncodings.append(method.name()); - } - } + try (var fileStream = Files.newInputStream(file.path())) { + var out = new HttpChunkedInput(new ChunkedStream(fileStream)); - try (var fileStream = Files.newInputStream(file.path())) { - fileStream.transferTo(out); - } + final HttpResponse response = new DefaultHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK); + HttpUtil.setKeepAlive(response, HttpUtil.isKeepAlive(msg)); + response.headers().set(HttpHeaderNames.CONTENT_TYPE, "application/octet-stream"); + response.headers().set("filename", file.name()); + response.headers().set(HttpHeaderNames.TRANSFER_ENCODING, HttpHeaderValues.CHUNKED); - usedEncodingsStr = usedEncodings.toString(); - out.flush(); - out.close(); + ctx.write(response); + ctx.writeAndFlush(out); } catch (IOException e) { throw new RuntimeException(e); } - - LOGGER.debug("Sending {} with {} compression ({} -> {} bytes)", file.name(), usedEncodingsStr.isEmpty() ? "no" : usedEncodingsStr, file.size(), content.writerIndex()); - - final FullHttpResponse response = new DefaultFullHttpResponse(HttpVersion.HTTP_1_1, HttpResponseStatus.OK, content); - HttpUtil.setKeepAlive(response, HttpUtil.isKeepAlive(msg)); - response.headers().set(HttpHeaderNames.CONTENT_TYPE, "application/octet-stream"); - response.headers().set("filename", file.name()); - if (!usedEncodingsStr.isEmpty()) { - response.headers().set("Content-Encoding", usedEncodingsStr); - } - HttpUtil.setContentLength(response, content.writerIndex()); - - ctx.writeAndFlush(response); } } diff --git a/src/main/java/net/forgecraft/serverpacklocator/server/SimpleHttpServer.java b/src/main/java/net/forgecraft/serverpacklocator/server/SimpleHttpServer.java index 2d39773..6fe6ae2 100644 --- a/src/main/java/net/forgecraft/serverpacklocator/server/SimpleHttpServer.java +++ b/src/main/java/net/forgecraft/serverpacklocator/server/SimpleHttpServer.java @@ -1,16 +1,21 @@ package net.forgecraft.serverpacklocator.server; -import net.forgecraft.serverpacklocator.secure.IConnectionSecurityManager; import com.google.common.util.concurrent.ThreadFactoryBuilder; import com.mojang.logging.LogUtils; import io.netty.bootstrap.ServerBootstrap; -import io.netty.channel.*; +import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelInboundHandlerAdapter; +import io.netty.channel.ChannelInitializer; +import io.netty.channel.ChannelOption; import io.netty.channel.nio.NioEventLoopGroup; import io.netty.channel.socket.ServerSocketChannel; import io.netty.channel.socket.SocketChannel; import io.netty.channel.socket.nio.NioServerSocketChannel; +import io.netty.handler.codec.http.HttpContentCompressor; import io.netty.handler.codec.http.HttpObjectAggregator; import io.netty.handler.codec.http.HttpServerCodec; +import io.netty.handler.stream.ChunkedWriteHandler; +import net.forgecraft.serverpacklocator.secure.IConnectionSecurityManager; import org.slf4j.Logger; /** @@ -53,6 +58,8 @@ public void channelActive(final ChannelHandlerContext ctx) { @Override protected void initChannel(final SocketChannel ch) { ch.pipeline().addLast("codec", new HttpServerCodec()); + ch.pipeline().addLast("deflater", new HttpContentCompressor()); + ch.pipeline().addLast("chunkedWriter", new ChunkedWriteHandler()); ch.pipeline().addLast("aggregator", new HttpObjectAggregator(MAX_CONTENT_LENGTH)); ch.pipeline().addLast("request", new RequestHandler( securityManager, fileManager