use crate::google::storage::v2::{
BidiWriteObjectRequest, ChecksummedData, bidi_write_object_request::Data,
};
use bytes::Bytes;
use std::collections::VecDeque;
pub const DEFAULT_REPLAY_BUFFER_SIZE: usize = 32 * 1024 * 1024;
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ReplayChunk {
pub write_offset: i64,
pub data: Bytes,
pub crc32c: u32,
}
impl ReplayChunk {
pub fn new(write_offset: i64, data: Bytes, crc32c: u32) -> Self {
Self {
write_offset,
data,
crc32c,
}
}
pub fn end_offset(&self) -> i64 {
self.write_offset + self.data.len() as i64
}
pub fn to_request(&self) -> BidiWriteObjectRequest {
BidiWriteObjectRequest {
write_offset: self.write_offset,
data: Some(Data::ChecksummedData(ChecksummedData {
content: self.data.clone(),
crc32c: Some(self.crc32c),
})),
..BidiWriteObjectRequest::default()
}
}
}
#[derive(Debug)]
pub struct ReplayBuffer {
queue: VecDeque<ReplayChunk>,
unpersisted_bytes: usize,
capacity: usize,
}
impl Default for ReplayBuffer {
fn default() -> Self {
Self::new()
}
}
impl ReplayBuffer {
pub fn new() -> Self {
Self::with_capacity(DEFAULT_REPLAY_BUFFER_SIZE)
}
pub fn with_capacity(capacity: usize) -> Self {
Self {
queue: VecDeque::new(),
unpersisted_bytes: 0,
capacity,
}
}
pub fn capacity(&self) -> usize {
self.capacity
}
pub fn push(&mut self, chunk: ReplayChunk) {
self.unpersisted_bytes += chunk.data.len();
self.queue.push_back(chunk);
}
pub fn ack(&mut self, persisted_size: i64) {
while let Some(front) = self.queue.front() {
if front.end_offset() <= persisted_size {
self.unpersisted_bytes -= front.data.len();
self.queue.pop_front();
} else {
break;
}
}
if let Some(front) = self.queue.front_mut()
&& front.write_offset < persisted_size
{
let trimmed_bytes = (persisted_size - front.write_offset) as usize;
front.data = front.data.slice(trimmed_bytes..);
front.write_offset = persisted_size;
front.crc32c = crc32c::crc32c(&front.data);
self.unpersisted_bytes -= trimmed_bytes;
}
}
pub fn is_full(&self) -> bool {
self.unpersisted_bytes >= self.capacity
}
pub fn is_empty(&self) -> bool {
self.queue.is_empty()
}
pub fn num_chunks(&self) -> usize {
self.queue.len()
}
pub fn unpersisted_bytes(&self) -> usize {
self.unpersisted_bytes
}
pub fn chunks_to_replay(&self) -> impl Iterator<Item = &ReplayChunk> {
self.queue.iter()
}
pub fn clear(&mut self) {
self.queue.clear();
self.unpersisted_bytes = 0;
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn new_buffer_state() {
let mut buf = ReplayBuffer::new();
assert!(buf.is_empty());
assert_eq!(buf.num_chunks(), 0);
assert_eq!(buf.unpersisted_bytes(), 0);
assert!(!buf.is_full());
buf.ack(100);
assert!(buf.is_empty());
}
#[test]
fn push_and_ack_full_chunks() {
let mut buf = ReplayBuffer::new();
let chunk1 = Bytes::from_static(b"hello ");
let chunk2 = Bytes::from_static(b"world!");
buf.push(ReplayChunk::new(0, chunk1.clone(), crc32c::crc32c(&chunk1)));
buf.push(ReplayChunk::new(6, chunk2.clone(), crc32c::crc32c(&chunk2)));
assert_eq!(buf.num_chunks(), 2);
assert_eq!(buf.unpersisted_bytes(), 12);
buf.ack(4);
assert_eq!(buf.num_chunks(), 2);
assert_eq!(buf.unpersisted_bytes(), 8);
let chunks: Vec<_> = buf.chunks_to_replay().cloned().collect();
assert_eq!(chunks[0].write_offset, 4);
assert_eq!(chunks[0].data.as_ref(), b"o ");
assert_eq!(chunks[0].crc32c, crc32c::crc32c(b"o "));
assert_eq!(chunks[1].write_offset, 6);
assert_eq!(chunks[1].data.as_ref(), b"world!");
assert_eq!(chunks[1].crc32c, crc32c::crc32c(b"world!"));
buf.ack(10);
assert_eq!(buf.num_chunks(), 1);
assert_eq!(buf.unpersisted_bytes(), 2);
let chunks: Vec<_> = buf.chunks_to_replay().cloned().collect();
assert_eq!(chunks[0].write_offset, 10);
assert_eq!(chunks[0].data.as_ref(), b"d!");
assert_eq!(chunks[0].crc32c, crc32c::crc32c(b"d!"));
buf.ack(12);
assert!(buf.is_empty());
assert_eq!(buf.unpersisted_bytes(), 0);
}
#[test]
fn ack_duplicate_or_earlier_offset() {
let mut buf = ReplayBuffer::new();
let chunk = Bytes::from_static(b"abcdef");
buf.push(ReplayChunk::new(10, chunk.clone(), crc32c::crc32c(&chunk)));
buf.ack(5);
assert_eq!(buf.num_chunks(), 1);
assert_eq!(buf.unpersisted_bytes(), 6);
buf.ack(10);
assert_eq!(buf.num_chunks(), 1);
assert_eq!(buf.unpersisted_bytes(), 6);
}
#[test]
fn is_full_threshold() {
let mut buf = ReplayBuffer::new();
let huge_chunk = Bytes::from(vec![0u8; DEFAULT_REPLAY_BUFFER_SIZE]);
buf.push(ReplayChunk::new(0, huge_chunk, 0));
assert!(buf.is_full());
buf.ack(1);
assert!(!buf.is_full());
}
#[test]
fn chunk_to_request_conversion() {
let data = Bytes::from_static(b"replay data");
let crc = crc32c::crc32c(&data);
let chunk = ReplayChunk::new(42, data.clone(), crc);
assert_eq!(chunk.end_offset(), 42 + data.len() as i64);
let req = chunk.to_request();
assert_eq!(req.write_offset, 42);
if let Some(Data::ChecksummedData(cd)) = req.data {
assert_eq!(cd.content, data);
assert_eq!(cd.crc32c, Some(crc));
} else {
panic!("expected ChecksummedData");
}
}
#[test]
fn clear_resets_size_and_queue() {
let mut buf = ReplayBuffer::new();
let chunk = Bytes::from_static(b"test data");
buf.push(ReplayChunk::new(0, chunk, 0));
assert!(!buf.is_empty());
assert!(buf.unpersisted_bytes() > 0);
buf.clear();
assert!(buf.is_empty());
assert_eq!(buf.unpersisted_bytes(), 0);
}
#[test]
fn with_capacity_default() {
let buf = ReplayBuffer::new();
assert_eq!(buf.capacity(), DEFAULT_REPLAY_BUFFER_SIZE);
}
#[test]
fn with_capacity_custom() {
const CUSTOM_CAPACITY: usize = 64 * 1024 * 1024;
let buf = ReplayBuffer::with_capacity(CUSTOM_CAPACITY);
assert_eq!(buf.capacity(), CUSTOM_CAPACITY);
}
#[test]
fn is_full_with_custom_capacity() {
let mut buf = ReplayBuffer::with_capacity(100);
let chunk = Bytes::from(vec![0u8; 100]);
buf.push(ReplayChunk::new(0, chunk, 0));
assert!(buf.is_full());
buf.ack(1);
assert!(!buf.is_full());
}
}