Add RevBuilder for reverse message encoding

Generate RevBuilder structs in codegen alongside ProtoBuilder.
RevBuilder writes fields into a buffer starting from the end,
allowing fields to be set in reverse declaration order. This
avoids reallocation and produces valid protobuf data matching
the forward layout.

Update generated interop files and add comprehensive tests for
wire format, nesting, map entries, overflow handling, and buffer
marking utilities.
This commit is contained in:
2026-06-19 19:39:21 -07:00
parent 3640af6e57
commit 229727340d
4 changed files with 644 additions and 134 deletions
+204 -129
View File
@@ -1,16 +1,8 @@
// @generated by protoc-gen-roto — do not edit
#![allow(
unused,
unused_imports,
unused_assignments,
unused_variables,
non_camel_case_types
)]
use bytes::{Buf, BufMut, Bytes, BytesMut};
#[allow(unused, unused_imports, unused_assignments, unused_variables, non_camel_case_types)]
use roto_runtime::{ProtoAccessor, ProtoBuilder, RevBuilder, Result, RotoError, read_varint, RepeatedFieldIterator, RotoMessage};
use core::str;
use roto_runtime::{
ProtoAccessor, ProtoBuilder, RepeatedFieldIterator, Result, RotoError, RotoMessage, read_varint,
};
use bytes::{Bytes, BytesMut, Buf, BufMut};
pub struct UnaryRequest<'a> {
accessor: roto_runtime::ProtoAccessor<'a>,
@@ -23,21 +15,17 @@ impl<'a> UnaryRequest<'a> {
let mut message_offset = None;
for item in accessor.fields() {
let (offset, tag, _) = item?;
if tag.field_number == 1 {
message_offset = Some(offset);
}
if tag.field_number == 1 { message_offset = Some(offset); }
}
Ok(Self {
accessor,
message_offset,
message_offset,
})
}
pub fn message(&self) -> roto_runtime::Result<&'a str> {
let offset = self
.message_offset
.ok_or(roto_runtime::RotoError::FieldNotFound)?;
let offset = self.message_offset.ok_or(roto_runtime::RotoError::FieldNotFound)?;
let (bytes, _) = self.accessor.get_value_at(offset)?;
core::str::from_utf8(bytes).map_err(|_| roto_runtime::RotoError::WireFormatViolation)
}
@@ -46,13 +34,12 @@ impl<'a> UnaryRequest<'a> {
self.message().or(Ok(""))
}
pub fn has_message(&self) -> bool {
self.message_offset.is_some()
}
pub fn has_message(&self) -> bool { self.message_offset.is_some() }
pub fn raw_fields(&self) -> roto_runtime::RawFieldIterator<'a> {
self.accessor.raw_fields()
}
}
pub struct UnaryRequestBuilder<'b> {
@@ -93,6 +80,45 @@ impl<'b> UnaryRequestBuilder<'b> {
}
}
pub struct UnaryRequestRevBuilder<'b> {
builder: roto_runtime::RevBuilder<'b>,
message_written: bool,
}
impl<'b> UnaryRequestRevBuilder<'b> {
pub fn builder(buf: &mut [u8]) -> UnaryRequestRevBuilder<'_> {
UnaryRequestRevBuilder {
builder: roto_runtime::RevBuilder::new(buf),
message_written: false,
}
}
pub fn message(mut self, value: &str) -> roto_runtime::Result<Self> {
self.builder.write_string(1, value)?;
self.message_written = true;
Ok(self)
}
pub fn with(mut self, msg: &UnaryRequest<'_>) -> roto_runtime::Result<Self> {
let fields: Vec<_> = msg.accessor.raw_fields().collect();
for item in fields.into_iter().rev() {
let (field_number, raw_bytes) = item?;
let is_written = match field_number {
1 => self.message_written,
_ => false,
};
if !is_written {
self.builder.write_raw(raw_bytes)?;
}
}
Ok(self)
}
pub fn finish(self) -> roto_runtime::Result<&'b mut [u8]> {
self.builder.finish()
}
}
pub struct OwnedUnaryRequest {
pub data: bytes::Bytes,
}
@@ -125,21 +151,17 @@ impl<'a> UnaryResponse<'a> {
let mut reply_offset = None;
for item in accessor.fields() {
let (offset, tag, _) = item?;
if tag.field_number == 1 {
reply_offset = Some(offset);
}
if tag.field_number == 1 { reply_offset = Some(offset); }
}
Ok(Self {
accessor,
reply_offset,
reply_offset,
})
}
pub fn reply(&self) -> roto_runtime::Result<&'a str> {
let offset = self
.reply_offset
.ok_or(roto_runtime::RotoError::FieldNotFound)?;
let offset = self.reply_offset.ok_or(roto_runtime::RotoError::FieldNotFound)?;
let (bytes, _) = self.accessor.get_value_at(offset)?;
core::str::from_utf8(bytes).map_err(|_| roto_runtime::RotoError::WireFormatViolation)
}
@@ -148,13 +170,12 @@ impl<'a> UnaryResponse<'a> {
self.reply().or(Ok(""))
}
pub fn has_reply(&self) -> bool {
self.reply_offset.is_some()
}
pub fn has_reply(&self) -> bool { self.reply_offset.is_some() }
pub fn raw_fields(&self) -> roto_runtime::RawFieldIterator<'a> {
self.accessor.raw_fields()
}
}
pub struct UnaryResponseBuilder<'b> {
@@ -195,6 +216,45 @@ impl<'b> UnaryResponseBuilder<'b> {
}
}
pub struct UnaryResponseRevBuilder<'b> {
builder: roto_runtime::RevBuilder<'b>,
reply_written: bool,
}
impl<'b> UnaryResponseRevBuilder<'b> {
pub fn builder(buf: &mut [u8]) -> UnaryResponseRevBuilder<'_> {
UnaryResponseRevBuilder {
builder: roto_runtime::RevBuilder::new(buf),
reply_written: false,
}
}
pub fn reply(mut self, value: &str) -> roto_runtime::Result<Self> {
self.builder.write_string(1, value)?;
self.reply_written = true;
Ok(self)
}
pub fn with(mut self, msg: &UnaryResponse<'_>) -> roto_runtime::Result<Self> {
let fields: Vec<_> = msg.accessor.raw_fields().collect();
for item in fields.into_iter().rev() {
let (field_number, raw_bytes) = item?;
let is_written = match field_number {
1 => self.reply_written,
_ => false,
};
if !is_written {
self.builder.write_raw(raw_bytes)?;
}
}
Ok(self)
}
pub fn finish(self) -> roto_runtime::Result<&'b mut [u8]> {
self.builder.finish()
}
}
pub struct OwnedUnaryResponse {
pub data: bytes::Bytes,
}
@@ -227,21 +287,17 @@ impl<'a> StreamingRequest<'a> {
let mut query_offset = None;
for item in accessor.fields() {
let (offset, tag, _) = item?;
if tag.field_number == 1 {
query_offset = Some(offset);
}
if tag.field_number == 1 { query_offset = Some(offset); }
}
Ok(Self {
accessor,
query_offset,
query_offset,
})
}
pub fn query(&self) -> roto_runtime::Result<&'a str> {
let offset = self
.query_offset
.ok_or(roto_runtime::RotoError::FieldNotFound)?;
let offset = self.query_offset.ok_or(roto_runtime::RotoError::FieldNotFound)?;
let (bytes, _) = self.accessor.get_value_at(offset)?;
core::str::from_utf8(bytes).map_err(|_| roto_runtime::RotoError::WireFormatViolation)
}
@@ -250,13 +306,12 @@ impl<'a> StreamingRequest<'a> {
self.query().or(Ok(""))
}
pub fn has_query(&self) -> bool {
self.query_offset.is_some()
}
pub fn has_query(&self) -> bool { self.query_offset.is_some() }
pub fn raw_fields(&self) -> roto_runtime::RawFieldIterator<'a> {
self.accessor.raw_fields()
}
}
pub struct StreamingRequestBuilder<'b> {
@@ -297,6 +352,45 @@ impl<'b> StreamingRequestBuilder<'b> {
}
}
pub struct StreamingRequestRevBuilder<'b> {
builder: roto_runtime::RevBuilder<'b>,
query_written: bool,
}
impl<'b> StreamingRequestRevBuilder<'b> {
pub fn builder(buf: &mut [u8]) -> StreamingRequestRevBuilder<'_> {
StreamingRequestRevBuilder {
builder: roto_runtime::RevBuilder::new(buf),
query_written: false,
}
}
pub fn query(mut self, value: &str) -> roto_runtime::Result<Self> {
self.builder.write_string(1, value)?;
self.query_written = true;
Ok(self)
}
pub fn with(mut self, msg: &StreamingRequest<'_>) -> roto_runtime::Result<Self> {
let fields: Vec<_> = msg.accessor.raw_fields().collect();
for item in fields.into_iter().rev() {
let (field_number, raw_bytes) = item?;
let is_written = match field_number {
1 => self.query_written,
_ => false,
};
if !is_written {
self.builder.write_raw(raw_bytes)?;
}
}
Ok(self)
}
pub fn finish(self) -> roto_runtime::Result<&'b mut [u8]> {
self.builder.finish()
}
}
pub struct OwnedStreamingRequest {
pub data: bytes::Bytes,
}
@@ -329,21 +423,17 @@ impl<'a> StreamingResponse<'a> {
let mut item_offset = None;
for item in accessor.fields() {
let (offset, tag, _) = item?;
if tag.field_number == 1 {
item_offset = Some(offset);
}
if tag.field_number == 1 { item_offset = Some(offset); }
}
Ok(Self {
accessor,
item_offset,
item_offset,
})
}
pub fn item(&self) -> roto_runtime::Result<&'a str> {
let offset = self
.item_offset
.ok_or(roto_runtime::RotoError::FieldNotFound)?;
let offset = self.item_offset.ok_or(roto_runtime::RotoError::FieldNotFound)?;
let (bytes, _) = self.accessor.get_value_at(offset)?;
core::str::from_utf8(bytes).map_err(|_| roto_runtime::RotoError::WireFormatViolation)
}
@@ -352,13 +442,12 @@ impl<'a> StreamingResponse<'a> {
self.item().or(Ok(""))
}
pub fn has_item(&self) -> bool {
self.item_offset.is_some()
}
pub fn has_item(&self) -> bool { self.item_offset.is_some() }
pub fn raw_fields(&self) -> roto_runtime::RawFieldIterator<'a> {
self.accessor.raw_fields()
}
}
pub struct StreamingResponseBuilder<'b> {
@@ -399,6 +488,45 @@ impl<'b> StreamingResponseBuilder<'b> {
}
}
pub struct StreamingResponseRevBuilder<'b> {
builder: roto_runtime::RevBuilder<'b>,
item_written: bool,
}
impl<'b> StreamingResponseRevBuilder<'b> {
pub fn builder(buf: &mut [u8]) -> StreamingResponseRevBuilder<'_> {
StreamingResponseRevBuilder {
builder: roto_runtime::RevBuilder::new(buf),
item_written: false,
}
}
pub fn item(mut self, value: &str) -> roto_runtime::Result<Self> {
self.builder.write_string(1, value)?;
self.item_written = true;
Ok(self)
}
pub fn with(mut self, msg: &StreamingResponse<'_>) -> roto_runtime::Result<Self> {
let fields: Vec<_> = msg.accessor.raw_fields().collect();
for item in fields.into_iter().rev() {
let (field_number, raw_bytes) = item?;
let is_written = match field_number {
1 => self.item_written,
_ => false,
};
if !is_written {
self.builder.write_raw(raw_bytes)?;
}
}
Ok(self)
}
pub fn finish(self) -> roto_runtime::Result<&'b mut [u8]> {
self.builder.finish()
}
}
pub struct OwnedStreamingResponse {
pub data: bytes::Bytes,
}
@@ -420,41 +548,26 @@ impl roto_runtime::RotoMessage for OwnedStreamingResponse {
}
}
use crate::{BufferPool, StatusBody};
use futures_util::StreamExt;
use http_body::Body;
use http_body_util::BodyExt;
use std::future::Future;
#[allow(unused, unused_imports, unused_assignments, unused_variables, non_camel_case_types)]
use tonic::{Request, Response, Status};
use tokio_stream::Stream;
use std::pin::Pin;
use std::sync::Arc;
use std::task::{Context, Poll};
use tokio_stream::Stream;
use std::future::Future;
use tonic::body::BoxBody;
#[allow(
unused,
unused_imports,
unused_assignments,
unused_variables,
non_camel_case_types
)]
use tonic::{Request, Response, Status};
use tower::Service;
use futures_util::StreamExt;
use http_body_util::BodyExt;
use http_body::Body;
use crate::{BufferPool, StatusBody};
#[async_trait::async_trait]
pub trait InteropService: Send + Sync + 'static {
async fn unary_call(
&self,
request: Request<OwnedUnaryRequest>,
) -> std::result::Result<Response<OwnedUnaryResponse>, Status>;
async fn streaming_call(
&self,
request: Request<OwnedStreamingRequest>,
) -> std::result::Result<
Response<
Pin<Box<dyn Stream<Item = std::result::Result<OwnedStreamingResponse, Status>> + Send>>,
>,
Status,
>;
async fn unary_call(&self, request: Request<OwnedUnaryRequest>) -> std::result::Result<Response<OwnedUnaryResponse>, Status>;
async fn streaming_call(&self, request: Request<OwnedStreamingRequest>) -> std::result::Result<Response<Pin<Box<dyn Stream<Item = std::result::Result<OwnedStreamingResponse, Status>> + Send>>>, Status>;
}
#[derive(Clone)]
@@ -476,8 +589,7 @@ impl tonic::server::NamedService for InteropServiceServer {
impl Service<http::Request<BoxBody>> for InteropServiceServer {
type Response = http::Response<BoxBody>;
type Error = std::convert::Infallible;
type Future =
Pin<Box<dyn Future<Output = std::result::Result<Self::Response, Self::Error>> + Send>>;
type Future = Pin<Box<dyn Future<Output = std::result::Result<Self::Response, Self::Error>> + Send>>;
fn poll_ready(&mut self, _: &mut Context<'_>) -> Poll<std::result::Result<(), Self::Error>> {
Poll::Ready(Ok(()))
@@ -502,14 +614,8 @@ impl Service<http::Request<BoxBody>> for InteropServiceServer {
let bytes_vec = buf.split_to(total_len).freeze();
pool.put(buf);
if bytes_vec.len() < 5 {
let res_body = BoxBody::new(StatusBody::new(
Some(Bytes::from_static(&[0, 0, 0, 0, 0])),
0,
));
return Ok(http::Response::builder()
.status(200)
.body(res_body)
.unwrap());
let res_body = BoxBody::new(StatusBody::new(Some(Bytes::from_static(&[0, 0, 0, 0, 0])), 0));
return Ok(http::Response::builder().status(200).body(res_body).unwrap());
}
let payload = bytes_vec.slice(5..);
@@ -519,28 +625,16 @@ impl Service<http::Request<BoxBody>> for InteropServiceServer {
let request_msg = match OwnedUnaryRequest::decode(payload) {
Ok(msg) => msg,
Err(_e) => {
let res_body = BoxBody::new(StatusBody::new(
Some(Bytes::from_static(&[0, 0, 0, 0, 0])),
0,
));
return Ok(http::Response::builder()
.status(200)
.body(res_body)
.unwrap());
let res_body = BoxBody::new(StatusBody::new(Some(Bytes::from_static(&[0, 0, 0, 0, 0])), 0));
return Ok(http::Response::builder().status(200).body(res_body).unwrap());
}
};
let response = match inner.unary_call(Request::new(request_msg)).await {
Ok(res) => res,
Err(_e) => {
let res_body = BoxBody::new(StatusBody::new(
Some(Bytes::from_static(&[0, 0, 0, 0, 0])),
0,
));
return Ok(http::Response::builder()
.status(200)
.body(res_body)
.unwrap());
let res_body = BoxBody::new(StatusBody::new(Some(Bytes::from_static(&[0, 0, 0, 0, 0])), 0));
return Ok(http::Response::builder().status(200).body(res_body).unwrap());
}
};
@@ -556,36 +650,17 @@ impl Service<http::Request<BoxBody>> for InteropServiceServer {
pool.put(res_buf);
let res_body = BoxBody::new(StatusBody::new(Some(frame), 0));
routed = true;
return Ok(http::Response::builder()
.status(200)
.header("content-type", "application/grpc")
.body(res_body)
.unwrap());
return Ok(http::Response::builder().status(200).header("content-type", "application/grpc").body(res_body).unwrap());
}
if path == "/interop.InteropService/StreamingCall" {
let res_body = BoxBody::new(StatusBody::new(
Some(Bytes::from_static(&[0, 0, 0, 0, 0])),
0,
));
return Ok(http::Response::builder()
.status(200)
.body(res_body)
.unwrap());
let res_body = BoxBody::new(StatusBody::new(Some(Bytes::from_static(&[0, 0, 0, 0, 0])), 0));
return Ok(http::Response::builder().status(200).body(res_body).unwrap());
}
if !routed {
let res_body = BoxBody::new(StatusBody::new(
Some(Bytes::from_static(&[0, 0, 0, 0, 0])),
0,
));
return Ok(http::Response::builder()
.status(200)
.body(res_body)
.unwrap());
let res_body = BoxBody::new(StatusBody::new(Some(Bytes::from_static(&[0, 0, 0, 0, 0])), 0));
return Ok(http::Response::builder().status(200).body(res_body).unwrap());
}
Ok(http::Response::builder()
.status(200)
.body(BoxBody::new(StatusBody::new(None, 0)))
.unwrap())
Ok(http::Response::builder().status(200).body(BoxBody::new(StatusBody::new(None, 0))).unwrap())
})
}
}