Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -198,6 +198,7 @@

<modules>
<module>vertx-grpc-common</module>
<module>vertx-grpc-compression</module>
<module>vertx-grpc-transcoding</module>
<module>vertx-grpc-server</module>
<module>vertx-grpc-reflection</module>
Expand Down Expand Up @@ -274,4 +275,4 @@
</plugins>
</build>

</project>
</project>
6 changes: 6 additions & 0 deletions vertx-grpc-client/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,12 @@
<type>test-jar</type>
<scope>test</scope>
</dependency>
<dependency>
<groupId>io.vertx</groupId>
<artifactId>vertx-grpc-compression</artifactId>
<version>${project.version}</version>
<scope>test</scope>
</dependency>
<dependency>
<groupId>io.grpc</groupId>
<artifactId>grpc-protobuf</artifactId>
Expand Down
Original file line number Diff line number Diff line change
@@ -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();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -11,14 +11,19 @@
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;

/**
* Options configuring a gRPC client.
*/
@Unstable
@DataObject
@JsonGen(publicConverter = false)
public class GrpcClientOptions {

/**
Expand All @@ -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;
}

/**
Expand All @@ -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);
}

/**
Expand Down Expand Up @@ -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();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -14,21 +14,17 @@
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;
import io.vertx.core.net.Address;
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;

/**
Expand All @@ -44,6 +40,10 @@ public class GrpcClientImpl implements GrpcClient {
private final int timeout;
private final TimeUnit timeoutUnit;

private final String compressionAlgorithm;
private final Map<String, GrpcCompressor> compressors;
private final Map<String, GrpcDecompressor> decompressors;

public GrpcClientImpl(Vertx vertx, HttpClient client) {
this(vertx, new GrpcClientOptions(), client, false);
}
Expand All @@ -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() {
Expand All @@ -70,8 +74,12 @@ public Future<GrpcClientRequest<Buffer, Buffer>> 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;
});
Expand Down Expand Up @@ -123,8 +131,12 @@ private <Req, Resp> Future<GrpcClientRequest<Req, Resp>> 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);
Expand All @@ -137,7 +149,7 @@ public Future<Void> close() {
if (closeClient) {
return client.close();
} else {
return ((VertxInternal)vertx).getOrCreateContext().succeededFuture();
return ((VertxInternal) vertx).getOrCreateContext().succeededFuture();
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -51,8 +51,10 @@ public GrpcClientRequestImpl(HttpClientRequest httpRequest,
long maxMessageSize,
boolean scheduleDeadline,
GrpcMessageEncoder<Req> messageEncoder,
GrpcMessageDecoder<Resp> messageDecoder) {
super( ((PromiseInternal<?>)httpRequest.response()).context(), "application/grpc", httpRequest, messageEncoder);
GrpcMessageDecoder<Resp> messageDecoder,
Map<String, GrpcCompressor> compressors,
Map<String, GrpcDecompressor> decompressors) {
super( ((PromiseInternal<?>)httpRequest.response()).context(), "application/grpc", httpRequest, messageEncoder, compressors, decompressors);
this.httpRequest = httpRequest;
this.scheduleDeadline = scheduleDeadline;
this.timeout = 0L;
Expand Down Expand Up @@ -83,7 +85,8 @@ public GrpcClientRequestImpl(HttpClientRequest httpRequest,
format,
status,
httpResponse,
messageDecoder);
messageDecoder,
decompressors);
grpcResponse.init(this, maxMessageSize);
grpcResponse.invalidMessageHandler(invalidMsg -> {
cancel();
Expand Down Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
import io.vertx.grpc.common.impl.Http2GrpcMessageDeframer;

import java.nio.charset.StandardCharsets;
import java.util.Map;

/**
* @author <a href="mailto:julien@julienviet.com">Julien Viet</a>
Expand All @@ -41,14 +42,18 @@ public GrpcClientResponseImpl(ContextInternal context,
GrpcClientRequestImpl<Req, Resp> request,
WireFormat format,
GrpcStatus status,
HttpClientResponse httpResponse, GrpcMessageDecoder<Resp> messageDecoder) {
HttpClientResponse httpResponse,
GrpcMessageDecoder<Resp> messageDecoder,
Map<String, GrpcDecompressor> 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;
Expand Down
12 changes: 9 additions & 3 deletions vertx-grpc-client/src/main/java/module-info.java
Original file line number Diff line number Diff line change
@@ -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;
}
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,8 @@ private void testEncode(TestContext should, String requestEncoding, GrpcMessage
}));
callRequest.endMessage(msg);
}));

test.await(5000L);
}

@Test
Expand Down
Loading