const SLACK_FOR_EIGHT_BYTE_HASHING: usize = 7;
const HEAD_ROOM: usize = 2;
const TAIL_SENTINEL: u8 = 241;
#[derive(Copy, Clone)]
pub(crate) struct Window<'a> {
pub(crate) data: &'a [u8],
pub(crate) mask: usize,
}
#[derive(Copy, Clone, Debug, Eq, PartialEq)]
pub(crate) struct BlockSpan {
pub(crate) position: u32,
pub(crate) bytes: u32,
}
pub(crate) struct RingBuffer {
size: usize,
mask: usize,
tail_size: usize,
total_size: usize,
cur_size: usize,
tail_start: usize,
pos: u32,
data: Vec<u8>,
}
impl RingBuffer {
pub(crate) fn new(rb_bits: usize, lgblock: usize) -> Self {
let size = 1usize << rb_bits;
let tail_size = 1usize << lgblock;
Self {
size,
mask: size - 1,
tail_size,
total_size: size + tail_size,
cur_size: 0,
tail_start: 0,
pos: 0,
data: Vec::new(),
}
}
pub(crate) const fn mask(&self) -> usize {
self.mask
}
pub(crate) fn buffer(&self) -> &[u8] {
match self.data.get(HEAD_ROOM..) {
Some(buffer) => buffer,
None => &[],
}
}
pub(crate) fn retained_bytes(&self) -> usize {
self.data.capacity()
}
pub(crate) const fn is_allocated(&self) -> bool {
self.cur_size != 0
}
pub(crate) fn reset(&mut self) {
self.cur_size = 0;
self.tail_start = 0;
self.pos = 0;
}
#[cfg_attr(feature = "hotpath", hotpath::measure)]
fn init_buffer(&mut self, buflen: usize) {
self.data
.resize(HEAD_ROOM + buflen + SLACK_FOR_EIGHT_BYTE_HASHING, 0);
self.cur_size = buflen;
for index in 0..SLACK_FOR_EIGHT_BYTE_HASHING {
if let Some(byte) = self.data.get_mut(HEAD_ROOM + self.cur_size + index) {
*byte = 0;
}
}
if let Some(byte) = self.data.get_mut(0) {
*byte = 0;
}
if let Some(byte) = self.data.get_mut(1) {
*byte = 0;
}
}
fn set(&mut self, index: usize, value: u8) {
if let Some(byte) = self.data.get_mut(HEAD_ROOM + index) {
*byte = value;
}
}
fn write_tail(&mut self, bytes: &[u8]) {
let masked_pos = (self.pos as usize) & self.mask;
if masked_pos >= self.tail_size {
return;
}
let count = bytes.len().min(self.tail_size - masked_pos);
let start = HEAD_ROOM + self.size + masked_pos;
if let Some(target) = self.data.get_mut(start..start + count)
&& let Some(source) = bytes.get(..count)
{
target.copy_from_slice(source);
}
}
#[cfg_attr(feature = "hotpath", hotpath::measure)]
pub(crate) fn write(&mut self, bytes: &[u8]) {
let end = self.pos as usize + bytes.len();
if self.cur_size < self.total_size && end < self.size.saturating_sub(8) {
if self.pos == 0 {
self.tail_start = if bytes.len() < self.tail_size {
bytes.len()
} else {
0
};
}
self.init_buffer(end);
let start = HEAD_ROOM + self.pos as usize;
if let Some(target) = self.data.get_mut(start..HEAD_ROOM + end) {
target.copy_from_slice(bytes);
}
self.pos = end as u32;
return;
}
if self.cur_size < self.total_size {
let mirrored_end = (self.pos as usize).min(self.tail_size);
self.init_buffer(self.total_size);
self.set(self.size - 2, 0);
self.set(self.size - 1, 0);
self.set(self.size, TAIL_SENTINEL);
if self.tail_start < mirrored_end {
self.data.copy_within(
HEAD_ROOM + self.tail_start..HEAD_ROOM + mirrored_end,
HEAD_ROOM + self.size + self.tail_start,
);
}
}
let masked_pos = (self.pos as usize) & self.mask;
self.write_tail(bytes);
if masked_pos + bytes.len() <= self.size {
let start = HEAD_ROOM + masked_pos;
if let Some(target) = self.data.get_mut(start..start + bytes.len()) {
target.copy_from_slice(bytes);
}
} else {
let head = (self.total_size - masked_pos).min(bytes.len());
let start = HEAD_ROOM + masked_pos;
if let Some(target) = self.data.get_mut(start..start + head)
&& let Some(source) = bytes.get(..head)
{
target.copy_from_slice(source);
}
let wrapped = self.size - masked_pos;
if let Some(source) = bytes.get(wrapped..) {
let count = source.len();
if let Some(target) = self.data.get_mut(HEAD_ROOM..HEAD_ROOM + count) {
target.copy_from_slice(source);
}
}
}
let not_first_lap = (self.pos & (1u32 << 31)) != 0;
let pos_mask = (1u32 << 31) - 1;
let last_but_one = self.buffer().get(self.size - 2).copied().unwrap_or(0);
let last = self.buffer().get(self.size - 1).copied().unwrap_or(0);
if let Some(byte) = self.data.get_mut(0) {
*byte = last_but_one;
}
if let Some(byte) = self.data.get_mut(1) {
*byte = last;
}
self.pos = (self.pos & pos_mask) + ((bytes.len() as u32) & pos_mask);
if not_first_lap {
self.pos |= 1u32 << 31;
}
}
pub(crate) fn clear_margin(&mut self) {
if self.pos as usize > self.mask {
return;
}
let start = HEAD_ROOM + self.pos as usize;
let end = (start + SLACK_FOR_EIGHT_BYTE_HASHING).min(self.data.len());
if let Some(target) = self.data.get_mut(start..end) {
target.fill(0);
}
}
}
pub(crate) const fn wrap_position(position: u64) -> u32 {
let result = position as u32;
let gb = position >> 30;
if gb > 2 {
(result & ((1u32 << 30) - 1)) | ((((gb - 1) & 1) as u32 + 1) << 30)
} else {
result
}
}
#[cfg(test)]
mod tests {
use super::*;
fn ring(lgwin: usize, lgblock: usize) -> RingBuffer {
RingBuffer::new(1 + lgwin.max(lgblock), lgblock)
}
#[test]
fn a_short_first_write_only_allocates_what_it_holds() {
let mut rb = ring(16, 16);
assert!(!rb.is_allocated());
rb.write(b"hello");
assert!(rb.is_allocated());
assert_eq!(&rb.buffer()[..5], b"hello");
assert_eq!(&rb.buffer()[5..12], &[0u8; 7]);
}
#[test]
fn writes_before_the_first_wrap_only_allocate_the_written_prefix() {
let mut rb = ring(16, 16);
let payload = vec![7u8; 1 << 16];
rb.write(&payload);
assert_eq!(rb.cur_size, payload.len());
assert_eq!(rb.buffer()[0], 7);
rb.write(&[8; 16]);
assert_eq!(rb.cur_size, payload.len() + 16);
assert_eq!(&rb.buffer()[payload.len()..payload.len() + 16], &[8; 16]);
assert_eq!(rb.mask(), rb.size - 1);
}
#[test]
fn the_tail_sentinel_survives_until_a_lap_writes_over_it() {
let mut rb = ring(16, 16);
let mut payload = vec![1u8; 1 << 16];
payload.truncate(1 << 16);
rb.write(&vec![2u8; rb.size]);
assert_eq!(rb.buffer()[rb.size], 2);
}
#[test]
fn materializing_the_tail_preserves_previous_writes_and_the_short_write_sentinel() {
for first in [3, 16] {
let mut rb = RingBuffer::new(6, 4);
rb.write(&vec![2; first]);
rb.write(&vec![3; 52 - first]);
rb.write(&[4; 12]);
assert_eq!(&rb.buffer()[..first], vec![2; first]);
assert_eq!(&rb.buffer()[52..64], &[4; 12]);
assert_eq!(&rb.data[..2], &[4; 2]);
assert_eq!(rb.buffer()[64], if first < 16 { TAIL_SENTINEL } else { 2 });
assert_eq!(
&rb.buffer()[64 + first.min(16)..80],
&rb.buffer()[first.min(16)..16]
);
rb.write(&[9; 8]);
assert_eq!(&rb.buffer()[..8], &[9; 8]);
assert_eq!(&rb.buffer()[64..72], &[9; 8]);
rb.reset();
rb.write(&[5; 16]);
rb.clear_margin();
assert_eq!(&rb.buffer()[..16], &[5; 16]);
assert_eq!(&rb.buffer()[16..23], &[0; 7]);
}
}
#[test]
fn writes_wrap_around_and_mirror_into_the_tail() {
let mut rb = ring(10, 16);
let window = rb.size;
rb.write(&vec![1u8; window - 4]);
rb.write(&[9, 9, 9, 9, 8, 8, 8, 8]);
assert_eq!(&rb.buffer()[window - 4..window], &[9, 9, 9, 9]);
assert_eq!(&rb.buffer()[..4], &[8, 8, 8, 8]);
assert_eq!(&rb.buffer()[window..window + 4], &[8, 8, 8, 8]);
}
#[test]
fn the_head_bytes_mirror_the_end_of_the_window() {
let mut rb = ring(10, 16);
let window = rb.size;
rb.write(&vec![5u8; window]);
assert_eq!(rb.data[0], 5);
assert_eq!(rb.data[1], 5);
assert_eq!(rb.buffer()[window - 1], 5);
}
#[test]
fn clearing_the_margin_only_touches_the_first_lap() {
let mut rb = ring(10, 16);
rb.write(&[3u8; 8]);
rb.clear_margin();
assert_eq!(&rb.buffer()[8..15], &[0u8; 7]);
}
#[test]
fn position_wrapping_keeps_the_lap_parity() {
assert_eq!(wrap_position(0), 0);
assert_eq!(wrap_position(1234), 1234);
assert_eq!(wrap_position((1u64 << 30) - 1), (1 << 30) - 1);
assert_eq!(wrap_position(3u64 << 30), 1 << 30);
assert_eq!(wrap_position(4u64 << 30), 2 << 30);
assert_eq!(wrap_position(5u64 << 30), 1 << 30);
assert_eq!(wrap_position((3u64 << 30) + 17), (1 << 30) + 17);
}
}