use kevy_resp::{ProtocolError, Reply, parse_reply};
#[derive(Debug)]
pub struct ReplyReadBuf {
buf: Vec<u8>,
pos: usize,
}
impl ReplyReadBuf {
pub fn with_capacity(cap: usize) -> Self {
Self {
buf: Vec::with_capacity(cap),
pos: 0,
}
}
pub fn parse_next(&mut self) -> Result<Option<Reply>, ProtocolError> {
match parse_reply(&self.buf[self.pos..])? {
Some((reply, used)) => {
self.pos += used;
self.reclaim_after_consume();
Ok(Some(reply))
}
None => Ok(None),
}
}
pub fn extend(&mut self, bytes: &[u8]) {
if self.pos > 0 && bytes.len() > self.buf.len() - self.pos {
self.compact();
}
self.buf.extend_from_slice(bytes);
}
fn reclaim_after_consume(&mut self) {
if self.pos == self.buf.len() {
self.buf.clear();
self.pos = 0;
} else if self.pos > self.buf.capacity() / 2 {
self.compact();
}
}
fn compact(&mut self) {
self.buf.drain(..self.pos);
self.pos = 0;
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parses_sequential_replies_via_cursor() {
let mut b = ReplyReadBuf::with_capacity(64);
b.extend(b"+A\r\n+B\r\n");
assert!(matches!(b.parse_next(), Ok(Some(Reply::Simple(s))) if s == b"A"));
assert!(matches!(b.parse_next(), Ok(Some(Reply::Simple(s))) if s == b"B"));
assert!(matches!(b.parse_next(), Ok(None)));
}
#[test]
fn partial_reply_returns_none_then_completes() {
let mut b = ReplyReadBuf::with_capacity(16);
b.extend(b"$5\r\nhel");
assert!(matches!(b.parse_next(), Ok(None)));
b.extend(b"lo\r\n");
assert!(matches!(b.parse_next(), Ok(Some(Reply::Bulk(v))) if v == b"hello"));
}
#[test]
fn survives_compaction_mid_stream() {
let mut b = ReplyReadBuf::with_capacity(8);
b.extend(b"+first-longish\r\n$3\r\nabc");
assert!(matches!(b.parse_next(), Ok(Some(Reply::Simple(s))) if s == b"first-longish"));
assert!(matches!(b.parse_next(), Ok(None)));
b.extend(b"\r\n");
assert!(matches!(b.parse_next(), Ok(Some(Reply::Bulk(v))) if v == b"abc"));
}
#[test]
fn malformed_frame_errors() {
let mut b = ReplyReadBuf::with_capacity(16);
b.extend(b"!bogus\r\n");
assert!(b.parse_next().is_err());
}
}