<?xml version="1.0" encoding="UTF-8"?>
<rss version="2.0" xmlns:atom="http://www.w3.org/2005/Atom" xmlns:dc="http://purl.org/dc/elements/1.1/">
  <channel>
    <title>DEV Community: Wang Lee</title>
    <description>The latest articles on DEV Community by Wang Lee (@qianwj).</description>
    <link>https://dev.to/qianwj</link>
    <image>
      <url>https://media2.dev.to/dynamic/image/width=90,height=90,fit=cover,gravity=auto,format=auto/https:%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Fuser%2Fprofile_image%2F317772%2F2129f777-227e-4884-a83f-e26df2c7d7fe.jpeg</url>
      <title>DEV Community: Wang Lee</title>
      <link>https://dev.to/qianwj</link>
    </image>
    <atom:link rel="self" type="application/rss+xml" href="https://dev.to/feed/qianwj"/>
    <language>en</language>
    <item>
      <title>Unary gRPC on Reactor Netty: Event Loop Serialization, Trailers, and Cancellation</title>
      <dc:creator>Wang Lee</dc:creator>
      <pubDate>Sun, 09 Aug 2026 03:18:25 +0000</pubDate>
      <link>https://dev.to/qianwj/unary-grpc-on-reactor-netty-event-loop-serialization-trailers-and-cancellation-50ld</link>
      <guid>https://dev.to/qianwj/unary-grpc-on-reactor-netty-event-loop-serialization-trailers-and-cancellation-50ld</guid>
      <description>&lt;p&gt;With protocol values and message framing complete, Stage 2 delivered the first end-to-end call: plaintext h2c unary RPC. This is already on &lt;code&gt;main&lt;/code&gt;, and Stage 3 and Stage 4 subsequently completed all four RPC cardinalities on the same transport primitive.&lt;/p&gt;

&lt;p&gt;Previous: &lt;a href="https://dev.to/qianwj/building-a-leak-safe-grpc-frame-decoder-on-reactor-netty-po7"&gt;Building a Leak-Safe gRPC Frame Decoder on Reactor Netty&lt;/a&gt;&lt;/p&gt;

&lt;h2&gt;
  
  
  Method Descriptor Is Where Protocol Meets Types
&lt;/h2&gt;

&lt;p&gt;A method requires a precise service name, method name, cardinality, and request/response marshallers:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="kt"&gt;var&lt;/span&gt; &lt;span class="n"&gt;echo&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="nc"&gt;GrpcMethod&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&amp;gt;(&lt;/span&gt;
    &lt;span class="s"&gt;"testing.EchoService"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt;
    &lt;span class="s"&gt;"Echo"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt;
    &lt;span class="nc"&gt;GrpcMethod&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;Cardinality&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;UNARY&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt;
    &lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="nc"&gt;ProtobufMarshaller&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&amp;gt;(&lt;/span&gt;&lt;span class="nc"&gt;StringValue&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;parser&lt;/span&gt;&lt;span class="o"&gt;()),&lt;/span&gt;
    &lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="nc"&gt;ProtobufMarshaller&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&amp;gt;(&lt;/span&gt;&lt;span class="nc"&gt;StringValue&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;parser&lt;/span&gt;&lt;span class="o"&gt;()));&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;The generated path must be:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;/testing.EchoService/Echo
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;The service registry matches by exact full path. An unknown path returns &lt;code&gt;UNIMPLEMENTED&lt;/code&gt;; registering the same path twice fails immediately when building the service definition.&lt;/p&gt;

&lt;h2&gt;
  
  
  Server Validates Protocol Before Subscribing to Business Logic
&lt;/h2&gt;

&lt;p&gt;&lt;code&gt;ReactorGrpcServer&lt;/code&gt; uses Reactor Netty h2c:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="nc"&gt;DisposableServer&lt;/span&gt; &lt;span class="n"&gt;bound&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nc"&gt;HttpServer&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;create&lt;/span&gt;&lt;span class="o"&gt;()&lt;/span&gt;
    &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;host&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;host&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
    &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;port&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;port&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
    &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;protocol&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;HttpProtocol&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;H2C&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
    &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;handle&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nl"&gt;handler:&lt;/span&gt;&lt;span class="o"&gt;:&lt;/span&gt;&lt;span class="n"&gt;handle&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
    &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;bindNow&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;Duration&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;ofSeconds&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;10&lt;/span&gt;&lt;span class="o"&gt;));&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Incoming requests are validated in order:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;HTTP method must be POST;&lt;/li&gt;
&lt;li&gt;content-type must be &lt;code&gt;application/grpc&lt;/code&gt; or &lt;code&gt;application/grpc+...&lt;/code&gt;;&lt;/li&gt;
&lt;li&gt;
&lt;code&gt;te&lt;/code&gt; must declare trailers;&lt;/li&gt;
&lt;li&gt;path must exist;&lt;/li&gt;
&lt;li&gt;currently only unary cardinality is allowed;&lt;/li&gt;
&lt;li&gt;metadata and message size must not exceed limits.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;Only after validation passes does it create a &lt;code&gt;GrpcCallContext&lt;/code&gt; and subscribe to the request body, preventing invalid requests from entering the business handler.&lt;/p&gt;

&lt;h2&gt;
  
  
  HTTP 200 Does Not Mean RPC Success
&lt;/h2&gt;

&lt;p&gt;The server writes a compatible content-type first; the final status comes from trailing headers:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="n"&gt;response&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;status&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;200&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
    &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;header&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;HttpHeaderNames&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;CONTENT_TYPE&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="s"&gt;"application/grpc+proto"&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;

&lt;span class="n"&gt;response&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;trailerHeaders&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;trailers&lt;/span&gt; &lt;span class="o"&gt;-&amp;gt;&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
  &lt;span class="nc"&gt;GrpcException&lt;/span&gt; &lt;span class="n"&gt;error&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;terminal&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;get&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt;
  &lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;error&lt;/span&gt; &lt;span class="o"&gt;==&lt;/span&gt; &lt;span class="kc"&gt;null&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
    &lt;span class="n"&gt;writeStatus&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;trailers&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="nc"&gt;GrpcStatus&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;OK&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="nc"&gt;GrpcMetadata&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;empty&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;
  &lt;span class="o"&gt;}&lt;/span&gt; &lt;span class="k"&gt;else&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
    &lt;span class="n"&gt;writeStatus&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;trailers&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;error&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;status&lt;/span&gt;&lt;span class="o"&gt;(),&lt;/span&gt; &lt;span class="n"&gt;error&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;trailers&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;
  &lt;span class="o"&gt;}&lt;/span&gt;
&lt;span class="o"&gt;});&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;A &lt;code&gt;GrpcException&lt;/code&gt; thrown by the application preserves an explicit status and trailing metadata. Ordinary application exceptions map to &lt;code&gt;UNKNOWN&lt;/code&gt;; protocol violations map to &lt;code&gt;INTERNAL&lt;/code&gt;. Fatal JVM errors should not be wrapped as ordinary business status.&lt;/p&gt;

&lt;p&gt;After receiving a response, the client must read the final &lt;code&gt;grpc-status&lt;/code&gt;. HTTP 200 without a status is &lt;code&gt;INTERNAL&lt;/code&gt;; a non-gRPC HTTP 404 maps to &lt;code&gt;UNIMPLEMENTED&lt;/code&gt; per the standard fallback, and 503 maps to &lt;code&gt;UNAVAILABLE&lt;/code&gt;.&lt;/p&gt;

&lt;h2&gt;
  
  
  Metadata Must Preserve Order and Allow Duplicates
&lt;/h2&gt;

&lt;p&gt;gRPC metadata is not a simple &lt;code&gt;Map&amp;lt;String, String&amp;gt;&lt;/code&gt;. The same key can repeat, order must be preserved, and the &lt;code&gt;-bin&lt;/code&gt; suffix indicates binary values encoded as no-padding Base64 on the wire.&lt;/p&gt;

&lt;p&gt;&lt;code&gt;GrpcMetadata&lt;/code&gt; therefore stores an immutable entry list internally:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="nc"&gt;GrpcMetadata&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;builder&lt;/span&gt;&lt;span class="o"&gt;()&lt;/span&gt;
    &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;addAscii&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"trace-id"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="s"&gt;"abc"&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
    &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;addBinary&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"span-bin"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;bytes&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
    &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;build&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;User metadata cannot overwrite reserved fields like &lt;code&gt;content-type&lt;/code&gt;, &lt;code&gt;te&lt;/code&gt;, &lt;code&gt;grpc-status&lt;/code&gt;, &lt;code&gt;grpc-timeout&lt;/code&gt;, and HTTP/2 pseudo headers. The total encoded size is also limited on read to prevent peers from consuming unbounded memory through headers.&lt;/p&gt;

&lt;h2&gt;
  
  
  Deadline Wire Values Cannot Be Simply Truncated
&lt;/h2&gt;

&lt;p&gt;&lt;code&gt;grpc-timeout&lt;/code&gt; consists of up to eight digits and a unit, such as &lt;code&gt;100m&lt;/code&gt; or &lt;code&gt;2S&lt;/code&gt;. Formatting must use ceiling division to ensure the wire deadline is never shorter than what the caller requested:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="nc"&gt;BigInteger&lt;/span&gt; &lt;span class="n"&gt;amount&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;ceilDivide&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;nanos&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;unit&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;nanos&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
&lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;amount&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;compareTo&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="no"&gt;MAX_VALUE&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;&amp;lt;=&lt;/span&gt; &lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
  &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="n"&gt;amount&lt;/span&gt; &lt;span class="o"&gt;+&lt;/span&gt; &lt;span class="nc"&gt;String&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;valueOf&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;unit&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;suffix&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
&lt;span class="o"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;In the Stage 2 commit covered by this post, the client only writes the timeout header and sets the Reactor Netty response timeout slightly longer to distinguish gRPC deadlines from underlying connection waits. Merely being able to parse the header does not constitute full support; Stage 5 later adds client and server Reactor deadline timers, handler cancellation, and &lt;code&gt;GrpcCallContext&lt;/code&gt; propagation. The current implementation still does not automatically derive shorter deadlines for nested calls on behalf of business code.&lt;/p&gt;

&lt;h2&gt;
  
  
  Unary Cardinality Is Enforced with single()
&lt;/h2&gt;

&lt;p&gt;The client uses &lt;code&gt;requests.single()&lt;/code&gt; for unary requests, and responses likewise require exactly one value. Empty requests, duplicate requests, empty responses, and duplicate responses are not silent defaults — they are explicit errors.&lt;/p&gt;

&lt;p&gt;Current transport test coverage:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;unary message and ASCII metadata round-trip;&lt;/li&gt;
&lt;li&gt;declared error and trailing metadata;&lt;/li&gt;
&lt;li&gt;unknown method;&lt;/li&gt;
&lt;li&gt;empty/duplicate unary data;&lt;/li&gt;
&lt;li&gt;non-gRPC HTTP response and missing status;&lt;/li&gt;
&lt;li&gt;client cancellation propagation to server publisher;&lt;/li&gt;
&lt;li&gt;connection failure mapping to &lt;code&gt;UNAVAILABLE&lt;/code&gt;.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;The unary vertical slice tests continue serving as regression baselines for streaming transport. Server, client, and bidirectional streaming have been completed in subsequent commits; TLS, GOAWAY, keepalive, compression, deadline, connection pool, and transport diagnostics were also completed in Stage 5 and merged to &lt;code&gt;main&lt;/code&gt;.&lt;/p&gt;

&lt;h2&gt;
  
  
  Full Lifecycle of a Client Request
&lt;/h2&gt;

&lt;p&gt;A &lt;code&gt;ReactorGrpcClient&lt;/code&gt; unary call ultimately maps to a complete HTTP/2 stream interaction. Here is the full path from user call to network response:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;client.unary(method, Mono.just(request))
    |
    +-- Flux.defer() -&amp;gt; lazy assembly, independent per subscription
    |
    +-- channel.select() -&amp;gt; pick subchannel (Stage 2: fixed single connection)
    |
    +-- httpClient.post(method.path())
    |   +-- Write HTTP/2 request headers
    |   |   content-type: application/grpc+proto
    |   |   te: trailers
    |   |   grpc-accept-encoding: gzip
    |   |   grpc-timeout: &amp;lt;if set&amp;gt;
    |   |   custom metadata
    |   |
    |   +-- onEventLoop() -&amp;gt; serialize on connection's event loop
    |   |
    |   +-- GrpcOutboundMessageEncoder.encode()
    |   |   +-- marshaller.serialize() -&amp;gt; Protobuf -&amp;gt; ByteBuf
    |   |   +-- validate size &amp;lt;= maxOutboundMessageSize
    |   |   +-- GrpcFrameCodec.encode() -&amp;gt; 5-byte header + payload
    |   |
    |   +-- outbound.send(requestFrames) -&amp;gt; HTTP/2 DATA frame
    |
    +-- .response((response, content) -&amp;gt; decodeResponse)
        +-- Check response headers for grpc-status (trailers-only error)
        +-- Validate HTTP 200 + content-type = application/grpc
        +-- Read grpc-encoding header
        +-- GrpcFrameCodec.decode(content) -&amp;gt; incremental frame parsing
        +-- marshaller.deserialize() -&amp;gt; ByteBuf -&amp;gt; Response protobuf
        +-- Read grpc-status + grpc-message from trailers
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Key design: the entire pipeline is lazy — no HTTP/2 connection is initiated until the &lt;code&gt;Mono&lt;/code&gt; is subscribed to. This differs from grpc-java's &lt;code&gt;ClientCall.start()&lt;/code&gt; which triggers connection immediately.&lt;/p&gt;

&lt;h2&gt;
  
  
  Why Serialize on the Event Loop
&lt;/h2&gt;

&lt;p&gt;Each Reactor Netty HTTP/2 connection is bound to a single Netty event loop thread. If marshalling happens on the user thread, the serialized ByteBuf must be transferred across threads to the event loop, incurring synchronization overhead.&lt;/p&gt;

&lt;p&gt;&lt;code&gt;onEventLoop()&lt;/code&gt; defers serialization to the connection's event loop:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="kd"&gt;private&lt;/span&gt; &lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="no"&gt;T&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&lt;/span&gt; &lt;span class="nc"&gt;Flux&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="no"&gt;T&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&lt;/span&gt; &lt;span class="nf"&gt;onEventLoop&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;NettyOutbound&lt;/span&gt; &lt;span class="n"&gt;outbound&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="nc"&gt;Flux&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="no"&gt;T&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&lt;/span&gt; &lt;span class="n"&gt;source&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
    &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="n"&gt;source&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;publishOn&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;
        &lt;span class="nc"&gt;Schedulers&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;fromExecutor&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;outbound&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;alloc&lt;/span&gt;&lt;span class="o"&gt;().&lt;/span&gt;&lt;span class="na"&gt;buffer&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="o"&gt;).&lt;/span&gt;&lt;span class="na"&gt;alloc&lt;/span&gt;&lt;span class="o"&gt;()&lt;/span&gt;
            &lt;span class="o"&gt;...&lt;/span&gt; &lt;span class="c1"&gt;// obtain event loop executor&lt;/span&gt;
        &lt;span class="o"&gt;));&lt;/span&gt;
&lt;span class="o"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Benefits:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;Serialization and network writes happen on the same thread, lock-free;&lt;/li&gt;
&lt;li&gt;The ByteBuf allocator is consistent with the connection, reducing cross-arena allocation;&lt;/li&gt;
&lt;li&gt;
&lt;code&gt;concatMap(..., 1)&lt;/code&gt; controls prefetch, ensuring messages are sent in strict order.&lt;/li&gt;
&lt;/ul&gt;

&lt;h2&gt;
  
  
  Error Mapping: From HTTP/2 Exceptions to gRPC Status
&lt;/h2&gt;

&lt;p&gt;Network-layer errors cannot be exposed directly to users — they need to be mapped to gRPC semantics:&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Network Exception&lt;/th&gt;
&lt;th&gt;gRPC Status&lt;/th&gt;
&lt;th&gt;Meaning&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;&lt;code&gt;Http2Exception.StreamException(REFUSED_STREAM)&lt;/code&gt;&lt;/td&gt;
&lt;td&gt;UNAVAILABLE&lt;/td&gt;
&lt;td&gt;Server refused this stream, safe to retry&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;GOAWAY frame&lt;/td&gt;
&lt;td&gt;UNAVAILABLE&lt;/td&gt;
&lt;td&gt;Server is shutting down gracefully&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;Connection closed / reset&lt;/td&gt;
&lt;td&gt;UNAVAILABLE&lt;/td&gt;
&lt;td&gt;Network disconnected&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;HTTP 404 (non-gRPC response)&lt;/td&gt;
&lt;td&gt;UNIMPLEMENTED&lt;/td&gt;
&lt;td&gt;Path doesn't exist (may have passed through a non-gRPC proxy)&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;HTTP 503&lt;/td&gt;
&lt;td&gt;UNAVAILABLE&lt;/td&gt;
&lt;td&gt;Service temporarily unavailable&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;Missing &lt;code&gt;grpc-status&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;INTERNAL&lt;/td&gt;
&lt;td&gt;Protocol violation&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;
&lt;code&gt;content-type&lt;/code&gt; mismatch&lt;/td&gt;
&lt;td&gt;INTERNAL&lt;/td&gt;
&lt;td&gt;Response is not gRPC&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;This mapping follows the &lt;a href="https://github.com/grpc/grpc/blob/master/doc/PROTOCOL-HTTP2.md" rel="noopener noreferrer"&gt;HTTP status code to gRPC status conversion table&lt;/a&gt; defined in the gRPC protocol spec. &lt;code&gt;UNAVAILABLE&lt;/code&gt; means the client can retry; &lt;code&gt;INTERNAL&lt;/code&gt; and &lt;code&gt;UNIMPLEMENTED&lt;/code&gt; are generally not worth retrying.&lt;/p&gt;

&lt;h2&gt;
  
  
  How Client Cancellation Propagates to Server
&lt;/h2&gt;

&lt;p&gt;Reactor's &lt;code&gt;dispose()&lt;/code&gt; needs to be mapped to HTTP/2's RST_STREAM, letting the server know this call is no longer needed:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;client: subscription.dispose()
    -&amp;gt; Reactor Netty observes subscriber cancel
    -&amp;gt; Sends HTTP/2 RST_STREAM frame
    -&amp;gt; Server connection receives RST_STREAM
    -&amp;gt; request body publisher is cancelled
    -&amp;gt; server handler's Mono/Flux receives cancel signal
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;The server side listens for stream close events through a connection handler and converts them into a &lt;code&gt;peerCancellation&lt;/code&gt; signal. The business handler uses &lt;code&gt;takeUntilOther(peerCancellation)&lt;/code&gt; to subscribe to this signal, ensuring the server stops unnecessary computation after client cancellation.&lt;/p&gt;

&lt;p&gt;This mechanism naturally aligns Reactor's backpressure semantics (cancel = I no longer need data) with gRPC's cancellation semantics (RST_STREAM = call aborted).&lt;/p&gt;

&lt;h2&gt;
  
  
  Interoperability Verification with grpc-java
&lt;/h2&gt;

&lt;p&gt;Stage 2's exit criteria are not "calling ourselves successfully" — it must interoperate bidirectionally with grpc-java. &lt;code&gt;UnaryInteroperabilityTest&lt;/code&gt; verifies both directions:&lt;/p&gt;

&lt;p&gt;&lt;strong&gt;Reactor client -&amp;gt; grpc-java server:&lt;/strong&gt;&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="nd"&gt;@Test&lt;/span&gt;
&lt;span class="kt"&gt;void&lt;/span&gt; &lt;span class="nf"&gt;reactorClientCallsGrpcJavaServer&lt;/span&gt;&lt;span class="o"&gt;()&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
    &lt;span class="k"&gt;try&lt;/span&gt; &lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="kt"&gt;var&lt;/span&gt; &lt;span class="n"&gt;fixture&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nc"&gt;GrpcJavaFixture&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;start&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt;
         &lt;span class="kt"&gt;var&lt;/span&gt; &lt;span class="n"&gt;client&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nc"&gt;ReactorGrpcClient&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;builder&lt;/span&gt;&lt;span class="o"&gt;().&lt;/span&gt;&lt;span class="na"&gt;port&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;fixture&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;port&lt;/span&gt;&lt;span class="o"&gt;()).&lt;/span&gt;&lt;span class="na"&gt;build&lt;/span&gt;&lt;span class="o"&gt;())&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
        &lt;span class="kt"&gt;var&lt;/span&gt; &lt;span class="n"&gt;response&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;stub&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;unary&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;Mono&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;just&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;request&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"hello"&lt;/span&gt;&lt;span class="o"&gt;))).&lt;/span&gt;&lt;span class="na"&gt;block&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt;
        &lt;span class="n"&gt;assertEquals&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"hello"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;response&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;getValue&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;

        &lt;span class="c1"&gt;// Verify error status and trailing metadata propagate correctly&lt;/span&gt;
        &lt;span class="kt"&gt;var&lt;/span&gt; &lt;span class="n"&gt;error&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;assertThrows&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;GrpcException&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;class&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt;
            &lt;span class="o"&gt;()&lt;/span&gt; &lt;span class="o"&gt;-&amp;gt;&lt;/span&gt; &lt;span class="n"&gt;stub&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;unary&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;Mono&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;just&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;request&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"error"&lt;/span&gt;&lt;span class="o"&gt;))).&lt;/span&gt;&lt;span class="na"&gt;block&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;
        &lt;span class="n"&gt;assertEquals&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;GrpcStatus&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;Code&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;INVALID_ARGUMENT&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;error&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;status&lt;/span&gt;&lt;span class="o"&gt;().&lt;/span&gt;&lt;span class="na"&gt;code&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;
        &lt;span class="n"&gt;assertEquals&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"rejected"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;error&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;status&lt;/span&gt;&lt;span class="o"&gt;().&lt;/span&gt;&lt;span class="na"&gt;message&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;
        &lt;span class="n"&gt;assertEquals&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"fixture"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;error&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;trailers&lt;/span&gt;&lt;span class="o"&gt;().&lt;/span&gt;&lt;span class="na"&gt;getLast&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"error-id"&lt;/span&gt;&lt;span class="o"&gt;).&lt;/span&gt;&lt;span class="na"&gt;orElseThrow&lt;/span&gt;&lt;span class="o"&gt;().&lt;/span&gt;&lt;span class="na"&gt;asciiValue&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;
    &lt;span class="o"&gt;}&lt;/span&gt;
&lt;span class="o"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;&lt;strong&gt;grpc-java client -&amp;gt; Reactor server:&lt;/strong&gt;&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="nd"&gt;@Test&lt;/span&gt;
&lt;span class="kt"&gt;void&lt;/span&gt; &lt;span class="nf"&gt;grpcJavaClientCallsReactorServer&lt;/span&gt;&lt;span class="o"&gt;()&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
    &lt;span class="kt"&gt;var&lt;/span&gt; &lt;span class="n"&gt;service&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nc"&gt;InteropTestServiceReactorGrpc&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;bindService&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="nc"&gt;Service&lt;/span&gt;&lt;span class="o"&gt;()&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
        &lt;span class="nd"&gt;@Override&lt;/span&gt;
        &lt;span class="kd"&gt;public&lt;/span&gt; &lt;span class="nc"&gt;Mono&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="nc"&gt;TestResponse&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&lt;/span&gt; &lt;span class="nf"&gt;unary&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;GrpcCallContext&lt;/span&gt; &lt;span class="n"&gt;context&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="nc"&gt;Mono&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="nc"&gt;TestRequest&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&lt;/span&gt; &lt;span class="n"&gt;requests&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
            &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="n"&gt;requests&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;map&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;req&lt;/span&gt; &lt;span class="o"&gt;-&amp;gt;&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
                &lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;req&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;getValue&lt;/span&gt;&lt;span class="o"&gt;().&lt;/span&gt;&lt;span class="na"&gt;equals&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"error"&lt;/span&gt;&lt;span class="o"&gt;))&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
                    &lt;span class="k"&gt;throw&lt;/span&gt; &lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="nf"&gt;GrpcException&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;
                        &lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="nf"&gt;GrpcStatus&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;GrpcStatus&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;Code&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;INVALID_ARGUMENT&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="s"&gt;"rejected"&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt;
                        &lt;span class="nc"&gt;GrpcMetadata&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;builder&lt;/span&gt;&lt;span class="o"&gt;().&lt;/span&gt;&lt;span class="na"&gt;addAscii&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"error-id"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="s"&gt;"reactor"&lt;/span&gt;&lt;span class="o"&gt;).&lt;/span&gt;&lt;span class="na"&gt;build&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;
                &lt;span class="o"&gt;}&lt;/span&gt;
                &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="nc"&gt;TestResponse&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;newBuilder&lt;/span&gt;&lt;span class="o"&gt;().&lt;/span&gt;&lt;span class="na"&gt;setValue&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;req&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;getValue&lt;/span&gt;&lt;span class="o"&gt;()).&lt;/span&gt;&lt;span class="na"&gt;build&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt;
            &lt;span class="o"&gt;});&lt;/span&gt;
        &lt;span class="o"&gt;}&lt;/span&gt;
    &lt;span class="o"&gt;});&lt;/span&gt;

    &lt;span class="k"&gt;try&lt;/span&gt; &lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="kt"&gt;var&lt;/span&gt; &lt;span class="n"&gt;server&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nc"&gt;ReactorGrpcServer&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;builder&lt;/span&gt;&lt;span class="o"&gt;().&lt;/span&gt;&lt;span class="na"&gt;addService&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;service&lt;/span&gt;&lt;span class="o"&gt;).&lt;/span&gt;&lt;span class="na"&gt;start&lt;/span&gt;&lt;span class="o"&gt;())&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
        &lt;span class="kt"&gt;var&lt;/span&gt; &lt;span class="n"&gt;stub&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nc"&gt;InteropTestServiceGrpc&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;newBlockingStub&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;channel&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
        &lt;span class="n"&gt;assertEquals&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"hello"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;stub&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;unary&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;request&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"hello"&lt;/span&gt;&lt;span class="o"&gt;)).&lt;/span&gt;&lt;span class="na"&gt;getValue&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;

        &lt;span class="kt"&gt;var&lt;/span&gt; &lt;span class="n"&gt;error&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;assertThrows&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;StatusRuntimeException&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;class&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt;
            &lt;span class="o"&gt;()&lt;/span&gt; &lt;span class="o"&gt;-&amp;gt;&lt;/span&gt; &lt;span class="n"&gt;stub&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;unary&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;request&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"error"&lt;/span&gt;&lt;span class="o"&gt;)));&lt;/span&gt;
        &lt;span class="n"&gt;assertEquals&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;Status&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;Code&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;INVALID_ARGUMENT&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;error&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;getStatus&lt;/span&gt;&lt;span class="o"&gt;().&lt;/span&gt;&lt;span class="na"&gt;getCode&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;
    &lt;span class="o"&gt;}&lt;/span&gt;
&lt;span class="o"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Both directions verify: normal responses, error status codes, percent-encoding round-trip of error messages, and custom trailing metadata propagation. Only when grpc-java can correctly parse Reactor server's response can we prove the wire protocol implementation is compatible.&lt;/p&gt;

&lt;h2&gt;
  
  
  Stage 2 Exit Criteria
&lt;/h2&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;code&gt;./gradlew clean spotlessCheck test --no-daemon&lt;/code&gt; passes in full&lt;/li&gt;
&lt;li&gt;Reactor client &amp;lt;-&amp;gt; grpc-java server bidirectional unary calls succeed&lt;/li&gt;
&lt;li&gt;grpc-java client &amp;lt;-&amp;gt; Reactor server bidirectional unary calls succeed&lt;/li&gt;
&lt;li&gt;Error status, message, and trailing metadata propagate correctly&lt;/li&gt;
&lt;li&gt;Client cancellation propagates to server via RST_STREAM&lt;/li&gt;
&lt;li&gt;Connection failure maps to UNAVAILABLE&lt;/li&gt;
&lt;li&gt;Unknown path maps to UNIMPLEMENTED&lt;/li&gt;
&lt;li&gt;All Stage 0 and Stage 1 tests remain green&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;The next post enters Stage 3 and Stage 4: implementing server streaming, client streaming, and bidirectional streaming on the same transport primitive, facing new challenges of multi-message backpressure coordination, bounded inbound buffering, and resource cleanup on premature stream termination.&lt;/p&gt;

</description>
      <category>java</category>
      <category>grpc</category>
      <category>eventloop</category>
      <category>netty</category>
    </item>
    <item>
      <title>Building a Leak-Safe gRPC Frame Decoder on Reactor Netty</title>
      <dc:creator>Wang Lee</dc:creator>
      <pubDate>Sat, 08 Aug 2026 06:09:36 +0000</pubDate>
      <link>https://dev.to/qianwj/building-a-leak-safe-grpc-frame-decoder-on-reactor-netty-po7</link>
      <guid>https://dev.to/qianwj/building-a-leak-safe-grpc-frame-decoder-on-reactor-netty-po7</guid>
      <description>&lt;p&gt;This is the second article in my &lt;a href="https://github.com/qianwj/grpc-reactor" rel="noopener noreferrer"&gt;grpc-reactor&lt;/a&gt; series. The &lt;a href="https://dev.to/qianwj/why-build-grpc-directly-on-reactor-netty-2g96"&gt;first article&lt;/a&gt; explains why I chose to build the runtime directly on Reactor Netty and where its compatibility boundary sits. This article moves one layer down into the Stage 1 protocol implementation: the frame decoder that every RPC shape relies on.&lt;/p&gt;

&lt;p&gt;gRPC protobuf messages are not written directly as raw bytes into HTTP/2 DATA frames. Every message starts with a five-byte envelope:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;byte 0      bit 0 indicates compression; bits 1-7 must be zero
bytes 1-4   unsigned big-endian payload length
byte 5..n   protobuf message, or its compressed representation
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Encoding this envelope is straightforward. The difficult part is decoding it without assuming that one input buffer contains one complete frame. HTTP/2, TCP, and Reactor Netty do not promise that buffer boundaries will line up with gRPC message boundaries.&lt;/p&gt;

&lt;p&gt;This post describes the Stage 1 protocol layer. The project has since progressed beyond it, but the ownership and bounded-decoding rules introduced here remain the foundation for the later transport stages.&lt;/p&gt;

&lt;h2&gt;
  
  
  Encoding Must Define Ownership
&lt;/h2&gt;

&lt;p&gt;The contract of &lt;a href="https://github.com/qianwj/grpc-reactor/blob/main/grpc-reactor-protocol/src/main/java/io/github/qianwj/grpc/reactor/protocol/GrpcFrameCodec.java" rel="noopener noreferrer"&gt;&lt;code&gt;GrpcFrameCodec.encode&lt;/code&gt;&lt;/a&gt; is deliberately explicit: the returned frame and the input message have independent lifetimes. Encoding must not move the input reader index or release the input buffer.&lt;/p&gt;

&lt;p&gt;The implementation currently copies the readable bytes into a byte array before applying compression:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="kd"&gt;public&lt;/span&gt; &lt;span class="kd"&gt;static&lt;/span&gt; &lt;span class="nc"&gt;ByteBuf&lt;/span&gt; &lt;span class="nf"&gt;encode&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;
        &lt;span class="nc"&gt;ByteBufAllocator&lt;/span&gt; &lt;span class="n"&gt;allocator&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt;
        &lt;span class="nc"&gt;ByteBuf&lt;/span&gt; &lt;span class="n"&gt;message&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt;
        &lt;span class="nc"&gt;GrpcCompression&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;Codec&lt;/span&gt; &lt;span class="n"&gt;compression&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
    &lt;span class="kt"&gt;boolean&lt;/span&gt; &lt;span class="n"&gt;compressed&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="o"&gt;!&lt;/span&gt;&lt;span class="n"&gt;compression&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;name&lt;/span&gt;&lt;span class="o"&gt;().&lt;/span&gt;&lt;span class="na"&gt;equals&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"identity"&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
    &lt;span class="kt"&gt;byte&lt;/span&gt;&lt;span class="o"&gt;[]&lt;/span&gt; &lt;span class="n"&gt;payload&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="kt"&gt;byte&lt;/span&gt;&lt;span class="o"&gt;[&lt;/span&gt;&lt;span class="n"&gt;message&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;readableBytes&lt;/span&gt;&lt;span class="o"&gt;()];&lt;/span&gt;
    &lt;span class="n"&gt;message&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;getBytes&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;message&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;readerIndex&lt;/span&gt;&lt;span class="o"&gt;(),&lt;/span&gt; &lt;span class="n"&gt;payload&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
    &lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;compressed&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
        &lt;span class="n"&gt;payload&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;compression&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;compress&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;payload&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
    &lt;span class="o"&gt;}&lt;/span&gt;
    &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="n"&gt;allocator&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;buffer&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;GrpcFrameCodec&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;HEADER_SIZE&lt;/span&gt; &lt;span class="o"&gt;+&lt;/span&gt; &lt;span class="n"&gt;payload&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;length&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
            &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;writeByte&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;compressed&lt;/span&gt; &lt;span class="o"&gt;?&lt;/span&gt; &lt;span class="mi"&gt;1&lt;/span&gt; &lt;span class="o"&gt;:&lt;/span&gt; &lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
            &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;writeInt&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;payload&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;length&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
            &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;writeBytes&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;payload&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
&lt;span class="o"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;This is not a zero-copy implementation, and it should not be presented as the fastest possible design. The copy makes the ownership boundary easy to reason about first. If a later optimization uses a slice or a composite buffer, cancellation and exception paths must be re-proven instead of assuming that the old ownership rules still hold.&lt;/p&gt;

&lt;h2&gt;
  
  
  The Decoder Is a Per-Subscription State Machine
&lt;/h2&gt;

&lt;p&gt;&lt;code&gt;decode&lt;/code&gt; uses &lt;code&gt;Flux.defer&lt;/code&gt; so every subscription receives an independent decoder instance:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="nc"&gt;Flux&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;defer&lt;/span&gt;&lt;span class="o"&gt;(()&lt;/span&gt; &lt;span class="o"&gt;-&amp;gt;&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
    &lt;span class="kt"&gt;var&lt;/span&gt; &lt;span class="n"&gt;decoder&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="nc"&gt;Decoder&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;
            &lt;span class="n"&gt;maxWireMessageSize&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt;
            &lt;span class="n"&gt;maxDecompressedMessageSize&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt;
            &lt;span class="n"&gt;compression&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt;
            &lt;span class="n"&gt;maxBufferedBytesPerStream&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
    &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="n"&gt;input&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;concatMap&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nl"&gt;decoder:&lt;/span&gt;&lt;span class="o"&gt;:&lt;/span&gt;&lt;span class="n"&gt;accept&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
            &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;concatWith&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;Flux&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;defer&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nl"&gt;decoder:&lt;/span&gt;&lt;span class="o"&gt;:&lt;/span&gt;&lt;span class="n"&gt;finish&lt;/span&gt;&lt;span class="o"&gt;))&lt;/span&gt;
            &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;publishOn&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;Schedulers&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;immediate&lt;/span&gt;&lt;span class="o"&gt;(),&lt;/span&gt; &lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
&lt;span class="o"&gt;});&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;The state is intentionally small: either the five-byte header is incomplete, or the header is complete and the decoder is collecting a payload of a known length. One input &lt;code&gt;ByteBuf&lt;/code&gt; may contain one byte of a header, or three complete messages back-to-back.&lt;/p&gt;

&lt;p&gt;&lt;code&gt;concatMap(..., 1)&lt;/code&gt; preserves source order and limits the number of source buffers being processed at once. The decoder still has to respect downstream demand when it emits decoded messages. Every source buffer is released in &lt;code&gt;doFinally&lt;/code&gt;, including success, failure, and cancellation:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="kd"&gt;private&lt;/span&gt; &lt;span class="nc"&gt;Flux&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="nc"&gt;ByteBuf&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&lt;/span&gt; &lt;span class="nf"&gt;accept&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;ByteBuf&lt;/span&gt; &lt;span class="n"&gt;source&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
    &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="nc"&gt;Flux&lt;/span&gt;&lt;span class="o"&gt;.&amp;lt;&lt;/span&gt;&lt;span class="nc"&gt;ByteBuf&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&lt;/span&gt;&lt;span class="n"&gt;generate&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;sink&lt;/span&gt; &lt;span class="o"&gt;-&amp;gt;&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
        &lt;span class="k"&gt;try&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
            &lt;span class="k"&gt;while&lt;/span&gt; &lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;source&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;isReadable&lt;/span&gt;&lt;span class="o"&gt;())&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
                &lt;span class="c1"&gt;// read the header, allocate a bounded payload,&lt;/span&gt;
                &lt;span class="c1"&gt;// and emit a complete message when available&lt;/span&gt;
            &lt;span class="o"&gt;}&lt;/span&gt;
            &lt;span class="n"&gt;sink&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;complete&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt;
        &lt;span class="o"&gt;}&lt;/span&gt; &lt;span class="k"&gt;catch&lt;/span&gt; &lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;Throwable&lt;/span&gt; &lt;span class="n"&gt;error&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
            &lt;span class="n"&gt;sink&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;error&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;error&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
        &lt;span class="o"&gt;}&lt;/span&gt;
    &lt;span class="o"&gt;}).&lt;/span&gt;&lt;span class="na"&gt;doFinally&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;ignored&lt;/span&gt; &lt;span class="o"&gt;-&amp;gt;&lt;/span&gt; &lt;span class="n"&gt;source&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;release&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;
&lt;span class="o"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;The complete implementation also tracks an explicit per-stream buffered-byte limit. That limit covers an incomplete header, a partial payload, and bytes still present in the current source buffer.&lt;/p&gt;

&lt;h2&gt;
  
  
  Validate Peer Input Before Allocating
&lt;/h2&gt;

&lt;p&gt;The length field comes from the peer, so it must be validated before allocating a payload array. The decoder first combines the unsigned big-endian bytes in a &lt;code&gt;long&lt;/code&gt;, then checks the wire-size limit:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="kt"&gt;long&lt;/span&gt; &lt;span class="n"&gt;length&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="o"&gt;((&lt;/span&gt;&lt;span class="kt"&gt;long&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;header&lt;/span&gt;&lt;span class="o"&gt;[&lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="o"&gt;]&lt;/span&gt; &lt;span class="o"&gt;&amp;amp;&lt;/span&gt; &lt;span class="mh"&gt;0xff&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;&amp;lt;&amp;lt;&lt;/span&gt; &lt;span class="mi"&gt;24&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;|&lt;/span&gt; &lt;span class="o"&gt;((&lt;/span&gt;&lt;span class="kt"&gt;long&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;header&lt;/span&gt;&lt;span class="o"&gt;[&lt;/span&gt;&lt;span class="mi"&gt;2&lt;/span&gt;&lt;span class="o"&gt;]&lt;/span&gt; &lt;span class="o"&gt;&amp;amp;&lt;/span&gt; &lt;span class="mh"&gt;0xff&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;&amp;lt;&amp;lt;&lt;/span&gt; &lt;span class="mi"&gt;16&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;|&lt;/span&gt; &lt;span class="o"&gt;((&lt;/span&gt;&lt;span class="kt"&gt;long&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;header&lt;/span&gt;&lt;span class="o"&gt;[&lt;/span&gt;&lt;span class="mi"&gt;3&lt;/span&gt;&lt;span class="o"&gt;]&lt;/span&gt; &lt;span class="o"&gt;&amp;amp;&lt;/span&gt; &lt;span class="mh"&gt;0xff&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;&amp;lt;&amp;lt;&lt;/span&gt; &lt;span class="mi"&gt;8&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;|&lt;/span&gt; &lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;header&lt;/span&gt;&lt;span class="o"&gt;[&lt;/span&gt;&lt;span class="mi"&gt;4&lt;/span&gt;&lt;span class="o"&gt;]&lt;/span&gt; &lt;span class="o"&gt;&amp;amp;&lt;/span&gt; &lt;span class="mh"&gt;0xff&lt;/span&gt;&lt;span class="no"&gt;L&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;

&lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;length&lt;/span&gt; &lt;span class="o"&gt;&amp;gt;&lt;/span&gt; &lt;span class="n"&gt;maxWireMessageSize&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
    &lt;span class="k"&gt;throw&lt;/span&gt; &lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="nf"&gt;GrpcProtocolException&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"wire message length exceeds limit"&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
&lt;span class="o"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Wire size and decompressed size are separate limits. A tiny gzip payload can expand into a huge message, so gzip decompression must enforce a second output limit to defend against decompression bombs.&lt;/p&gt;

&lt;p&gt;Only the lowest compression-flag bit is valid. Any reserved bit is a protocol error. A compressed frame received while the negotiated codec is still &lt;code&gt;identity&lt;/code&gt; is also rejected; the decoder must not guess which algorithm the peer intended.&lt;/p&gt;

&lt;h2&gt;
  
  
  Test Every Header Split Point
&lt;/h2&gt;

&lt;p&gt;Testing one arbitrary two-buffer split is not enough. A five-byte header has six representative split positions, including before the first byte and after the complete header. The test suite uses a dynamic test for every split from 0 through 5:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="nc"&gt;IntStream&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;rangeClosed&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="nc"&gt;GrpcFrameCodec&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;HEADER_SIZE&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;mapToObj&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;split&lt;/span&gt; &lt;span class="o"&gt;-&amp;gt;&lt;/span&gt; &lt;span class="nc"&gt;DynamicTest&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;dynamicTest&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;
                &lt;span class="s"&gt;"split after byte "&lt;/span&gt; &lt;span class="o"&gt;+&lt;/span&gt; &lt;span class="n"&gt;split&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt;
                &lt;span class="o"&gt;()&lt;/span&gt; &lt;span class="o"&gt;-&amp;gt;&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
                    &lt;span class="kt"&gt;byte&lt;/span&gt;&lt;span class="o"&gt;[]&lt;/span&gt; &lt;span class="n"&gt;wire&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;wireBytes&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"hello"&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
                    &lt;span class="nc"&gt;Flux&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="nc"&gt;ByteBuf&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&lt;/span&gt; &lt;span class="n"&gt;chunks&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nc"&gt;Flux&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;just&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;
                            &lt;span class="n"&gt;wrapped&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;wire&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;split&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt;
                            &lt;span class="n"&gt;wrapped&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;wire&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;split&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;wire&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;length&lt;/span&gt; &lt;span class="o"&gt;-&lt;/span&gt; &lt;span class="n"&gt;split&lt;/span&gt;&lt;span class="o"&gt;));&lt;/span&gt;
                    &lt;span class="nc"&gt;StepVerifier&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;create&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;GrpcFrameCodec&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;decode&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;chunks&lt;/span&gt;&lt;span class="o"&gt;))&lt;/span&gt;
                            &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;assertNext&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;message&lt;/span&gt; &lt;span class="o"&gt;-&amp;gt;&lt;/span&gt; &lt;span class="n"&gt;assertMessage&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;message&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="s"&gt;"hello"&lt;/span&gt;&lt;span class="o"&gt;))&lt;/span&gt;
                            &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;verifyComplete&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt;
                &lt;span class="o"&gt;}));&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;The same test class covers arbitrary body fragmentation, multiple messages coalesced into one buffer, empty messages, gzip, reserved flags, truncated frames, wire/decompressed limits, and the fact that encoding does not consume the input buffer. See &lt;a href="https://github.com/qianwj/grpc-reactor/blob/main/grpc-reactor-protocol/src/test/java/io/github/qianwj/grpc/reactor/protocol/GrpcFrameCodecTest.java" rel="noopener noreferrer"&gt;&lt;code&gt;GrpcFrameCodecTest&lt;/code&gt;&lt;/a&gt; for the executable cases.&lt;/p&gt;

&lt;h2&gt;
  
  
  Cancellation Must Release Undelivered Data
&lt;/h2&gt;

&lt;p&gt;Suppose one source &lt;code&gt;ByteBuf&lt;/code&gt; contains three messages: &lt;code&gt;one&lt;/code&gt;, &lt;code&gt;two&lt;/code&gt;, and &lt;code&gt;three&lt;/code&gt;. The downstream requests two messages and then cancels. The test must assert not only the values it received, but also that the source buffer was released:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="nc"&gt;StepVerifier&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;create&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;GrpcFrameCodec&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;decode&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;Flux&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;just&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;source&lt;/span&gt;&lt;span class="o"&gt;)),&lt;/span&gt; &lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;thenRequest&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;assertNext&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;message&lt;/span&gt; &lt;span class="o"&gt;-&amp;gt;&lt;/span&gt; &lt;span class="n"&gt;assertMessage&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;message&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="s"&gt;"one"&lt;/span&gt;&lt;span class="o"&gt;))&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;thenRequest&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;assertNext&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;message&lt;/span&gt; &lt;span class="o"&gt;-&amp;gt;&lt;/span&gt; &lt;span class="n"&gt;assertMessage&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;message&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="s"&gt;"two"&lt;/span&gt;&lt;span class="o"&gt;))&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;thenCancel&lt;/span&gt;&lt;span class="o"&gt;()&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;verify&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt;

&lt;span class="n"&gt;assertEquals&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;source&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;refCnt&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;That assertion is more important than a happy-path content check. Network code often behaves correctly under normal completion; leaks tend to appear during cancellation, size-limit failures, truncated frames, or competing terminal signals.&lt;/p&gt;

&lt;p&gt;Cancellation can also arrive before a complete message exists. In that case there is no decoded value for the subscriber to release, so the decoder itself must release the partially accumulated source buffer:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="nd"&gt;@Test&lt;/span&gt;
&lt;span class="kt"&gt;void&lt;/span&gt; &lt;span class="nf"&gt;releasesPartialFrameInputWhenCancelled&lt;/span&gt;&lt;span class="o"&gt;()&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
    &lt;span class="nc"&gt;ByteBuf&lt;/span&gt; &lt;span class="n"&gt;partial&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nc"&gt;Unpooled&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;buffer&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;8&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
            &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;writeByte&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
            &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;writeInt&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;16&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
            &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;writeBytes&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="kt"&gt;byte&lt;/span&gt;&lt;span class="o"&gt;[]{&lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="mi"&gt;2&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="mi"&gt;3&lt;/span&gt;&lt;span class="o"&gt;});&lt;/span&gt;

    &lt;span class="nc"&gt;StepVerifier&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;create&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;
                    &lt;span class="nc"&gt;GrpcFrameCodec&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;decode&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;
                            &lt;span class="nc"&gt;Flux&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;just&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;partial&lt;/span&gt;&lt;span class="o"&gt;).&lt;/span&gt;&lt;span class="na"&gt;concatWith&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;Flux&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;never&lt;/span&gt;&lt;span class="o"&gt;())),&lt;/span&gt;
                    &lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
            &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;thenRequest&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
            &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;thenAwait&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;java&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;time&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;Duration&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;ofMillis&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;10&lt;/span&gt;&lt;span class="o"&gt;))&lt;/span&gt;
            &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;thenCancel&lt;/span&gt;&lt;span class="o"&gt;()&lt;/span&gt;
            &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;verify&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt;

    &lt;span class="n"&gt;assertEquals&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;partial&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;refCnt&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;
&lt;span class="o"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;The full executable case is &lt;a href="https://github.com/qianwj/grpc-reactor/blob/main/grpc-reactor-protocol/src/test/java/io/github/qianwj/grpc/reactor/protocol/GrpcFrameCodecTest.java" rel="noopener noreferrer"&gt;&lt;code&gt;releasesPartialFrameInputWhenCancelled&lt;/code&gt;&lt;/a&gt;. It covers the lifecycle edge that a normal decode-complete test cannot exercise.&lt;/p&gt;

&lt;p&gt;Run only the frame codec suite from the repository root:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight shell"&gt;&lt;code&gt;./gradlew :grpc-reactor-protocol:test &lt;span class="se"&gt;\&lt;/span&gt;
  &lt;span class="nt"&gt;--tests&lt;/span&gt; io.github.qianwj.grpc.reactor.protocol.GrpcFrameCodecTest &lt;span class="se"&gt;\&lt;/span&gt;
  &lt;span class="nt"&gt;--no-daemon&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;On JDK 25, the Gradle build and generated JUnit report produced:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;GrpcFrameCodecTest: 15 tests, 0 failures, 0 errors, 0 skipped
BUILD SUCCESSFUL in 3s
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;At the Stage 1 boundary, the frame decoder verifies message-level demand but does not yet implement the two-level flow-control problem of the streaming transport. Reactive Streams counts messages, while HTTP/2 flow control counts bytes. They cannot be treated as the same quantity. Stage 3 later adds bounded inbound buffering and demand-aware delivery, and Stage 4 extends those rules to bidirectional streaming. Those transport and stress tests are covered in later posts.&lt;/p&gt;

&lt;h2&gt;
  
  
  Metadata: Ordering, Duplicates, and Binary Values
&lt;/h2&gt;

&lt;p&gt;gRPC metadata is not a simple &lt;code&gt;Map&amp;lt;String, String&amp;gt;&lt;/code&gt;. It must satisfy all of these rules:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;The same key may occur more than once, and insertion order matters.&lt;/li&gt;
&lt;li&gt;Keys may contain only lowercase letters, digits, &lt;code&gt;_&lt;/code&gt;, &lt;code&gt;.&lt;/code&gt;, and &lt;code&gt;-&lt;/code&gt;, validated by &lt;code&gt;[0-9a-z_.-]+&lt;/code&gt;.&lt;/li&gt;
&lt;li&gt;Keys ending in &lt;code&gt;-bin&lt;/code&gt; carry binary values and use unpadded Base64 on the wire.&lt;/li&gt;
&lt;li&gt;Applications cannot set reserved fields such as &lt;code&gt;content-type&lt;/code&gt;, &lt;code&gt;te&lt;/code&gt;, &lt;code&gt;grpc-status&lt;/code&gt;, or &lt;code&gt;grpc-timeout&lt;/code&gt;.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;&lt;code&gt;GrpcMetadata&lt;/code&gt; stores an immutable entry list so it can be safely shared across asynchronous boundaries:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="nc"&gt;GrpcMetadata&lt;/span&gt; &lt;span class="n"&gt;metadata&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nc"&gt;GrpcMetadata&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;builder&lt;/span&gt;&lt;span class="o"&gt;()&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;addAscii&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"trace-id"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="s"&gt;"abc123"&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;addAscii&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"trace-id"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="s"&gt;"def456"&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="c1"&gt;// duplicates are allowed&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;addBinary&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"auth-token-bin"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;tokenBytes&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;build&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt;

&lt;span class="c1"&gt;// Order is preserved when reading.&lt;/span&gt;
&lt;span class="nc"&gt;List&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&lt;/span&gt;&lt;span class="nc"&gt;GrpcMetadata&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;Entry&lt;/span&gt;&lt;span class="o"&gt;&amp;gt;&lt;/span&gt; &lt;span class="n"&gt;all&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;metadata&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;getAll&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"trace-id"&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt; &lt;span class="c1"&gt;// [abc123, def456]&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;When parsing HTTP/2 headers, a binary value may be comma-joined by header handling. The implementation splits it on commas and decodes each Base64 segment independently. The total encoded size is bounded at 8 KiB by default, preventing a peer from exhausting memory with oversized headers.&lt;/p&gt;

&lt;h2&gt;
  
  
  Status: 17 Codes and Percent-Encoding
&lt;/h2&gt;

&lt;p&gt;gRPC defines &lt;a href="https://github.com/grpc/grpc/blob/HEAD/doc/statuscodes.md" rel="noopener noreferrer"&gt;17 standard status codes&lt;/a&gt;, each with a specific meaning for client error handling and future retry policies. &lt;code&gt;GrpcStatus&lt;/code&gt; is a record containing a code and a human-readable message:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="kd"&gt;public&lt;/span&gt; &lt;span class="n"&gt;record&lt;/span&gt; &lt;span class="nf"&gt;GrpcStatus&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;Code&lt;/span&gt; &lt;span class="n"&gt;code&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="nc"&gt;String&lt;/span&gt; &lt;span class="n"&gt;message&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
    &lt;span class="kd"&gt;public&lt;/span&gt; &lt;span class="kd"&gt;enum&lt;/span&gt; &lt;span class="nc"&gt;Code&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
        &lt;span class="no"&gt;OK&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt; &lt;span class="no"&gt;CANCELLED&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt; &lt;span class="no"&gt;UNKNOWN&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;2&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt; &lt;span class="no"&gt;INVALID_ARGUMENT&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;3&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt;
        &lt;span class="no"&gt;DEADLINE_EXCEEDED&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;4&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt; &lt;span class="no"&gt;NOT_FOUND&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;5&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt; &lt;span class="no"&gt;ALREADY_EXISTS&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;6&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt;
        &lt;span class="no"&gt;PERMISSION_DENIED&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;7&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt; &lt;span class="no"&gt;RESOURCE_EXHAUSTED&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;8&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt;
        &lt;span class="no"&gt;FAILED_PRECONDITION&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;9&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt; &lt;span class="no"&gt;ABORTED&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;10&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt; &lt;span class="no"&gt;OUT_OF_RANGE&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;11&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt;
        &lt;span class="no"&gt;UNIMPLEMENTED&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;12&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt; &lt;span class="no"&gt;INTERNAL&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;13&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt; &lt;span class="no"&gt;UNAVAILABLE&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;14&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt;
        &lt;span class="no"&gt;DATA_LOSS&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;15&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt; &lt;span class="no"&gt;UNAUTHENTICATED&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;16&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
    &lt;span class="o"&gt;}&lt;/span&gt;
&lt;span class="o"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Several details matter in practice:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;code&gt;DEADLINE_EXCEEDED&lt;/code&gt; may be returned even after the operation completed successfully. If the successful response crosses the deadline in transit, the client can still observe a timeout.&lt;/li&gt;
&lt;li&gt;
&lt;code&gt;UNAVAILABLE&lt;/code&gt; indicates a transient failure for which a client may later retry safely; &lt;code&gt;INTERNAL&lt;/code&gt; generally describes a server-side bug and should not be blindly retried.&lt;/li&gt;
&lt;li&gt;
&lt;code&gt;UNIMPLEMENTED&lt;/code&gt; carries the semantic meaning of an unsupported method, commonly surfaced through an HTTP 404 response at the protocol boundary.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;Text in the &lt;code&gt;grpc-message&lt;/code&gt; trailer uses percent-encoding: printable ASCII characters other than &lt;code&gt;%&lt;/code&gt; can pass through, while other bytes become &lt;code&gt;%HH&lt;/code&gt;. This allows UTF-8 error descriptions to travel through ASCII HTTP/2 headers safely.&lt;/p&gt;

&lt;p&gt;Unknown numeric status codes are mapped to &lt;code&gt;UNKNOWN&lt;/code&gt; instead of causing a parse failure. That preserves forward compatibility when a peer adopts a newer gRPC specification.&lt;/p&gt;

&lt;h2&gt;
  
  
  Timeout: Eight Digits and a Unit
&lt;/h2&gt;

&lt;p&gt;The gRPC &lt;code&gt;grpc-timeout&lt;/code&gt; header carries a relative duration, not an absolute timestamp. By the time the server receives a request, part of the caller's original time budget has already been consumed by transport latency.&lt;/p&gt;

&lt;p&gt;The wire format is compact: at most eight decimal digits followed by a unit suffix:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;100m       -&amp;gt; 100 milliseconds
2S         -&amp;gt; 2 seconds
99999999H  -&amp;gt; roughly 11,415 years (the maximum value)
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;The six units are &lt;code&gt;H&lt;/code&gt; (hours), &lt;code&gt;M&lt;/code&gt; (minutes), &lt;code&gt;S&lt;/code&gt; (seconds), &lt;code&gt;m&lt;/code&gt; (milliseconds), &lt;code&gt;u&lt;/code&gt; (microseconds), and &lt;code&gt;n&lt;/code&gt; (nanoseconds).&lt;/p&gt;

&lt;p&gt;When formatting a &lt;code&gt;Duration&lt;/code&gt;, the implementation rounds upward using ceiling division. The encoded deadline must never be shorter than the caller's requested duration:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="nc"&gt;BigInteger&lt;/span&gt; &lt;span class="n"&gt;amount&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;nanos&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;add&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;unitNanos&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;subtract&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;BigInteger&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;ONE&lt;/span&gt;&lt;span class="o"&gt;))&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;divide&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;unitNanos&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt; &lt;span class="c1"&gt;// ceiling division&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;&lt;code&gt;BigInteger&lt;/code&gt; avoids overflow during nanosecond arithmetic. The formatter scans from nanoseconds upward and selects the first unit whose value fits within 99,999,999.&lt;/p&gt;

&lt;h2&gt;
  
  
  Compression: Identity by Default, Gzip Built In
&lt;/h2&gt;

&lt;p&gt;&lt;code&gt;GrpcCompression&lt;/code&gt; manages codec registration and negotiation:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="nc"&gt;GrpcCompression&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;Registry&lt;/span&gt; &lt;span class="n"&gt;registry&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nc"&gt;GrpcCompression&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;Registry&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;builder&lt;/span&gt;&lt;span class="o"&gt;()&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;add&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;GrpcCompression&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;GZIP&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;build&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt;

&lt;span class="c1"&gt;// Produces grpc-accept-encoding: gzip&lt;/span&gt;
&lt;span class="nc"&gt;String&lt;/span&gt; &lt;span class="n"&gt;advertised&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;registry&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;advertisedEncodings&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;&lt;code&gt;identity&lt;/code&gt; is always implicit and appears first in the registry. The codec interface has only three operations: a name, byte-array compression, and bounded byte-array decompression.&lt;/p&gt;

&lt;p&gt;Gzip decompression reads in 8 KiB chunks and uses &lt;code&gt;Math.addExact()&lt;/code&gt; while accumulating the output size. It can stop immediately after exceeding &lt;code&gt;maxDecompressedSize&lt;/code&gt;, and arithmetic overflow cannot silently wrap the counter. A 100-byte gzip payload can expand to gigabytes, so decompression-bomb protection is part of the protocol contract rather than an optional optimization.&lt;/p&gt;

&lt;h2&gt;
  
  
  GrpcMethod: One Description for Four Cardinalities
&lt;/h2&gt;

&lt;p&gt;Each RPC is described by one &lt;code&gt;GrpcMethod&lt;/code&gt; record:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="kt"&gt;var&lt;/span&gt; &lt;span class="n"&gt;method&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="nc"&gt;GrpcMethod&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&amp;gt;(&lt;/span&gt;
        &lt;span class="s"&gt;"testing.InteropTestService"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="c1"&gt;// full service name&lt;/span&gt;
        &lt;span class="s"&gt;"Unary"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt;                      &lt;span class="c1"&gt;// method name&lt;/span&gt;
        &lt;span class="nc"&gt;GrpcMethod&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;Cardinality&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;UNARY&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt;
        &lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="nc"&gt;ProtobufMarshaller&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&amp;gt;(&lt;/span&gt;&lt;span class="nc"&gt;TestRequest&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;parser&lt;/span&gt;&lt;span class="o"&gt;()),&lt;/span&gt;
        &lt;span class="k"&gt;new&lt;/span&gt; &lt;span class="nc"&gt;ProtobufMarshaller&lt;/span&gt;&lt;span class="o"&gt;&amp;lt;&amp;gt;(&lt;/span&gt;&lt;span class="nc"&gt;TestResponse&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;parser&lt;/span&gt;&lt;span class="o"&gt;()));&lt;/span&gt;

&lt;span class="n"&gt;method&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;path&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt;           &lt;span class="c1"&gt;// /testing.InteropTestService/Unary&lt;/span&gt;
&lt;span class="n"&gt;method&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;fullMethodName&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt; &lt;span class="c1"&gt;// testing.InteropTestService/Unary&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;The &lt;code&gt;Cardinality&lt;/code&gt; enum exposes &lt;code&gt;singleRequest()&lt;/code&gt; and &lt;code&gt;singleResponse()&lt;/code&gt;. The transport uses those flags to insert &lt;code&gt;single()&lt;/code&gt; at the API boundary, turning cardinality violations into explicit errors instead of silently dropping values.&lt;/p&gt;

&lt;h2&gt;
  
  
  ProtobufMarshaller: Serialization Does Not Own the Input
&lt;/h2&gt;

&lt;p&gt;&lt;code&gt;ProtobufMarshaller&lt;/code&gt; wraps a protobuf &lt;code&gt;Parser&amp;lt;T&amp;gt;&lt;/code&gt;:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="kd"&gt;public&lt;/span&gt; &lt;span class="nc"&gt;ByteBuf&lt;/span&gt; &lt;span class="nf"&gt;serialize&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;ByteBufAllocator&lt;/span&gt; &lt;span class="n"&gt;allocator&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="no"&gt;T&lt;/span&gt; &lt;span class="n"&gt;value&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
    &lt;span class="kt"&gt;var&lt;/span&gt; &lt;span class="n"&gt;result&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;allocator&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;buffer&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;value&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;getSerializedSize&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;
    &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="n"&gt;result&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;writeBytes&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;value&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;toByteArray&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;
&lt;span class="o"&gt;}&lt;/span&gt;

&lt;span class="kd"&gt;public&lt;/span&gt; &lt;span class="no"&gt;T&lt;/span&gt; &lt;span class="nf"&gt;deserialize&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;ByteBuf&lt;/span&gt; &lt;span class="n"&gt;message&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
    &lt;span class="nc"&gt;ByteBuffer&lt;/span&gt; &lt;span class="n"&gt;bytes&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;message&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;nioBuffer&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;message&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;readerIndex&lt;/span&gt;&lt;span class="o"&gt;(),&lt;/span&gt; &lt;span class="n"&gt;message&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;readableBytes&lt;/span&gt;&lt;span class="o"&gt;());&lt;/span&gt;
    &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="n"&gt;parser&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;parseFrom&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;bytes&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
&lt;span class="o"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;The important contract is that &lt;code&gt;deserialize&lt;/code&gt; reads through a NIO &lt;code&gt;ByteBuffer&lt;/code&gt; view. It does not move the input reader index and does not release the input. Ownership remains with the caller, allowing the frame decoder's &lt;code&gt;doFinally(release)&lt;/code&gt; to manage the buffer lifecycle uniformly whether parsing succeeds or fails.&lt;/p&gt;

&lt;h2&gt;
  
  
  Stage 1 Exit Criteria
&lt;/h2&gt;

&lt;p&gt;The protocol module does not need Reactor Netty on its classpath. That dependency boundary is itself part of the verification. The Stage 1 exit criteria are:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;every header and payload split point decodes correctly;&lt;/li&gt;
&lt;li&gt;cancellation leaves no unreleased &lt;code&gt;ByteBuf&lt;/code&gt; (&lt;code&gt;refCnt&lt;/code&gt; assertions);&lt;/li&gt;
&lt;li&gt;metadata preserves order, supports binary values, and enforces its size limit;&lt;/li&gt;
&lt;li&gt;unknown status codes do not throw;&lt;/li&gt;
&lt;li&gt;all timeout units parse correctly, formatting rounds upward, and arithmetic is overflow-safe;&lt;/li&gt;
&lt;li&gt;gzip decompression limits reject bombs;&lt;/li&gt;
&lt;li&gt;the marshaller does not change the input buffer state.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;These tests use JUnit 5 &lt;code&gt;@TestFactory&lt;/code&gt; and &lt;code&gt;DynamicTest&lt;/code&gt; for parameterization, together with Reactor Test's &lt;code&gt;StepVerifier&lt;/code&gt; for asynchronous behavior. The protocol layer is the reason the later transport stages can focus on HTTP/2 lifecycle instead of rediscovering framing and ownership rules.&lt;/p&gt;

&lt;p&gt;The complete implementation is in &lt;a href="https://github.com/qianwj/grpc-reactor/blob/main/grpc-reactor-protocol/src/main/java/io/github/qianwj/grpc/reactor/protocol/GrpcFrameCodec.java" rel="noopener noreferrer"&gt;&lt;code&gt;GrpcFrameCodec.java&lt;/code&gt;&lt;/a&gt;. The related protocol primitives are &lt;a href="https://github.com/qianwj/grpc-reactor/blob/main/grpc-reactor-protocol/src/main/java/io/github/qianwj/grpc/reactor/protocol/GrpcMetadata.java" rel="noopener noreferrer"&gt;&lt;code&gt;GrpcMetadata&lt;/code&gt;&lt;/a&gt;, &lt;a href="https://github.com/qianwj/grpc-reactor/blob/main/grpc-reactor-protocol/src/main/java/io/github/qianwj/grpc/reactor/protocol/GrpcStatus.java" rel="noopener noreferrer"&gt;&lt;code&gt;GrpcStatus&lt;/code&gt;&lt;/a&gt;, &lt;a href="https://github.com/qianwj/grpc-reactor/blob/main/grpc-reactor-protocol/src/main/java/io/github/qianwj/grpc/reactor/protocol/GrpcTimeout.java" rel="noopener noreferrer"&gt;&lt;code&gt;GrpcTimeout&lt;/code&gt;&lt;/a&gt;, &lt;a href="https://github.com/qianwj/grpc-reactor/blob/main/grpc-reactor-protocol/src/main/java/io/github/qianwj/grpc/reactor/protocol/GrpcCompression.java" rel="noopener noreferrer"&gt;&lt;code&gt;GrpcCompression&lt;/code&gt;&lt;/a&gt;, and &lt;a href="https://github.com/qianwj/grpc-reactor/blob/main/grpc-reactor-protocol/src/main/java/io/github/qianwj/grpc/reactor/protocol/ProtobufMarshaller.java" rel="noopener noreferrer"&gt;&lt;code&gt;ProtobufMarshaller&lt;/code&gt;&lt;/a&gt;.&lt;/p&gt;

</description>
      <category>grpc</category>
      <category>netty</category>
      <category>java</category>
      <category>bytebuf</category>
    </item>
    <item>
      <title>Why Build gRPC Directly on Reactor Netty</title>
      <dc:creator>Wang Lee</dc:creator>
      <pubDate>Fri, 07 Aug 2026 14:23:22 +0000</pubDate>
      <link>https://dev.to/qianwj/why-build-grpc-directly-on-reactor-netty-2g96</link>
      <guid>https://dev.to/qianwj/why-build-grpc-directly-on-reactor-netty-2g96</guid>
      <description>&lt;p&gt;I recently started building &lt;a href="https://github.com/qianwj/grpc-reactor" rel="noopener noreferrer"&gt;grpc-reactor&lt;/a&gt;: an experimental gRPC implementation built directly on Reactor Netty HTTP/2. The project uses Protobuf but has zero runtime dependency on grpc-java transport, &lt;code&gt;ClientCall&lt;/code&gt;, &lt;code&gt;ServerCall&lt;/code&gt;, or &lt;code&gt;StreamObserver&lt;/code&gt;.&lt;/p&gt;

&lt;p&gt;This is not about proving grpc-java is bad. grpc-java is mature, stable, and covers load balancing, NameResolver, retries, and rich observability. This project answers a different question: if your programming model is already Reactor, can &lt;code&gt;Mono&lt;/code&gt; and &lt;code&gt;Flux&lt;/code&gt; flow all the way from the generated API down to the HTTP/2 stream, without adapting between two async abstractions?&lt;/p&gt;

&lt;p&gt;This post is based on JDK 25, Gradle 9.2.1, Reactor 3.8.6, Reactor Netty 1.3.6, and Netty 4.2.15.Final. The &lt;code&gt;main&lt;/code&gt; branch has completed Stage 0 through Stage 10: beyond protocol, four RPC cardinalities, production transport, DNS, codegen, standard services, and Stage 9 hardening, Stage 10 adds an optional load-balancing extensions artifact for grpclb, RLS, load reporting, and ORCA.&lt;/p&gt;

&lt;h2&gt;
  
  
  The API Goal Is Not Wrapping StreamObserver
&lt;/h2&gt;

&lt;p&gt;Protobuf methods have four cardinalities. The target API maps them directly:&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Protobuf Method&lt;/th&gt;
&lt;th&gt;Reactor Signature&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;unary&lt;/td&gt;
&lt;td&gt;&lt;code&gt;Mono&amp;lt;Resp&amp;gt; method(Mono&amp;lt;Req&amp;gt;)&lt;/code&gt;&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;server streaming&lt;/td&gt;
&lt;td&gt;&lt;code&gt;Flux&amp;lt;Resp&amp;gt; method(Mono&amp;lt;Req&amp;gt;)&lt;/code&gt;&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;client streaming&lt;/td&gt;
&lt;td&gt;&lt;code&gt;Mono&amp;lt;Resp&amp;gt; method(Flux&amp;lt;Req&amp;gt;)&lt;/code&gt;&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;bidirectional streaming&lt;/td&gt;
&lt;td&gt;&lt;code&gt;Flux&amp;lt;Resp&amp;gt; method(Flux&amp;lt;Req&amp;gt;)&lt;/code&gt;&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;Why does gRPC define four instead of just unary? It's essentially a 2x2 combination where request and response each independently choose "single value or stream":&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;&lt;/th&gt;
&lt;th&gt;Single Response&lt;/th&gt;
&lt;th&gt;Streaming Response&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;Single Request&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;unary&lt;/td&gt;
&lt;td&gt;server streaming&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;strong&gt;Streaming Request&lt;/strong&gt;&lt;/td&gt;
&lt;td&gt;client streaming&lt;/td&gt;
&lt;td&gt;bidirectional&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;They solve different scenarios: unary covers classic request-response; server streaming handles server push (event subscriptions, large paginated pulls); client streaming handles bulk uploads (file chunks, batch writes); bidirectional streaming handles real-time two-way communication (chat, collaborative editing). These aren't invented patterns — HTTP/2 streams are inherently full-duplex. A single connection can multiplex hundreds of concurrent streams, with both request and response sending multiple data frames independently. gRPC elevates this transport capability to first-class API semantics: not a variant of chunked transfer encoding, but compiler-level type checking that callers and implementers match.&lt;/p&gt;

&lt;p&gt;The low-level transport retains a single unified model:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;Flux&amp;lt;Req&amp;gt; -&amp;gt; HTTP/2 stream -&amp;gt; Flux&amp;lt;Resp&amp;gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Generated code applies cardinality checks like &lt;code&gt;single()&lt;/code&gt; at the boundary. This avoids implementing four separate network logic paths for four RPC types, while explicitly rejecting empty requests, duplicate requests, or duplicate responses in unary calls.&lt;/p&gt;

&lt;h2&gt;
  
  
  Compatibility Targets the Wire Protocol
&lt;/h2&gt;

&lt;p&gt;The project's compatibility target is the public &lt;a href="https://github.com/grpc/grpc/blob/master/doc/PROTOCOL-HTTP2.md" rel="noopener noreferrer"&gt;gRPC over HTTP/2 protocol spec&lt;/a&gt;, not grpc-java internal APIs:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;generated Reactor API
        |
call dispatcher / service registry
        |
marshaller + message framer/deframer
        |
metadata + status + deadline
        |
Reactor Netty HTTP/2 stream
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Each RPC corresponds to one HTTP/2 stream. Requests start with a HEADERS frame containing at minimum:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;:method      POST
:path        /&amp;lt;package.Service&amp;gt;/&amp;lt;Method&amp;gt;
content-type application/grpc+proto
te           trailers
grpc-timeout &amp;lt;relative timeout, e.g. 100m for 100ms&amp;gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Message bodies use Length-Prefixed-Message format encapsulated in DATA frames — each Protobuf message is preceded by a five-byte envelope (1 byte compression flag + 4 bytes big-endian length). A single DATA frame may contain multiple gRPC messages, and a large message may span multiple DATA frames.&lt;/p&gt;

&lt;p&gt;&lt;strong&gt;Why does gRPC require HTTP trailers?&lt;/strong&gt; This is one of the protocol's most counterintuitive designs. HTTP status codes are nearly useless for gRPC — the protocol requires the HTTP layer to always return 200, with the actual call result (&lt;code&gt;grpc-status&lt;/code&gt; and &lt;code&gt;grpc-message&lt;/code&gt;) placed in trailers. The reason: in streaming scenarios, the server may have already sent thousands of messages, and whether it ultimately succeeded or failed can only be determined after processing the last piece of data. HTTP headers are sent before the body and cannot carry this posterior result. Trailers are the only mechanism in HTTP/2 that can append metadata after the body.&lt;/p&gt;

&lt;p&gt;The protocol also defines &lt;strong&gt;trailers-only&lt;/strong&gt; mode: when the server can determine failure before reading the body (e.g., path not found, authentication failed), &lt;code&gt;grpc-status&lt;/code&gt; is returned directly in response headers without sending a body, saving one round-trip. Clients must check both headers and trailers to correctly extract the final status.&lt;/p&gt;

&lt;h2&gt;
  
  
  Modules Split Along Protocol Boundaries
&lt;/h2&gt;

&lt;p&gt;The project currently has nine modules:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;grpc-reactor-protocol      → Transport-independent protocol primitives
grpc-reactor-transport     → Reactor Netty HTTP/2 mapping
grpc-reactor-codegen       → Protoc plugin, generates Reactor stubs
grpc-reactor-gradle-plugin → Gradle integration
grpc-reactor-maven-plugin  → Maven generate-sources integration
grpc-reactor-services      → Optional Health, Reflection &amp;amp; Channelz standard services
grpc-reactor-binlog        → Optional, bounded canonical binary logging
grpc-reactor-lb-extensions → Optional grpclb, RLS, load reporting &amp;amp; ORCA
grpc-reactor-interop-test  → grpc-java compatibility tests
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;&lt;code&gt;protocol&lt;/code&gt; handles only transport-independent values and codecs: message framing, metadata, status, timeout, compression, and protobuf marshallers. It depends on Reactor Core and Netty Buffer but not Reactor Netty — meaning the protocol layer can be tested independently without real HTTP/2 connections.&lt;/p&gt;

&lt;p&gt;&lt;code&gt;transport&lt;/code&gt; maps the protocol onto Reactor Netty HTTP/2, handling client, server, service registry, and per-call context.&lt;/p&gt;

&lt;p&gt;&lt;code&gt;codegen&lt;/code&gt; consumes Protobuf &lt;code&gt;CodeGeneratorRequest&lt;/code&gt; and generates type-safe Reactor client and service binders.&lt;/p&gt;

&lt;p&gt;&lt;code&gt;services&lt;/code&gt; uses the same codegen to generate canonical gRPC service bindings, implementing Health v1, Reflection v1, and Channelz v1 on top of transport's immutable descriptor/diagnostics snapshots. This module requires explicit registration — adding the dependency alone won't expose management endpoints.&lt;/p&gt;

&lt;p&gt;&lt;code&gt;binlog&lt;/code&gt; is also an explicitly-enabled standalone module. It captures canonical binary-log v1 events via a transport interceptor, controlling information exposure and memory limits through metadata/message truncation, sensitive key redaction, and fixed-capacity sinks.&lt;/p&gt;

&lt;p&gt;&lt;code&gt;interop-test&lt;/code&gt; introduces grpc-java in test scope. grpc-java serves as the compatibility oracle here, never entering the project runtime.&lt;/p&gt;

&lt;h2&gt;
  
  
  Dependency Selection and Version Pinning
&lt;/h2&gt;

&lt;p&gt;The runtime dependency chain is intentionally kept short:&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Dependency&lt;/th&gt;
&lt;th&gt;Version&lt;/th&gt;
&lt;th&gt;Purpose&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;Reactor Core&lt;/td&gt;
&lt;td&gt;3.8.6&lt;/td&gt;
&lt;td&gt;Mono/Flux programming model&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;Reactor Netty&lt;/td&gt;
&lt;td&gt;1.3.6&lt;/td&gt;
&lt;td&gt;HTTP/2 client/server&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;Netty&lt;/td&gt;
&lt;td&gt;4.2.15&lt;/td&gt;
&lt;td&gt;ByteBuf, HTTP/2 codec&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;Protobuf-java&lt;/td&gt;
&lt;td&gt;4.35.1&lt;/td&gt;
&lt;td&gt;Message serialization&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;The project does not introduce Spring, Micrometer, or any DI framework. Tests use JUnit 6.1.2 and Reactor Test's &lt;code&gt;StepVerifier&lt;/code&gt;. grpc-java 1.82.2 appears only in &lt;code&gt;interop-test&lt;/code&gt;'s test classpath with zero runtime intrusion.&lt;/p&gt;

&lt;p&gt;Versions are centrally managed via &lt;code&gt;gradle/libs.versions.toml&lt;/code&gt; with Gradle dependency locking generating lock files, ensuring fully reproducible builds across machines.&lt;/p&gt;

&lt;h2&gt;
  
  
  Why the Protocol Layer Must Come First
&lt;/h2&gt;

&lt;p&gt;It's tempting to start with "spin up an HTTP/2 Server" — you quickly get an echo demo, but it pushes the truly difficult problems to later:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;DATA may split in the middle of the five-byte gRPC header;&lt;/li&gt;
&lt;li&gt;A single DATA buffer may contain multiple messages;&lt;/li&gt;
&lt;li&gt;Metadata allows duplicate keys, and binary values require Base64;&lt;/li&gt;
&lt;li&gt;Timeout wire values are at most eight digits with a unit suffix;&lt;/li&gt;
&lt;li&gt;Final status comes from trailers;&lt;/li&gt;
&lt;li&gt;ByteBuf must be released on success, failure, and cancellation paths;&lt;/li&gt;
&lt;li&gt;Reactive Streams demand counts messages, HTTP/2 flow control counts bytes.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;Therefore the project progresses by Stage: first lock down the build and interop fixtures, then complete the protocol layer, then implement unary transport, streaming, production features, and codegen. Each Stage has executable exit criteria — "classes have been created" is not a completion standard.&lt;/p&gt;

&lt;h2&gt;
  
  
  Stage 0: Build Baseline and Interop Fixture
&lt;/h2&gt;

&lt;p&gt;Before writing any protocol code, Stage 0 solves "how to prove code is correct":&lt;/p&gt;

&lt;p&gt;&lt;strong&gt;JDK 25 compilation&lt;/strong&gt; — The project uses &lt;code&gt;-Xlint:all -parameters -encoding UTF-8&lt;/code&gt; for strict compilation. All warnings are compilation errors; silent suppression is not allowed. JDK 25 was chosen to validate Netty and Protobuf compatibility on the latest JVM early.&lt;/p&gt;

&lt;p&gt;&lt;strong&gt;Spotless formatting&lt;/strong&gt; — Unified Eclipse formatter config plus ktlint. &lt;code&gt;spotlessCheck&lt;/code&gt; is the first gate in CI. This eliminates all code review discussions about formatting.&lt;/p&gt;

&lt;p&gt;&lt;strong&gt;interop.proto test fixture&lt;/strong&gt; — Defines a test service covering all four RPC cardinalities:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight protobuf"&gt;&lt;code&gt;&lt;span class="kd"&gt;service&lt;/span&gt; &lt;span class="n"&gt;InteropTestService&lt;/span&gt; &lt;span class="p"&gt;{&lt;/span&gt;
  &lt;span class="k"&gt;rpc&lt;/span&gt; &lt;span class="n"&gt;Unary&lt;/span&gt; &lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;TestRequest&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; &lt;span class="k"&gt;returns&lt;/span&gt; &lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;TestResponse&lt;/span&gt;&lt;span class="p"&gt;);&lt;/span&gt;
  &lt;span class="k"&gt;rpc&lt;/span&gt; &lt;span class="n"&gt;ServerStreaming&lt;/span&gt; &lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;TestRequest&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; &lt;span class="k"&gt;returns&lt;/span&gt; &lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;stream&lt;/span&gt; &lt;span class="n"&gt;TestResponse&lt;/span&gt;&lt;span class="p"&gt;);&lt;/span&gt;
  &lt;span class="k"&gt;rpc&lt;/span&gt; &lt;span class="n"&gt;ClientStreaming&lt;/span&gt; &lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;stream&lt;/span&gt; &lt;span class="n"&gt;TestRequest&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; &lt;span class="k"&gt;returns&lt;/span&gt; &lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;TestResponse&lt;/span&gt;&lt;span class="p"&gt;);&lt;/span&gt;
  &lt;span class="k"&gt;rpc&lt;/span&gt; &lt;span class="n"&gt;BidirectionalStreaming&lt;/span&gt; &lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;stream&lt;/span&gt; &lt;span class="n"&gt;TestRequest&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; &lt;span class="k"&gt;returns&lt;/span&gt; &lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;stream&lt;/span&gt; &lt;span class="n"&gt;TestResponse&lt;/span&gt;&lt;span class="p"&gt;);&lt;/span&gt;
&lt;span class="p"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;&lt;strong&gt;GrpcJavaFixture&lt;/strong&gt; — In the &lt;code&gt;interop-test&lt;/code&gt; module, a test utility class starts both a grpc-java server and client, providing &lt;code&gt;start()&lt;/code&gt; / &lt;code&gt;close()&lt;/code&gt; lifecycle. Through it, bidirectional verification is possible: Reactor client calls grpc-java server, and grpc-java client calls Reactor server. Both use the same &lt;code&gt;.proto&lt;/code&gt; generated code, ensuring wire compatibility.&lt;/p&gt;

&lt;p&gt;&lt;strong&gt;Utilities&lt;/strong&gt;:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;code&gt;FreePorts&lt;/code&gt;: Allocates independent ports for each test, avoiding parallel test conflicts;&lt;/li&gt;
&lt;li&gt;
&lt;code&gt;TlsTestCertificates&lt;/code&gt;: Pre-generates self-signed certificates for subsequent TLS tests;&lt;/li&gt;
&lt;li&gt;
&lt;code&gt;LeakDetection&lt;/code&gt;: Integrates Netty's &lt;code&gt;ResourceLeakDetector&lt;/code&gt;, ensuring ByteBuf leaks are immediately exposed in tests.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;&lt;strong&gt;Stage 0 exit criteria&lt;/strong&gt;: &lt;code&gt;./gradlew clean test&lt;/code&gt; passes from fresh checkout, CI is green on Linux + Java 25, protoc generation is deterministically reproducible, grpc-java fixture communicates bidirectionally.&lt;/p&gt;

&lt;h2&gt;
  
  
  Staged Verification Strategy
&lt;/h2&gt;

&lt;p&gt;The project progresses through 12 Stages, each with clear goals, executable exit criteria, and regression coverage:&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Stage&lt;/th&gt;
&lt;th&gt;Goal&lt;/th&gt;
&lt;th&gt;Key Deliverable&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;0&lt;/td&gt;
&lt;td&gt;Build baseline&lt;/td&gt;
&lt;td&gt;CI, formatting, interop fixture&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;1&lt;/td&gt;
&lt;td&gt;Protocol foundation&lt;/td&gt;
&lt;td&gt;Frame codec, metadata, status, timeout, compression&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;2&lt;/td&gt;
&lt;td&gt;Unary transport&lt;/td&gt;
&lt;td&gt;End-to-end h2c unary call&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;3&lt;/td&gt;
&lt;td&gt;Server streaming&lt;/td&gt;
&lt;td&gt;Multi-message response stream&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;4&lt;/td&gt;
&lt;td&gt;Full cardinality&lt;/td&gt;
&lt;td&gt;Client streaming + bidirectional&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;5&lt;/td&gt;
&lt;td&gt;Production transport&lt;/td&gt;
&lt;td&gt;TLS, gzip, deadline, GOAWAY, connection pool, keepalive&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;6&lt;/td&gt;
&lt;td&gt;Name resolution&lt;/td&gt;
&lt;td&gt;DNS, subchannel, pick_first, round_robin&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;7&lt;/td&gt;
&lt;td&gt;Codegen &amp;amp; build integration&lt;/td&gt;
&lt;td&gt;Protoc plugin, descriptor registry, Gradle/Maven plugins&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;8&lt;/td&gt;
&lt;td&gt;Standard services&lt;/td&gt;
&lt;td&gt;Health v1, Reflection v1, Channelz v1&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;9&lt;/td&gt;
&lt;td&gt;Operations &amp;amp; hardening&lt;/td&gt;
&lt;td&gt;Interceptor, observer, binlog, canonical smoke, fuzz/churn&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;10&lt;/td&gt;
&lt;td&gt;Load-balancing extensions&lt;/td&gt;
&lt;td&gt;grpclb, RLS, load reporter, ORCA&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;11&lt;/td&gt;
&lt;td&gt;Diagnostics (planned)&lt;/td&gt;
&lt;td&gt;Channelz v2 over the bounded diagnostics registry&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;Each Stage completion requires: &lt;code&gt;./gradlew clean spotlessCheck test --no-daemon&lt;/code&gt; passes in full, and all prior Stage tests continue running as regression. This means when Stage 4 completes, Stage 2's unary tests are still green.&lt;/p&gt;

&lt;h2&gt;
  
  
  Where Reactor semantics differ from grpc-java
&lt;/h2&gt;

&lt;p&gt;The wire contract is shared, but the application contracts are not. The tests therefore compare the two implementations at the protocol boundary and test each runtime's lifecycle rules separately:&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Concern&lt;/th&gt;
&lt;th&gt;Reactor contract&lt;/th&gt;
&lt;th&gt;grpc-java contract&lt;/th&gt;
&lt;th&gt;What the tests must prove&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;Demand&lt;/td&gt;
&lt;td&gt;
&lt;code&gt;Subscription.request(n)&lt;/code&gt; counts decoded messages; transport demand must also respect HTTP/2 byte windows&lt;/td&gt;
&lt;td&gt;inbound flow is controlled through &lt;code&gt;ClientCall.request(n)&lt;/code&gt; / readiness callbacks&lt;/td&gt;
&lt;td&gt;no response DATA is delivered before downstream demand, and coalesced frames do not bypass demand&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;Cancellation&lt;/td&gt;
&lt;td&gt;disposing a &lt;code&gt;Mono&lt;/code&gt;/&lt;code&gt;Flux&lt;/code&gt; cancels both application publishers and resets an open stream&lt;/td&gt;
&lt;td&gt;
&lt;code&gt;ClientCall.cancel&lt;/code&gt; or &lt;code&gt;Context&lt;/code&gt; cancellation reaches &lt;code&gt;StreamObserver&lt;/code&gt; callbacks&lt;/td&gt;
&lt;td&gt;both sides stop, exactly one terminal signal wins, and late DATA/trailers are ignored&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;Trailers&lt;/td&gt;
&lt;td&gt;success trailers are retained in the call result; non-OK trailers become &lt;code&gt;GrpcException.trailers()&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;trailers are exposed through &lt;code&gt;ClientCall.Listener#onClose&lt;/code&gt; or &lt;code&gt;StatusRuntimeException#getTrailers()&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;
&lt;code&gt;grpc-status&lt;/code&gt;, &lt;code&gt;grpc-message&lt;/code&gt;, and custom trailers survive in both directions, including trailers-only errors&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;Buffer ownership&lt;/td&gt;
&lt;td&gt;Netty &lt;code&gt;ByteBuf&lt;/code&gt; is reference-counted and must be released on success, error, and cancel; decoded protobuf values cross the API boundary&lt;/td&gt;
&lt;td&gt;generated protobuf messages hide transport buffers from application code&lt;/td&gt;
&lt;td&gt;every rejected, partial, compressed, and cancelled frame reaches &lt;code&gt;refCnt() == 0&lt;/code&gt;
&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;For example, the Reactor demand test starts with zero demand and asks for one response at a time:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="nc"&gt;StepVerifier&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;create&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;client&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;bidirectionalStreaming&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="no"&gt;BIDI&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;requests&lt;/span&gt;&lt;span class="o"&gt;),&lt;/span&gt; &lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;thenRequest&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;expectNext&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;first&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;thenRequest&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;2&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;expectNext&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;second&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="n"&gt;third&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;expectComplete&lt;/span&gt;&lt;span class="o"&gt;()&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;verify&lt;/span&gt;&lt;span class="o"&gt;();&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;The corresponding grpc-java test uses &lt;code&gt;StreamObserver&lt;/code&gt; callbacks and completion signals. The messages and trailers must be wire-compatible, but it would be incorrect to describe the two tests as asserting the same backpressure mechanism. This distinction is why the project keeps both cross-runtime interop and runtime-specific failure tests.&lt;/p&gt;

&lt;p&gt;The repository now also has two server-streaming cancellation interop cases in &lt;code&gt;ServerStreamingInteroperabilityTest&lt;/code&gt;: a Reactor client receives one response from a grpc-java server and cancels, while a grpc-java client cancels a Reactor server after its first response. Both sides wait for the peer's cancellation callback under paranoid leak detection. The lower-level &lt;code&gt;GrpcFrameCodecTest&lt;/code&gt; keeps the direct ownership assertion by checking that cancelled and partial inputs finish with &lt;code&gt;refCnt() == 0&lt;/code&gt;.&lt;/p&gt;

&lt;p&gt;The essential cancellation paths are deliberately small:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight java"&gt;&lt;code&gt;&lt;span class="c1"&gt;// Reactor client: cancel the HTTP/2 call after the first response.&lt;/span&gt;
&lt;span class="nc"&gt;TestResponse&lt;/span&gt; &lt;span class="n"&gt;first&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;reactorStub&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;serverStreaming&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;Mono&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;just&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="n"&gt;request&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;5&lt;/span&gt;&lt;span class="o"&gt;)))&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;take&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;single&lt;/span&gt;&lt;span class="o"&gt;()&lt;/span&gt;
        &lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;block&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;Duration&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;ofSeconds&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="mi"&gt;5&lt;/span&gt;&lt;span class="o"&gt;));&lt;/span&gt;

&lt;span class="c1"&gt;// grpc-java client: cancel from the response callback.&lt;/span&gt;
&lt;span class="nd"&gt;@Override&lt;/span&gt;
&lt;span class="kd"&gt;public&lt;/span&gt; &lt;span class="kt"&gt;void&lt;/span&gt; &lt;span class="nf"&gt;onNext&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="nc"&gt;TestResponse&lt;/span&gt; &lt;span class="n"&gt;response&lt;/span&gt;&lt;span class="o"&gt;)&lt;/span&gt; &lt;span class="o"&gt;{&lt;/span&gt;
    &lt;span class="n"&gt;requestStream&lt;/span&gt;&lt;span class="o"&gt;.&lt;/span&gt;&lt;span class="na"&gt;cancel&lt;/span&gt;&lt;span class="o"&gt;(&lt;/span&gt;&lt;span class="s"&gt;"cancel after first response"&lt;/span&gt;&lt;span class="o"&gt;,&lt;/span&gt; &lt;span class="kc"&gt;null&lt;/span&gt;&lt;span class="o"&gt;);&lt;/span&gt;
&lt;span class="o"&gt;}&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;On the grpc-java server side, the test installs &lt;code&gt;setOnCancelHandler&lt;/code&gt;; on the Reactor server side, the response &lt;code&gt;Flux&lt;/code&gt; uses &lt;code&gt;doOnCancel&lt;/code&gt;. The complete test, including latches, peer status assertions, shutdown, and leak-detection scope, is in &lt;a href="https://github.com/qianwj/grpc-reactor/blob/db1a2fb/grpc-reactor-interop-test/src/test/java/io/github/qianwj/grpc/reactor/testing/ServerStreamingInteroperabilityTest.java#L97-L183" rel="noopener noreferrer"&gt;&lt;code&gt;ServerStreamingInteroperabilityTest.java&lt;/code&gt;&lt;/a&gt;.&lt;/p&gt;

&lt;h2&gt;
  
  
  Current Verification Results
&lt;/h2&gt;

&lt;p&gt;The project uses this unified gate that simultaneously checks formatting, compilation, protocol tests, transport tests, and grpc-java interop tests:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight shell"&gt;&lt;code&gt;./gradlew clean spotlessCheck &lt;span class="nb"&gt;test&lt;/span&gt; &lt;span class="nt"&gt;--no-daemon&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;This command passes as of this writing. The gate now includes the optional Stage 10 load-balancing extension tests in addition to protocol, transport, codegen, services, binary-log, and grpc-java interop tests. The dependency baseline is Protobuf 4.35.1, grpc-java 1.82.2, and JUnit 6.1.2. Under JDK 25, you'll still see Protobuf &lt;code&gt;Unsafe&lt;/code&gt;, Netty/Gradle native access, and Gradle deprecated feature warnings — they don't affect test results but need ongoing tracking through future JDK and Gradle upgrades.&lt;/p&gt;

&lt;p&gt;Subsequent posts will continue discussing the five-byte gRPC message envelope, ByteBuf ownership, four RPC cardinalities, production transport, DNS, code generation, standard management services, and how to integrate interceptors, observation, and binary logging without breaking Reactive Streams semantics.&lt;/p&gt;

</description>
      <category>java</category>
      <category>grpc</category>
      <category>netty</category>
      <category>reactive</category>
    </item>
  </channel>
</rss>
