#[cfg(test)]
mod tests {
use std::sync::Arc;
use std::sync::Barrier;
use std::thread;
use crate::{
buffer::REGISTERED_BUFFER_CAPACITY, buffer::REGISTERED_BUFFER_COUNT, Buffer, BufferPool,
Completed, Owned, Submitted,
};
#[test]
fn pool_rejects_zero_capacity() {
assert!(BufferPool::new(0, 1).is_err());
}
#[test]
fn pool_rejects_zero_buffers() {
assert!(BufferPool::new(4, 0).is_err());
}
#[test]
fn acquire_returns_capacity_sized_buffer() {
let pool = BufferPool::new(32, 1).expect("pool");
let buffer = pool.acquire().expect("buffer");
assert_eq!(buffer.capacity(), 32);
}
#[test]
fn owned_to_submitted_to_completed_cycle() {
let pool = BufferPool::new(16, 1).expect("pool");
let mut owned = pool.acquire().expect("owned");
owned.as_mut_slice()[..3].copy_from_slice(b"abc");
let owned = owned.set_filled_len(3).expect("filled");
let submitted: Buffer<Submitted> = owned.into_submitted();
let completed: Buffer<Completed> = submitted.into_completed(3).expect("done");
assert_eq!(completed.filled(), b"abc");
let released: Buffer<Owned> = completed.release();
assert_eq!(released.filled_len(), 0);
}
#[test]
fn set_filled_len_rejects_oversize() {
let pool = BufferPool::new(8, 1).expect("pool");
let owned = pool.acquire().expect("owned");
assert!(owned.set_filled_len(9).is_err());
}
#[test]
fn completed_rejects_oversize_length() {
let pool = BufferPool::new(8, 1).expect("pool");
let owned = pool.acquire().expect("owned");
let submitted = owned.into_submitted();
assert!(submitted.into_completed(9).is_err());
}
#[test]
fn concurrent_acquire_respects_max_buffers() {
let pool = Arc::new(BufferPool::new(32, 2).expect("pool"));
let barrier = Arc::new(Barrier::new(8));
let mut joins = Vec::new();
for _ in 0..8 {
let pool = pool.clone();
let barrier = barrier.clone();
joins.push(thread::spawn(move || {
let acquired = pool.acquire().ok();
barrier.wait();
acquired.is_some()
}));
}
let successes = joins
.into_iter()
.map(|join: thread::JoinHandle<bool>| join.join().expect("join"))
.filter(|success| *success)
.count();
assert_eq!(successes, 2);
}
#[test]
fn registered_pool_round_trip_preserves_shape() {
let pool = BufferPool::registered().expect("pool");
let completed = pool
.acquire()
.expect("owned")
.into_submitted()
.into_completed(0)
.expect("completed");
let released = completed.release();
assert_eq!(released.capacity(), REGISTERED_BUFFER_CAPACITY);
assert_eq!(pool.max_buffers(), REGISTERED_BUFFER_COUNT);
}
#[test]
fn registered_pool_exhaustion_falls_back_to_heap_buffer() {
let pool = BufferPool::registered().expect("pool");
let mut buffers = Vec::new();
for _ in 0..REGISTERED_BUFFER_COUNT {
buffers.push(pool.acquire().expect("pooled acquire"));
}
let overflow = pool.acquire().expect("overflow acquire");
assert_eq!(overflow.capacity(), REGISTERED_BUFFER_CAPACITY);
drop(overflow);
drop(buffers);
for _ in 0..REGISTERED_BUFFER_COUNT {
let buffer = pool.acquire().expect("reacquire");
assert_eq!(buffer.capacity(), REGISTERED_BUFFER_CAPACITY);
}
}
#[test]
fn registered_pool_four_threads_get_and_release() {
let pool = Arc::new(BufferPool::registered().expect("pool"));
let mut joins = Vec::new();
for _ in 0..4 {
let pool = pool.clone();
joins.push(thread::spawn(move || {
for _ in 0..512 {
let buffer = pool.acquire().expect("acquire");
drop(buffer);
}
}));
}
for join in joins {
join.join().expect("join");
}
}
}