#![allow(
clippy::indexing_slicing,
clippy::expect_used,
clippy::arithmetic_side_effects,
reason = "bench harness over hand-built buffers: a panic on a malformed \
fixture is the report, and adding error handling would put \
branches in the code being measured"
)]
use crate::{
Error, ErrorKind, Result,
client::{BufferConfig, PubSubMessage},
network::PubSubPush,
resp::{
BufferDecoder, Command, CommandEncoder, ParsedFrame, RespBuf, RespFrameParser,
RespResponse, RespTapeMut,
},
};
use bytes::{Bytes, BytesMut};
use serde::de::DeserializeOwned;
use smallvec::SmallVec;
use tokio_util::codec::{Decoder, Encoder as _};
#[inline]
pub fn bench_decode_to<T: DeserializeOwned>(bytes: &[u8]) -> Result<T> {
let mut tape = RespTapeMut::default();
let (frame, frame_len) = RespFrameParser::new(bytes, &mut tape).parse()?;
let buf = RespBuf::from(Bytes::copy_from_slice(&bytes[..frame_len]));
RespResponse::new(buf, frame).to()
}
pub struct BenchPubSubPush {
buf: RespBuf,
tape: RespTapeMut,
}
impl BenchPubSubPush {
pub fn new(bytes: &[u8]) -> Result<Self> {
let mut tape = RespTapeMut::default();
let (_, frame_len) = RespFrameParser::new(bytes, &mut tape).parse()?;
Ok(Self {
buf: RespBuf::from(Bytes::copy_from_slice(&bytes[..frame_len])),
tape,
})
}
fn response(&mut self) -> Result<RespResponse> {
let (frame, _) = RespFrameParser::new(&self.buf, &mut self.tape).parse()?;
Ok(RespResponse::new(self.buf.clone(), frame))
}
fn segments(response: &RespResponse) -> Result<(&[u8], &[u8], &[u8])> {
match PubSubPush::try_from(response) {
Ok(PubSubPush::Message(channel, payload) | PubSubPush::SMessage(channel, payload)) => {
Ok((&[], channel, payload))
}
Ok(PubSubPush::PMessage(pattern, channel, payload)) => Ok((pattern, channel, payload)),
_ => Err(Error::from(ErrorKind::EOF)),
}
}
#[inline(never)]
pub fn deliver(&mut self) -> Result<PubSubMessage> {
let response = self.response()?;
PubSubMessage::try_from(&response)
}
#[inline(never)]
pub fn deliver_owned(&mut self) -> Result<(Vec<u8>, Vec<u8>, Vec<u8>)> {
let response = self.response()?;
let (pattern, channel, payload) = Self::segments(&response)?;
Ok((pattern.to_vec(), channel.to_vec(), payload.to_vec()))
}
#[inline(never)]
pub fn deliver_inline(&mut self) -> Result<(SmallVec<[u8; 64]>, usize, usize)> {
let response = self.response()?;
let (pattern, channel, payload) = Self::segments(&response)?;
let channel_start = pattern.len();
let payload_start = channel_start + channel.len();
let mut buf = SmallVec::with_capacity(payload_start + payload.len());
buf.extend_from_slice(pattern);
buf.extend_from_slice(channel);
buf.extend_from_slice(payload);
Ok((buf, channel_start, payload_start))
}
}
#[inline]
pub fn bench_decode_chunked<T: DeserializeOwned>(chunks: &[&[u8]]) -> Result<T> {
let mut decoder = BufferDecoder::new();
let mut buf = BytesMut::new();
for chunk in chunks {
buf.extend_from_slice(chunk);
if let Some(resp) = decoder.decode(&mut buf)? {
return resp.to();
}
}
match decoder.decode(&mut buf)? {
Some(resp) => resp.to(),
None => Err(Error::from(ErrorKind::EOF)),
}
}
#[inline(never)]
pub fn bench_decode_stream_grow(data: &[u8], chunk: usize) -> Result<usize> {
drive_stream(data, chunk, false)
}
#[inline(never)]
pub fn bench_decode_stream_prereserve(data: &[u8], chunk: usize) -> Result<usize> {
drive_stream(data, chunk, true)
}
#[inline(always)]
fn drive_stream(data: &[u8], chunk: usize, prereserve: bool) -> Result<usize> {
let mut decoder = BufferDecoder::new();
let mut src = BytesMut::with_capacity(BufferConfig::DEFAULT.read_capacity);
let mut pos = 0usize;
let mut reserved = false;
loop {
if let Some(resp) = decoder.decode(&mut src)? {
let len = match std::hint::black_box(&resp) {
RespResponse::Frame { buf, .. } => buf.as_ref().len(),
_ => 0,
};
return Ok(len);
}
if pos >= data.len() {
return Err(Error::from(ErrorKind::EOF));
}
if prereserve && !reserved && !src.is_empty() {
src.reserve(data.len().saturating_sub(src.len()));
reserved = true;
}
src.reserve(1);
let spare = src.capacity() - src.len();
let take = spare.min(chunk).min(data.len() - pos);
src.extend_from_slice(&data[pos..pos + take]);
pos += take;
}
}
#[derive(Default)]
pub struct BenchTape(RespTapeMut);
impl BenchTape {
pub fn new() -> Self {
Self::default()
}
}
#[inline(never)]
pub fn bench_parse_only(bytes: &[u8], tape: &mut BenchTape) {
let tape = &mut tape.0;
let (frame, frame_len) = RespFrameParser::new(bytes, tape)
.parse()
.expect("bench_parse_only fed a valid frame");
std::hint::black_box((&frame, frame_len));
}
#[inline(never)]
pub fn bench_tape_footprint(bytes: &[u8]) -> (usize, usize) {
let mut tape = RespTapeMut::default();
let (frame, frame_len) = RespFrameParser::new(bytes, &mut tape)
.parse()
.expect("bench_tape_footprint fed a valid frame");
let tape_bytes = match &frame {
ParsedFrame::Collection(tape) => tape.byte_len(),
_ => 0,
};
(frame_len, tape_bytes)
}
pub struct BenchDecoder {
decoder: BufferDecoder,
src: BytesMut,
}
impl Default for BenchDecoder {
fn default() -> Self {
Self::new()
}
}
impl BenchDecoder {
pub fn new() -> Self {
Self {
decoder: BufferDecoder::new(),
src: BytesMut::with_capacity(BufferConfig::DEFAULT.read_capacity),
}
}
pub fn feed(&mut self, reply: &[u8]) -> Result<usize> {
self.src.extend_from_slice(reply);
let Some(resp) = self.decoder.decode(&mut self.src)? else {
return Err(Error::from(ErrorKind::EOF));
};
let tape_bytes = match &resp {
RespResponse::Frame { tape, .. } => tape.byte_len(),
_ => 0,
};
drop(resp);
Ok(tape_bytes)
}
pub fn retained_tape_capacity(&self) -> usize {
self.decoder.tape_capacity()
}
}
#[inline(never)]
pub fn bench_encode_command(command: &Command, buf: &mut BytesMut) {
buf.clear();
CommandEncoder
.encode(command, buf)
.expect("bench_encode_command fed a valid command");
std::hint::black_box(&buf);
}