From 1f8f490e2fafe0fd1a6b96a4f54ae853c0621014 Mon Sep 17 00:00:00 2001 From: Daniel Fiala Date: Sun, 16 Aug 2026 12:08:52 +0200 Subject: [PATCH 1/5] Enable JSON transcoding for server-streaming RPCs and refine response formatting. Motivation: - Extend JSON transcoding support for server-streaming RPCs to enhance compatibility with streaming APIs. - Improve response flexibility with support for JSON Array, NDJSON, and SSE formats based on the `Accept` header. Changes: - Added `streaming` and `StreamFormat` fields to `TranscodingGrpcOutboundStream` for response format negotiation. - Implemented JSON Array, NDJSON, and SSE response formatting for server-streaming. - Derived the streaming cardinality from `ServiceMethod#serverStreaming()`. - Updated `MessageWeaver` to handle `JsonObject` outputs for request/response weaving. - Modified tests to validate new streaming response formats and HTTP transcoding logic. Signed-off-by: Daniel Fiala --- .../src/main/asciidoc/transcoding.adoc | 34 ++++ .../grpc/transcoding/impl/MessageWeaver.java | 37 ++-- .../impl/TranscodingGrpcOutboundStream.java | 122 ++++++++++- .../impl/TranscodingServiceMethodImpl.java | 11 +- .../benchmarks/MessageWeaverBenchmark.java | 12 +- .../transcoding/tests/MessageWeaverTest.java | 146 ++++++++----- .../tests/ServerTranscodingTest.java | 191 ++++++++++++++++++ 7 files changed, 464 insertions(+), 89 deletions(-) diff --git a/vertx-grpc-docs/src/main/asciidoc/transcoding.adoc b/vertx-grpc-docs/src/main/asciidoc/transcoding.adoc index 1ef302817..bb5a433b3 100644 --- a/vertx-grpc-docs/src/main/asciidoc/transcoding.adoc +++ b/vertx-grpc-docs/src/main/asciidoc/transcoding.adoc @@ -212,6 +212,40 @@ 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 per-request from the HTTP `Accept` header so clients can pick the encoding that fits their consumption pattern without any server configuration: + +|=== +| `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"} +---- + +If the RPC terminates with a non-OK gRPC status before any message has been written, the response is finished with the corresponding HTTP status code and no body. If messages have already been written, the stream ends (with the closing `]` for JSON array mode). Mid-stream errors cannot be signalled on the body since the HTTP status was already sent. + +Client-streaming and bidirectional-streaming RPCs are not supported by transcoding. + === 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-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); @@ -83,4 +138,53 @@ public Future writeMessage(GrpcMessageFrame frame) { } 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()); + } + return super.writeEnd(); + } } 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) From 05d2bceba98421057269352f0a6074216c9ffd4d Mon Sep 17 00:00:00 2001 From: Daniel Fiala Date: Sun, 16 Aug 2026 12:13:04 +0200 Subject: [PATCH 2/5] Add integration tests for server-streaming JSON transcoding Motivation: - The server-streaming transcoding formats were only covered by unit tests using hand-built service methods, leaving the protoc plugin path untested end to end. Changes: - Exercise the generated `StreamingTranscodingGreeter` service method over HTTP for the JSON array, NDJSON and SSE formats. - Cover the empty stream and the coexistence of the transcoded and gRPC routes for the same method. Signed-off-by: Daniel Fiala --- .../vertx/grpc/it/tests/TranscodingTest.java | 169 ++++++++++++++++++ 1 file changed, 169 insertions(+) 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..09f30f845 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,36 @@ 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.stub.StreamObserver; +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.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 +357,162 @@ public void testUnaryCollisionWithoutOption() throws TimeoutException { assertEquals("Hello Julien", reply.getMessage()); } + @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(); From 8ce579942c4216ef3ca120c8fe43142a995561b2 Mon Sep 17 00:00:00 2001 From: Daniel Fiala Date: Sun, 16 Aug 2026 12:16:59 +0200 Subject: [PATCH 3/5] Register the mount point paths of the methods of a service Motivation: - `GrpcServer#addService` only registered the canonical `/package.Service/Method` path, so the HTTP rules of a transcoded service method were ignored and the transcoded routes returned a 500. `GrpcServer#callHandler` already mounts these paths. Changes: - Mount the `MountPoint` paths of a service method in `addService`, like `callHandler` does. - Add integration tests binding a transcoded unary and server-streaming service with `addService`. Signed-off-by: Daniel Fiala --- .../vertx/grpc/it/tests/TranscodingTest.java | 58 +++++++++++++++++++ .../grpc/server/impl/GrpcServerImpl.java | 9 ++- 2 files changed, 66 insertions(+), 1 deletion(-) 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 09f30f845..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 @@ -6,7 +6,9 @@ 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.*; @@ -14,6 +16,7 @@ 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; @@ -357,6 +360,61 @@ 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(); 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); From afb6a69e5b1b952f81895905585ebb56de494b22 Mon Sep 17 00:00:00 2001 From: Daniel Fiala Date: Wed, 19 Aug 2026 19:15:02 +0200 Subject: [PATCH 4/5] Resolve the transcoding head promise when a unary call ends without a message. Motivation: Transcoding defers the HTTP head until the response body is known, so writeHead returns a promise resolved by writeUnaryMessage. A unary call that ends without ever writing a message never resolves it and leaves the caller waiting. Changes: Resolve the head promise from writeEnd when it is still pending. Fix the indentation of the SSE branch. --- .../impl/TranscodingGrpcOutboundStream.java | 18 ++++++++++++++---- 1 file changed, 14 insertions(+), 4 deletions(-) diff --git a/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingGrpcOutboundStream.java b/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingGrpcOutboundStream.java index 95a3d5fd2..d754cea47 100644 --- a/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingGrpcOutboundStream.java +++ b/vertx-grpc-transcoding/src/main/java/io/vertx/grpc/transcoding/impl/TranscodingGrpcOutboundStream.java @@ -133,8 +133,10 @@ private Future writeUnaryMessage(Buffer payload) { 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; } @@ -158,7 +160,7 @@ private Future writeStreamingMessage(Buffer payload) { chunk = Buffer.buffer(transcoded.length() + NEWLINE.length()).appendBuffer(transcoded).appendBuffer(NEWLINE); break; } - case SSE: { + case SSE: { chunk = Buffer.buffer(SSE_PREFIX.length() + transcoded.length() + SSE_SUFFIX.length()) .appendBuffer(SSE_PREFIX).appendBuffer(transcoded).appendBuffer(SSE_SUFFIX); break; @@ -185,6 +187,14 @@ public Future writeEnd() { if (status != GrpcStatus.OK) { httpResponse.setStatusCode(GrpcTranscodingError.fromHttp2Code(status.code).getHttpStatusCode()); } - return super.writeEnd(); + 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; } } From fd3fe780ec9298e3bebd227d49dd1e7961ef06e4 Mon Sep 17 00:00:00 2001 From: Daniel Fiala Date: Thu, 3 Sep 2026 19:05:39 +0200 Subject: [PATCH 5/5] Reword the server-streaming transcoding documentation Motivation: - Review feedback on the server-streaming transcoding section asked for a more concise wording of the content negotiation and error reporting paragraphs. Changes: - Describe the wire format as selected according to the HTTP `accept` header. - Split the error reporting paragraph into the trailers-only and trailers cases. - Reword the unsupported streaming modes sentence. Signed-off-by: Daniel Fiala --- vertx-grpc-docs/src/main/asciidoc/transcoding.adoc | 9 ++++++--- 1 file changed, 6 insertions(+), 3 deletions(-) diff --git a/vertx-grpc-docs/src/main/asciidoc/transcoding.adoc b/vertx-grpc-docs/src/main/asciidoc/transcoding.adoc index bb5a433b3..2f335fec9 100644 --- a/vertx-grpc-docs/src/main/asciidoc/transcoding.adoc +++ b/vertx-grpc-docs/src/main/asciidoc/transcoding.adoc @@ -214,7 +214,7 @@ 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 per-request from the HTTP `Accept` header so clients can pick the encoding that fits their consumption pattern without any server configuration: +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 @@ -242,9 +242,12 @@ curl -X POST -H "Accept: text/event-stream" -H "Content-Type: application/json" # data: {"payload":"second"} ---- -If the RPC terminates with a non-OK gRPC status before any message has been written, the response is finished with the corresponding HTTP status code and no body. If messages have already been written, the stream ends (with the closing `]` for JSON array mode). Mid-stream errors cannot be signalled on the body since the HTTP status was already sent. +A gRPC error is reported to the client depending on when it happens: -Client-streaming and bidirectional-streaming RPCs are not supported by transcoding. +- 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