diff --git a/pom.xml b/pom.xml index 7d7ab6ec6..51c8f0268 100644 --- a/pom.xml +++ b/pom.xml @@ -198,6 +198,7 @@ vertx-grpc-common + vertx-grpc-compression vertx-grpc-transcoding vertx-grpc-server vertx-grpc-reflection @@ -274,4 +275,4 @@ - \ No newline at end of file + diff --git a/vertx-grpc-client/pom.xml b/vertx-grpc-client/pom.xml index 933ce3d05..1521ee43b 100644 --- a/vertx-grpc-client/pom.xml +++ b/vertx-grpc-client/pom.xml @@ -42,6 +42,12 @@ test-jar test + + io.vertx + vertx-grpc-compression + ${project.version} + test + io.grpc grpc-protobuf diff --git a/vertx-grpc-client/src/main/java/io/vertx/grpc/client/GrpcClientCompressionOptions.java b/vertx-grpc-client/src/main/java/io/vertx/grpc/client/GrpcClientCompressionOptions.java new file mode 100644 index 000000000..6f2eadad5 --- /dev/null +++ b/vertx-grpc-client/src/main/java/io/vertx/grpc/client/GrpcClientCompressionOptions.java @@ -0,0 +1,75 @@ +package io.vertx.grpc.client; + +import io.vertx.codegen.annotations.DataObject; +import io.vertx.codegen.annotations.Unstable; +import io.vertx.codegen.json.annotations.JsonGen; +import io.vertx.core.json.JsonObject; +import io.vertx.grpc.common.GrpcCompressionOptions; + +@Unstable +@DataObject +@JsonGen(publicConverter = false) +public class GrpcClientCompressionOptions extends GrpcCompressionOptions { + + /** + * The default compression algorithm accepted by the client = {@code gzip} + */ + public static final String DEFAULT_COMPRESSION_ALGORITHM = "identity"; + + private String compressionAlgorithm; + + /** + * Default options. + */ + public GrpcClientCompressionOptions() { + this.compressionAlgorithm = DEFAULT_COMPRESSION_ALGORITHM; + } + + /** + * Copy constructor. + */ + public GrpcClientCompressionOptions(GrpcClientCompressionOptions other) { + super(other); + this.compressionAlgorithm = other.compressionAlgorithm; + } + + /** + * Creates options from JSON. + */ + public GrpcClientCompressionOptions(JsonObject json) { + this(); + GrpcClientCompressionOptionsConverter.fromJson(json, this); + } + + /** + * @return the compression algorithm accepted by the client + */ + public String getCompressionAlgorithm() { + return compressionAlgorithm; + } + + /** + * Set the compression algorithm accepted by the client. + * + * @param compressionAlgorithm the compression algorithm + * @return a reference to this, so the API can be used fluently + */ + public GrpcClientCompressionOptions setCompressionAlgorithm(String compressionAlgorithm) { + this.compressionAlgorithm = compressionAlgorithm; + return this; + } + + /** + * @return a JSON representation of options + */ + public JsonObject toJson() { + JsonObject json = new JsonObject(); + GrpcClientCompressionOptionsConverter.toJson(this, json); + return json; + } + + @Override + public String toString() { + return toJson().encode(); + } +} diff --git a/vertx-grpc-client/src/main/java/io/vertx/grpc/client/GrpcClientOptions.java b/vertx-grpc-client/src/main/java/io/vertx/grpc/client/GrpcClientOptions.java index bb979a92d..71df181c0 100644 --- a/vertx-grpc-client/src/main/java/io/vertx/grpc/client/GrpcClientOptions.java +++ b/vertx-grpc-client/src/main/java/io/vertx/grpc/client/GrpcClientOptions.java @@ -11,6 +11,9 @@ package io.vertx.grpc.client; import io.vertx.codegen.annotations.DataObject; +import io.vertx.codegen.annotations.Unstable; +import io.vertx.codegen.json.annotations.JsonGen; +import io.vertx.core.json.JsonObject; import java.util.Objects; import java.util.concurrent.TimeUnit; @@ -18,7 +21,9 @@ /** * Options configuring a gRPC client. */ +@Unstable @DataObject +@JsonGen(publicConverter = false) public class GrpcClientOptions { /** @@ -41,19 +46,26 @@ public class GrpcClientOptions { */ public static final long DEFAULT_MAX_MESSAGE_SIZE = 256 * 1024; + /** + * The default compression options + */ + public static final GrpcClientCompressionOptions DEFAULT_COMPRESSION = new GrpcClientCompressionOptions(); + private boolean scheduleDeadlineAutomatically; private int timeout; private TimeUnit timeoutUnit; private long maxMessageSize; + private GrpcClientCompressionOptions compression; /** * Default constructor. */ public GrpcClientOptions() { - scheduleDeadlineAutomatically = DEFAULT_SCHEDULE_DEADLINE_AUTOMATICALLY; - timeout = DEFAULT_TIMEOUT; - timeoutUnit = DEFAULT_TIMEOUT_UNIT; + this.scheduleDeadlineAutomatically = DEFAULT_SCHEDULE_DEADLINE_AUTOMATICALLY; + this.timeout = DEFAULT_TIMEOUT; + this.timeoutUnit = DEFAULT_TIMEOUT_UNIT; this.maxMessageSize = DEFAULT_MAX_MESSAGE_SIZE; + this.compression = DEFAULT_COMPRESSION; } /** @@ -62,10 +74,21 @@ public GrpcClientOptions() { * @param other the options to copy */ public GrpcClientOptions(GrpcClientOptions other) { - scheduleDeadlineAutomatically = other.scheduleDeadlineAutomatically; - timeout = other.timeout; - timeoutUnit = other.timeoutUnit; - maxMessageSize = other.maxMessageSize; + this.scheduleDeadlineAutomatically = other.scheduleDeadlineAutomatically; + this.timeout = other.timeout; + this.timeoutUnit = other.timeoutUnit; + this.maxMessageSize = other.maxMessageSize; + this.compression = new GrpcClientCompressionOptions(other.compression); + } + + /** + * Create a client options from JSON. + * + * @param json the JSON + */ + public GrpcClientOptions(JsonObject json) { + this(); + GrpcClientOptionsConverter.fromJson(json, this); } /** @@ -158,4 +181,36 @@ public GrpcClientOptions setMaxMessageSize(long maxMessageSize) { this.maxMessageSize = maxMessageSize; return this; } + + /** + * @return the compression options + */ + public GrpcClientCompressionOptions getCompression() { + return compression; + } + + /** + * Set the compression options. + * + * @param compression the compression options + * @return a reference to this, so the API can be used fluently + */ + public GrpcClientOptions setCompression(GrpcClientCompressionOptions compression) { + this.compression = Objects.requireNonNull(compression); + return this; + } + + /** + * @return a JSON representation of options + */ + public JsonObject toJson() { + JsonObject json = new JsonObject(); + GrpcClientOptionsConverter.toJson(this, json); + return json; + } + + @Override + public String toString() { + return toJson().encode(); + } } diff --git a/vertx-grpc-client/src/main/java/io/vertx/grpc/client/impl/GrpcClientImpl.java b/vertx-grpc-client/src/main/java/io/vertx/grpc/client/impl/GrpcClientImpl.java index 04b6b59ee..6da9597e3 100644 --- a/vertx-grpc-client/src/main/java/io/vertx/grpc/client/impl/GrpcClientImpl.java +++ b/vertx-grpc-client/src/main/java/io/vertx/grpc/client/impl/GrpcClientImpl.java @@ -14,9 +14,7 @@ import io.vertx.core.Vertx; import io.vertx.core.buffer.Buffer; import io.vertx.core.http.HttpClient; -import io.vertx.core.http.HttpClientOptions; import io.vertx.core.http.HttpMethod; -import io.vertx.core.http.HttpVersion; import io.vertx.core.http.RequestOptions; import io.vertx.core.internal.ContextInternal; import io.vertx.core.internal.VertxInternal; @@ -24,11 +22,9 @@ import io.vertx.grpc.client.GrpcClient; import io.vertx.grpc.client.GrpcClientOptions; import io.vertx.grpc.client.GrpcClientRequest; -import io.vertx.grpc.common.ServiceMethod; -import io.vertx.grpc.common.GrpcMessageDecoder; -import io.vertx.grpc.common.GrpcMessageEncoder; -import io.vertx.grpc.common.GrpcLocal; +import io.vertx.grpc.common.*; +import java.util.Map; import java.util.concurrent.TimeUnit; /** @@ -44,6 +40,10 @@ public class GrpcClientImpl implements GrpcClient { private final int timeout; private final TimeUnit timeoutUnit; + private final String compressionAlgorithm; + private final Map compressors; + private final Map decompressors; + public GrpcClientImpl(Vertx vertx, HttpClient client) { this(vertx, new GrpcClientOptions(), client, false); } @@ -52,10 +52,14 @@ protected GrpcClientImpl(Vertx vertx, GrpcClientOptions grpcOptions, HttpClient this.vertx = vertx; this.client = client; this.scheduleDeadlineAutomatically = grpcOptions.getScheduleDeadlineAutomatically(); - this.maxMessageSize = grpcOptions.getMaxMessageSize();; + this.maxMessageSize = grpcOptions.getMaxMessageSize(); this.timeout = grpcOptions.getTimeout(); this.timeoutUnit = grpcOptions.getTimeoutUnit(); this.closeClient = close; + + this.compressionAlgorithm = grpcOptions.getCompression().getCompressionAlgorithm(); + this.compressors = grpcOptions.getCompression().getCompressors(); + this.decompressors = grpcOptions.getCompression().getDecompressors(); } public Vertx vertx() { @@ -70,8 +74,12 @@ public Future> request(RequestOptions options) maxMessageSize, scheduleDeadlineAutomatically, GrpcMessageEncoder.IDENTITY, - GrpcMessageDecoder.IDENTITY); + GrpcMessageDecoder.IDENTITY, + this.compressors, + this.decompressors + ); grpcRequest.init(); + grpcRequest.encoding(this.compressionAlgorithm); configureTimeout(grpcRequest); return grpcRequest; }); @@ -123,8 +131,12 @@ private Future> request(RequestOptions maxMessageSize, scheduleDeadlineAutomatically, method.encoder(), - method.decoder()); + method.decoder(), + this.compressors, + this.decompressors + ); call.init(); + call.encoding(this.compressionAlgorithm); call.serviceName(method.serviceName()); call.methodName(method.methodName()); configureTimeout(call); @@ -137,7 +149,7 @@ public Future close() { if (closeClient) { return client.close(); } else { - return ((VertxInternal)vertx).getOrCreateContext().succeededFuture(); + return ((VertxInternal) vertx).getOrCreateContext().succeededFuture(); } } } diff --git a/vertx-grpc-client/src/main/java/io/vertx/grpc/client/impl/GrpcClientRequestImpl.java b/vertx-grpc-client/src/main/java/io/vertx/grpc/client/impl/GrpcClientRequestImpl.java index fcbf5f100..4c7ea5b71 100644 --- a/vertx-grpc-client/src/main/java/io/vertx/grpc/client/impl/GrpcClientRequestImpl.java +++ b/vertx-grpc-client/src/main/java/io/vertx/grpc/client/impl/GrpcClientRequestImpl.java @@ -51,8 +51,10 @@ public GrpcClientRequestImpl(HttpClientRequest httpRequest, long maxMessageSize, boolean scheduleDeadline, GrpcMessageEncoder messageEncoder, - GrpcMessageDecoder messageDecoder) { - super( ((PromiseInternal)httpRequest.response()).context(), "application/grpc", httpRequest, messageEncoder); + GrpcMessageDecoder messageDecoder, + Map compressors, + Map decompressors) { + super( ((PromiseInternal)httpRequest.response()).context(), "application/grpc", httpRequest, messageEncoder, compressors, decompressors); this.httpRequest = httpRequest; this.scheduleDeadline = scheduleDeadline; this.timeout = 0L; @@ -83,7 +85,8 @@ public GrpcClientRequestImpl(HttpClientRequest httpRequest, format, status, httpResponse, - messageDecoder); + messageDecoder, + decompressors); grpcResponse.init(this, maxMessageSize); grpcResponse.invalidMessageHandler(invalidMsg -> { cancel(); @@ -185,7 +188,7 @@ protected void setHeaders(String contentType, MultiMap headers) { if (encoding != null) { httpRequest.putHeader(GrpcHeaderNames.GRPC_ENCODING, encoding); } - httpRequest.putHeader(GrpcHeaderNames.GRPC_ACCEPT_ENCODING, "gzip"); + httpRequest.putHeader(GrpcHeaderNames.GRPC_ACCEPT_ENCODING, String.join(",", GrpcDecompressor.getSupportedEncodings())); httpRequest.putHeader(HttpHeaderNames.TE, "trailers"); httpRequest.setChunked(true); httpRequest.setURI(uri); diff --git a/vertx-grpc-client/src/main/java/io/vertx/grpc/client/impl/GrpcClientResponseImpl.java b/vertx-grpc-client/src/main/java/io/vertx/grpc/client/impl/GrpcClientResponseImpl.java index e47a2d2c4..956ad6d18 100644 --- a/vertx-grpc-client/src/main/java/io/vertx/grpc/client/impl/GrpcClientResponseImpl.java +++ b/vertx-grpc-client/src/main/java/io/vertx/grpc/client/impl/GrpcClientResponseImpl.java @@ -26,6 +26,7 @@ import io.vertx.grpc.common.impl.Http2GrpcMessageDeframer; import java.nio.charset.StandardCharsets; +import java.util.Map; /** * @author Julien Viet @@ -41,14 +42,18 @@ public GrpcClientResponseImpl(ContextInternal context, GrpcClientRequestImpl request, WireFormat format, GrpcStatus status, - HttpClientResponse httpResponse, GrpcMessageDecoder messageDecoder) { + HttpClientResponse httpResponse, + GrpcMessageDecoder messageDecoder, + Map decompressors) { super( context, httpResponse, httpResponse.headers().get(GrpcHeaderNames.GRPC_ENCODING), format, new Http2GrpcMessageDeframer(httpResponse.headers().get(GrpcHeaderNames.GRPC_ENCODING), format), - messageDecoder); + messageDecoder, + decompressors + ); this.request = request; this.httpResponse = httpResponse; this.status = status; diff --git a/vertx-grpc-client/src/main/java/module-info.java b/vertx-grpc-client/src/main/java/module-info.java index 62f52e0dd..9f440bbfe 100644 --- a/vertx-grpc-client/src/main/java/module-info.java +++ b/vertx-grpc-client/src/main/java/module-info.java @@ -1,15 +1,21 @@ module io.vertx.grpc.client{ + + requires static io.vertx.docgen; + requires static io.vertx.codegen.api; + requires static io.vertx.codegen.json; + requires io.netty.buffer; requires io.netty.codec.http; requires io.netty.codec; requires io.vertx.core.logging; requires io.vertx.core; requires io.vertx.grpc.common; - requires static io.vertx.docgen; - requires static io.vertx.codegen.api; - requires static io.vertx.codegen.json; requires com.google.protobuf; requires com.google.common; + + uses io.vertx.grpc.common.GrpcCompressor; + uses io.vertx.grpc.common.GrpcDecompressor; + exports io.vertx.grpc.client; exports io.vertx.grpc.client.impl to io.vertx.tests.client; } diff --git a/vertx-grpc-client/src/test/java/io/vertx/tests/client/ClientMessageEncodingTest.java b/vertx-grpc-client/src/test/java/io/vertx/tests/client/ClientMessageEncodingTest.java index d6a8927a2..e3f69c478 100644 --- a/vertx-grpc-client/src/test/java/io/vertx/tests/client/ClientMessageEncodingTest.java +++ b/vertx-grpc-client/src/test/java/io/vertx/tests/client/ClientMessageEncodingTest.java @@ -93,6 +93,8 @@ private void testEncode(TestContext should, String requestEncoding, GrpcMessage })); callRequest.endMessage(msg); })); + + test.await(5000L); } @Test diff --git a/vertx-grpc-client/src/test/java/io/vertx/tests/client/ClientTest.java b/vertx-grpc-client/src/test/java/io/vertx/tests/client/ClientTest.java index c49e76893..5f71fa40e 100644 --- a/vertx-grpc-client/src/test/java/io/vertx/tests/client/ClientTest.java +++ b/vertx-grpc-client/src/test/java/io/vertx/tests/client/ClientTest.java @@ -67,6 +67,11 @@ public void testUnaryCompression(TestContext should) throws IOException { testUnary(should, "gzip", "identity"); } + @Test + public void testUnaryCompressionDecompression(TestContext should) throws IOException { + testUnary(should, "gzip", "gzip"); + } + protected void testUnary(TestContext should, String requestEncoding, String responseEncoding) throws IOException { TestServiceGrpc.TestServiceImplBase called = new TestServiceGrpc.TestServiceImplBase() { @Override diff --git a/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcCompressionOptions.java b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcCompressionOptions.java new file mode 100644 index 000000000..398f27621 --- /dev/null +++ b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcCompressionOptions.java @@ -0,0 +1,193 @@ +package io.vertx.grpc.common; + +import io.vertx.codegen.annotations.DataObject; +import io.vertx.codegen.annotations.Unstable; +import io.vertx.codegen.json.annotations.JsonGen; +import io.vertx.core.json.JsonObject; + +import java.util.Map; +import java.util.ServiceLoader; +import java.util.Set; +import java.util.stream.Collectors; + +@Unstable +@DataObject +@JsonGen(publicConverter = false) +public class GrpcCompressionOptions { + + private static final Set DEFAULT_COMPRESSORS = ServiceLoader.load(GrpcCompressor.class).stream().map(ServiceLoader.Provider::get) + .collect(Collectors.toUnmodifiableSet()); + private static final Set DEFAULT_DECOMPRESSORS = ServiceLoader.load(GrpcDecompressor.class).stream().map(ServiceLoader.Provider::get) + .collect(Collectors.toUnmodifiableSet()); + + /** + * Whether compression is enabled, by default = {@code false} + */ + public static final boolean DEFAULT_COMPRESSION_ENABLED = false; + + /** + * The default compression algorithms = {@code empty} + */ + public static final Set DEFAULT_COMPRESSION_ALGORITHMS = DEFAULT_COMPRESSORS.stream().map(GrpcCompressor::encoding).collect(Collectors.toUnmodifiableSet()); + + /** + * The default decompression algorithms + */ + public static final Set DEFAULT_DECOMPRESSION_ALGORITHMS = DEFAULT_DECOMPRESSORS.stream().map(GrpcDecompressor::encoding).collect(Collectors.toUnmodifiableSet()); + + private boolean compressionEnabled; + private Set compressionAlgorithms; + private Set decompressionAlgorithms; + + /** + * Default options. + */ + public GrpcCompressionOptions() { + compressionEnabled = DEFAULT_COMPRESSION_ENABLED; + compressionAlgorithms = DEFAULT_COMPRESSION_ALGORITHMS; + decompressionAlgorithms = DEFAULT_DECOMPRESSION_ALGORITHMS; + } + + /** + * Copy constructor. + */ + public GrpcCompressionOptions(GrpcCompressionOptions other) { + compressionEnabled = other.compressionEnabled; + compressionAlgorithms = other.compressionAlgorithms; + decompressionAlgorithms = other.decompressionAlgorithms; + } + + /** + * Creates options from JSON. + */ + public GrpcCompressionOptions(JsonObject json) { + this(); + GrpcCompressionOptionsConverter.fromJson(json, this); + } + + /** + * @return whether compression is enabled + */ + public boolean isCompressionEnabled() { + return compressionEnabled; + } + + /** + * Set whether compression is enabled. + * + * @param compressionEnabled whether to enable compression + * @return a reference to this, so the API can be used fluently + */ + public GrpcCompressionOptions setCompressionEnabled(boolean compressionEnabled) { + this.compressionEnabled = compressionEnabled; + return this; + } + + /** + * @return the supported compression algorithms + */ + public Set getCompressionAlgorithms() { + return compressionAlgorithms; + } + + /** + * Set the supported compression algorithms. + * + * @param compressionAlgorithms the compression algorithms + * @return a reference to this, so the API can be used fluently + */ + public GrpcCompressionOptions setCompressionAlgorithms(Set compressionAlgorithms) { + this.compressionAlgorithms = compressionAlgorithms; + return this; + } + + /** + * Add a supported compression algorithm. + * + * @param compressionAlgorithm the compression algorithm + * @return a reference to this, so the API can be used fluently + */ + public GrpcCompressionOptions addCompressionAlgorithm(String compressionAlgorithm) { + this.compressionAlgorithms.add(compressionAlgorithm); + return this; + } + + /** + * Remove a supported compression algorithm. + * + * @param compressionAlgorithm the compression algorithm + * @return a reference to this, so the API can be used fluently + */ + public GrpcCompressionOptions removeCompressionAlgorithm(String compressionAlgorithm) { + this.compressionAlgorithms.remove(compressionAlgorithm); + return this; + } + + /** + * @return the supported decompression algorithms + */ + public Set getDecompressionAlgorithms() { + return decompressionAlgorithms; + } + + /** + * Set the supported decompression algorithms. + * + * @param decompressionAlgorithms the decompression algorithms + * @return a reference to this, so the API can be used fluently + */ + public GrpcCompressionOptions setDecompressionAlgorithms(Set decompressionAlgorithms) { + this.decompressionAlgorithms = decompressionAlgorithms; + return this; + } + + /** + * Add a supported decompression algorithm. + * + * @param decompressionAlgorithm the decompression algorithm + * @return a reference to this, so the API can be used fluently + */ + public GrpcCompressionOptions addDecompressionAlgorithm(String decompressionAlgorithm) { + this.decompressionAlgorithms.add(decompressionAlgorithm); + return this; + } + + /** + * Remove a supported decompression algorithm. + * + * @param decompressionAlgorithm the decompression algorithm + * @return a reference to this, so the API can be used fluently + */ + public GrpcCompressionOptions removeDecompressionAlgorithm(String decompressionAlgorithm) { + this.decompressionAlgorithms.remove(decompressionAlgorithm); + return this; + } + + public Map getCompressors() { + return ServiceLoader.load(GrpcCompressor.class) + .stream().map(ServiceLoader.Provider::get) + .filter(compressor -> compressionAlgorithms.contains(compressor.encoding())) + .collect(Collectors.toUnmodifiableMap(GrpcCompressor::encoding, c -> c)); + } + + public Map getDecompressors() { + return ServiceLoader.load(GrpcDecompressor.class) + .stream().map(ServiceLoader.Provider::get) + .filter(decompressor -> decompressionAlgorithms.contains(decompressor.encoding())) + .collect(Collectors.toUnmodifiableMap(GrpcDecompressor::encoding, d -> d)); + } + + /** + * @return a JSON representation of options + */ + public JsonObject toJson() { + JsonObject json = new JsonObject(); + GrpcCompressionOptionsConverter.toJson(this, json); + return json; + } + + @Override + public String toString() { + return toJson().encode(); + } +} diff --git a/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcCompressor.java b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcCompressor.java new file mode 100644 index 000000000..5e277d83d --- /dev/null +++ b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcCompressor.java @@ -0,0 +1,52 @@ +/* + * Copyright (c) 2011-2025 Contributors to the Eclipse Foundation + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0 + * which is available at https://www.apache.org/licenses/LICENSE-2.0. + * + * SPDX-License-Identifier: EPL-2.0 OR Apache-2.0 + */ +package io.vertx.grpc.common; + +import io.vertx.codegen.annotations.VertxGen; +import io.vertx.core.buffer.Buffer; + +import java.util.ServiceLoader; +import java.util.Set; +import java.util.stream.Collectors; + +/** + * A compressor for gRPC messages. + */ +@VertxGen +public interface GrpcCompressor { + + Set COMPRESSORS = ServiceLoader.load(GrpcCompressor.class).stream().map(ServiceLoader.Provider::get).collect(Collectors.toUnmodifiableSet()); + + static Set getDefaultCompressors() { + return COMPRESSORS; + } + + static GrpcCompressor lookupCompressor(String encoding) { + return getDefaultCompressors().stream() + .filter(compressor -> compressor.encoding().equals(encoding)) + .findFirst() + .orElse(null); + } + + /** + * @return the encoding name of this compressor (e.g., "gzip") + */ + String encoding(); + + /** + * Compresses the given buffer. + * + * @param data the buffer to compress + * @return the compressed buffer + * @throws CodecException if compression fails + */ + Buffer compress(Buffer data) throws CodecException; +} diff --git a/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcDecompressor.java b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcDecompressor.java new file mode 100644 index 000000000..8050575fe --- /dev/null +++ b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcDecompressor.java @@ -0,0 +1,56 @@ +/* + * Copyright (c) 2011-2025 Contributors to the Eclipse Foundation + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0 + * which is available at https://www.apache.org/licenses/LICENSE-2.0. + * + * SPDX-License-Identifier: EPL-2.0 OR Apache-2.0 + */ +package io.vertx.grpc.common; + +import io.vertx.codegen.annotations.VertxGen; +import io.vertx.core.buffer.Buffer; + +import java.util.ServiceLoader; +import java.util.Set; +import java.util.stream.Collectors; + +/** + * A decompressor for gRPC messages. + */ +@VertxGen +public interface GrpcDecompressor { + + Set DECOMPRESSORS = ServiceLoader.load(GrpcDecompressor.class).stream().map(ServiceLoader.Provider::get).collect(Collectors.toUnmodifiableSet()); + + static Set getDefaultDecompressors() { + return DECOMPRESSORS; + } + + static Set getSupportedEncodings() { + return getDefaultDecompressors().stream().map(GrpcDecompressor::encoding).collect(Collectors.toUnmodifiableSet()); + } + + static GrpcDecompressor lookupDecompressor(String encoding) { + return getDefaultDecompressors().stream() + .filter(decompressor -> decompressor.encoding().equals(encoding)) + .findFirst() + .orElse(null); + } + + /** + * @return the encoding name of this decompressor (e.g., "gzip") + */ + String encoding(); + + /** + * Decompresses the given buffer. + * + * @param data the buffer to decompress + * @return the decompressed buffer + * @throws CodecException if decompression fails + */ + Buffer decompress(Buffer data) throws CodecException; +} diff --git a/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcRequestTransformer.java b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcRequestTransformer.java new file mode 100644 index 000000000..c9313f568 --- /dev/null +++ b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcRequestTransformer.java @@ -0,0 +1,7 @@ +package io.vertx.grpc.common; + +import io.vertx.core.streams.Pipe; + +public interface GrpcRequestTransformer extends Pipe { + +} diff --git a/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcResponseTransformer.java b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcResponseTransformer.java new file mode 100644 index 000000000..da9c4ad4b --- /dev/null +++ b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/GrpcResponseTransformer.java @@ -0,0 +1,6 @@ +package io.vertx.grpc.common; + +import io.vertx.core.streams.Pipe; + +public interface GrpcResponseTransformer extends Pipe { +} diff --git a/vertx-grpc-common/src/main/java/io/vertx/grpc/common/impl/GrpcReadStreamBase.java b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/impl/GrpcReadStreamBase.java index 1f297ef30..da1af69b0 100644 --- a/vertx-grpc-common/src/main/java/io/vertx/grpc/common/impl/GrpcReadStreamBase.java +++ b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/impl/GrpcReadStreamBase.java @@ -22,6 +22,8 @@ import io.vertx.core.streams.ReadStream; import io.vertx.grpc.common.*; +import java.util.Map; + import static io.vertx.grpc.common.GrpcError.mapHttp2ErrorCode; /** @@ -47,6 +49,8 @@ public Buffer payload() { }; protected final ContextInternal context; + protected final Map decompressors; + private final String encoding; private final WireFormat format; private final ReadStream stream; @@ -66,7 +70,8 @@ protected GrpcReadStreamBase(Context context, String encoding, WireFormat format, GrpcMessageDeframer messageDeframer, - GrpcMessageDecoder messageDecoder) { + GrpcMessageDecoder messageDecoder, + Map decompressors) { ContextInternal ctx = (ContextInternal) context; this.context = ctx; this.encoding = encoding; @@ -91,6 +96,7 @@ protected void handleMessage(GrpcMessage msg) { } }; this.messageDecoder = messageDecoder; + this.decompressors = decompressors; this.end = ctx.promise(); this.deframer = messageDeframer; } @@ -116,16 +122,13 @@ public void init(GrpcWriteStreamBase ws, long maxMessageSize) { } protected final T decodeMessage(GrpcMessage msg) throws CodecException { - switch (msg.encoding()) { - case "identity": - // Nothing to do - break; - case "gzip": { - msg = GrpcMessage.message("identity", msg.format(), Utils.GZIP_DECODER.apply(msg.payload())); - break; + String encoding = msg.encoding(); + if (!encoding.equals("identity")) { + GrpcDecompressor decompressor = this.decompressors.get(encoding); + if (decompressor == null) { + throw new UnsupportedOperationException("Unsupported encoding: " + encoding); } - default: - throw new UnsupportedOperationException(); + msg = GrpcMessage.message("identity", msg.format(), decompressor.decompress(msg.payload())); } return messageDecoder.decode(msg); } diff --git a/vertx-grpc-common/src/main/java/io/vertx/grpc/common/impl/GrpcWriteStreamBase.java b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/impl/GrpcWriteStreamBase.java index acf9dbafe..bbd8ad2bd 100644 --- a/vertx-grpc-common/src/main/java/io/vertx/grpc/common/impl/GrpcWriteStreamBase.java +++ b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/impl/GrpcWriteStreamBase.java @@ -9,7 +9,9 @@ import io.vertx.core.streams.WriteStream; import io.vertx.grpc.common.*; +import java.util.Map; import java.util.Objects; +import java.util.Optional; import static io.vertx.grpc.common.GrpcError.mapHttp2ErrorCode; @@ -19,7 +21,10 @@ public abstract class GrpcWriteStreamBase, T private final GrpcMessageEncoder messageEncoder; private final WriteStream writeStream; - protected String mediaType; + protected final String mediaType; + protected final Map compressors; + protected final Map decompressors; + protected String encoding; protected WireFormat format; private boolean headersSent; @@ -33,11 +38,13 @@ public abstract class GrpcWriteStreamBase, T private Handler exceptionHandler; private Handler errorHandler; - public GrpcWriteStreamBase(ContextInternal context, String mediaType, WriteStream writeStream, GrpcMessageEncoder messageEncoder) { + public GrpcWriteStreamBase(ContextInternal context, String mediaType, WriteStream writeStream, GrpcMessageEncoder messageEncoder, Map compressors, Map decompressors) { this.context = context; this.writeStream = writeStream; this.messageEncoder = messageEncoder; this.mediaType = mediaType; + this.compressors = compressors; + this.decompressors = decompressors; } public void init() { @@ -175,11 +182,7 @@ public final Future end(T message) { } private GrpcMessage encodeMessage(T message) { - WireFormat f = format; - if (f == null) { - f = WireFormat.PROTOBUF; - } - return messageEncoder.encode(message, f); + return messageEncoder.encode(message, Optional.ofNullable(format).orElse(WireFormat.PROTOBUF)); } @Override @@ -242,39 +245,51 @@ protected Future writeMessage(GrpcMessage message, boolean end) { boolean compressed; if (message != null) { if (encoding != null) { - switch (encoding) { - case "gzip": - compressed = true; - if (message.encoding().equals("identity")) { - try { - payload = Utils.GZIP_ENCODER.apply(message.payload()); - } catch (CodecException e) { - return Future.failedFuture(e); - } - } else { - if (!message.encoding().equals("gzip")) { - return Future.failedFuture("Encoding " + message.encoding() + " is not supported"); - } - payload = message.payload(); - } - break; - case "identity": + if (message.encoding().equals(encoding)) { + // Message is already in the desired encoding + compressed = !encoding.equals("identity"); + payload = message.payload(); + } else if (message.encoding().equals("identity")) { + // Message is in identity encoding, need to compress + GrpcCompressor compressor = compressors.get(encoding); + if (compressor == null) { + return Future.failedFuture("Encoding " + encoding + " is not supported"); + } + compressed = !encoding.equals("identity"); + try { + payload = compressor.compress(message.payload()); + } catch (CodecException e) { + return Future.failedFuture(e); + } + } else { + // Message is in some other encoding, need to decompress first then compress + GrpcDecompressor decompressor = decompressors.get(message.encoding()); + if (decompressor == null) { + return Future.failedFuture("Encoding " + message.encoding() + " is not supported"); + } + + Buffer decompressed; + try { + decompressed = decompressor.decompress(message.payload()); + } catch (CodecException e) { + return Future.failedFuture(e); + } + + if (encoding.equals("identity")) { compressed = false; - if (!message.encoding().equals("identity")) { - if (!message.encoding().equals("gzip")) { - return Future.failedFuture("Encoding " + message.encoding() + " is not supported"); - } - try { - payload = Utils.GZIP_DECODER.apply(message.payload()); - } catch (CodecException e) { - return Future.failedFuture(e); - } - } else { - payload = message.payload(); + payload = decompressed; + } else { + GrpcCompressor compressor = compressors.get(encoding); + if (compressor == null) { + return Future.failedFuture("Encoding " + encoding + " is not supported"); } - break; - default: - return Future.failedFuture("Encoding " + encoding + " is not supported"); + compressed = true; + try { + payload = compressor.compress(decompressed); + } catch (CodecException e) { + return Future.failedFuture(e); + } + } } } else { compressed = !message.encoding().equals("identity"); diff --git a/vertx-grpc-common/src/main/java/io/vertx/grpc/common/impl/IdentityCompressor.java b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/impl/IdentityCompressor.java new file mode 100644 index 000000000..b9441174a --- /dev/null +++ b/vertx-grpc-common/src/main/java/io/vertx/grpc/common/impl/IdentityCompressor.java @@ -0,0 +1,34 @@ +/* + * Copyright (c) 2011-2025 Contributors to the Eclipse Foundation + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0 + * which is available at https://www.apache.org/licenses/LICENSE-2.0. + * + * SPDX-License-Identifier: EPL-2.0 OR Apache-2.0 + */ +package io.vertx.grpc.common.impl; + +import io.vertx.core.buffer.Buffer; +import io.vertx.grpc.common.CodecException; +import io.vertx.grpc.common.GrpcCompressor; +import io.vertx.grpc.common.GrpcDecompressor; + +public class IdentityCompressor implements GrpcCompressor, GrpcDecompressor { + + @Override + public String encoding() { + return "identity"; + } + + @Override + public Buffer compress(Buffer data) throws CodecException { + return data; + } + + @Override + public Buffer decompress(Buffer data) throws CodecException { + return data; + } +} diff --git a/vertx-grpc-common/src/main/java/module-info.java b/vertx-grpc-common/src/main/java/module-info.java index 553878629..c9d50ad7e 100644 --- a/vertx-grpc-common/src/main/java/module-info.java +++ b/vertx-grpc-common/src/main/java/module-info.java @@ -1,6 +1,8 @@ module io.vertx.grpc.common { requires static io.vertx.codegen.api; + requires static io.vertx.codegen.json; + requires static io.vertx.docgen; requires io.vertx.core; requires io.netty.common; @@ -11,8 +13,13 @@ requires com.google.protobuf; requires com.google.protobuf.util; + uses io.vertx.grpc.common.GrpcCompressor; + uses io.vertx.grpc.common.GrpcDecompressor; + exports io.vertx.grpc.common; exports io.vertx.grpc.common.impl to io.vertx.tests.common, io.vertx.grpc.server, io.vertx.grpc.client, io.vertx.grpc.transcoding, io.vertx.tests.server, io.vertx.tests.client; provides io.vertx.core.spi.VertxServiceProvider with io.vertx.grpc.common.impl.GrpcRequestLocalRegistration; + provides io.vertx.grpc.common.GrpcCompressor with io.vertx.grpc.common.impl.IdentityCompressor; + provides io.vertx.grpc.common.GrpcDecompressor with io.vertx.grpc.common.impl.IdentityCompressor; } diff --git a/vertx-grpc-common/src/main/resources/META-INF/services/io.vertx.grpc.common.GrpcCompressor b/vertx-grpc-common/src/main/resources/META-INF/services/io.vertx.grpc.common.GrpcCompressor new file mode 100644 index 000000000..5a54e1853 --- /dev/null +++ b/vertx-grpc-common/src/main/resources/META-INF/services/io.vertx.grpc.common.GrpcCompressor @@ -0,0 +1 @@ +io.vertx.grpc.common.impl.IdentityCompressor diff --git a/vertx-grpc-common/src/main/resources/META-INF/services/io.vertx.grpc.common.GrpcDecompressor b/vertx-grpc-common/src/main/resources/META-INF/services/io.vertx.grpc.common.GrpcDecompressor new file mode 100644 index 000000000..5a54e1853 --- /dev/null +++ b/vertx-grpc-common/src/main/resources/META-INF/services/io.vertx.grpc.common.GrpcDecompressor @@ -0,0 +1 @@ +io.vertx.grpc.common.impl.IdentityCompressor diff --git a/vertx-grpc-common/src/test/java/io/vertx/tests/common/GrpcTestBase.java b/vertx-grpc-common/src/test/java/io/vertx/tests/common/GrpcTestBase.java index 368a81c42..210bfa095 100644 --- a/vertx-grpc-common/src/test/java/io/vertx/tests/common/GrpcTestBase.java +++ b/vertx-grpc-common/src/test/java/io/vertx/tests/common/GrpcTestBase.java @@ -11,7 +11,6 @@ package io.vertx.tests.common; import io.vertx.core.Vertx; -import io.vertx.core.VertxOptions; import io.vertx.core.buffer.Buffer; import io.vertx.ext.unit.TestContext; import io.vertx.ext.unit.junit.VertxUnitRunner; @@ -25,6 +24,9 @@ import java.util.zip.GZIPInputStream; import java.util.zip.GZIPOutputStream; +import io.vertx.grpc.common.GrpcCompressor; +import io.vertx.grpc.common.GrpcDecompressor; + /** * @author Julien Viet */ @@ -72,4 +74,14 @@ public static Buffer zip(Buffer buffer) { } return Buffer.buffer(ret.toByteArray()); } + + public static Buffer snappyCompress(Buffer buffer) { + GrpcCompressor compressor = GrpcCompressor.lookupCompressor("snappy"); + return compressor.compress(buffer); + } + + public static Buffer snappyDecompress(Buffer buffer) { + GrpcDecompressor decompressor = GrpcDecompressor.lookupDecompressor("snappy"); + return decompressor.decompress(buffer); + } } diff --git a/vertx-grpc-common/src/test/java/module-info.java b/vertx-grpc-common/src/test/java/module-info.java index ef679790b..b1934a534 100644 --- a/vertx-grpc-common/src/test/java/module-info.java +++ b/vertx-grpc-common/src/test/java/module-info.java @@ -8,6 +8,10 @@ requires io.grpc; requires io.grpc.protobuf; requires io.grpc.stub; + + uses io.vertx.grpc.common.GrpcCompressor; + uses io.vertx.grpc.common.GrpcDecompressor; + exports io.vertx.tests.common; exports io.vertx.tests.common.grpc; } diff --git a/vertx-grpc-compression/pom.xml b/vertx-grpc-compression/pom.xml new file mode 100644 index 000000000..da506fc4c --- /dev/null +++ b/vertx-grpc-compression/pom.xml @@ -0,0 +1,38 @@ + + + + + 4.0.0 + + + io.vertx + vertx-grpc-aggregator + 5.1.0-SNAPSHOT + ../pom.xml + + + vertx-grpc-compression + + Vert.x gRPC Compression + + + + io.vertx + vertx-grpc-common + + + diff --git a/vertx-grpc-compression/src/main/java/io/vertx/grpc/compression/CompressionUtils.java b/vertx-grpc-compression/src/main/java/io/vertx/grpc/compression/CompressionUtils.java new file mode 100644 index 000000000..4eba936ab --- /dev/null +++ b/vertx-grpc-compression/src/main/java/io/vertx/grpc/compression/CompressionUtils.java @@ -0,0 +1,98 @@ +/* + * Copyright (c) 2011-2025 Contributors to the Eclipse Foundation + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0 + * which is available at https://www.apache.org/licenses/LICENSE-2.0. + * + * SPDX-License-Identifier: EPL-2.0 OR Apache-2.0 + */ +package io.vertx.grpc.compression; + +import io.netty.buffer.ByteBuf; +import io.netty.buffer.CompositeByteBuf; +import io.netty.buffer.Unpooled; +import io.netty.channel.ChannelFuture; +import io.netty.channel.ChannelHandler; +import io.netty.channel.embedded.EmbeddedChannel; +import io.vertx.core.buffer.Buffer; +import io.vertx.core.internal.buffer.BufferInternal; +import io.vertx.grpc.common.CodecException; + +import java.util.Queue; + +public final class CompressionUtils { + + private CompressionUtils() { + } + + /** + * Helper method to decode data using a specified decoder + * + * @param data the data to decode + * @param decoder the decoder to use + * @param errorMessage the error message to use if decoding fails + * @return the decoded buffer + * @throws CodecException if decoding fails + */ + public static Buffer decode(Buffer data, ChannelHandler decoder, String errorMessage) { + if (data.length() == 0) { + return BufferInternal.buffer(); + } + + EmbeddedChannel channel = new EmbeddedChannel(decoder); + channel.config().setAllocator(BufferInternal.buffer().getByteBuf().alloc()); + try { + ChannelFuture fut = channel.writeOneInbound(((BufferInternal) data).getByteBuf()); + if (fut.isSuccess()) { + Buffer decoded = null; + while (true) { + ByteBuf buf = channel.readInbound(); + if (buf == null) { + break; + } + if (decoded == null) { + decoded = BufferInternal.buffer(buf); + } else { + decoded.appendBuffer(BufferInternal.buffer(buf)); + } + } + if (decoded == null) { + throw new CodecException(errorMessage); + } + return decoded; + } else { + throw new CodecException(fut.cause()); + } + } finally { + channel.close(); + } + } + + /** + * Helper method to encode data using a specified encoder + * + * @param data the data to encode + * @param encoder the encoder to use + * @return the encoded buffer + */ + public static Buffer encode(Buffer data, ChannelHandler encoder) { + if (data.length() == 0) { + return BufferInternal.buffer(); + } + + CompositeByteBuf composite = Unpooled.compositeBuffer(); + EmbeddedChannel channel = new EmbeddedChannel(encoder); + channel.config().setAllocator(BufferInternal.buffer().getByteBuf().alloc()); + channel.writeOutbound(((BufferInternal) data).getByteBuf()); + channel.finish(); + Queue messages = channel.outboundMessages(); + ByteBuf a; + while ((a = (ByteBuf) messages.poll()) != null) { + composite.addComponent(true, a); + } + channel.close(); + return BufferInternal.buffer(composite); + } +} diff --git a/vertx-grpc-compression/src/main/java/io/vertx/grpc/compression/GzipCompressor.java b/vertx-grpc-compression/src/main/java/io/vertx/grpc/compression/GzipCompressor.java new file mode 100644 index 000000000..37e836509 --- /dev/null +++ b/vertx-grpc-compression/src/main/java/io/vertx/grpc/compression/GzipCompressor.java @@ -0,0 +1,44 @@ +/* + * Copyright (c) 2011-2025 Contributors to the Eclipse Foundation + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0 + * which is available at https://www.apache.org/licenses/LICENSE-2.0. + * + * SPDX-License-Identifier: EPL-2.0 OR Apache-2.0 + */ +package io.vertx.grpc.compression; + +import io.netty.handler.codec.compression.*; +import io.vertx.core.buffer.Buffer; +import io.vertx.grpc.common.CodecException; +import io.vertx.grpc.common.GrpcCompressor; +import io.vertx.grpc.common.GrpcDecompressor; + +import java.util.function.Function; + +public class GzipCompressor implements GrpcCompressor, GrpcDecompressor { + + public static final Function GZIP_DECODER = data -> CompressionUtils.decode(data, ZlibCodecFactory.newZlibDecoder(ZlibWrapper.GZIP), "Invalid GZIP input"); + public static final Function GZIP_ENCODER = data -> { + GzipOptions options = StandardCompressionOptions.gzip(); + ZlibEncoder encoder = ZlibCodecFactory.newZlibEncoder(ZlibWrapper.GZIP, options.compressionLevel(), options.windowBits(), options.memLevel()); + return CompressionUtils.encode(data, encoder); + }; + + @Override + public String encoding() { + return "gzip"; + } + + @Override + public Buffer compress(Buffer data) throws CodecException { + return GZIP_ENCODER.apply(data); + } + + @Override + public Buffer decompress(Buffer data) throws CodecException { + return GZIP_DECODER.apply(data); + } +} diff --git a/vertx-grpc-compression/src/main/java/io/vertx/grpc/compression/SnappyCompressor.java b/vertx-grpc-compression/src/main/java/io/vertx/grpc/compression/SnappyCompressor.java new file mode 100644 index 000000000..c072f792b --- /dev/null +++ b/vertx-grpc-compression/src/main/java/io/vertx/grpc/compression/SnappyCompressor.java @@ -0,0 +1,41 @@ +/* + * Copyright (c) 2011-2025 Contributors to the Eclipse Foundation + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0 + * which is available at https://www.apache.org/licenses/LICENSE-2.0. + * + * SPDX-License-Identifier: EPL-2.0 OR Apache-2.0 + */ +package io.vertx.grpc.compression; + +import io.netty.handler.codec.compression.SnappyFrameDecoder; +import io.netty.handler.codec.compression.SnappyFrameEncoder; +import io.vertx.core.buffer.Buffer; +import io.vertx.grpc.common.CodecException; +import io.vertx.grpc.common.GrpcCompressor; +import io.vertx.grpc.common.GrpcDecompressor; + +import java.util.function.Function; + +public class SnappyCompressor implements GrpcCompressor, GrpcDecompressor { + + public static final Function SNAPPY_DECODER = data -> CompressionUtils.decode(data, new SnappyFrameDecoder(), "Invalid Snappy input"); + public static final Function SNAPPY_ENCODER = data -> CompressionUtils.encode(data, new SnappyFrameEncoder()); + + @Override + public String encoding() { + return "snappy"; + } + + @Override + public Buffer compress(Buffer data) throws CodecException { + return SNAPPY_ENCODER.apply(data); + } + + @Override + public Buffer decompress(Buffer data) throws CodecException { + return SNAPPY_DECODER.apply(data); + } +} diff --git a/vertx-grpc-compression/src/main/java/module-info.java b/vertx-grpc-compression/src/main/java/module-info.java new file mode 100644 index 000000000..e2b5b4cc0 --- /dev/null +++ b/vertx-grpc-compression/src/main/java/module-info.java @@ -0,0 +1,19 @@ +module io.vertx.grpc.compression { + + requires static io.vertx.codegen.api; + + requires io.netty.common; + requires io.netty.buffer; + requires io.netty.codec; + requires io.netty.codec.compression; + requires io.netty.transport; + requires com.google.protobuf; + requires com.google.protobuf.util; + requires io.vertx.core; + requires io.vertx.grpc.common; + + exports io.vertx.grpc.compression to io.vertx.grpc.client, io.vertx.grpc.server, io.vertx.grpc.transcoding, io.vertx.tests.client, io.vertx.tests.common, io.vertx.tests.server; + + provides io.vertx.grpc.common.GrpcCompressor with io.vertx.grpc.compression.GzipCompressor, io.vertx.grpc.compression.SnappyCompressor; + provides io.vertx.grpc.common.GrpcDecompressor with io.vertx.grpc.compression.GzipCompressor, io.vertx.grpc.compression.SnappyCompressor; +} diff --git a/vertx-grpc-compression/src/main/resources/META-INF/services/io.vertx.grpc.common.GrpcCompressor b/vertx-grpc-compression/src/main/resources/META-INF/services/io.vertx.grpc.common.GrpcCompressor new file mode 100644 index 000000000..4a989d24b --- /dev/null +++ b/vertx-grpc-compression/src/main/resources/META-INF/services/io.vertx.grpc.common.GrpcCompressor @@ -0,0 +1,3 @@ +io.vertx.grpc.compression.GzipCompressor +io.vertx.grpc.common.impl.IdentityCompressor +io.vertx.grpc.compression.SnappyCompressor diff --git a/vertx-grpc-compression/src/main/resources/META-INF/services/io.vertx.grpc.common.GrpcDecompressor b/vertx-grpc-compression/src/main/resources/META-INF/services/io.vertx.grpc.common.GrpcDecompressor new file mode 100644 index 000000000..4a989d24b --- /dev/null +++ b/vertx-grpc-compression/src/main/resources/META-INF/services/io.vertx.grpc.common.GrpcDecompressor @@ -0,0 +1,3 @@ +io.vertx.grpc.compression.GzipCompressor +io.vertx.grpc.common.impl.IdentityCompressor +io.vertx.grpc.compression.SnappyCompressor diff --git a/vertx-grpc-compression/src/test/java/io/vertx/tests/compression/CompressorTestBase.java b/vertx-grpc-compression/src/test/java/io/vertx/tests/compression/CompressorTestBase.java new file mode 100644 index 000000000..b4ee6ff3e --- /dev/null +++ b/vertx-grpc-compression/src/test/java/io/vertx/tests/compression/CompressorTestBase.java @@ -0,0 +1,110 @@ +/* + * Copyright (c) 2011-2025 Contributors to the Eclipse Foundation + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0 + * which is available at https://www.apache.org/licenses/LICENSE-2.0. + * + * SPDX-License-Identifier: EPL-2.0 OR Apache-2.0 + */ +package io.vertx.tests.compression; + +import io.vertx.core.buffer.Buffer; +import io.vertx.ext.unit.TestContext; +import io.vertx.ext.unit.junit.VertxUnitRunner; +import io.vertx.grpc.common.GrpcCompressor; +import io.vertx.grpc.common.GrpcDecompressor; +import io.vertx.grpc.common.GrpcMessage; +import io.vertx.grpc.common.WireFormat; +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; + +@RunWith(VertxUnitRunner.class) +public abstract class CompressorTestBase { + + /** + * @return the encoding name of the compressor being tested (e.g., "gzip", "snappy") + */ + protected abstract String getEncodingName(); + + protected boolean shouldReduceSize() { + return true; + } + + @Test + public void testCompressor(TestContext should) { + // Create a message with some content + String original = "Hello, World!"; + Buffer originalBuffer = Buffer.buffer(original); + GrpcMessage message = GrpcMessage.message("identity", WireFormat.PROTOBUF, originalBuffer); + + // Compress the message using the specified compressor + GrpcCompressor compressor = GrpcCompressor.lookupCompressor(getEncodingName()); + should.assertNotNull(compressor, getEncodingName() + " compressor should be registered"); + + Buffer compressed = compressor.compress(message.payload()); + GrpcMessage compressedMessage = GrpcMessage.message(getEncodingName(), message.format(), compressed); + + // Decompress the message + GrpcDecompressor decompressor = GrpcDecompressor.lookupDecompressor(getEncodingName()); + should.assertNotNull(decompressor, getEncodingName() + " decompressor should be registered"); + + Buffer decompressed = decompressor.decompress(compressedMessage.payload()); + should.assertEquals(originalBuffer, decompressed, "Decompressed message should match original"); + } + + @Test + public void testCompressorWithLargePayload(TestContext should) { + // Create a message with a large content + StringBuilder sb = new StringBuilder(); + for (int i = 0; i < 1000; i++) { + sb.append("This is a test message with some repetitive content. "); + } + String original = sb.toString(); + Buffer originalBuffer = Buffer.buffer(original); + GrpcMessage message = GrpcMessage.message("identity", WireFormat.PROTOBUF, originalBuffer); + + // Compress the message using the specified compressor + GrpcCompressor compressor = GrpcCompressor.lookupCompressor(getEncodingName()); + should.assertNotNull(compressor, getEncodingName() + " compressor should be registered"); + + Buffer compressed = compressor.compress(message.payload()); + GrpcMessage compressedMessage = GrpcMessage.message(getEncodingName(), message.format(), compressed); + + // Verify compression ratio if the compressor should reduce size + if (shouldReduceSize()) { + should.assertTrue(compressed.length() < originalBuffer.length(), "Compressed message should be smaller than original"); + } + + // Decompress the message + GrpcDecompressor decompressor = GrpcDecompressor.lookupDecompressor(getEncodingName()); + should.assertNotNull(decompressor, getEncodingName() + " decompressor should be registered"); + + Buffer decompressed = decompressor.decompress(compressedMessage.payload()); + should.assertEquals(originalBuffer, decompressed, "Decompressed message should match original"); + } + + @Test + public void testCompressorWithEmptyPayload(TestContext should) { + // Create a message with an empty content + String original = ""; + Buffer originalBuffer = Buffer.buffer(original); + GrpcMessage message = GrpcMessage.message("identity", WireFormat.PROTOBUF, originalBuffer); + + // Compress the message using the specified compressor + GrpcCompressor compressor = GrpcCompressor.lookupCompressor(getEncodingName()); + should.assertNotNull(compressor, getEncodingName() + " compressor should be registered"); + + Buffer compressed = compressor.compress(message.payload()); + GrpcMessage compressedMessage = GrpcMessage.message(getEncodingName(), message.format(), compressed); + + // Decompress the message + GrpcDecompressor decompressor = GrpcDecompressor.lookupDecompressor(getEncodingName()); + should.assertNotNull(decompressor, getEncodingName() + " decompressor should be registered"); + + Buffer decompressed = decompressor.decompress(compressedMessage.payload()); + should.assertEquals(originalBuffer, decompressed, "Decompressed message should match original"); + } +} diff --git a/vertx-grpc-compression/src/test/java/io/vertx/tests/compression/CustomCompressorTest.java b/vertx-grpc-compression/src/test/java/io/vertx/tests/compression/CustomCompressorTest.java new file mode 100644 index 000000000..cc916623e --- /dev/null +++ b/vertx-grpc-compression/src/test/java/io/vertx/tests/compression/CustomCompressorTest.java @@ -0,0 +1,55 @@ +/* + * Copyright (c) 2011-2025 Contributors to the Eclipse Foundation + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0 + * which is available at https://www.apache.org/licenses/LICENSE-2.0. + * + * SPDX-License-Identifier: EPL-2.0 OR Apache-2.0 + */ +package io.vertx.tests.compression; + +import io.vertx.core.buffer.Buffer; +import io.vertx.grpc.common.CodecException; +import io.vertx.grpc.common.GrpcCompressor; +import io.vertx.grpc.common.GrpcDecompressor; + +public class CustomCompressorTest extends CompressorTestBase { + + /** + * A simple custom compressor that reverses the bytes in the buffer. + */ + public static class ReverseCompressor implements GrpcCompressor, GrpcDecompressor { + @Override + public String encoding() { + return "reverse"; + } + + @Override + public Buffer compress(Buffer data) throws CodecException { + byte[] bytes = data.getBytes(); + byte[] reversed = new byte[bytes.length]; + for (int i = 0; i < bytes.length; i++) { + reversed[i] = bytes[bytes.length - 1 - i]; + } + return Buffer.buffer(reversed); + } + + @Override + public Buffer decompress(Buffer data) throws CodecException { + // For this simple example, compression and decompression are the same operation + return compress(data); + } + } + + @Override + protected String getEncodingName() { + return "reverse"; + } + + @Override + protected boolean shouldReduceSize() { + return false; + } +} diff --git a/vertx-grpc-compression/src/test/java/io/vertx/tests/compression/GzipCompressorTest.java b/vertx-grpc-compression/src/test/java/io/vertx/tests/compression/GzipCompressorTest.java new file mode 100644 index 000000000..b0762d5d0 --- /dev/null +++ b/vertx-grpc-compression/src/test/java/io/vertx/tests/compression/GzipCompressorTest.java @@ -0,0 +1,19 @@ +/* + * Copyright (c) 2011-2025 Contributors to the Eclipse Foundation + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0 + * which is available at https://www.apache.org/licenses/LICENSE-2.0. + * + * SPDX-License-Identifier: EPL-2.0 OR Apache-2.0 + */ +package io.vertx.tests.compression; + +public class GzipCompressorTest extends CompressorTestBase { + + @Override + protected String getEncodingName() { + return "gzip"; + } +} diff --git a/vertx-grpc-compression/src/test/java/io/vertx/tests/compression/SnappyCompressorTest.java b/vertx-grpc-compression/src/test/java/io/vertx/tests/compression/SnappyCompressorTest.java new file mode 100644 index 000000000..3da570bcb --- /dev/null +++ b/vertx-grpc-compression/src/test/java/io/vertx/tests/compression/SnappyCompressorTest.java @@ -0,0 +1,19 @@ +/* + * Copyright (c) 2011-2025 Contributors to the Eclipse Foundation + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0 + * which is available at https://www.apache.org/licenses/LICENSE-2.0. + * + * SPDX-License-Identifier: EPL-2.0 OR Apache-2.0 + */ +package io.vertx.tests.compression; + +public class SnappyCompressorTest extends CompressorTestBase { + + @Override + protected String getEncodingName() { + return "snappy"; + } +} diff --git a/vertx-grpc-compression/src/test/java/module-info.java b/vertx-grpc-compression/src/test/java/module-info.java new file mode 100644 index 000000000..2ee563745 --- /dev/null +++ b/vertx-grpc-compression/src/test/java/module-info.java @@ -0,0 +1,16 @@ +import io.vertx.tests.compression.CustomCompressorTest; + +open module io.vertx.tests.compression { + requires io.vertx.core; + requires io.vertx.grpc.common; + requires io.vertx.testing.unit; + requires junit; + + uses io.vertx.grpc.common.GrpcCompressor; + uses io.vertx.grpc.common.GrpcDecompressor; + + exports io.vertx.tests.compression; + + provides io.vertx.grpc.common.GrpcCompressor with CustomCompressorTest.ReverseCompressor; + provides io.vertx.grpc.common.GrpcDecompressor with CustomCompressorTest.ReverseCompressor; +} diff --git a/vertx-grpc-docs/src/main/asciidoc/client.adoc b/vertx-grpc-docs/src/main/asciidoc/client.adoc index a3654fe4d..052d0e85d 100644 --- a/vertx-grpc-docs/src/main/asciidoc/client.adoc +++ b/vertx-grpc-docs/src/main/asciidoc/client.adoc @@ -391,7 +391,7 @@ You can also specify the JSON wire format when creating an idiomatic client === Compression -You can compress request messages by setting the request encoding *prior* before sending any message +You can compress request messages by setting a compression algorithm when creating the client, or per message. [source,java] ---- @@ -400,7 +400,8 @@ You can compress request messages by setting the request encoding *prior* before === Decompression -Decompression is achieved transparently by the client when the server sends encoded responses. +Decompression is achieved by client if it has access to proper decompressors and if the proper decompression algorithm is +registered in `GrpcClientCompressionOptions`. === Message level API diff --git a/vertx-grpc-docs/src/main/asciidoc/server.adoc b/vertx-grpc-docs/src/main/asciidoc/server.adoc index c3a580569..0e2eab554 100644 --- a/vertx-grpc-docs/src/main/asciidoc/server.adoc +++ b/vertx-grpc-docs/src/main/asciidoc/server.adoc @@ -316,19 +316,47 @@ Anemic JSON is also supported with Vert.x `JsonObject` === Compression -You can compress response messages by setting the response encoding *prior* before sending any message +You can compress response messages by setting the compression algorithm when creating the server, or per message. [source,java] ---- {@link examples.GrpcServerExamples#responseCompression} ---- +By default, Vert.x gRPC supports three compression types: +- `identity` (no compression) +- `gzip` compression +- `snappy` compression + +You can add support for custom encoding by implementing the `GrpcCompressor` interface and registering it with the `GrpcCompressorRegistry`: + +[source,java] +---- +{@link examples.GrpcServerExamples#customCompressor} +---- + +After registering your custom compressor, the server will be able to compress messages that use your custom encoding. + NOTE: Compression is not supported over the gRPC-Web protocol. === Decompression Decompression is done transparently by the server when the client send encoded requests. +By default, Vert.x gRPC supports three compression types: +- `identity` (no decompression) +- `gzip` decompression +- `snappy` decompression + +You can add support for custom encoding by implementing the `GrpcDecompressor` interface and registering it with the `GrpcDecompressorRegistry`: + +[source,java] +---- +{@link examples.GrpcServerExamples#customDecompressor} +---- + +After registering your custom decompressor, the server will be able to decompress messages that use your custom encoding. + NOTE: Decompression is not supported over the gRPC-Web protocol. === Message level API diff --git a/vertx-grpc-docs/src/main/java/examples/GrpcClientExamples.java b/vertx-grpc-docs/src/main/java/examples/GrpcClientExamples.java index e0c6dc040..69a775571 100644 --- a/vertx-grpc-docs/src/main/java/examples/GrpcClientExamples.java +++ b/vertx-grpc-docs/src/main/java/examples/GrpcClientExamples.java @@ -216,13 +216,22 @@ public void jsonWireFormat02(GrpcClient client, SocketAddress server) { }); } - public void requestCompression(GrpcClientRequest request) { - request.encoding("gzip"); + public void requestCompression(Vertx vertx) { + GrpcClientCompressionOptions compressionOptions = new GrpcClientCompressionOptions(); - // Write items after encoding has been defined - request.write(Item.newBuilder().setValue("item-1").build()); - request.write(Item.newBuilder().setValue("item-2").build()); - request.write(Item.newBuilder().setValue("item-3").build()); + compressionOptions.setCompressionEnabled(true); + compressionOptions.addCompressionAlgorithm("gzip"); + compressionOptions.addDecompressionAlgorithm("gzip"); + compressionOptions.setCompressionAlgorithm("gzip"); + + GrpcClient client = GrpcClient.client(vertx, new GrpcClientOptions().setCompression(compressionOptions)); + + SocketAddress server = SocketAddress.inetSocketAddress(443, "example.com"); + Future> fut = client.request(server, GreeterGrpcClient.SayHello); + fut.onSuccess(request -> { + // The end method calls the service + request.encoding("gzip").end(HelloRequest.newBuilder().setName("Bob").build()); + }); } public void requestWithDeadline(Vertx vertx) { diff --git a/vertx-grpc-docs/src/main/java/examples/GrpcServerExamples.java b/vertx-grpc-docs/src/main/java/examples/GrpcServerExamples.java index 7e4872b3d..2e6f2431f 100644 --- a/vertx-grpc-docs/src/main/java/examples/GrpcServerExamples.java +++ b/vertx-grpc-docs/src/main/java/examples/GrpcServerExamples.java @@ -169,13 +169,22 @@ public void anemicJson(GrpcServer server) { }); } - public void responseCompression(GrpcServerResponse response) { - response.encoding("gzip"); + public void responseCompression(Vertx vertx) { + GrpcCompressionOptions compressionOptions = new GrpcCompressionOptions(); - // Write items after encoding has been defined - response.write(Item.newBuilder().setValue("item-1").build()); - response.write(Item.newBuilder().setValue("item-2").build()); - response.write(Item.newBuilder().setValue("item-3").build()); + compressionOptions.setCompressionEnabled(true); + compressionOptions.addCompressionAlgorithm("gzip"); + compressionOptions.addDecompressionAlgorithm("gzip"); + + GrpcServer server = GrpcServer.server(vertx, new GrpcServerOptions().setCompressionOptions(compressionOptions)); + + // Create a response + server.callHandler(GreeterGrpcService.SayHello, request -> { + request.handler(hello -> { + // Handle hello message + request.response().encoding("gzip").end(HelloReply.newBuilder().setMessage("Hello " + hello.getName()).build()); + }); + }); } public void protobufLevelAPI(GrpcServer server) { @@ -341,4 +350,44 @@ public void healthServiceExample(Vertx vertx, HttpServerOptions options) { .requestHandler(grpcServer) .listen(); } + + public void customDecompressor() { + // Create a custom decompressor + GrpcDecompressor myDecompressor = new GrpcDecompressor() { + @Override + public String encoding() { + return "my-custom-encoding"; + } + + @Override + public Buffer decompress(Buffer data) throws CodecException { + // Implement your custom decompression logic here + // This is just a placeholder example + return Buffer.buffer(data.getBytes()); + } + }; + + // Then you need to add: + // provides io.vertx.grpc.common.GrpcDecompressor with . to your module-info.java + } + + public void customCompressor() { + // Create a custom compressor + GrpcCompressor myCompressor = new GrpcCompressor() { + @Override + public String encoding() { + return "my-custom-encoding"; + } + + @Override + public Buffer compress(Buffer data) throws CodecException { + // Implement your custom compression logic here + // This is just a placeholder example + return Buffer.buffer(data.getBytes()); + } + }; + + // Then you need to add: + // provides io.vertx.grpc.common.GrpcCompressor with . to your module-info.java + } } diff --git a/vertx-grpc-it/pom.xml b/vertx-grpc-it/pom.xml index 80594e53c..c65dc4ddd 100644 --- a/vertx-grpc-it/pom.xml +++ b/vertx-grpc-it/pom.xml @@ -55,6 +55,11 @@ vertx-grpc-reflection ${project.version} + + io.vertx + vertx-grpc-compression + ${project.version} + io.vertx vertx-grpc-transcoding diff --git a/vertx-grpc-it/src/test/java/io/vertx/grpc/it/CompressionTest.java b/vertx-grpc-it/src/test/java/io/vertx/grpc/it/CompressionTest.java new file mode 100644 index 000000000..fdf376356 --- /dev/null +++ b/vertx-grpc-it/src/test/java/io/vertx/grpc/it/CompressionTest.java @@ -0,0 +1,140 @@ +/* + * Copyright (c) 2011-2025 Contributors to the Eclipse Foundation + * + * This program and the accompanying materials are made available under the + * terms of the Eclipse Public License 2.0 which is available at + * http://www.eclipse.org/legal/epl-2.0, or the Apache License, Version 2.0 + * which is available at https://www.apache.org/licenses/LICENSE-2.0. + * + * SPDX-License-Identifier: EPL-2.0 OR Apache-2.0 + */ +package io.vertx.grpc.it; + +import io.grpc.examples.helloworld.*; +import io.grpc.testing.integration.*; +import io.vertx.core.Future; +import io.vertx.core.http.HttpServer; +import io.vertx.core.net.SocketAddress; +import io.vertx.ext.unit.Async; +import io.vertx.ext.unit.TestContext; +import io.vertx.grpc.client.GrpcClient; +import io.vertx.grpc.client.GrpcClientCompressionOptions; +import io.vertx.grpc.client.GrpcClientOptions; +import io.vertx.grpc.server.GrpcServer; +import io.vertx.grpc.server.GrpcServerOptions; +import io.vertx.grpc.server.Service; +import org.junit.Test; + +public class CompressionTest extends ProtocPluginTestBase { + + @Test + public void testIdentityCompression(TestContext should) throws Exception { + testCompression(should, "identity"); + } + + @Test + public void testGzipCompression(TestContext should) throws Exception { + testCompression(should, "gzip"); + } + + @Test + public void testSnappyCompression(TestContext should) throws Exception { + testCompression(should, "snappy"); + } + + @Test + public void testInvalidCompression(TestContext should) throws Exception { + // Create gRPC Server with invalid compression + GrpcServer grpcServer = GrpcServer.server(vertx, new GrpcServerOptions()); + + grpcServer.addService(greeterService(new GreeterService() { + @Override + public Future sayHello(HelloRequest request) { + return Future.succeededFuture(HelloReply.newBuilder() + .setMessage("Hello " + request.getName()) + .build()); + } + })); + + HttpServer httpServer = vertx.createHttpServer(); + httpServer.requestHandler(grpcServer).listen(8080).toCompletionStage().toCompletableFuture().get(20, java.util.concurrent.TimeUnit.SECONDS); + + // Create gRPC Client with invalid compression + GrpcClient grpcClient = GrpcClient.client(vertx, new GrpcClientOptions().setCompression(new GrpcClientCompressionOptions().setCompressionAlgorithm("invalid"))); + GreeterClient client = greeterClient(grpcClient, SocketAddress.inetSocketAddress(port, "localhost")); + + Async test = should.async(); + client.sayHello(HelloRequest.newBuilder().setName("World").build()) + .onFailure(exception -> test.complete()) + .onSuccess(reply -> should.fail("Expected failure due to invalid compression")); + + test.awaitSuccess(20_000); + + // Close the server + httpServer.close(); + } + + private void testCompression(TestContext should, String compressionType) throws Exception { + // Create gRPC Server with the specified compression + GrpcServer grpcServer = GrpcServer.server(vertx); + + grpcServer.addService(greeterService(new GreeterService() { + @Override + public Future sayHello(HelloRequest request) { + return Future.succeededFuture(HelloReply.newBuilder() + .setMessage("Hello " + request.getName()) + .build()); + } + })); + + HttpServer httpServer = vertx.createHttpServer(); + httpServer.requestHandler(grpcServer).listen(8080).toCompletionStage().toCompletableFuture().get(20, java.util.concurrent.TimeUnit.SECONDS); + + // Create gRPC Client with the specified compression + GrpcClient grpcClient = GrpcClient.client(vertx, new GrpcClientOptions().setCompression(new GrpcClientCompressionOptions().setCompressionAlgorithm(compressionType))); + GreeterClient client = greeterClient(grpcClient, SocketAddress.inetSocketAddress(port, "localhost")); + + Async test = should.async(); + client.sayHello(HelloRequest.newBuilder() + .setName("World") + .build()) + .onComplete(should.asyncAssertSuccess(reply -> { + should.assertEquals("Hello World", reply.getMessage()); + test.complete(); + })); + test.awaitSuccess(20_000); + + // Close the server + httpServer.close(); + } + + @Override + protected GrpcServer grpcServer() { + return GrpcServer.server(vertx); + } + + @Override + protected GrpcClient grpcClient() { + return GrpcClient.client(vertx); + } + + @Override + protected Service greeterService(GreeterService service) { + return GreeterGrpcService.of(service); + } + + @Override + protected GreeterClient greeterClient(GrpcClient grpcClient, SocketAddress socketAddress) { + return GreeterGrpcClient.create(grpcClient, socketAddress); + } + + @Override + protected Service testService(TestServiceService service) { + return TestServiceGrpcService.of(service); + } + + @Override + protected TestServiceClient testClient(GrpcClient client, SocketAddress socketAddress) { + return TestServiceGrpcClient.create(client, socketAddress); + } +} diff --git a/vertx-grpc-it/src/test/java/io/vertx/grpc/it/ProtocPluginTest.java b/vertx-grpc-it/src/test/java/io/vertx/grpc/it/ProtocPluginTest.java index bbad857be..bd9beef90 100644 --- a/vertx-grpc-it/src/test/java/io/vertx/grpc/it/ProtocPluginTest.java +++ b/vertx-grpc-it/src/test/java/io/vertx/grpc/it/ProtocPluginTest.java @@ -13,6 +13,7 @@ import io.grpc.examples.helloworld.*; import io.grpc.testing.integration.*; import io.vertx.core.net.SocketAddress; +import io.vertx.grpc.client.GrpcClientOptions; import io.vertx.grpc.server.GrpcServer; import io.vertx.grpc.client.GrpcClient; import io.vertx.grpc.server.Service; diff --git a/vertx-grpc-server/pom.xml b/vertx-grpc-server/pom.xml index d692d57e0..0c0ee5a2a 100644 --- a/vertx-grpc-server/pom.xml +++ b/vertx-grpc-server/pom.xml @@ -42,6 +42,12 @@ test-jar test + + io.vertx + vertx-grpc-compression + ${project.version} + test + org.bouncycastle diff --git a/vertx-grpc-server/src/main/generated/io/vertx/grpc/server/GrpcServerOptionsConverter.java b/vertx-grpc-server/src/main/generated/io/vertx/grpc/server/GrpcServerOptionsConverter.java index 1bcc0de6d..a7aa06e28 100644 --- a/vertx-grpc-server/src/main/generated/io/vertx/grpc/server/GrpcServerOptionsConverter.java +++ b/vertx-grpc-server/src/main/generated/io/vertx/grpc/server/GrpcServerOptionsConverter.java @@ -20,6 +20,11 @@ static void fromJson(Iterable> json, GrpcSer }); } break; + case "compressionOptions": + if (member.getValue() instanceof JsonObject) { + obj.setCompressionOptions(new io.vertx.grpc.common.GrpcCompressionOptions((io.vertx.core.json.JsonObject)member.getValue())); + } + break; case "scheduleDeadlineAutomatically": if (member.getValue() instanceof Boolean) { obj.setScheduleDeadlineAutomatically((Boolean)member.getValue()); @@ -49,6 +54,9 @@ static void toJson(GrpcServerOptions obj, java.util.Map json) { obj.getEnabledProtocols().forEach(item -> array.add(item.name())); json.put("enabledProtocols", array); } + if (obj.getCompressionOptions() != null) { + json.put("compressionOptions", obj.getCompressionOptions().toJson()); + } json.put("scheduleDeadlineAutomatically", obj.getScheduleDeadlineAutomatically()); json.put("deadlinePropagation", obj.getDeadlinePropagation()); json.put("maxMessageSize", obj.getMaxMessageSize()); diff --git a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/GrpcServerOptions.java b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/GrpcServerOptions.java index cc342992d..fb8095f49 100644 --- a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/GrpcServerOptions.java +++ b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/GrpcServerOptions.java @@ -14,17 +14,19 @@ import io.vertx.codegen.annotations.Unstable; import io.vertx.codegen.json.annotations.JsonGen; import io.vertx.core.json.JsonObject; +import io.vertx.grpc.common.GrpcCompressionOptions; import java.util.Collections; import java.util.EnumSet; +import java.util.Objects; import java.util.Set; /** * Configuration for a {@link GrpcServer}. */ +@Unstable @DataObject @JsonGen(publicConverter = false) -@Unstable public class GrpcServerOptions { /** @@ -32,6 +34,11 @@ public class GrpcServerOptions { */ public static final Set DEFAULT_ENABLED_PROTOCOLS = Collections.unmodifiableSet(EnumSet.allOf(GrpcProtocol.class)); + /** + * The default compression options + */ + public static final GrpcCompressionOptions DEFAULT_COMPRESSION = new GrpcCompressionOptions(); + /** * Whether the server schedule deadline automatically when a request carrying a timeout is received, by default = {@code false} */ @@ -48,6 +55,7 @@ public class GrpcServerOptions { public static final long DEFAULT_MAX_MESSAGE_SIZE = 256 * 1024; private Set enabledProtocols; + private GrpcCompressionOptions compressionOptions; private boolean scheduleDeadlineAutomatically; private boolean deadlinePropagation; private long maxMessageSize; @@ -57,6 +65,7 @@ public class GrpcServerOptions { */ public GrpcServerOptions() { enabledProtocols = EnumSet.copyOf(DEFAULT_ENABLED_PROTOCOLS); + compressionOptions = DEFAULT_COMPRESSION; scheduleDeadlineAutomatically = DEFAULT_SCHEDULE_DEADLINE_AUTOMATICALLY; deadlinePropagation = DEFAULT_PROPAGATE_DEADLINE; maxMessageSize = DEFAULT_MAX_MESSAGE_SIZE; @@ -67,6 +76,7 @@ public GrpcServerOptions() { */ public GrpcServerOptions(GrpcServerOptions other) { enabledProtocols = EnumSet.copyOf(other.enabledProtocols); + compressionOptions = new GrpcCompressionOptions(other.compressionOptions); scheduleDeadlineAutomatically = other.scheduleDeadlineAutomatically; deadlinePropagation = other.deadlinePropagation; maxMessageSize = other.maxMessageSize; @@ -121,6 +131,15 @@ public Set getEnabledProtocols() { return enabledProtocols; } + public GrpcCompressionOptions getCompressionOptions() { + return compressionOptions; + } + + public GrpcServerOptions setCompressionOptions(GrpcCompressionOptions compressionOptions) { + this.compressionOptions = Objects.requireNonNull(compressionOptions, "compressionOptions must not be null"); + return this; + } + /** * @return whether the server will automatically schedule a deadline when a request carrying a timeout is received. */ @@ -165,7 +184,6 @@ public GrpcServerOptions setDeadlinePropagation(boolean deadlinePropagation) { return this; } - /** * @return the maximum message size in bytes accepted by the server */ @@ -175,6 +193,7 @@ public long getMaxMessageSize() { /** * Set the maximum message size in bytes accepted from a client, the maximum value is {@code 0xFFFFFFFF} + * * @param maxMessageSize the size * @return a reference to this, so the API can be used fluently */ diff --git a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerImpl.java b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerImpl.java index 6e5c0ea1e..6ca246e2d 100644 --- a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerImpl.java +++ b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerImpl.java @@ -26,7 +26,6 @@ import io.vertx.grpc.server.*; import java.util.*; -import java.util.regex.Pattern; import java.util.stream.Collectors; /** @@ -34,23 +33,28 @@ */ public class GrpcServerImpl implements GrpcServer, Closeable { - private static final Pattern CONTENT_TYPE_PATTERN = Pattern.compile("application/grpc(-web(-text)?)?(\\+(json|proto))?"); - private static final Logger log = LoggerFactory.getLogger(GrpcServer.class); private final GrpcServerOptions options; + private Handler> requestHandler; + private final List invokers; private final List services = new ArrayList<>(); private final Map>> methodCallHandlers = new HashMap<>(); - private final List invokers; + private final Map compressors; + private final Map decompressors; private boolean closing; public GrpcServerImpl(Vertx vertx, GrpcServerOptions options) { - ServiceLoader loader = ServiceLoader.load(GrpcHttpInvoker.class); - this.invokers = loader.stream().map(ServiceLoader.Provider::get).collect(Collectors.toList()); + ServiceLoader invokerServiceLoader = ServiceLoader.load(GrpcHttpInvoker.class); + + this.invokers = invokerServiceLoader.stream().map(ServiceLoader.Provider::get).collect(Collectors.toList()); + this.compressors = options.getCompressionOptions().getCompressors(); + this.decompressors = options.getCompressionOptions().getDecompressors(); + this.options = new GrpcServerOptions(Objects.requireNonNull(options, "options is null")); } @@ -126,6 +130,18 @@ private int validate(GrpcServerRequestInspector.RequestInspectionDetails details return 415; } + // Check encoding + if (!decompressors.containsKey(details.encoding)) { + log.trace("Compression algorithm " + details.encoding + " is not implemented, sending error 404"); + return 404; + } + + // Check if we support at least one of the accepted encodings + if (!details.acceptEncodings.isEmpty() && details.acceptEncodings.stream().noneMatch(decompressors::containsKey)) { + log.trace("None of the accepted encodings " + details.acceptEncodings + " is implemented, sending error 404"); + return 404; + } + return -1; } @@ -145,13 +161,18 @@ private boolean handle(MethodCallHandler method, HttpServ format, httpRequest, method.messageDecoder, - methodCall); + this.decompressors, + methodCall + ); grpcResponse = new Http2GrpcServerResponse<>( context, grpcRequest, protocol, httpRequest.response(), - method.messageEncoder); + method.messageEncoder, + this.compressors, + this.decompressors + ); break; case WEB: case WEB_TEXT: @@ -165,13 +186,18 @@ private boolean handle(MethodCallHandler method, HttpServ options.getMaxMessageSize(), httpRequest, method.messageDecoder, - methodCall); + this.decompressors, + methodCall + ); grpcResponse = new WebGrpcServerResponse<>( context, grpcRequest, protocol, httpRequest.response(), - method.messageEncoder); + method.messageEncoder, + this.compressors, + this.decompressors + ); break; case TRANSCODING: grpcRequest = null; diff --git a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerRequestImpl.java b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerRequestImpl.java index d828e42ed..d6bc819f1 100644 --- a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerRequestImpl.java +++ b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerRequestImpl.java @@ -75,8 +75,9 @@ public GrpcServerRequestImpl(ContextInternal context, HttpServerRequest httpRequest, GrpcMessageDeframer messageDeframer, GrpcMessageDecoder messageDecoder, - GrpcMethodCall methodCall) { - super(context, httpRequest, httpRequest.headers().get(GrpcHeaderNames.GRPC_ENCODING), format, messageDeframer, messageDecoder); + GrpcMethodCall methodCall, + Map decompressors) { + super(context, httpRequest, httpRequest.headers().get(GrpcHeaderNames.GRPC_ENCODING), format, messageDeframer, messageDecoder, decompressors); String timeoutHeader = httpRequest.getHeader(GrpcHeaderNames.GRPC_TIMEOUT); long timeout = timeoutHeader != null ? parseTimeout(timeoutHeader) : 0L; diff --git a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerRequestInspector.java b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerRequestInspector.java index 17afd2a20..2359f5841 100644 --- a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerRequestInspector.java +++ b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerRequestInspector.java @@ -3,15 +3,21 @@ import io.vertx.core.http.HttpHeaders; import io.vertx.core.http.HttpServerRequest; import io.vertx.core.http.HttpVersion; +import io.vertx.grpc.common.GrpcHeaderNames; import io.vertx.grpc.common.WireFormat; import io.vertx.grpc.server.GrpcProtocol; +import java.util.Arrays; +import java.util.Set; import java.util.regex.Matcher; import java.util.regex.Pattern; +import java.util.stream.Collectors; final class GrpcServerRequestInspector { private static final Pattern CONTENT_TYPE_PATTERN = Pattern.compile("application/grpc(-web(-text)?)?(\\+(json|proto))?"); + private static final String DEFAULT_ENCODING = "identity"; + private static final String DEFAULT_ACCEPT_ENCODING = "identity"; private GrpcServerRequestInspector() { } @@ -22,6 +28,8 @@ public static RequestInspectionDetails inspect(HttpServerRequest request) { return null; } + determineEncoding(request, builder); + return builder.build(); } @@ -61,15 +69,46 @@ private static boolean determineContentType(String contentType, RequestInspectio return false; } + private static void determineEncoding(HttpServerRequest request, RequestInspectionDetailsBuilder builder) { + String encoding = request.getHeader(HttpHeaders.CONTENT_ENCODING); + if (request.getHeader(GrpcHeaderNames.GRPC_ENCODING) != null) { + encoding = request.getHeader(GrpcHeaderNames.GRPC_ENCODING); + } + + String acceptEncoding = request.getHeader(HttpHeaders.ACCEPT_ENCODING); + if (request.getHeader(GrpcHeaderNames.GRPC_ACCEPT_ENCODING) != null) { + acceptEncoding = request.getHeader(GrpcHeaderNames.GRPC_ACCEPT_ENCODING); + } + + if (encoding == null) { + encoding = DEFAULT_ENCODING; + } + + if (acceptEncoding == null) { + acceptEncoding = DEFAULT_ACCEPT_ENCODING; + } + + if(!acceptEncoding.contains(encoding)) { + acceptEncoding += "," + encoding; + } + + builder.encoding(encoding); + builder.acceptEncodings(Arrays.stream(acceptEncoding.split(",")).map(String::trim).collect(Collectors.toUnmodifiableSet())); + } + static final class RequestInspectionDetails { final HttpVersion version; final GrpcProtocol protocol; final WireFormat format; + final String encoding; + final Set acceptEncodings; - private RequestInspectionDetails(HttpVersion version, GrpcProtocol protocol, WireFormat format) { + RequestInspectionDetails(HttpVersion version, GrpcProtocol protocol, WireFormat format, String encoding, Set acceptEncodings) { this.version = version; this.protocol = protocol; this.format = format; + this.encoding = encoding; + this.acceptEncodings = acceptEncodings; } } @@ -77,6 +116,8 @@ static final class RequestInspectionDetailsBuilder { private HttpVersion version; private GrpcProtocol protocol; private WireFormat format; + private String encoding; + private Set acceptEncodings; RequestInspectionDetailsBuilder() { } @@ -96,8 +137,18 @@ RequestInspectionDetailsBuilder format(WireFormat format) { return this; } + RequestInspectionDetailsBuilder encoding(String encoding) { + this.encoding = encoding; + return this; + } + + RequestInspectionDetailsBuilder acceptEncodings(Set acceptEncodings) { + this.acceptEncodings = acceptEncodings; + return this; + } + RequestInspectionDetails build() { - return new RequestInspectionDetails(version, protocol, format); + return new RequestInspectionDetails(version, protocol, format, encoding, acceptEncodings); } } } diff --git a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerResponseImpl.java b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerResponseImpl.java index af65296f7..4488e4bda 100644 --- a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerResponseImpl.java +++ b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/GrpcServerResponseImpl.java @@ -10,18 +10,19 @@ */ package io.vertx.grpc.server.impl; +import com.google.common.net.UrlEscapers; import io.vertx.core.Future; import io.vertx.core.MultiMap; import io.vertx.core.buffer.Buffer; import io.vertx.core.http.HttpServerResponse; import io.vertx.core.internal.ContextInternal; +import io.vertx.grpc.common.*; import io.vertx.grpc.common.GrpcHeaderNames; import io.vertx.grpc.common.GrpcMessage; import io.vertx.grpc.common.GrpcMessageEncoder; import io.vertx.grpc.common.GrpcStatus; import io.vertx.grpc.common.impl.GrpcMessageImpl; import io.vertx.grpc.common.impl.GrpcWriteStreamBase; -import io.vertx.grpc.common.impl.Utils; import io.vertx.grpc.server.GrpcProtocol; import io.vertx.grpc.server.GrpcServerResponse; import io.vertx.grpc.server.StatusException; @@ -44,8 +45,10 @@ public GrpcServerResponseImpl(ContextInternal context, GrpcServerRequestImpl request, GrpcProtocol protocol, HttpServerResponse httpResponse, - GrpcMessageEncoder encoder) { - super(context, protocol.mediaType(), httpResponse, encoder); + GrpcMessageEncoder encoder, + Map compressors, + Map decompressors) { + super(context, protocol.mediaType(), httpResponse, encoder, compressors, decompressors); this.request = request; this.httpResponse = httpResponse; } @@ -157,7 +160,7 @@ protected void encodeGrpcStatus(MultiMap entries) { if (status != GrpcStatus.OK) { String msg = statusMessage; if (msg != null && !entries.contains(GrpcHeaderNames.GRPC_MESSAGE)) { - entries.set(GrpcHeaderNames.GRPC_MESSAGE, Utils.utf8PercentEncode(msg)); + entries.set(GrpcHeaderNames.GRPC_MESSAGE, UrlEscapers.urlFragmentEscaper().escape(msg)); } } else { entries.remove(GrpcHeaderNames.GRPC_MESSAGE); diff --git a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/Http2GrpcServerRequest.java b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/Http2GrpcServerRequest.java index dcf09e3bd..93de617db 100644 --- a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/Http2GrpcServerRequest.java +++ b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/Http2GrpcServerRequest.java @@ -12,6 +12,7 @@ import io.vertx.core.http.HttpServerRequest; import io.vertx.core.internal.ContextInternal; +import io.vertx.grpc.common.GrpcDecompressor; import io.vertx.grpc.common.GrpcHeaderNames; import io.vertx.grpc.common.GrpcMessageDecoder; import io.vertx.grpc.common.WireFormat; @@ -19,9 +20,11 @@ import io.vertx.grpc.common.impl.Http2GrpcMessageDeframer; import io.vertx.grpc.server.GrpcProtocol; +import java.util.Map; + public class Http2GrpcServerRequest extends GrpcServerRequestImpl { - public Http2GrpcServerRequest(ContextInternal context, GrpcProtocol protocol, WireFormat format, HttpServerRequest httpRequest, GrpcMessageDecoder messageDecoder, GrpcMethodCall methodCall) { - super(context, protocol, format, httpRequest, new Http2GrpcMessageDeframer(httpRequest.headers().get(GrpcHeaderNames.GRPC_ENCODING), format), messageDecoder, methodCall); + public Http2GrpcServerRequest(ContextInternal context, GrpcProtocol protocol, WireFormat format, HttpServerRequest httpRequest, GrpcMessageDecoder messageDecoder, Map decompressors, GrpcMethodCall methodCall) { + super(context, protocol, format, httpRequest, new Http2GrpcMessageDeframer(httpRequest.headers().get(GrpcHeaderNames.GRPC_ENCODING), format), messageDecoder, methodCall, decompressors); } } diff --git a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/Http2GrpcServerResponse.java b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/Http2GrpcServerResponse.java index 81d68b950..9a1f920c9 100644 --- a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/Http2GrpcServerResponse.java +++ b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/Http2GrpcServerResponse.java @@ -13,20 +13,30 @@ import io.vertx.core.MultiMap; import io.vertx.core.http.HttpServerResponse; import io.vertx.core.internal.ContextInternal; +import io.vertx.grpc.common.GrpcCompressor; +import io.vertx.grpc.common.GrpcDecompressor; import io.vertx.grpc.common.GrpcHeaderNames; import io.vertx.grpc.common.GrpcMessageEncoder; import io.vertx.grpc.server.GrpcProtocol; -public class Http2GrpcServerResponse extends GrpcServerResponseImpl { +import java.util.Map; - public Http2GrpcServerResponse(ContextInternal context, GrpcServerRequestImpl request, GrpcProtocol protocol, HttpServerResponse httpResponse, GrpcMessageEncoder encoder) { - super(context, request, protocol, httpResponse, encoder); +public class Http2GrpcServerResponse extends GrpcServerResponseImpl { + + public Http2GrpcServerResponse(ContextInternal context, + GrpcServerRequestImpl request, + GrpcProtocol protocol, + HttpServerResponse httpResponse, + GrpcMessageEncoder encoder, + Map compressors, + Map decompressors) { + super(context, request, protocol, httpResponse, encoder, compressors, decompressors); } @Override protected void encodeGrpcHeaders(MultiMap grpcHeaders, MultiMap httpHeaders) { super.encodeGrpcHeaders(grpcHeaders, httpHeaders); httpHeaders.set(GrpcHeaderNames.GRPC_ENCODING, encoding); - httpHeaders.set(GrpcHeaderNames.GRPC_ACCEPT_ENCODING, "gzip"); + httpHeaders.set(GrpcHeaderNames.GRPC_ACCEPT_ENCODING, String.join(",", decompressors.keySet())); } } diff --git a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/WebGrpcServerRequest.java b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/WebGrpcServerRequest.java index 43f8369d6..5c3c7b2ca 100644 --- a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/WebGrpcServerRequest.java +++ b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/WebGrpcServerRequest.java @@ -22,6 +22,8 @@ import io.vertx.grpc.common.impl.Http2GrpcMessageDeframer; import io.vertx.grpc.server.GrpcProtocol; +import java.util.Map; + import static io.vertx.core.http.HttpHeaders.CONTENT_TYPE; public class WebGrpcServerRequest extends GrpcServerRequestImpl { @@ -88,7 +90,7 @@ public Object next() { } } - public WebGrpcServerRequest(ContextInternal context, GrpcProtocol protocol, WireFormat format, long maxMessageSize, HttpServerRequest httpRequest, GrpcMessageDecoder messageDecoder, GrpcMethodCall methodCall) { - super(context, protocol, format, httpRequest, httpRequest.version() != HttpVersion.HTTP_2 && GrpcMediaType.isGrpcWebText(httpRequest.getHeader(CONTENT_TYPE)) ? new TextMessageDeframer() : new Http2GrpcMessageDeframer(httpRequest.headers().get(GrpcHeaderNames.GRPC_ENCODING), format), messageDecoder, methodCall); + public WebGrpcServerRequest(ContextInternal context, GrpcProtocol protocol, WireFormat format, long maxMessageSize, HttpServerRequest httpRequest, GrpcMessageDecoder messageDecoder, Map decompressors, GrpcMethodCall methodCall) { + super(context, protocol, format, httpRequest, httpRequest.version() != HttpVersion.HTTP_2 && GrpcMediaType.isGrpcWebText(httpRequest.getHeader(CONTENT_TYPE)) ? new TextMessageDeframer() : new Http2GrpcMessageDeframer(httpRequest.headers().get(GrpcHeaderNames.GRPC_ENCODING), format), messageDecoder, methodCall, decompressors); } } diff --git a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/WebGrpcServerResponse.java b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/WebGrpcServerResponse.java index e189ae5f2..e2d8f0cdb 100644 --- a/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/WebGrpcServerResponse.java +++ b/vertx-grpc-server/src/main/java/io/vertx/grpc/server/impl/WebGrpcServerResponse.java @@ -18,6 +18,8 @@ import io.vertx.core.http.HttpServerResponse; import io.vertx.core.internal.ContextInternal; import io.vertx.core.internal.buffer.BufferInternal; +import io.vertx.grpc.common.GrpcCompressor; +import io.vertx.grpc.common.GrpcDecompressor; import io.vertx.grpc.common.GrpcMessageEncoder; import io.vertx.grpc.server.GrpcProtocol; @@ -31,8 +33,14 @@ public class WebGrpcServerResponse extends GrpcServerResponseImpl request, GrpcProtocol protocol, HttpServerResponse httpResponse, GrpcMessageEncoder encoder) { - super(context, request, protocol, httpResponse, encoder); + public WebGrpcServerResponse(ContextInternal context, + GrpcServerRequestImpl request, + GrpcProtocol protocol, + HttpServerResponse httpResponse, + GrpcMessageEncoder encoder, + Map compressors, + Map decompressors) { + super(context, request, protocol, httpResponse, encoder, compressors, decompressors); this.protocol = protocol; this.httpResponse = httpResponse; diff --git a/vertx-grpc-server/src/main/java/module-info.java b/vertx-grpc-server/src/main/java/module-info.java index c5b86fa49..205ec2bc3 100644 --- a/vertx-grpc-server/src/main/java/module-info.java +++ b/vertx-grpc-server/src/main/java/module-info.java @@ -10,7 +10,10 @@ requires io.netty.codec; requires io.netty.buffer; requires com.google.protobuf; + requires com.google.common; + uses io.vertx.grpc.common.GrpcCompressor; + uses io.vertx.grpc.common.GrpcDecompressor; uses io.vertx.grpc.server.impl.GrpcHttpInvoker; exports io.vertx.grpc.server; diff --git a/vertx-grpc-server/src/test/java/io/vertx/tests/server/ServerMessageEncodingTest.java b/vertx-grpc-server/src/test/java/io/vertx/tests/server/ServerMessageEncodingTest.java index f8e055a95..2a19bf160 100644 --- a/vertx-grpc-server/src/test/java/io/vertx/tests/server/ServerMessageEncodingTest.java +++ b/vertx-grpc-server/src/test/java/io/vertx/tests/server/ServerMessageEncodingTest.java @@ -69,6 +69,21 @@ public void testIdentityRequestPassThrough(TestContext should) { testEncode(should, "identity", GrpcMessage.message("identity", Buffer.buffer("Hello World")), false); } + @Test + public void testSnappyResponseCompress(TestContext should) { + testEncode(should, "snappy", GrpcMessage.message("identity", Buffer.buffer("Hello World")), true); + } + + @Test + public void testSnappyResponsePassThrough(TestContext should) { + testEncode(should, "snappy", GrpcMessage.message("snappy", GrpcTestBase.snappyCompress(Buffer.buffer("Hello World"))), true); + } + + @Test + public void testIdentityResponseUnsnappy(TestContext should) { + testEncode(should, "identity", GrpcMessage.message("snappy", GrpcTestBase.snappyCompress(Buffer.buffer("Hello World"))), false); + } + private void testEncode(TestContext should, String encoding, GrpcMessage msg, boolean compressed) { Buffer expected = Buffer.buffer("Hello World"); @@ -104,13 +119,19 @@ private void testEncode(TestContext should, String encoding, GrpcMessage msg, bo int len = body.getInt(1); Buffer received = body.slice(5, 5 + len); if (compressed) { - received = GrpcTestBase.unzip(received); + if (encoding.equals("gzip")) { + received = GrpcTestBase.unzip(received); + } else if (encoding.equals("snappy")) { + received = GrpcTestBase.snappyDecompress(received); + } } should.assertEquals(expected, received); done.complete(); })); })); })); + + done.awaitSuccess(5000); } @Test diff --git a/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingGrpcServerRequest.java b/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingGrpcServerRequest.java index 7c49cee72..282db84d1 100644 --- a/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingGrpcServerRequest.java +++ b/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingGrpcServerRequest.java @@ -19,11 +19,13 @@ import io.vertx.grpc.common.GrpcMessageDecoder; import io.vertx.grpc.common.WireFormat; import io.vertx.grpc.common.impl.GrpcMethodCall; +import io.vertx.grpc.common.impl.IdentityCompressor; import io.vertx.grpc.server.GrpcProtocol; import io.vertx.grpc.server.impl.GrpcServerRequestImpl; import io.vertx.grpc.transcoding.impl.config.HttpVariableBinding; import java.util.List; +import java.util.Map; public class TranscodingGrpcServerRequest extends GrpcServerRequestImpl { @@ -50,7 +52,7 @@ public Req decode(GrpcMessage msg) throws CodecException { public boolean accepts(WireFormat format) { return messageDecoder.accepts(format); } - }, methodCall); + }, methodCall, Map.of("identity", new IdentityCompressor())); this.transcodingRequestBody = transcodingRequestBody; this.bindings = bindings; diff --git a/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingGrpcServerResponse.java b/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingGrpcServerResponse.java index e402c5aa7..c26483439 100644 --- a/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingGrpcServerResponse.java +++ b/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingGrpcServerResponse.java @@ -21,10 +21,13 @@ import io.vertx.grpc.common.GrpcMessageEncoder; import io.vertx.grpc.common.GrpcStatus; import io.vertx.grpc.common.WireFormat; +import io.vertx.grpc.common.impl.IdentityCompressor; import io.vertx.grpc.server.GrpcProtocol; import io.vertx.grpc.server.impl.GrpcServerRequestImpl; import io.vertx.grpc.server.impl.GrpcServerResponseImpl; +import java.util.Map; + public class TranscodingGrpcServerResponse extends GrpcServerResponseImpl { private final TranscodingGrpcServerRequest request; @@ -33,7 +36,7 @@ public class TranscodingGrpcServerResponse extends GrpcServerResponse private Promise head; public TranscodingGrpcServerResponse(ContextInternal context, GrpcServerRequestImpl request, GrpcProtocol protocol, HttpServerResponse httpResponse, String transcodingResponseBody, GrpcMessageEncoder encoder) { - super(context, request, protocol, httpResponse, encoder); + super(context, request, protocol, httpResponse, encoder, Map.of("identity", new IdentityCompressor()), Map.of("identity", new IdentityCompressor())); this.request = (TranscodingGrpcServerRequest) request; this.httpResponse = httpResponse; diff --git a/vertx-grpcio-client/pom.xml b/vertx-grpcio-client/pom.xml index 24dae6c84..330f65f41 100644 --- a/vertx-grpcio-client/pom.xml +++ b/vertx-grpcio-client/pom.xml @@ -55,6 +55,12 @@ test-jar test + + io.vertx + vertx-grpc-compression + ${project.version} + test + io.vertx vertx-grpc-client diff --git a/vertx-grpcio-server/pom.xml b/vertx-grpcio-server/pom.xml index c95f4abf0..04b77256a 100644 --- a/vertx-grpcio-server/pom.xml +++ b/vertx-grpcio-server/pom.xml @@ -50,6 +50,12 @@ ${project.version} test + + io.vertx + vertx-grpc-compression + ${project.version} + test + io.vertx vertx-grpc-common