grpc_webnext/
grpc_framing.rs1use bytes::{Buf, BufMut, Bytes, BytesMut};
9
10const HEADER_LEN: usize = 5;
11
12pub fn frame(message: &[u8]) -> Bytes {
14 let mut buf = BytesMut::with_capacity(HEADER_LEN + message.len());
15 buf.put_u8(0); buf.put_u32(message.len() as u32);
17 buf.put_slice(message);
18 buf.freeze()
19}
20
21#[derive(Default)]
23pub struct Deframer {
24 buf: BytesMut,
25}
26
27impl Deframer {
28 pub fn new() -> Self {
29 Self::default()
30 }
31
32 pub fn push(&mut self, chunk: &[u8]) {
34 self.buf.put_slice(chunk);
35 }
36
37 pub fn next_message(&mut self) -> Option<Bytes> {
39 if self.buf.len() < HEADER_LEN {
40 return None;
41 }
42 let len = u32::from_be_bytes([self.buf[1], self.buf[2], self.buf[3], self.buf[4]]) as usize;
44 if self.buf.len() < HEADER_LEN + len {
45 return None;
46 }
47 self.buf.advance(HEADER_LEN);
48 Some(self.buf.split_to(len).freeze())
49 }
50
51 pub fn is_empty(&self) -> bool {
53 self.buf.is_empty()
54 }
55}
56
57pub fn deframe_all(body: &[u8]) -> Vec<Bytes> {
59 let mut d = Deframer::new();
60 d.push(body);
61 let mut out = Vec::new();
62 while let Some(m) = d.next_message() {
63 out.push(m);
64 }
65 out
66}