diff --git a/Cargo.toml b/Cargo.toml index bb427a10..f3de20b3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -144,6 +144,12 @@ harness = false [[bench]] name = "multiline_framing" harness = false +required-features = ["framing-bench"] + +[[bench]] +name = "retained_buffer" +harness = false +required-features = ["framing-bench"] [[bench]] name = "yenc_decoding" diff --git a/benches/retained_buffer.rs b/benches/retained_buffer.rs new file mode 100644 index 00000000..0382eaef --- /dev/null +++ b/benches/retained_buffer.rs @@ -0,0 +1,34 @@ +//! Retained-buffer append policy benchmarks. +//! +//! Run with: +//! `nix develop -c cargo bench --features framing-bench --bench retained_buffer` + +use divan::{Bencher, black_box}; +use nntp_proxy::pool::buffer::retained_append_benchmark; +use tokio::runtime::Builder; + +fn main() { + divan::main(); +} + +fn bench_tail(bencher: Bencher, physical_tail: usize) { + let runtime = Builder::new_current_thread().build().unwrap(); + bencher + .with_inputs(|| retained_append_benchmark(physical_tail)) + .bench_local_values(|case| black_box(runtime.block_on(case.drain()))); +} + +#[divan::bench(sample_count = 100, sample_size = 100)] +fn exact_full(bencher: Bencher) { + bench_tail(bencher, 0); +} + +#[divan::bench(sample_count = 100, sample_size = 100)] +fn one_byte_tail(bencher: Bencher) { + bench_tail(bencher, 1); +} + +#[divan::bench(sample_count = 100, sample_size = 100)] +fn roomy_tail(bencher: Bencher) { + bench_tail(bencher, 4096); +} diff --git a/src/pool/buffer.rs b/src/pool/buffer.rs index 62ff03d5..9e8cb4b3 100644 --- a/src/pool/buffer.rs +++ b/src/pool/buffer.rs @@ -39,6 +39,79 @@ pub struct PooledBuffer { counts_toward_pool: bool, } +/// One permission to append a backend read to a retained buffer. +pub(crate) struct RetainedAppendPermit<'a> { + buffer: &'a mut PooledBuffer, +} + +pub(crate) struct AppendedRead<'a> { + buffer: &'a mut PooledBuffer, + previous_len: usize, +} + +pub(crate) enum AppendOutcome<'a> { + Data(AppendedRead<'a>), + Eof(&'a mut PooledBuffer), +} + +#[cfg(feature = "framing-bench")] +#[doc(hidden)] +pub struct RetainedAppendBenchmark { + buffer: PooledBuffer, + reader: std::io::Cursor>, +} + +#[cfg(feature = "framing-bench")] +#[doc(hidden)] +#[must_use] +pub fn retained_append_benchmark(physical_tail: usize) -> RetainedAppendBenchmark { + const CAPACITY: usize = 64 * 1024; + const RETAINED: usize = 32; + + assert!(physical_tail <= CAPACITY - RETAINED); + let initialized = CAPACITY - physical_tail; + let pool = BufferPool::new( + BufferSize::try_new(CAPACITY).expect("valid benchmark size"), + 1, + ); + let mut buffer = pool.acquire(); + buffer.copy_from_slice(&vec![b'x'; initialized]); + buffer.expose_initialized_range_without_copying(initialized - RETAINED..initialized); + + RetainedAppendBenchmark { + buffer, + reader: std::io::Cursor::new(vec![b'y'; CAPACITY - RETAINED]), + } +} + +#[cfg(feature = "framing-bench")] +impl RetainedAppendBenchmark { + pub async fn drain(mut self) -> (usize, usize) { + let Some(permit) = self.buffer.retained_append_permit() else { + panic!("benchmark buffer should have logical room to append"); + }; + let AppendOutcome::Data(appended) = permit + .read(&mut self.reader) + .await + .expect("benchmark cursor read should succeed") + else { + panic!("benchmark cursor should contain data"); + }; + let first_read = appended.as_new_bytes().len(); + let buffer = appended.into_inner(); + let second_read = if self.reader.position() < self.reader.get_ref().len() as u64 { + buffer + .read_from(&mut self.reader) + .await + .expect("benchmark cursor read should succeed") + } else { + 0 + }; + assert_eq!(first_read + second_read, 64 * 1024 - 32); + (first_read, second_read) + } +} + #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum PooledBufferKind { Regular, @@ -57,54 +130,39 @@ struct BufferAcquisition { counts_toward_pool: bool, } -#[derive(Clone, Copy, Debug, PartialEq, Eq)] -enum ReadMode { - Reset, - Append, -} - #[derive(Default)] struct BufferStorage { bytes: BytesMut, - hidden_prefix: Option, - allocation_capacity: usize, + visible_from: usize, } impl BufferStorage { fn new(bytes: BytesMut) -> Self { - let allocation_capacity = bytes.capacity(); Self { bytes, - hidden_prefix: None, - allocation_capacity, + visible_from: 0, } } #[inline] fn capacity(&self) -> usize { - self.allocation_capacity + self.bytes.capacity() } #[inline] fn initialized_len(&self) -> usize { - self.bytes.len() + self.bytes.len() - self.visible_from } #[inline] fn clear(&mut self) { - self.restore_hidden_prefix(); - self.bytes.clear(); - } - - #[inline] - fn clear_for_fresh_read(&mut self) { - self.restore_hidden_prefix(); self.bytes.clear(); + self.visible_from = 0; } #[inline] fn as_slice(&self) -> &[u8] { - &self.bytes + &self.bytes[self.visible_from..] } fn copy_from_slice(&mut self, data: &[u8]) { @@ -115,7 +173,6 @@ impl BufferStorage { fn extend_from_slice(&mut self, data: &[u8]) { self.compact_visible(); self.bytes.extend_from_slice(data); - self.allocation_capacity = self.bytes.capacity(); } fn retain_range(&mut self, range: Range) { @@ -123,34 +180,31 @@ impl BufferStorage { range.start < range.end && range.end <= self.initialized_len(), "exposed range must be non-empty and inside initialized bytes" ); - let mut visible = self.bytes.split_off(range.start); - visible.truncate(range.len()); - let newly_hidden_prefix = std::mem::replace(&mut self.bytes, visible); - if let Some(prefix) = &mut self.hidden_prefix { - prefix.unsplit(newly_hidden_prefix); - } else { - self.hidden_prefix = Some(newly_hidden_prefix); - } + let start = self.visible_from + range.start; + let end = self.visible_from + range.end; + self.bytes.truncate(end); + self.visible_from = start; } fn compact_visible(&mut self) { - let Some(mut prefix) = self.hidden_prefix.take() else { + if self.visible_from == 0 { return; - }; - let prefix_len = prefix.len(); - let visible_len = self.bytes.len(); - prefix.unsplit(std::mem::take(&mut self.bytes)); - prefix.copy_within(prefix_len..prefix_len + visible_len, 0); - prefix.truncate(visible_len); - self.bytes = prefix; + } + + let len = self.initialized_len(); + self.bytes.copy_within(self.visible_from.., 0); + self.bytes.truncate(len); + self.visible_from = 0; } - fn restore_hidden_prefix(&mut self) { - let Some(mut prefix) = self.hidden_prefix.take() else { - return; - }; - prefix.unsplit(std::mem::take(&mut self.bytes)); - self.bytes = prefix; + fn compact_visible_if_tail_full(&mut self) { + if self.visible_from != 0 && self.bytes.capacity() == self.bytes.len() { + self.compact_visible(); + } + } + + fn has_retained_prefix(&self) -> bool { + self.visible_from != 0 } fn freeze(&mut self) -> Bytes { @@ -158,22 +212,24 @@ impl BufferStorage { } fn take(&mut self) -> BytesMut { - drop(self.hidden_prefix.take()); - self.allocation_capacity = 0; - std::mem::take(&mut self.bytes) + let visible_from = std::mem::take(&mut self.visible_from); + let mut bytes = std::mem::take(&mut self.bytes); + if visible_from == 0 { + bytes + } else { + bytes.split_off(visible_from) + } } fn restore(&mut self, bytes: BytesMut) { - self.allocation_capacity = bytes.capacity(); self.bytes = bytes; - self.hidden_prefix = None; + self.visible_from = 0; } async fn read_from(&mut self, reader: &mut R, read_len: usize) -> std::io::Result where R: AsyncRead + Unpin, { - debug_assert!(self.hidden_prefix.is_none()); if read_len == 0 { return Ok(0); } @@ -252,7 +308,8 @@ impl PooledBuffer { where R: AsyncRead + Unpin, { - self.read_into_spare(reader, ReadMode::Reset).await + self.buffer.clear(); + self.read_into_spare(reader, self.read_limit()).await } /// Read more data at the current initialized offset, accumulating bytes. @@ -268,24 +325,21 @@ impl PooledBuffer { where R: AsyncRead + Unpin, { - self.read_into_spare(reader, ReadMode::Append).await + self.buffer.compact_visible_if_tail_full(); + let read_len = self + .read_limit() + .saturating_sub(self.buffer.initialized_len()); + self.read_into_spare(reader, read_len).await } - async fn read_into_spare(&mut self, reader: &mut R, mode: ReadMode) -> std::io::Result + async fn read_into_spare( + &mut self, + reader: &mut R, + read_len: usize, + ) -> std::io::Result where R: AsyncRead + Unpin, { - let read_len = match mode { - ReadMode::Reset => { - self.buffer.clear_for_fresh_read(); - self.read_limit() - } - ReadMode::Append => { - self.buffer.compact_visible(); - self.read_limit() - .saturating_sub(self.buffer.initialized_len()) - } - }; self.buffer.read_from(reader, read_len).await } @@ -294,22 +348,17 @@ impl PooledBuffer { self.buffer.retain_range(range); } - #[must_use] - pub(crate) fn is_exposed_range_view(&self) -> bool { - self.buffer.hidden_prefix.is_some() - } - - #[must_use] - pub(crate) fn has_remaining_fixed_writable_region(&self) -> bool { - self.buffer.initialized_len() < self.read_limit() + pub(crate) fn retained_append_permit(&mut self) -> Option> { + if !self.buffer.has_retained_prefix() || self.buffer.initialized_len() >= self.read_limit() + { + return None; + } + Some(RetainedAppendPermit { buffer: self }) } #[cfg(test)] pub(crate) fn allocation_ptr(&self) -> *const u8 { - self.buffer - .hidden_prefix - .as_ref() - .map_or_else(|| self.buffer.bytes.as_ptr(), |prefix| prefix.as_ptr()) + self.buffer.bytes.as_ptr() } /// Copy data into buffer and mark as initialized @@ -378,6 +427,45 @@ impl PooledBuffer { } } +impl<'a> RetainedAppendPermit<'a> { + pub(crate) async fn read(self, reader: &mut R) -> std::io::Result> + where + R: AsyncRead + Unpin, + { + let Self { buffer } = self; + let previous_len = buffer.initialized(); + let read = buffer.read_more(reader).await?; + Ok(match read { + 0 => AppendOutcome::Eof(buffer), + _ => AppendOutcome::Data(AppendedRead { + buffer, + previous_len, + }), + }) + } +} + +impl<'a> AppendedRead<'a> { + #[must_use] + pub(crate) fn previous_len(&self) -> usize { + self.previous_len + } + + #[must_use] + pub(crate) fn as_new_bytes(&self) -> &[u8] { + &self.buffer.as_ref()[self.previous_len..] + } + + #[must_use] + pub(crate) fn total_len(&self) -> usize { + self.buffer.initialized() + } + + pub(crate) fn into_inner(self) -> &'a mut PooledBuffer { + self.buffer + } +} + impl Deref for PooledBuffer { type Target = [u8]; @@ -1479,15 +1567,17 @@ mod tests { let mut buffer = pool.acquire(); buffer.copy_from_slice(b"discardkeepdiscard"); let allocation = buffer.allocation_ptr(); + let capacity = buffer.capacity(); buffer.expose_initialized_range_without_copying(7..11); assert_eq!(buffer.as_ref(), b"keep"); assert_eq!(buffer.allocation_ptr(), allocation); + assert_eq!(buffer.capacity(), capacity); } #[tokio::test] - async fn read_more_compacts_a_retained_prefix_before_appending() { + async fn read_more_appends_after_a_retained_prefix_without_compacting() { let pool = BufferPool::new(BufferSize::try_new(1024).unwrap(), 1); let mut buffer = pool.acquire(); buffer.copy_from_slice(b"discard22unused"); @@ -1497,11 +1587,168 @@ mod tests { writer.write_all(b"0 ready\r\n").await.unwrap(); drop(writer); - let read = buffer.read_more(&mut reader).await.unwrap(); + let Some(appendable) = buffer.retained_append_permit() else { + panic!("retained prefix should have spare tail capacity"); + }; + let AppendOutcome::Data(appended) = appendable.read(&mut reader).await.unwrap() else { + panic!("the reader has data available"); + }; + let read = appended.as_new_bytes().len(); + let buffer = appended.into_inner(); assert_eq!(read, 9); assert_eq!(buffer.as_ref(), b"220 ready\r\n"); assert_eq!(buffer.allocation_ptr(), allocation); + assert!(buffer.buffer.has_retained_prefix()); + } + + #[tokio::test] + async fn retained_append_result_owns_the_new_logical_window() { + let pool = BufferPool::new(BufferSize::try_new(8).unwrap(), 1); + let mut buffer = pool.acquire(); + buffer.copy_from_slice(b"discard22unused"); + buffer.expose_initialized_range_without_copying(7..9); + + let mut reader = std::io::Cursor::new(b"abcdefghi"); + let Some(appendable) = buffer.retained_append_permit() else { + panic!("retained prefix should have logical room to append"); + }; + let AppendOutcome::Data(appended) = appendable.read(&mut reader).await.unwrap() else { + panic!("the reader has data available"); + }; + + assert_eq!(appended.previous_len(), 2); + assert_eq!(appended.as_new_bytes(), b"abcdef"); + let buffer = appended.into_inner(); + assert_eq!(buffer.as_ref(), b"22abcdef"); + assert!(buffer.retained_append_permit().is_none()); + } + + #[tokio::test] + async fn retained_append_result_distinguishes_eof_from_logical_exhaustion() { + let pool = BufferPool::new(BufferSize::try_new(8).unwrap(), 1); + let mut buffer = pool.acquire(); + buffer.copy_from_slice(b"discard22unused"); + buffer.expose_initialized_range_without_copying(7..9); + + let mut reader = std::io::Cursor::new(Vec::::new()); + let Some(appendable) = buffer.retained_append_permit() else { + panic!("retained prefix should have logical room to append"); + }; + + let AppendOutcome::Eof(buffer) = appendable.read(&mut reader).await.unwrap() else { + panic!("an exhausted reader should produce the typed EOF outcome"); + }; + assert_eq!(buffer.initialized(), 2); + } + + #[tokio::test] + async fn read_more_compacts_a_retained_prefix_when_tail_is_full() { + let pool = BufferPool::new(BufferSize::try_new(4096).unwrap(), 1); + let mut buffer = pool.acquire(); + buffer.copy_from_slice(&vec![b'x'; 4096]); + let allocation = buffer.allocation_ptr(); + buffer.expose_initialized_range_without_copying(4094..4096); + let (mut writer, mut reader) = tokio::io::duplex(64); + writer.write_all(b"y").await.unwrap(); + drop(writer); + + let Some(appendable) = buffer.retained_append_permit() else { + panic!("retained prefix should have logical room to append"); + }; + let AppendOutcome::Data(appended) = appendable.read(&mut reader).await.unwrap() else { + panic!("the reader has data available"); + }; + let read = appended.as_new_bytes().len(); + assert_eq!(appended.as_new_bytes(), b"y"); + let buffer = appended.into_inner(); + + assert_eq!(read, 1); + assert_eq!(buffer.as_ref(), b"xxy"); + assert_eq!(buffer.allocation_ptr(), allocation); + } + + #[tokio::test] + async fn retained_append_uses_a_tiny_physical_tail_without_compacting() { + let pool = BufferPool::new(BufferSize::try_new(4096).unwrap(), 1); + let mut buffer = pool.acquire(); + buffer.copy_from_slice(&vec![b'x'; 4095]); + let allocation = buffer.allocation_ptr(); + buffer.expose_initialized_range_without_copying(4063..4095); + let mut reader = std::io::Cursor::new(vec![b'y'; 4096]); + + let Some(permit) = buffer.retained_append_permit() else { + panic!("retained prefix should have logical room to append"); + }; + let AppendOutcome::Data(appended) = permit.read(&mut reader).await.unwrap() else { + panic!("the reader has data available"); + }; + + assert_eq!(appended.as_new_bytes(), b"y"); + let buffer = appended.into_inner(); + assert_eq!(buffer.initialized(), 33); + assert_eq!(buffer.allocation_ptr(), allocation); + assert_eq!(buffer.buffer.visible_from, 4063); + } + + #[tokio::test] + async fn cancelling_retained_append_after_compaction_preserves_visible_bytes() { + let pool = BufferPool::new(BufferSize::try_new(4096).unwrap(), 1); + let mut buffer = pool.acquire(); + buffer.copy_from_slice(&vec![b'x'; 4096]); + let allocation = buffer.allocation_ptr(); + buffer.expose_initialized_range_without_copying(4094..4096); + let (_writer, mut reader) = tokio::io::duplex(64); + + let Some(permit) = buffer.retained_append_permit() else { + panic!("retained prefix should have logical room to append"); + }; + let mut read = Box::pin(permit.read(&mut reader)); + match futures::poll!(read.as_mut()) { + std::task::Poll::Pending => {} + std::task::Poll::Ready(_) => panic!("empty open reader should remain pending"), + } + drop(read); + + assert_eq!(buffer.as_ref(), b"xx"); + assert_eq!(buffer.allocation_ptr(), allocation); + assert_eq!(buffer.buffer.visible_from, 0); + } + + #[tokio::test] + async fn retained_append_error_after_compaction_preserves_visible_bytes() { + struct ErrorReader; + + impl AsyncRead for ErrorReader { + fn poll_read( + self: Pin<&mut Self>, + _cx: &mut std::task::Context<'_>, + _buf: &mut ReadBuf<'_>, + ) -> std::task::Poll> { + std::task::Poll::Ready(Err(std::io::Error::from( + std::io::ErrorKind::ConnectionReset, + ))) + } + } + + let pool = BufferPool::new(BufferSize::try_new(4096).unwrap(), 1); + let mut buffer = pool.acquire(); + buffer.copy_from_slice(&vec![b'x'; 4096]); + let allocation = buffer.allocation_ptr(); + buffer.expose_initialized_range_without_copying(4094..4096); + + let Some(permit) = buffer.retained_append_permit() else { + panic!("retained prefix should have logical room to append"); + }; + let error = match permit.read(&mut ErrorReader).await { + Err(error) => error, + Ok(_) => panic!("error reader should fail the append"), + }; + + assert_eq!(error.kind(), std::io::ErrorKind::ConnectionReset); + assert_eq!(buffer.as_ref(), b"xx"); + assert_eq!(buffer.allocation_ptr(), allocation); + assert_eq!(buffer.buffer.visible_from, 0); } #[tokio::test] diff --git a/src/pool/mod.rs b/src/pool/mod.rs index 733e05bb..0ecba585 100644 --- a/src/pool/mod.rs +++ b/src/pool/mod.rs @@ -10,6 +10,7 @@ pub mod health_check; pub mod prewarming; pub mod provider; +pub(crate) use buffer::{AppendOutcome, RetainedAppendPermit}; pub use buffer::{ BufferPool, ChunkedResponse, HotPathAllocationMetricsSnapshot, PooledBuffer, hot_path_allocation_metrics_snapshot, reset_hot_path_allocation_metrics, diff --git a/src/session/multiline_framing.rs b/src/session/multiline_framing.rs index 8db4f953..d22bd110 100644 --- a/src/session/multiline_framing.rs +++ b/src/session/multiline_framing.rs @@ -273,12 +273,11 @@ struct IncompleteMultilineWireChunk { impl IncompleteMultilineWireChunk { #[allow(clippy::too_many_arguments)] - async fn compact_packed_prefix_and_write_with_next_backend_chunk( + async fn append_to_packed_prefix_and_write_with_next_backend_chunk( self, writer: &mut W, - current_len: usize, framer: &mut MultilineFramer, - io_buffer: &mut crate::pool::PooledBuffer, + io_buffer: crate::pool::RetainedAppendPermit<'_>, conn: &mut crate::stream::ConnectionStream, pool: &crate::pool::BufferPool, backend_id: crate::types::BackendId, @@ -287,26 +286,34 @@ impl IncompleteMultilineWireChunk { W: AsyncWrite + Unpin, { let initial_response_start = self.response.start; - let n = io_buffer.read_more(conn).await.map_err(|e| { + let appended = io_buffer.read(conn).await.map_err(|e| { crate::session::response_transfer::ResponseTransferError::Io( anyhow::Error::from(e).context("Failed to read remaining response body"), ) })?; - if n == 0 { - return Err( - crate::session::response_transfer::ResponseTransferError::BackendEof { - backend_id, - bytes_received: self.response.len() as u64, - }, - ); - } - let total_len = current_len + n; + let appended = match appended { + crate::pool::AppendOutcome::Data(appended) => appended, + crate::pool::AppendOutcome::Eof(buffer) => { + let _ = buffer; + return Err( + crate::session::response_transfer::ResponseTransferError::BackendEof { + backend_id, + bytes_received: self.response.len() as u64, + }, + ); + } + }; + let previous_len = appended.previous_len(); + let total_len = appended.total_len(); + let frame = framer.frame_next_multiline_chunk(self, appended.as_new_bytes()); + let io_buffer = appended.into_inner(); - match framer.frame_next_multiline_chunk(self, &io_buffer[current_len..total_len]) { + match frame { FramedMultilineChunk::Complete(complete) => { - let combined_response = initial_response_start..current_len + complete.response.end; - let combined_next_response = current_len + complete.next_response_input.start - ..current_len + complete.next_response_input.end; + let combined_response = + initial_response_start..previous_len + complete.response.end; + let combined_next_response = previous_len + complete.next_response_input.start + ..previous_len + complete.next_response_input.end; write_response_chunk_preserving_suffix_on_error( writer, io_buffer, @@ -320,7 +327,7 @@ impl IncompleteMultilineWireChunk { } FramedMultilineChunk::Incomplete(incomplete) => { let combined_response = - initial_response_start..current_len + incomplete.response.end; + initial_response_start..previous_len + incomplete.response.end; let bytes_received = combined_response.len() as u64; if let Err(error) = writer .write_all(&io_buffer[combined_response.clone()]) @@ -1351,12 +1358,10 @@ where .await? } FramedMultilineChunk::Incomplete(incomplete) => { - if io_buffer.is_exposed_range_view() && io_buffer.has_remaining_fixed_writable_region() - { + if let Some(io_buffer) = io_buffer.retained_append_permit() { incomplete - .compact_packed_prefix_and_write_with_next_backend_chunk( + .append_to_packed_prefix_and_write_with_next_backend_chunk( writer, - initial_len, &mut framer, io_buffer, conn,