<?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>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://qianwj.github.io/en/2026/06/01/grpc-reactor-series-1/" rel="noopener noreferrer"&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>
