grpc-reactor 实战(二):五字节消息帧、分片解码与协议原语

grpc-reactor 实战(二):五字节消息帧、分片解码与协议原语

Elvis Lv2

gRPC 的 Protobuf 消息不会直接裸写进 HTTP/2 DATA。每条消息前面都有五字节 envelope:

1
2
3
byte 0      bit 0 表示是否压缩,bit 1-7 必须为 0
byte 1-4 payload 长度,无符号大端序
byte 5..n Protobuf 或压缩后的消息体

编码并不复杂,真正困难的是解码不能假设一次收到完整 frame。HTTP/2、TCP 和 Reactor Netty 都不保证 buffer 边界与 gRPC 消息边界一致。

编码必须定义所有权

GrpcFrameCodec.encode 的契约是:返回 frame 与输入 message 各自拥有独立生命周期,编码过程不移动输入的 reader index,也不释放输入。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
public static ByteBuf encode(
ByteBufAllocator allocator,
ByteBuf message,
GrpcCompression.Codec compression) {
byte[] payload = new byte[message.readableBytes()];
message.getBytes(message.readerIndex(), payload);
if (!compression.name().equals("identity")) {
payload = compression.compress(payload);
}
return allocator.buffer(5 + payload.length)
.writeByte(compressed ? 1 : 0)
.writeInt(payload.length)
.writeBytes(payload);
}

当前实现为了先把 ownership 做清楚,会复制为 byte array。它不是零拷贝实现,也不应该被描述成性能最优;后续若改用 slice 或 composite buffer,必须重新证明取消和异常路径没有泄漏。

解码器是每次订阅独立的状态机

decode 使用 Flux.defer 为每个订阅创建 Decoder:

1
2
3
4
5
6
7
8
return Flux.defer(() -> {
var decoder = new Decoder(
maxWireMessageSize,
maxDecompressedMessageSize,
compression);
return input.concatMap(decoder::accept, 1)
.concatWith(Flux.defer(decoder::finish));
});

状态只有两部分:尚未收满的五字节 header,以及已知长度但尚未收满的 payload。一个输入 ByteBuf 可能只包含 header 的一个字节,也可能连续包含三条完整消息。

concatMap(..., 1) 保持输入次序,同时让 decoder 按下游 demand 交付消息。每个 source buffer 最终通过 doFinally 释放:

1
2
3
4
5
6
7
8
9
10
11
12
private Flux<ByteBuf> accept(ByteBuf source) {
return Flux.<ByteBuf>generate(sink -> {
try {
while (source.isReadable()) {
// read header, allocate bounded payload, emit complete message
}
sink.complete();
} catch (Throwable error) {
sink.error(error);
}
}).doFinally(ignored -> source.release());
}

在分配前验证 peer 输入

长度字段来自远端,不能读取后立即分配。实现先用 long 组合无符号大端整数,再检查 wire size:

1
2
3
4
5
6
7
8
long length = ((long) (header[1] & 0xff) << 24)
| ((long) (header[2] & 0xff) << 16)
| ((long) (header[3] & 0xff) << 8)
| (header[4] & 0xffL);

if (length > maxWireMessageSize) {
throw new GrpcProtocolException("wire message length exceeds limit");
}

wire size 与 decompressed size 是两个不同限制。一个很小的 gzip payload 可能解压为巨大数据,因此 gzip 解压过程中还必须再次限制输出,防止 decompression bomb。

压缩标志只有最低位合法;保留位不为零直接视为协议错误。收到 compressed frame 但协商结果仍是 identity,同样不能猜测算法继续解码。

测试所有 header 分割点

固定测“拆成两个 buffer”还不够,五字节 header 有六个代表性分割位置。测试使用动态用例覆盖 0 到 5:

1
2
3
4
5
6
7
8
9
10
11
12
IntStream.rangeClosed(0, GrpcFrameCodec.HEADER_SIZE)
.mapToObj(split -> DynamicTest.dynamicTest(
"split after byte " + split,
() -> {
byte[] wire = wireBytes("hello");
Flux<ByteBuf> chunks = Flux.just(
wrapped(wire, 0, split),
wrapped(wire, split, wire.length - split));
StepVerifier.create(GrpcFrameCodec.decode(chunks))
.assertNext(message -> assertMessage(message, "hello"))
.verifyComplete();
}));

另外还覆盖:body 任意分片、多 frame 合并、空消息、gzip、保留位、截断 frame、wire/decompressed limit,以及编码不消费输入 buffer。

取消时也要释放未消费数据

一个 source ByteBuf 中合并了 onetwothree 三条消息。下游只请求两条后取消,测试最后断言 source 的引用计数归零:

1
2
3
4
5
6
7
StepVerifier.create(GrpcFrameCodec.decode(Flux.just(source)), 0)
.thenRequest(1).assertNext(message -> assertMessage(message, "one"))
.thenRequest(1).assertNext(message -> assertMessage(message, "two"))
.thenCancel()
.verify();

assertEquals(0, source.refCnt());

这条断言比“结果内容正确”更重要。网络库中的 happy path 很容易正确,泄漏往往发生在取消、超限、截断或重复 terminal signal 中。

在本篇对应的 Stage 1 提交中,frame decoder 只覆盖了消息级 demand,尚未完成 streaming transport 的两级 flow control。Reactive Streams 计算消息数量,HTTP/2 window 计算字节,两者不能混为一谈。这个边界后来由 Stage 3 的有界入站缓冲与按需交付补齐,并在 Stage 4 扩展到双向流;对应实现和压力测试会在系列第五、六篇展开。

Metadata:保序、重复、二进制三合一

gRPC metadata 不是简单的 Map<String, String>。它必须满足:

  • 同一个 key 可以出现多次,顺序敏感;
  • key 只允许小写字母、数字、_.-,由正则 [0-9a-z_.-]+ 校验;
  • -bin 结尾的 key 表示二进制值,wire 上使用无 padding Base64 编码;
  • 用户不能设置保留字段(content-typetegrpc-statusgrpc-timeout 等)。

GrpcMetadata 内部是一个不可变的 entry list,保证跨异步边界安全共享:

1
2
3
4
5
6
7
8
GrpcMetadata metadata = GrpcMetadata.builder()
.addAscii("trace-id", "abc123")
.addAscii("trace-id", "def456") // 允许重复
.addBinary("auth-token-bin", tokenBytes)
.build();

// 读取时保留顺序
List<GrpcMetadata.Entry> all = metadata.getAll("trace-id"); // [abc123, def456]

从 HTTP/2 headers 解析时,二进制值可能被逗号拼接(HTTP/2 合并同名 header 的行为),实现会按逗号拆分后逐段 Base64 解码。编码后总大小有上限(默认 8KB),防止 peer 用超大 headers 耗尽内存。

Status:17 种状态码与 percent-encoding

gRPC 定义了 17 种标准状态码,每种对应一个明确的语义,用于指导客户端的重试和错误处理策略。GrpcStatus 是一个 record,包含 status code 和人类可读 message:

1
2
3
4
5
6
7
8
9
10
public record GrpcStatus(Code code, String message) {
public enum Code {
OK(0), CANCELLED(1), UNKNOWN(2), INVALID_ARGUMENT(3),
DEADLINE_EXCEEDED(4), NOT_FOUND(5), ALREADY_EXISTS(6),
PERMISSION_DENIED(7), RESOURCE_EXHAUSTED(8),
FAILED_PRECONDITION(9), ABORTED(10), OUT_OF_RANGE(11),
UNIMPLEMENTED(12), INTERNAL(13), UNAVAILABLE(14),
DATA_LOSS(15), UNAUTHENTICATED(16);
}
}

几个值得注意的设计细节:

  • DEADLINE_EXCEEDED 即使操作已成功完成也可能返回——当 server 的成功响应在网络传输中超过了 deadline,client 仍然看到超时;
  • UNAVAILABLE 表示暂时性故障,client 可以安全重试;而 INTERNAL 表示 server 内部错误,盲目重试通常无意义;
  • UNIMPLEMENTED 对应 HTTP 404 的语义——路径存在但方法不被支持。

grpc-message trailer 中的文本需要 percent-encoding:可打印 ASCII 中除 % 外直接传输,其余字符转为 %HH。这使得 UTF-8 错误描述能安全地通过纯 ASCII 的 HTTP/2 header 传输。

未知的 status code 数值不会抛异常,而是映射到 UNKNOWN——这保证了向前兼容,当 peer 使用了更新的 gRPC spec 定义的新状态码时不会导致解析失败。

Timeout:八位数字 + 单位,向上取整

gRPC 协议规定 grpc-timeout 传输的是相对超时而非绝对时间戳。Server 收到请求时,deadline 已经比 client 设定的少了一次网络传输延迟,必须用剩余时间窗口判断是否应该继续处理。

wire 格式非常紧凑:最多八位十进制数字加一个单位后缀:

1
2
3
100m   → 100 毫秒
2S → 2 秒
99999999H → 约 11415 年(上限值)

支持六种单位:H(小时)、M(分钟)、S(秒)、m(毫秒)、u(微秒)、n(纳秒)。

格式化 Duration 为 wire value 时,必须向上取整(ceiling division),确保 wire deadline 不会比调用者设定的更短:

1
2
BigInteger amount = nanos.add(unitNanos.subtract(BigInteger.ONE))
.divide(unitNanos); // ceiling division

使用 BigInteger 而非 long 避免了 nanosecond 乘法溢出。选择单位的策略是从纳秒向上扫描,取第一个使数值不超过 99,999,999 的单位。

Compression:identity 是默认,gzip 是唯一内置

GrpcCompression 管理 codec 注册和协商:

1
2
3
4
5
6
GrpcCompression.Registry registry = GrpcCompression.Registry.builder()
.add(GrpcCompression.GZIP)
.build();

// 产生 grpc-accept-encoding: gzip
String advertised = registry.advertisedEncodings();

identity 始终隐含存在且排在第一位。Codec 接口只有两个方法:

1
2
3
4
5
interface Codec {
String name();
byte[] compress(byte[] input);
byte[] decompress(byte[] input, int maxSize);
}

Gzip 解压使用 8KB chunk 流式读取,每个 chunk 后用 Math.addExact() 累加已解压字节数——这既能在超出 maxSize 时立即中断,也能通过算术溢出检查防止 int wrap-around。一个 100 字节的 gzip payload 可能解压出数 GB 数据,所以 decompression bomb 防护不是可选项。

GrpcMethod:四种 cardinality 的统一描述

每个 RPC 方法用一个 GrpcMethod record 完整描述:

1
2
3
4
5
6
7
8
9
var method = new GrpcMethod<>(
"testing.InteropTestService", // fullServiceName
"Unary", // methodName
GrpcMethod.Cardinality.UNARY, // cardinality
new ProtobufMarshaller<>(TestRequest.parser()),
new ProtobufMarshaller<>(TestResponse.parser()));

method.path(); // "/testing.InteropTestService/Unary"
method.fullMethodName(); // "testing.InteropTestService/Unary"

Cardinality 枚举提供 singleRequest()singleResponse() 查询,transport 层用它在边界插入 single() 操作符,把多值违规转为明确异常而非默默丢弃。

ProtobufMarshaller:序列化不拥有输入

ProtobufMarshaller 包装 Protobuf 的 Parser<T>

1
2
3
4
5
6
7
8
public ByteBuf serialize(ByteBufAllocator allocator, T value) {
byte[] bytes = value.toByteArray();
return allocator.buffer(bytes.length).writeBytes(bytes);
}

public T deserialize(ByteBuf message) {
return parser.parseFrom(message.nioBuffer());
}

关键契约:deserialize 使用 NIO ByteBuffer 视图读取,不移动 ByteBuf 的 reader index,也不释放输入。所有权始终留给调用者。这个设计让 frame decoder 的 doFinally(release) 可以统一管理所有 buffer 生命周期,不论反序列化成功还是失败。

Stage 1 退出条件

协议模块的测试不需要 Reactor Netty 在 classpath 上——这验证了模块边界的正确性。退出条件:

  • 所有 header 分割点 + payload 分割的 frame 解码通过
  • 取消路径无 ByteBuf 泄漏(refCnt 断言)
  • Metadata 保序、二进制编解码、大小限制生效
  • Status 未知 code 不抛异常
  • Timeout 所有单位解析、格式化向上取整、溢出安全
  • Gzip 解压限制生效(decompression bomb 被拒绝)
  • Marshaller 不改变输入 buffer 状态

这些测试使用 JUnit 5 的 @TestFactory + DynamicTest 参数化,配合 Reactor Test 的 StepVerifier 验证异步流行为。协议层完成后,下一篇进入 Stage 2:在 Reactor Netty HTTP/2 上跑通第一条 unary 调用。

  • 标题: grpc-reactor 实战(二):五字节消息帧、分片解码与协议原语
  • 作者: Elvis
  • 创建于 : 2026-06-02 15:40:00
  • 更新于 : 2026-06-03 10:00:00
  • 链接: https://qianwj.github.io/2026/06/02/grpc-reactor-series-2/
  • 版权声明: 本文章采用 CC BY-NC-SA 4.0 进行许可。