2026-05-11 22:31:04 -07:00
|
|
|
use std::marker::PhantomData;
|
|
|
|
|
use tonic::codec::{Codec, Decoder, Encoder, DecodeBuf, EncodeBuf};
|
2026-05-13 23:08:21 -07:00
|
|
|
use bytes::{Buf, BufMut, Bytes, BytesMut};
|
2026-05-11 22:31:04 -07:00
|
|
|
use roto_runtime::RotoMessage;
|
2026-05-13 23:08:21 -07:00
|
|
|
use std::sync::{Arc, Mutex};
|
|
|
|
|
use std::pin::Pin;
|
|
|
|
|
use std::future::Future;
|
|
|
|
|
use std::task::{Context, Poll};
|
|
|
|
|
use http_body::Body;
|
2026-05-11 22:31:04 -07:00
|
|
|
|
|
|
|
|
pub struct RotoCodec<T, U> {
|
|
|
|
|
_phantom: PhantomData<(T, U)>,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl<T, U> Default for RotoCodec<T, U> {
|
|
|
|
|
fn default() -> Self {
|
|
|
|
|
Self {
|
|
|
|
|
_phantom: PhantomData,
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl<T, U> Codec for RotoCodec<T, U>
|
|
|
|
|
where
|
|
|
|
|
T: RotoMessage + Send + 'static,
|
|
|
|
|
U: RotoMessage + Send + 'static,
|
|
|
|
|
{
|
|
|
|
|
type Encode = U;
|
|
|
|
|
type Decode = T;
|
|
|
|
|
type Encoder = RotoEncoder<U>;
|
|
|
|
|
type Decoder = RotoDecoder<T>;
|
|
|
|
|
|
|
|
|
|
fn encoder(&mut self) -> Self::Encoder {
|
|
|
|
|
RotoEncoder(PhantomData)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn decoder(&mut self) -> Self::Decoder {
|
|
|
|
|
RotoDecoder(PhantomData)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub struct RotoEncoder<U>(PhantomData<U>);
|
|
|
|
|
|
|
|
|
|
impl<U> Encoder for RotoEncoder<U>
|
|
|
|
|
where
|
|
|
|
|
U: RotoMessage,
|
|
|
|
|
{
|
|
|
|
|
type Item = U;
|
|
|
|
|
type Error = tonic::Status;
|
|
|
|
|
|
|
|
|
|
fn encode(&mut self, message: Self::Item, buf: &mut EncodeBuf<'_>) -> Result<(), Self::Error> {
|
|
|
|
|
buf.put_slice(&message.bytes());
|
|
|
|
|
Ok(())
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub struct RotoDecoder<T>(PhantomData<T>);
|
|
|
|
|
|
|
|
|
|
impl<T> Decoder for RotoDecoder<T>
|
|
|
|
|
where
|
|
|
|
|
T: RotoMessage,
|
|
|
|
|
{
|
|
|
|
|
type Item = T;
|
|
|
|
|
type Error = tonic::Status;
|
|
|
|
|
|
|
|
|
|
fn decode(&mut self, buf: &mut DecodeBuf<'_>) -> Result<Option<Self::Item>, Self::Error> {
|
|
|
|
|
if buf.remaining() == 0 {
|
|
|
|
|
return Ok(None);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
let bytes = buf.copy_to_bytes(buf.remaining());
|
|
|
|
|
match T::decode(bytes) {
|
|
|
|
|
Ok(msg) => Ok(Some(msg)),
|
|
|
|
|
Err(e) => Err(tonic::Status::internal(format!("Roto decode error: {}", e))),
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-05-13 23:08:21 -07:00
|
|
|
|
|
|
|
|
pub struct BufferPool {
|
|
|
|
|
pool: Mutex<Vec<BytesMut>>,
|
|
|
|
|
default_capacity: usize,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
impl BufferPool {
|
|
|
|
|
pub fn new(default_capacity: usize) -> Self {
|
|
|
|
|
Self {
|
|
|
|
|
pool: Mutex::new(Vec::new()),
|
|
|
|
|
default_capacity,
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn get(&self) -> BytesMut {
|
|
|
|
|
self.pool.lock().unwrap().pop().unwrap_or_else(|| BytesMut::with_capacity(self.default_capacity))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub fn put(&self, mut buf: BytesMut) {
|
|
|
|
|
buf.clear();
|
|
|
|
|
if buf.capacity() >= self.default_capacity {
|
|
|
|
|
self.pool.lock().unwrap().push(buf);
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
pub struct StatusBody(pub(crate) Option<Bytes>);
|
|
|
|
|
|
|
|
|
|
impl Body for StatusBody {
|
|
|
|
|
type Data = Bytes;
|
|
|
|
|
type Error = tonic::Status;
|
|
|
|
|
|
|
|
|
|
fn poll_frame(
|
|
|
|
|
mut self: Pin<&mut Self>,
|
|
|
|
|
cx: &mut Context<'_>,
|
|
|
|
|
) -> Poll<Option<Result<http_body::Frame<Self::Data>, Self::Error>>> {
|
|
|
|
|
if let Some(data) = self.0.take() {
|
|
|
|
|
Poll::Ready(Some(Ok(http_body::Frame::data(data))))
|
|
|
|
|
} else {
|
|
|
|
|
Poll::Ready(None)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|