use super::MAX_WRITE_CHUNK_SIZE;
use bytes::{Bytes, BytesMut};
pub const COALESCING_CHUNK_SIZE: usize = MAX_WRITE_CHUNK_SIZE;
#[derive(Debug)]
pub struct CoalescingBuffer {
buffer: BytesMut,
}
impl Default for CoalescingBuffer {
fn default() -> Self {
Self::new()
}
}
impl CoalescingBuffer {
pub fn new() -> Self {
Self {
buffer: BytesMut::new(),
}
}
pub fn push(&mut self, mut chunk: Bytes) -> Vec<Bytes> {
let mut ready_chunks = Vec::new();
if !self.buffer.is_empty() {
let needed = COALESCING_CHUNK_SIZE - self.buffer.len();
if chunk.len() < needed {
self.buffer.extend_from_slice(&chunk);
return ready_chunks;
}
self.buffer.extend_from_slice(&chunk[..needed]);
ready_chunks.push(self.buffer.split().freeze());
chunk = chunk.slice(needed..);
}
while chunk.len() >= COALESCING_CHUNK_SIZE {
ready_chunks.push(chunk.slice(..COALESCING_CHUNK_SIZE));
chunk = chunk.slice(COALESCING_CHUNK_SIZE..);
}
if !chunk.is_empty() {
if self.buffer.capacity() < COALESCING_CHUNK_SIZE {
self.buffer
.reserve(COALESCING_CHUNK_SIZE - self.buffer.len());
}
self.buffer.extend_from_slice(&chunk);
}
ready_chunks
}
pub fn flush(&mut self) -> Option<Bytes> {
if self.buffer.is_empty() {
return None;
}
let residual = self.buffer.split().freeze();
Some(residual)
}
pub fn is_empty(&self) -> bool {
self.buffer.is_empty()
}
pub fn len(&self) -> usize {
self.buffer.len()
}
pub fn capacity(&self) -> usize {
self.buffer.capacity()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn new_buffer_is_empty() {
let mut buf = CoalescingBuffer::new();
assert!(buf.is_empty());
assert_eq!(buf.len(), 0);
assert_eq!(buf.flush(), None);
}
#[test]
fn push_empty_chunk() {
let mut buf = CoalescingBuffer::new();
let ready = buf.push(Bytes::new());
assert!(ready.is_empty());
assert!(buf.is_empty());
assert_eq!(buf.len(), 0);
}
#[test]
fn small_appends_coalesce_into_2mib() {
let mut buf = CoalescingBuffer::new();
let payload_1mib = Bytes::from(vec![1u8; 1024 * 1024]);
let ready1 = buf.push(payload_1mib.clone());
assert!(ready1.is_empty());
assert_eq!(buf.len(), 1024 * 1024);
assert!(!buf.is_empty());
let ready2 = buf.push(payload_1mib.clone());
assert_eq!(ready2.len(), 1);
assert_eq!(ready2[0].len(), COALESCING_CHUNK_SIZE);
assert_eq!(&ready2[0][..1024 * 1024], payload_1mib.as_ref());
assert_eq!(&ready2[0][1024 * 1024..], payload_1mib.as_ref());
assert!(buf.is_empty());
assert_eq!(buf.len(), 0);
}
#[test]
fn large_append_fast_path_zero_copy() {
let mut buf = CoalescingBuffer::new();
let payload_5mib = Bytes::from(vec![42u8; 5 * 1024 * 1024]);
let ready = buf.push(payload_5mib.clone());
assert_eq!(ready.len(), 2);
assert_eq!(ready[0].len(), COALESCING_CHUNK_SIZE);
assert_eq!(ready[1].len(), COALESCING_CHUNK_SIZE);
assert_eq!(
ready[0].as_ptr(),
payload_5mib[..COALESCING_CHUNK_SIZE].as_ptr()
);
assert_eq!(
ready[1].as_ptr(),
payload_5mib[COALESCING_CHUNK_SIZE..2 * COALESCING_CHUNK_SIZE].as_ptr()
);
assert_eq!(buf.len(), 1024 * 1024);
assert!(!buf.is_empty());
let residual = buf.flush();
assert_eq!(residual.map(|b| b.len()), Some(1024 * 1024));
assert!(buf.is_empty());
}
#[test]
fn lazy_allocation_and_flush_capacity() {
let mut buf = CoalescingBuffer::new();
assert_eq!(buf.capacity(), 0);
let exact = Bytes::from(vec![1u8; COALESCING_CHUNK_SIZE]);
let ready = buf.push(exact);
assert_eq!(ready.len(), 1);
assert_eq!(buf.capacity(), 0);
let partial = Bytes::from_static(b"hello world");
let ready = buf.push(partial);
assert!(ready.is_empty());
assert!(buf.capacity() >= COALESCING_CHUNK_SIZE);
let flushed = buf.flush();
assert!(flushed.is_some());
assert_eq!(buf.len(), 0);
}
#[test]
fn default_buffer() {
let buf = CoalescingBuffer::default();
assert!(buf.is_empty());
assert_eq!(buf.len(), 0);
assert_eq!(buf.capacity(), 0);
}
#[test]
fn mixed_small_and_large_appends() {
let mut buf = CoalescingBuffer::new();
let small = Bytes::from(vec![0xAAu8; 512 * 1024]); let large = Bytes::from(vec![0xBBu8; 3 * 1024 * 1024]);
let ready1 = buf.push(small);
assert!(ready1.is_empty());
assert_eq!(buf.len(), 512 * 1024);
let ready2 = buf.push(large);
assert_eq!(ready2.len(), 1);
assert_eq!(ready2[0].len(), COALESCING_CHUNK_SIZE);
assert_eq!(buf.len(), (1024 + 512) * 1024);
let residual = buf.flush().unwrap();
assert_eq!(residual.len(), (1024 + 512) * 1024);
assert!(buf.is_empty());
}
#[test]
fn exact_chunk_size_append() {
let mut buf = CoalescingBuffer::new();
let exact = Bytes::from(vec![0x77u8; COALESCING_CHUNK_SIZE]);
let ready = buf.push(exact.clone());
assert_eq!(ready.len(), 1);
assert_eq!(ready[0], exact);
assert!(buf.is_empty());
assert_eq!(buf.flush(), None);
}
#[test]
fn multiple_flushes() {
let mut buf = CoalescingBuffer::new();
let payload = Bytes::from_static(b"hello world");
buf.push(payload.clone());
assert_eq!(buf.flush(), Some(payload));
assert_eq!(buf.flush(), None);
assert_eq!(buf.flush(), None);
}
#[test]
fn flush_empty_buffer() {
let mut buf = CoalescingBuffer::new();
assert_eq!(buf.flush(), None);
}
}