diff --git a/vertx-grpc-docs/src/main/asciidoc/transcoding.adoc b/vertx-grpc-docs/src/main/asciidoc/transcoding.adoc index 1ef302817..2f335fec9 100644 --- a/vertx-grpc-docs/src/main/asciidoc/transcoding.adoc +++ b/vertx-grpc-docs/src/main/asciidoc/transcoding.adoc @@ -212,6 +212,43 @@ message HelloResponse { } ---- +=== Server-streaming responses + +Server-streaming RPCs (`rpc Foo (Req) returns (stream Resp)`) are transcoded by emitting the response over chunked transfer encoding. The wire format is selected according to the HTTP `accept` header, letting the client choose the most appropriate encoding: + +|=== +| `Accept` value | Response `Content-Type` | Framing + +| `application/json` (default) | `application/json` | Top-level JSON array: `[m1, m2, m3]`. Envoy-compatible. The body is not valid JSON until the final `]` arrives, so `JSON.parse(body)` only works after the stream completes. +| `application/x-ndjson` | `application/x-ndjson` | One JSON object per line, `\n`-separated. Each line is independently parseable as it arrives. +| `text/event-stream` | `text/event-stream` | Server-Sent Events: `data: \n\n` per message. Consumable by the browser `EventSource` API. +|=== + +[source] +---- +# JSON array (default) +curl -X POST -H "Content-Type: application/json" -d '...' http://localhost:8080/stream +# [{"payload":"first"},{"payload":"second"}] + +# NDJSON +curl -X POST -H "Accept: application/x-ndjson" -H "Content-Type: application/json" -d '...' http://localhost:8080/stream +# {"payload":"first"} +# {"payload":"second"} + +# Server-Sent Events +curl -X POST -H "Accept: text/event-stream" -H "Content-Type: application/json" -d '...' http://localhost:8080/stream +# data: {"payload":"first"} +# +# data: {"payload":"second"} +---- + +A gRPC error is reported to the client depending on when it happens: + +- a gRPC trailers-only response sends the corresponding HTTP status error without content +- a gRPC trailers response terminates the response with `]`, such error cannot be reported to the client + +Transcoding does not support client-streaming and bidirectional-streaming RPCs. + === Transcoding error handling If an error occurs during transcoding, the server will return an HTTP error response with the appropriate status code. diff --git a/vertx-grpc-it/src/test/java/io/vertx/grpc/it/tests/TranscodingTest.java b/vertx-grpc-it/src/test/java/io/vertx/grpc/it/tests/TranscodingTest.java index 0530aa609..c2de979fe 100644 --- a/vertx-grpc-it/src/test/java/io/vertx/grpc/it/tests/TranscodingTest.java +++ b/vertx-grpc-it/src/test/java/io/vertx/grpc/it/tests/TranscodingTest.java @@ -1,23 +1,39 @@ package io.vertx.grpc.it.tests; import io.grpc.examples.helloworld.*; +import io.grpc.examples.streamingtranscoding.StreamingHelloReply; +import io.grpc.examples.streamingtranscoding.StreamingHelloRequest; +import io.grpc.examples.streamingtranscoding.StreamingTranscodingGreeterClient; +import io.grpc.examples.streamingtranscoding.StreamingTranscodingGreeterGrpcClient; +import io.grpc.examples.streamingtranscoding.StreamingTranscodingGreeterGrpcService; +import io.grpc.examples.streamingtranscoding.StreamingTranscodingGreeterService; import io.grpc.stub.StreamObserver; +import io.vertx.core.Future; +import io.vertx.core.Promise; import io.vertx.core.buffer.Buffer; import io.vertx.core.http.*; import io.vertx.core.json.Json; +import io.vertx.core.json.JsonArray; import io.vertx.core.json.JsonObject; import io.vertx.core.net.SocketAddress; +import io.vertx.core.streams.WriteStream; import io.vertx.grpc.client.GrpcClient; import io.vertx.grpc.server.GrpcServer; +import io.vertx.grpc.server.GrpcServerResponse; import io.vertx.grpcio.server.GrpcIoServer; import org.junit.Test; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; import java.util.Map; import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import java.util.function.Consumer; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; public class TranscodingTest extends ProxyTestBase { @@ -344,6 +360,217 @@ public void testUnaryCollisionWithoutOption() throws TimeoutException { assertEquals("Hello Julien", reply.getMessage()); } + @Test + public void testUnaryAddService() throws TimeoutException { + HttpClient client = vertx.createHttpClient(); + + vertx.createHttpServer() + .requestHandler(GrpcServer.server(vertx).addService(GreeterGrpcService.of(new GreeterService() { + @Override + public Future sayHello(HelloRequest request) { + return Future.succeededFuture(HelloReply.newBuilder().setMessage("Hello " + request.getName()).build()); + } + }))).listen(8080, "localhost").await(10, TimeUnit.SECONDS); + + RequestOptions options = new RequestOptions().setHost("localhost").setPort(8080).setURI("/v1/hello/Julien").setMethod(HttpMethod.GET); + + Buffer body = client.request(options).compose(req -> { + req.putHeader(HttpHeaders.CONTENT_TYPE, "application/json"); + req.putHeader(HttpHeaders.ACCEPT, "application/json"); + return req.send(); + }).expecting(HttpResponseExpectation.SC_OK) + .expecting(HttpResponseExpectation.JSON) + .compose(HttpClientResponse::body) + .await(10, TimeUnit.SECONDS); + assertEquals("Hello Julien", getMessage(body.toString())); + } + + @Test + public void testServerStreamingAddService() throws TimeoutException { + HttpClient client = vertx.createHttpClient(); + + vertx.createHttpServer() + .requestHandler(GrpcServer.server(vertx).addService(StreamingTranscodingGreeterGrpcService.of(new StreamingTranscodingGreeterService() { + @Override + protected void sayHelloStreaming(StreamingHelloRequest request, WriteStream response) { + response.write(StreamingHelloReply.newBuilder().setMessage("Hello " + request.getName() + " 1").build()); + response.write(StreamingHelloReply.newBuilder().setMessage("Hello " + request.getName() + " 2").build()); + response.end(); + } + }))).listen(8080, "localhost").await(10, TimeUnit.SECONDS); + + RequestOptions options = new RequestOptions().setHost("localhost").setPort(8080).setURI("/v1/hello/stream/Julien").setMethod(HttpMethod.GET); + + Buffer body = client.request(options).compose(req -> { + req.putHeader(HttpHeaders.CONTENT_TYPE, "application/json"); + req.putHeader(HttpHeaders.ACCEPT, "application/json"); + return req.send(); + }).expecting(HttpResponseExpectation.SC_OK) + .compose(HttpClientResponse::body) + .await(10, TimeUnit.SECONDS); + + JsonArray array = new JsonArray(body); + assertEquals(2, array.size()); + assertEquals("Hello Julien 1", array.getJsonObject(0).getString("message")); + assertEquals("Hello Julien 2", array.getJsonObject(1).getString("message")); + } + + @Test + public void testServerStreaming() throws TimeoutException { + HttpClient client = vertx.createHttpClient(); + + vertx.createHttpServer() + .requestHandler(GrpcServer.server(vertx).callHandler(StreamingTranscodingGreeterGrpcService.SayHelloStreaming, call -> call.handler(request -> { + GrpcServerResponse response = call.response(); + response.write(StreamingHelloReply.newBuilder().setMessage("Hello " + request.getName() + " 1").build()); + response.write(StreamingHelloReply.newBuilder().setMessage("Hello " + request.getName() + " 2").build()); + response.write(StreamingHelloReply.newBuilder().setMessage("Hello " + request.getName() + " 3").build()); + response.end(); + }))).listen(8080, "localhost").await(10, TimeUnit.SECONDS); + + RequestOptions options = new RequestOptions().setHost("localhost").setPort(8080).setURI("/v1/hello/stream/Julien").setMethod(HttpMethod.GET); + + HttpClientResponse response = client.request(options).compose(req -> { + req.putHeader(HttpHeaders.CONTENT_TYPE, "application/json"); + req.putHeader(HttpHeaders.ACCEPT, "application/json"); + return req.send(); + }).expecting(HttpResponseExpectation.SC_OK) + .compose(resp -> resp.body().map(resp)) + .await(10, TimeUnit.SECONDS); + + assertTrue(response.headers().contains(HttpHeaders.CONTENT_TYPE, "application/json", true)); + // Streaming responses are chunked, so the length is not known up-front. + assertFalse(response.headers().contains(HttpHeaders.CONTENT_LENGTH)); + JsonArray array = new JsonArray(response.body().result()); + assertEquals(3, array.size()); + assertEquals("Hello Julien 1", array.getJsonObject(0).getString("message")); + assertEquals("Hello Julien 2", array.getJsonObject(1).getString("message")); + assertEquals("Hello Julien 3", array.getJsonObject(2).getString("message")); + } + + @Test + public void testServerStreamingNdjson() throws TimeoutException { + HttpClient client = vertx.createHttpClient(); + + vertx.createHttpServer() + .requestHandler(GrpcServer.server(vertx).callHandler(StreamingTranscodingGreeterGrpcService.SayHelloStreaming, call -> call.handler(request -> { + GrpcServerResponse response = call.response(); + response.write(StreamingHelloReply.newBuilder().setMessage("Hello " + request.getName() + " 1").build()); + response.write(StreamingHelloReply.newBuilder().setMessage("Hello " + request.getName() + " 2").build()); + response.end(); + }))).listen(8080, "localhost").await(10, TimeUnit.SECONDS); + + RequestOptions options = new RequestOptions().setHost("localhost").setPort(8080).setURI("/v1/hello/stream/Julien").setMethod(HttpMethod.GET); + + HttpClientResponse response = client.request(options).compose(req -> { + req.putHeader(HttpHeaders.CONTENT_TYPE, "application/json"); + req.putHeader(HttpHeaders.ACCEPT, "application/x-ndjson"); + return req.send(); + }).expecting(HttpResponseExpectation.SC_OK) + .compose(resp -> resp.body().map(resp)) + .await(10, TimeUnit.SECONDS); + + assertTrue(response.headers().contains(HttpHeaders.CONTENT_TYPE, "application/x-ndjson", true)); + String[] lines = response.body().result().toString().split("\n"); + assertEquals(2, lines.length); + assertEquals("Hello Julien 1", new JsonObject(lines[0]).getString("message")); + assertEquals("Hello Julien 2", new JsonObject(lines[1]).getString("message")); + } + + @Test + public void testServerStreamingSse() throws TimeoutException { + HttpClient client = vertx.createHttpClient(); + + vertx.createHttpServer() + .requestHandler(GrpcServer.server(vertx).callHandler(StreamingTranscodingGreeterGrpcService.SayHelloStreaming, call -> call.handler(request -> { + GrpcServerResponse response = call.response(); + response.write(StreamingHelloReply.newBuilder().setMessage("Hello " + request.getName() + " 1").build()); + response.write(StreamingHelloReply.newBuilder().setMessage("Hello " + request.getName() + " 2").build()); + response.end(); + }))).listen(8080, "localhost").await(10, TimeUnit.SECONDS); + + RequestOptions options = new RequestOptions().setHost("localhost").setPort(8080).setURI("/v1/hello/stream/Julien").setMethod(HttpMethod.GET); + + HttpClientResponse response = client.request(options).compose(req -> { + req.putHeader(HttpHeaders.CONTENT_TYPE, "application/json"); + req.putHeader(HttpHeaders.ACCEPT, "text/event-stream"); + return req.send(); + }).expecting(HttpResponseExpectation.SC_OK) + .compose(resp -> resp.body().map(resp)) + .await(10, TimeUnit.SECONDS); + + assertTrue(response.headers().contains(HttpHeaders.CONTENT_TYPE, "text/event-stream", true)); + String[] events = response.body().result().toString().split("\n\n"); + assertEquals(2, events.length); + assertTrue(events[0].startsWith("data: ")); + assertTrue(events[1].startsWith("data: ")); + assertEquals("Hello Julien 1", new JsonObject(events[0].substring("data: ".length())).getString("message")); + assertEquals("Hello Julien 2", new JsonObject(events[1].substring("data: ".length())).getString("message")); + } + + @Test + public void testServerStreamingEmpty() throws TimeoutException { + HttpClient client = vertx.createHttpClient(); + + vertx.createHttpServer() + .requestHandler(GrpcServer.server(vertx).callHandler(StreamingTranscodingGreeterGrpcService.SayHelloStreaming, call -> call.handler(request -> { + GrpcServerResponse response = call.response(); + response.end(); + }))).listen(8080, "localhost").await(10, TimeUnit.SECONDS); + + RequestOptions options = new RequestOptions().setHost("localhost").setPort(8080).setURI("/v1/hello/stream/Julien").setMethod(HttpMethod.GET); + + Buffer body = client.request(options).compose(req -> { + req.putHeader(HttpHeaders.CONTENT_TYPE, "application/json"); + req.putHeader(HttpHeaders.ACCEPT, "application/json"); + return req.send(); + }).expecting(HttpResponseExpectation.SC_OK) + .compose(HttpClientResponse::body) + .await(10, TimeUnit.SECONDS); + + assertEquals(0, new JsonArray(body).size()); + } + + @Test + public void testServerStreamingGrpcCollision() throws TimeoutException { + HttpClient httpClient = vertx.createHttpClient(); + + vertx.createHttpServer() + .requestHandler(GrpcServer.server(vertx).callHandler(StreamingTranscodingGreeterGrpcService.SayHelloStreaming, call -> call.handler(request -> { + GrpcServerResponse response = call.response(); + response.write(StreamingHelloReply.newBuilder().setMessage("Hello " + request.getName()).build()); + response.end(); + }))).listen(8080, "localhost").await(10, TimeUnit.SECONDS); + + RequestOptions options = new RequestOptions().setHost("localhost").setPort(8080).setURI("/v1/hello/stream/Julien").setMethod(HttpMethod.GET); + + Buffer httpBody = httpClient.request(options).compose(req -> { + req.putHeader(HttpHeaders.CONTENT_TYPE, "application/json"); + req.putHeader(HttpHeaders.ACCEPT, "application/json"); + return req.send(); + }).expecting(HttpResponseExpectation.SC_OK) + .compose(HttpClientResponse::body) + .await(10, TimeUnit.SECONDS); + assertEquals("Hello Julien", new JsonArray(httpBody).getJsonObject(0).getString("message")); + + // The same service method is still reachable over plain gRPC. + GrpcClient grpcClient = GrpcClient.client(vertx); + StreamingTranscodingGreeterClient greeterClient = StreamingTranscodingGreeterGrpcClient.create(grpcClient, SocketAddress.inetSocketAddress(8080, "localhost")); + + List received = greeterClient + .sayHelloStreaming(StreamingHelloRequest.newBuilder().setName("Julien").build()) + .compose(stream -> { + Promise> promise = Promise.promise(); + List replies = new ArrayList<>(); + stream.handler(reply -> replies.add(reply.getMessage())); + stream.endHandler(v -> promise.tryComplete(replies)); + stream.exceptionHandler(promise::tryFail); + return promise.future(); + }) + .await(10, TimeUnit.SECONDS); + assertEquals(Collections.singletonList("Hello Julien"), received); + } + @Test public void testUnknownService() throws TimeoutException { HttpClient client = vertx.createHttpClient(); 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 60b9c34de..eb9bde649 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 @@ -254,7 +254,14 @@ public GrpcServer addService(Service service) { } for (ServiceMethod method : service.methods()) { Handler handler = service.handler(method); - registerMethodCallHandler(service.pathOfMethod(method.methodName()), new ServiceMethodCallHandler(method, handler)); + ServiceMethodCallHandler smch = new ServiceMethodCallHandler<>(method, handler); + if (method instanceof MountPoint) { + MountPoint mountPoint = (MountPoint) method; + for (String path : mountPoint.paths()) { + registerMethodCallHandler(path, smch); + } + } + registerMethodCallHandler(service.pathOfMethod(method.methodName()), smch); } this.services.add(service); diff --git a/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/MessageWeaver.java b/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/MessageWeaver.java index 95c2a19e1..dd3952cc4 100644 --- a/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/MessageWeaver.java +++ b/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/MessageWeaver.java @@ -35,21 +35,22 @@ private MessageWeaver() { } /** - * Weaves HTTP variable bindings and request body into a gRPC message. + * Weaves HTTP variable bindings and request body into a gRPC message. The return type is always a {@link JsonObject} + * because a protobuf {@code Message} is always object-shaped in JSON form. * * @param message The original message buffer * @param bindings The HTTP variable bindings * @param transcodingRequestBody The transcoding request body path * @param descriptor The protobuf message descriptor, used to identify repeated fields - * @return The modified buffer with weaved content + * @return The woven message as a JsonObject * @throws DecodeException If JSON decoding fails */ - public static Buffer weaveRequestMessage(Buffer message, List bindings, String transcodingRequestBody, Descriptors.Descriptor descriptor) throws DecodeException { + public static JsonObject weaveRequestMessage(Buffer message, List bindings, String transcodingRequestBody, Descriptors.Descriptor descriptor) throws DecodeException { boolean hasBindings = bindings != null && !bindings.isEmpty(); boolean hasBody = transcodingRequestBody != null && !transcodingRequestBody.isEmpty(); if (!hasBindings && !hasBody) { - return new JsonObject().toBuffer(); + return new JsonObject(); } JsonObject result = new JsonObject(); @@ -72,7 +73,7 @@ public static Buffer weaveRequestMessage(Buffer message, List head; private final ContextInternal context; private final HttpServerResponse httpResponse; private final String transcodingResponseBody; + private final boolean streaming; + private final StreamFormat streamFormat; + private boolean firstMessageWritten; - public TranscodingGrpcOutboundStream(ContextInternal context, HttpServerRequest httpRequest, - String transcodingResponseBody, GrpcMessageDeframer deframer) { + public TranscodingGrpcOutboundStream( + ContextInternal context, + HttpServerRequest httpRequest, + String transcodingResponseBody, + GrpcMessageDeframer deframer, + boolean streaming + ) { super(httpRequest, GrpcProtocol.TRANSCODING, deframer); this.context = context; this.httpResponse = httpRequest.response(); this.transcodingResponseBody = transcodingResponseBody; + this.streaming = streaming; + this.streamFormat = streaming ? negotiateStreamFormat(httpRequest.getHeader(HttpHeaders.ACCEPT)) : StreamFormat.JSON_ARRAY; + } + + private static StreamFormat negotiateStreamFormat(String acceptHeader) { + if (acceptHeader == null) { + return StreamFormat.JSON_ARRAY; + } + String accept = acceptHeader.toLowerCase(); + if (accept.contains("text/event-stream")) { + return StreamFormat.SSE; + } + if (accept.contains("application/x-ndjson") || accept.contains("application/jsonl")) { + return StreamFormat.NDJSON; + } + return StreamFormat.JSON_ARRAY; } @Override protected String contentType(WireFormat wireFormat) { if (wireFormat instanceof JsonWireFormat) { - return protocol.mediaType(); + return streaming ? streamFormat.mediaType : protocol.mediaType(); } throw new UnsupportedOperationException(); } @@ -47,15 +92,18 @@ protected void encodeGrpcHeaders(MultiMap grpcHeaders, MultiMap httpHeaders, Str } @Override - public Future writeEnd() { - if (status != GrpcStatus.OK) { - httpResponse.setStatusCode(GrpcTranscodingError.fromHttp2Code(status.code).getHttpStatusCode()); + protected Future writeHeaders(GrpcHeadersFrame frame) { + if (streaming) { + httpResponse.setChunked(true); } - return super.writeEnd(); + return super.writeHeaders(frame); } @Override public Future writeHead() { + if (streaming) { + return httpResponse.writeHead(); + } head = context.promise(); return head.future(); } @@ -68,9 +116,16 @@ public Future writeMessage(GrpcMessageFrame frame) { } catch (CodecException e) { return context.failedFuture(e); } + if (streaming) { + return writeStreamingMessage(payload); + } + return writeUnaryMessage(payload); + } + + private Future writeUnaryMessage(Buffer payload) { Future res; try { - BufferInternal transcoded = (BufferInternal) MessageWeaver.weaveResponseMessage(payload, transcodingResponseBody); + Buffer transcoded = Json.encodeToBuffer(MessageWeaver.weaveResponseMessage(payload, transcodingResponseBody)); httpResponse.putHeader(HttpHeaders.CONTENT_LENGTH, Integer.toString(transcoded.length())); httpResponse.putHeader(HttpHeaders.CONTENT_TYPE, GrpcProtocol.TRANSCODING.mediaType()); res = httpResponse.write(transcoded); @@ -78,8 +133,67 @@ public Future writeMessage(GrpcMessageFrame frame) { httpResponse.setStatusCode(500).end(); res = context.failedFuture(e); } - if (head != null) { - res.onComplete(head); + Promise h = head; + if (h != null) { + head = null; + res.onComplete(h); + } + return res; + } + + private Future writeStreamingMessage(Buffer payload) { + Buffer transcoded; + try { + transcoded = Json.encodeToBuffer(MessageWeaver.weaveResponseMessage(payload, transcodingResponseBody)); + } catch (Exception e) { + httpResponse.setStatusCode(500).end(); + return context.failedFuture(e); + } + Buffer chunk; + switch (streamFormat) { + case JSON_ARRAY: { + Buffer prefix = firstMessageWritten ? COMMA : OPEN_ARRAY; + chunk = Buffer.buffer(prefix.length() + transcoded.length()).appendBuffer(prefix).appendBuffer(transcoded); + break; + } + case NDJSON: { + chunk = Buffer.buffer(transcoded.length() + NEWLINE.length()).appendBuffer(transcoded).appendBuffer(NEWLINE); + break; + } + case SSE: { + chunk = Buffer.buffer(SSE_PREFIX.length() + transcoded.length() + SSE_SUFFIX.length()) + .appendBuffer(SSE_PREFIX).appendBuffer(transcoded).appendBuffer(SSE_SUFFIX); + break; + } + default: + throw new AssertionError(streamFormat); + } + firstMessageWritten = true; + return httpResponse.write(chunk); + } + + @Override + public Future writeEnd() { + if (streaming) { + if (!firstMessageWritten && status != null && status != GrpcStatus.OK) { + httpResponse.setStatusCode(GrpcTranscodingError.fromHttp2Code(status.code).getHttpStatusCode()); + return httpResponse.end(); + } + if (streamFormat == StreamFormat.JSON_ARRAY) { + return httpResponse.end(firstMessageWritten ? CLOSE_ARRAY : EMPTY_ARRAY); + } + return httpResponse.end(); + } + if (status != GrpcStatus.OK) { + httpResponse.setStatusCode(GrpcTranscodingError.fromHttp2Code(status.code).getHttpStatusCode()); + } + Future res = super.writeEnd(); + // A unary call can end without ever writing a message, e.g. a failed call: the head promise + // is only resolved by writeUnaryMessage, so resolve it here instead of leaving it pending. + Promise h = head; + if (h != null) { + head = null; + res.onComplete(h); } return res; } diff --git a/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingServiceMethodImpl.java b/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingServiceMethodImpl.java index d183faa05..275ed8637 100644 --- a/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingServiceMethodImpl.java +++ b/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingServiceMethodImpl.java @@ -5,6 +5,7 @@ import io.vertx.core.http.HttpServerRequest; import io.vertx.core.internal.http.HttpServerRequestInternal; import io.vertx.core.json.DecodeException; +import io.vertx.core.json.JsonObject; import io.vertx.grpc.common.*; import io.vertx.grpc.server.GrpcProtocol; import io.vertx.grpc.server.impl.GrpcInvocation; @@ -95,6 +96,8 @@ public GrpcInvocation accept(HttpServerRequest httpRequest, WireFormat format) { return null; } + boolean streaming = cardinality == MethodCardinality.SERVER_STREAMING || cardinality == MethodCardinality.BIDI_STREAMING; + PathMatcherLookupResult res = pathMatcher == null ? null : pathMatcher.lookup(httpRequest.method().name(), httpRequest.path(), httpRequest.query()); if (res != null) { List bindings = new ArrayList<>(res.getVariableBindings()); @@ -102,21 +105,21 @@ public GrpcInvocation accept(HttpServerRequest httpRequest, WireFormat format) { TranscodingMessageDeframer deframer = new TranscodingMessageDeframer(format) { @Override protected Buffer decode(Buffer buffer) throws InvalidMessageException { - Buffer transcoded; + JsonObject transcoded; try { transcoded = MessageWeaver.weaveRequestMessage(buffer, bindings, res.getBodyFieldPath(), decoder.messageDescriptor()); } catch (DecodeException e) { throw new TranscodingInvalidMessageException(e); } - return transcoded; + return transcoded.toBuffer(); } }; - HttpGrpcOutboundStream protocolHandler = new TranscodingGrpcOutboundStream(context, httpRequest, options.getResponseBody(), deframer); + HttpGrpcOutboundStream protocolHandler = new TranscodingGrpcOutboundStream(context, httpRequest, options.getResponseBody(), deframer, streaming); return new GrpcInvocation(deframer, protocolHandler); } else if (options == null) { io.vertx.core.internal.ContextInternal context = ((HttpServerRequestInternal) httpRequest).context(); TranscodingMessageDeframer deframer = new TranscodingMessageDeframer(format); - HttpGrpcOutboundStream protocolHandler = new TranscodingGrpcOutboundStream(context, httpRequest, null, deframer); + HttpGrpcOutboundStream protocolHandler = new TranscodingGrpcOutboundStream(context, httpRequest, null, deframer, streaming); return new GrpcInvocation(deframer, protocolHandler); } diff --git a/vertx-grpc-transcoding/src/test/java/io/vertx/grpc/transcoding/benchmarks/MessageWeaverBenchmark.java b/vertx-grpc-transcoding/src/test/java/io/vertx/grpc/transcoding/benchmarks/MessageWeaverBenchmark.java index 06fc6c6eb..db1e2b16c 100644 --- a/vertx-grpc-transcoding/src/test/java/io/vertx/grpc/transcoding/benchmarks/MessageWeaverBenchmark.java +++ b/vertx-grpc-transcoding/src/test/java/io/vertx/grpc/transcoding/benchmarks/MessageWeaverBenchmark.java @@ -96,7 +96,7 @@ private HttpVariableBinding createBinding(String path, String value) { @Benchmark public void benchmarkWeaveRequestSimpleMessageNoBindings(Blackhole blackhole) { - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( simpleMessage, null, null, @@ -107,7 +107,7 @@ public void benchmarkWeaveRequestSimpleMessageNoBindings(Blackhole blackhole) { @Benchmark public void benchmarkWeaveRequestSimpleMessageWithBindings(Blackhole blackhole) { - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( simpleMessage, simpleBindings, null, @@ -118,7 +118,7 @@ public void benchmarkWeaveRequestSimpleMessageWithBindings(Blackhole blackhole) @Benchmark public void benchmarkWeaveRequestComplexMessageWithBindings(Blackhole blackhole) { - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( complexMessage, complexBindings, null, @@ -129,7 +129,7 @@ public void benchmarkWeaveRequestComplexMessageWithBindings(Blackhole blackhole) @Benchmark public void benchmarkWeaveRequestWithTranscodingPath(Blackhole blackhole) { - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( complexMessage, complexBindings, complexTranscodingPath, @@ -140,7 +140,7 @@ public void benchmarkWeaveRequestWithTranscodingPath(Blackhole blackhole) { @Benchmark public void benchmarkWeaveResponseSimpleMessage(Blackhole blackhole) { - Buffer result = MessageWeaver.weaveResponseMessage( + Object result = MessageWeaver.weaveResponseMessage( simpleMessage, simpleTranscodingPath ); @@ -149,7 +149,7 @@ public void benchmarkWeaveResponseSimpleMessage(Blackhole blackhole) { @Benchmark public void benchmarkWeaveResponseComplexMessage(Blackhole blackhole) { - Buffer result = MessageWeaver.weaveResponseMessage( + Object result = MessageWeaver.weaveResponseMessage( complexMessage, complexTranscodingPath ); diff --git a/vertx-grpc-transcoding/src/test/java/io/vertx/grpc/transcoding/tests/MessageWeaverTest.java b/vertx-grpc-transcoding/src/test/java/io/vertx/grpc/transcoding/tests/MessageWeaverTest.java index 1694d415e..4de246cea 100644 --- a/vertx-grpc-transcoding/src/test/java/io/vertx/grpc/transcoding/tests/MessageWeaverTest.java +++ b/vertx-grpc-transcoding/src/test/java/io/vertx/grpc/transcoding/tests/MessageWeaverTest.java @@ -15,8 +15,7 @@ import java.util.Arrays; import java.util.List; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertThrows; +import static org.junit.Assert.*; public class MessageWeaverTest { @@ -82,14 +81,14 @@ public void testPassThroughRequest() { .put("B", new JsonObject() .put("y", "b"))); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(message.encode()), new ArrayList<>(), "*", TEST_DESCRIPTOR ); - assertEquals(message, new JsonObject(result.toString())); + assertEquals(message, result); } @Test @@ -109,14 +108,14 @@ public void testLevel0Bindings() { .put("_y", "b") .put("_z", "c"); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(original.encode()), bindings, "*", TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test @@ -140,14 +139,14 @@ public void testLevel1Bindings() { .put("z", "f") .put("_x", "c")); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(original.encode()), bindings, "*", TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test @@ -179,14 +178,14 @@ public void testLevel2Bindings() { .put("u", "g") .put("_x", "c"))); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(original.encode()), bindings, "*", TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test @@ -195,14 +194,14 @@ public void testTranscodingRequestBodyWildcard() { .put("field1", "value1") .put("field2", "value2"); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(message.encode()), new ArrayList<>(), "*", TEST_DESCRIPTOR ); - assertEquals(message, new JsonObject(result.toString())); + assertEquals(message, result); } @Test @@ -213,14 +212,14 @@ public void testTranscodingRequestBodyPath() { JsonObject expected = new JsonObject().put("nested", new JsonObject().put("data", message)); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(message.encode()), new ArrayList<>(), "nested.data", TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test @@ -229,12 +228,12 @@ public void testResponsePassThrough() { .put("field1", "value1") .put("field2", "value2"); - Buffer result = MessageWeaver.weaveResponseMessage( + Object result = MessageWeaver.weaveResponseMessage( Buffer.buffer(message.encode()), null ); - assertEquals(message, new JsonObject(result.toString())); + assertEquals(message, result); } @Test @@ -243,12 +242,12 @@ public void testResponseBodyWildcard() { .put("field1", "value1") .put("field2", "value2"); - Buffer result = MessageWeaver.weaveResponseMessage( + Object result = MessageWeaver.weaveResponseMessage( Buffer.buffer(message.encode()), "*" ); - assertEquals(message, new JsonObject(result.toString())); + assertEquals(message, result); } @Test @@ -261,12 +260,53 @@ public void testResponseBodyPath() { .put("response", new JsonObject() .put("data", nested)); - Buffer result = MessageWeaver.weaveResponseMessage( + Object result = MessageWeaver.weaveResponseMessage( Buffer.buffer(message.encode()), "response.data" ); - assertEquals(nested, new JsonObject(result.toString())); + assertEquals(nested, result); + } + + @Test + public void testResponseBodyRepeatedField() { + JsonArray users = new JsonArray() + .add(new JsonObject().put("id", 1).put("name", "alice")) + .add(new JsonObject().put("id", 2).put("name", "bob")); + + JsonObject message = new JsonObject().put("users", users); + + Object result = MessageWeaver.weaveResponseMessage( + Buffer.buffer(message.encode()), + "users" + ); + + assertEquals(users, result); + } + + @Test + public void testResponseBodyScalarField() { + JsonObject message = new JsonObject() + .put("response", new JsonObject().put("token", "abc123")); + + Object result = MessageWeaver.weaveResponseMessage( + Buffer.buffer(message.encode()), + "response.token" + ); + + assertEquals("abc123", result); + } + + @Test + public void testResponseBodyNullLeaf() { + JsonObject message = new JsonObject().put("response", new JsonObject().putNull("token")); + + Object result = MessageWeaver.weaveResponseMessage( + Buffer.buffer(message.encode()), + "response.token" + ); + + assertNull(result); } @Test @@ -276,7 +316,7 @@ public void testResponseBodyInvalidPath() { assertThrows(IllegalArgumentException.class, () -> MessageWeaver.weaveResponseMessage( Buffer.buffer(message.encode()), - "invalid.path" + "field1.invalid" )); } @@ -289,14 +329,14 @@ public void testRepeatedBindingsLevel0() { JsonObject expected = new JsonObject() .put("x", new JsonArray().add("a").add("b").add("c")); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( null, bindings, null, TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test @@ -309,14 +349,14 @@ public void testRepeatedBindingsLevel1() { .put("A", new JsonObject() .put("x", new JsonArray().add("b").add("c").add("d"))); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( null, bindings, null, TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test @@ -333,14 +373,14 @@ public void testRepeatedBindingsWithBody() { .put("x", new JsonArray().add("b").add("c").add("d")) .put("y", "e")); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(body.encode()), bindings, "*", TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test @@ -353,14 +393,14 @@ public void testRepeatedBindingsMixedWithSingle() { .put("x", new JsonArray().add("a").add("b")) .put("y", "c"); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( null, bindings, null, TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test @@ -371,38 +411,38 @@ public void testRepeatedBindingsTwoValues() { JsonObject expected = new JsonObject() .put("keys", new JsonArray().add("A").add("B")); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( null, bindings, null, TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test public void testEmptyMessageReturnsEmptyObject() { - Buffer fromNull = MessageWeaver.weaveRequestMessage(null, new ArrayList<>(), null, TEST_DESCRIPTOR); - assertEquals(new JsonObject(), new JsonObject(fromNull.toString())); + JsonObject fromNull = MessageWeaver.weaveRequestMessage(null, new ArrayList<>(), null, TEST_DESCRIPTOR); + assertEquals(new JsonObject(), fromNull); - Buffer fromEmpty = MessageWeaver.weaveRequestMessage(Buffer.buffer(), new ArrayList<>(), null, TEST_DESCRIPTOR); - assertEquals(new JsonObject(), new JsonObject(fromEmpty.toString())); + JsonObject fromEmpty = MessageWeaver.weaveRequestMessage(Buffer.buffer(), new ArrayList<>(), null, TEST_DESCRIPTOR); + assertEquals(new JsonObject(), fromEmpty); - Buffer fromEmptyBody = MessageWeaver.weaveRequestMessage(Buffer.buffer(), new ArrayList<>(), "", TEST_DESCRIPTOR); - assertEquals(new JsonObject(), new JsonObject(fromEmptyBody.toString())); + JsonObject fromEmptyBody = MessageWeaver.weaveRequestMessage(Buffer.buffer(), new ArrayList<>(), "", TEST_DESCRIPTOR); + assertEquals(new JsonObject(), fromEmptyBody); } @Test public void testUnsetBodyIgnoresHttpBody() { JsonObject httpBody = new JsonObject().put("field1", "value1"); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(httpBody.encode()), new ArrayList<>(), null, TEST_DESCRIPTOR ); - assertEquals(new JsonObject(), new JsonObject(result.toString())); + assertEquals(new JsonObject(), result); } @Test @@ -417,14 +457,14 @@ public void testBindingOverridesBodyScalar() { .put("y", "from-binding") .put("A", new JsonObject().put("y", "nested-body")); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(body.encode()), bindings, "*", TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test @@ -437,14 +477,14 @@ public void testBindingOverridesNestedBodyScalar() { JsonObject expected = new JsonObject() .put("A", new JsonObject().put("y", "from-binding").put("x", new JsonArray().add("k"))); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(body.encode()), bindings, "*", TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test @@ -457,14 +497,14 @@ public void testBindingCreatesSubtreeAbsentFromBody() { .put("y", "root-body") .put("A", new JsonObject().put("y", "from-binding")); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(body.encode()), bindings, "*", TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test @@ -477,14 +517,14 @@ public void testRepeatedBindingAppendedToBodyValues() { JsonObject expected = new JsonObject() .put("x", new JsonArray().add("a").add("b").add("c")); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(body.encode()), bindings, "*", TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test @@ -498,14 +538,14 @@ public void testBindingAndBodyAtSpecificPath() { .put("x", new JsonArray().add("b")) .put("y", "from-binding")); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(body.encode()), bindings, "A", TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test @@ -521,14 +561,14 @@ public void testBindingsDoNotDescendIntoArrayElements() { .add(new JsonObject().put("A", new JsonObject().put("x", "in-array")))) .put("A", new JsonObject().put("y", "from-binding")); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(body.encode()), bindings, "*", TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test @@ -540,14 +580,14 @@ public void testDeepBodyPath() { .put("b", new JsonObject() .put("c", body))); - Buffer result = MessageWeaver.weaveRequestMessage( + JsonObject result = MessageWeaver.weaveRequestMessage( Buffer.buffer(body.encode()), new ArrayList<>(), "a.b.c", TEST_DESCRIPTOR ); - assertEquals(expected, new JsonObject(result.toString())); + assertEquals(expected, result); } @Test diff --git a/vertx-grpc-transcoding/src/test/java/io/vertx/grpc/transcoding/tests/ServerTranscodingTest.java b/vertx-grpc-transcoding/src/test/java/io/vertx/grpc/transcoding/tests/ServerTranscodingTest.java index 69698763a..15b206bd7 100644 --- a/vertx-grpc-transcoding/src/test/java/io/vertx/grpc/transcoding/tests/ServerTranscodingTest.java +++ b/vertx-grpc-transcoding/src/test/java/io/vertx/grpc/transcoding/tests/ServerTranscodingTest.java @@ -20,6 +20,7 @@ import io.vertx.core.buffer.Buffer; import io.vertx.core.http.*; import io.vertx.core.internal.buffer.BufferInternal; +import io.vertx.core.json.JsonArray; import io.vertx.core.json.JsonObject; import io.vertx.ext.unit.TestContext; import io.vertx.grpc.common.*; @@ -59,6 +60,8 @@ public static MethodTranscodingOptions create(String selector, HttpMethod httpMe public static GrpcMessageDecoder ECHO_REQUEST_BODY_DECODER = GrpcMessageDecoder.decoder(EchoRequestBody.newBuilder()); public static GrpcMessageEncoder ECHO_RESPONSE_ENCODER = GrpcMessageEncoder.encoder(); public static GrpcMessageEncoder ECHO_RESPONSE_BODY_ENCODER = GrpcMessageEncoder.encoder(); + public static GrpcMessageDecoder STREAMING_REQUEST_DECODER = GrpcMessageDecoder.decoder(StreamingRequest.newBuilder()); + public static GrpcMessageEncoder STREAMING_RESPONSE_ENCODER = GrpcMessageEncoder.encoder(); public static final ServiceName TEST_SERVICE_NAME = ServiceName.create(TestServiceGrpc.SERVICE_NAME); @@ -91,6 +94,15 @@ public static MethodTranscodingOptions create(String selector, HttpMethod httpMe public static final TranscodingServiceMethod UNARY_CALL_WITH_REPEATED_QUERY = TranscodingServiceMethod.server(TEST_SERVICE_NAME, "UnaryCallWithRepeatedQuery", ECHO_RESPONSE_ENCODER, ECHO_REQUEST_DECODER, UNARY_TRANSCODING_WITH_REPEATED_QUERY); + public static final MethodTranscodingOptions STREAMING_TRANSCODING = new MethodTranscodingOptions().setHttpMethod(HttpMethod.POST).setPath("/stream").setBody("*"); + public static final TranscodingServiceMethod STREAMING_CALL = TranscodingServiceMethod.server(TEST_SERVICE_NAME, "StreamingCall", + MethodCardinality.SERVER_STREAMING, STREAMING_RESPONSE_ENCODER, STREAMING_REQUEST_DECODER, STREAMING_TRANSCODING); + + public static final MethodTranscodingOptions UNARY_TRANSCODING_WITH_ARRAY_RESPONSE_BODY = new MethodTranscodingOptions().setHttpMethod(HttpMethod.POST).setPath("/keylist").setBody("*").setResponseBody("keys"); + public static GrpcMessageEncoder ECHO_REQUEST_ENCODER = GrpcMessageEncoder.encoder(); + public static final TranscodingServiceMethod UNARY_CALL_WITH_ARRAY_RESPONSE_BODY = TranscodingServiceMethod.server(TEST_SERVICE_NAME, "UnaryCallWithArrayResponseBody", + ECHO_REQUEST_ENCODER, ECHO_REQUEST_DECODER, UNARY_TRANSCODING_WITH_ARRAY_RESPONSE_BODY); + private static final CharSequence USER_AGENT = HttpHeaders.createOptimized("X-User-Agent"); private static final String CONTENT_TYPE = "application/json"; @@ -239,6 +251,36 @@ public void setUp(TestContext should) { response.end(responseMsg); }); }); + grpcServer.callHandler(UNARY_CALL_WITH_ARRAY_RESPONSE_BODY, request -> { + request.handler(requestMsg -> { + GrpcServerResponse response = request.response(); + EchoRequest responseMsg = EchoRequest.newBuilder() + .addKeys("alice") + .addKeys("bob") + .addKeys("carol") + .build(); + response.end(responseMsg); + }); + }); + grpcServer.callHandler(STREAMING_CALL, request -> { + request.handler(requestMsg -> { + GrpcServerResponse response = request.response(); + if (requestMsg.getResponseSizeList().isEmpty()) { + response.end(); + return; + } + for (int size : requestMsg.getResponseSizeList()) { + if (size < 0) { + response.status(GrpcStatus.INVALID_ARGUMENT).end(); + return; + } + char[] value = new char[size]; + Arrays.fill(value, 'a'); + response.write(StreamingResponse.newBuilder().setPayload(new String(value)).build()); + } + response.end(); + }); + }); httpServer = vertx.createHttpServer(new HttpServerOptions().setPort(port)).requestHandler(grpcServer); httpServer.listen().onComplete(should.asyncAssertSuccess()); } @@ -534,6 +576,155 @@ public void testRepeatedQueryParamsSingleValue(TestContext should) { }))); } + @Test + public void testServerStreaming(TestContext should) { + httpClient.request(HttpMethod.POST, "/stream").compose(req -> { + String body = encode(StreamingRequest.newBuilder().addResponseSize(1).addResponseSize(2).addResponseSize(3).build()).toString(); + req.headers().addAll(HEADERS); + return req.send(body).compose(response -> response.body().map(response)); + }).onComplete(should.asyncAssertSuccess(response -> should.verify(v -> { + assertEquals(200, response.statusCode()); + MultiMap headers = response.headers(); + assertTrue(headers.contains(HttpHeaders.CONTENT_TYPE, CONTENT_TYPE, true)); + assertFalse(headers.contains(HttpHeaders.CONTENT_LENGTH)); + JsonArray array = new JsonArray(response.body().result()); + assertEquals(3, array.size()); + assertEquals("a", array.getJsonObject(0).getString("payload")); + assertEquals("aa", array.getJsonObject(1).getString("payload")); + assertEquals("aaa", array.getJsonObject(2).getString("payload")); + }))); + } + + @Test + public void testServerStreamingEmpty(TestContext should) { + httpClient.request(HttpMethod.POST, "/stream").compose(req -> { + String body = encode(StreamingRequest.newBuilder().build()).toString(); + req.headers().addAll(HEADERS); + return req.send(body).compose(response -> response.body().map(response)); + }).onComplete(should.asyncAssertSuccess(response -> should.verify(v -> { + assertEquals(200, response.statusCode()); + JsonArray array = new JsonArray(response.body().result()); + assertEquals(0, array.size()); + }))); + } + + @Test + public void testServerStreamingErrorBeforeFirstMessage(TestContext should) { + httpClient.request(HttpMethod.POST, "/stream").compose(req -> { + String body = encode(StreamingRequest.newBuilder().addResponseSize(-1).build()).toString(); + req.headers().addAll(HEADERS); + return req.send(body).compose(response -> response.body().map(response)); + }).onComplete(should.asyncAssertSuccess(response -> should.verify(v -> { + assertEquals(400, response.statusCode()); + }))); + } + + @Test + public void testServerStreamingNdjson(TestContext should) { + httpClient.request(HttpMethod.POST, "/stream").compose(req -> { + String body = encode(StreamingRequest.newBuilder().addResponseSize(1).addResponseSize(2).addResponseSize(3).build()).toString(); + req.headers().addAll(HEADERS); + req.headers().set(HttpHeaders.ACCEPT, "application/x-ndjson"); + return req.send(body).compose(response -> response.body().map(response)); + }).onComplete(should.asyncAssertSuccess(response -> should.verify(v -> { + assertEquals(200, response.statusCode()); + assertTrue(response.headers().contains(HttpHeaders.CONTENT_TYPE, "application/x-ndjson", true)); + String[] lines = response.body().result().toString().split("\n"); + assertEquals(3, lines.length); + assertEquals("a", new JsonObject(lines[0]).getString("payload")); + assertEquals("aa", new JsonObject(lines[1]).getString("payload")); + assertEquals("aaa", new JsonObject(lines[2]).getString("payload")); + }))); + } + + @Test + public void testServerStreamingSse(TestContext should) { + httpClient.request(HttpMethod.POST, "/stream").compose(req -> { + String body = encode(StreamingRequest.newBuilder().addResponseSize(1).addResponseSize(2).addResponseSize(3).build()).toString(); + req.headers().addAll(HEADERS); + req.headers().set(HttpHeaders.ACCEPT, "text/event-stream"); + return req.send(body).compose(response -> response.body().map(response)); + }).onComplete(should.asyncAssertSuccess(response -> should.verify(v -> { + assertEquals(200, response.statusCode()); + assertTrue(response.headers().contains(HttpHeaders.CONTENT_TYPE, "text/event-stream", true)); + String[] events = response.body().result().toString().split("\n\n"); + assertEquals(3, events.length); + for (int i = 0; i < events.length; i++) { + assertTrue("event[" + i + "] should start with 'data: ': " + events[i], events[i].startsWith("data: ")); + String json = events[i].substring("data: ".length()); + assertEquals("a".repeat(i + 1), new JsonObject(json).getString("payload")); + } + }))); + } + + @Test + public void testServerStreamingErrorMidStream(TestContext should) { + httpClient.request(HttpMethod.POST, "/stream").compose(req -> { + String body = encode(StreamingRequest.newBuilder().addResponseSize(1).addResponseSize(-1).build()).toString(); + req.headers().addAll(HEADERS); + return req.send(body).compose(response -> response.body().map(response)); + }).onComplete(should.asyncAssertSuccess(response -> should.verify(v -> { + assertEquals(200, response.statusCode()); + JsonArray array = new JsonArray(response.body().result()); + assertEquals(1, array.size()); + assertEquals("a", array.getJsonObject(0).getString("payload")); + }))); + } + + @Test + public void testServerStreamingUnknownAcceptFallsBackToJsonArray(TestContext should) { + httpClient.request(HttpMethod.POST, "/stream").compose(req -> { + String body = encode(StreamingRequest.newBuilder().addResponseSize(1).build()).toString(); + req.headers().addAll(HEADERS); + req.headers().set(HttpHeaders.ACCEPT, "application/xml"); + return req.send(body).compose(response -> response.body().map(response)); + }).onComplete(should.asyncAssertSuccess(response -> should.verify(v -> { + assertEquals(200, response.statusCode()); + assertTrue(response.headers().contains(HttpHeaders.CONTENT_TYPE, CONTENT_TYPE, true)); + JsonArray array = new JsonArray(response.body().result()); + assertEquals(1, array.size()); + }))); + } + + @Test + public void testServerStreamingWildcardAcceptFallsBackToJsonArray(TestContext should) { + httpClient.request(HttpMethod.POST, "/stream").compose(req -> { + String body = encode(StreamingRequest.newBuilder().addResponseSize(1).build()).toString(); + req.headers().addAll(HEADERS); + req.headers().set(HttpHeaders.ACCEPT, "*/*"); + return req.send(body).compose(response -> response.body().map(response)); + }).onComplete(should.asyncAssertSuccess(response -> should.verify(v -> { + assertEquals(200, response.statusCode()); + assertTrue(response.headers().contains(HttpHeaders.CONTENT_TYPE, CONTENT_TYPE, true)); + JsonArray array = new JsonArray(response.body().result()); + assertEquals(1, array.size()); + }))); + } + + @Test + public void testResponseBodyRepeatedFieldUnwrapped(TestContext should) { + httpClient.request(HttpMethod.POST, "/keylist").compose(req -> { + String body = encode(EchoRequest.newBuilder().build()).toString(); + req.headers().addAll(HEADERS); + return req.send(body).compose(response -> response.body().map(response)); + }).onComplete(should.asyncAssertSuccess(response -> should.verify(v -> { + assertEquals(200, response.statusCode()); + assertTrue(response.headers().contains(HttpHeaders.CONTENT_TYPE, CONTENT_TYPE, true)); + JsonArray expected = new JsonArray().add("alice").add("bob").add("carol"); + assertEquals(expected, new JsonArray(response.body().result())); + }))); + } + + @Test + public void testServerStreamingMalformedBody(TestContext should) { + httpClient.request(HttpMethod.POST, "/stream").compose(req -> { + req.headers().addAll(HEADERS); + return req.send("not-json").compose(response -> response.body().map(response)); + }).onComplete(should.asyncAssertSuccess(response -> should.verify(v -> { + assertEquals(400, response.statusCode()); + }))); + } + @Test public void testHttp2Proto() { ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", port)