runsync-transfer 0.1.0

High-throughput P2P file transfer engine: adaptive compression, end-to-end AEAD, parallel chunked pipeline over QUIC or any async transport.
Documentation
//! Chunk buffer recycling.
//!
//! A 100 GB transfer at 1 MiB chunks is 100k chunks, each needing a plaintext
//! buffer and an output buffer. Allocating and freeing 200k multi-megabyte
//! buffers hands the allocator a workload it will happily serve and then return
//! to the OS, so the process spends its time in `madvise`/page faults instead of
//! moving bytes. Recycling keeps a small fixed set of hot buffers instead.

use parking_lot::Mutex;
use std::sync::Arc;

/// Fixed-capacity stack of reusable byte buffers.
pub struct BufPool {
    free: Mutex<Vec<Vec<u8>>>,
    capacity: usize,
    buf_hint: usize,
}

impl BufPool {
    /// `capacity` buffers retained at most, each pre-sized to `buf_hint`.
    pub fn new(capacity: usize, buf_hint: usize) -> Arc<Self> {
        Arc::new(Self {
            free: Mutex::new(Vec::with_capacity(capacity)),
            capacity,
            buf_hint,
        })
    }

    /// Take a buffer. Always returned empty, with `buf_hint` capacity available.
    pub fn take(&self) -> Vec<u8> {
        match self.free.lock().pop() {
            Some(mut b) => {
                b.clear();
                b
            }
            None => Vec::with_capacity(self.buf_hint),
        }
    }

    /// Return a buffer for reuse.
    pub fn put(&self, buf: Vec<u8>) {
        // Do not retain a buffer that grew far past the hint: a single
        // pathological chunk should not pin an oversized allocation forever.
        if buf.capacity() > self.buf_hint * 4 {
            return;
        }
        let mut free = self.free.lock();
        if free.len() < self.capacity {
            free.push(buf);
        }
    }
}

/// A buffer that returns itself to its pool when dropped.
pub struct Pooled {
    buf: Option<Vec<u8>>,
    pool: Arc<BufPool>,
}

impl Pooled {
    pub fn new(pool: Arc<BufPool>) -> Self {
        let buf = pool.take();
        Self {
            buf: Some(buf),
            pool,
        }
    }
}

impl std::ops::Deref for Pooled {
    type Target = Vec<u8>;
    fn deref(&self) -> &Vec<u8> {
        self.buf.as_ref().expect("buffer taken only on drop")
    }
}

impl std::ops::DerefMut for Pooled {
    fn deref_mut(&mut self) -> &mut Vec<u8> {
        self.buf.as_mut().expect("buffer taken only on drop")
    }
}

impl Drop for Pooled {
    fn drop(&mut self) {
        if let Some(b) = self.buf.take() {
            self.pool.put(b);
        }
    }
}

/// A small stack of reusable objects, handed out to workers and returned when
/// they finish.
///
/// The sealer this holds carries a per-file subkey cache. Sharing one behind a
/// mutex would serialise every encode on a stream — including the file read
/// that happens under the same critical section. Pooling gives each concurrent
/// worker its own, with the caches surviving between chunks.
pub struct ObjPool<T> {
    free: Mutex<Vec<T>>,
    capacity: usize,
}

impl<T> ObjPool<T> {
    pub fn new(capacity: usize) -> Arc<Self> {
        Arc::new(Self {
            free: Mutex::new(Vec::with_capacity(capacity)),
            capacity,
        })
    }

    /// Take an object, or build one with `make` if the pool is empty.
    pub fn take_or<F: FnOnce() -> T>(&self, make: F) -> T {
        match self.free.lock().pop() {
            Some(o) => o,
            None => make(),
        }
    }

    pub fn put(&self, obj: T) {
        let mut free = self.free.lock();
        if free.len() < self.capacity {
            free.push(obj);
        }
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn obj_pool_reuses_and_bounds() {
        let pool: Arc<ObjPool<Vec<u32>>> = ObjPool::new(2);
        let made = std::sync::atomic::AtomicUsize::new(0);
        let mk = || {
            made.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
            Vec::new()
        };
        let a = pool.take_or(mk);
        let b = pool.take_or(mk);
        assert_eq!(made.load(std::sync::atomic::Ordering::SeqCst), 2);
        pool.put(a);
        pool.put(b);
        // Both come back from the pool rather than being rebuilt.
        let _ = pool.take_or(mk);
        let _ = pool.take_or(mk);
        assert_eq!(made.load(std::sync::atomic::Ordering::SeqCst), 2);
        // Over capacity, extras are dropped rather than retained.
        for _ in 0..8 {
            pool.put(Vec::new());
        }
        assert_eq!(pool.free.lock().len(), 2);
    }

    #[test]
    fn buffers_are_recycled_not_reallocated() {
        let pool = BufPool::new(4, 1024);
        let ptr = {
            let mut b = Pooled::new(pool.clone());
            b.extend_from_slice(&[1u8; 512]);
            b.as_ptr()
        };
        let b2 = Pooled::new(pool.clone());
        assert!(b2.is_empty(), "recycled buffers come back empty");
        assert_eq!(b2.as_ptr(), ptr, "the same allocation came back");
    }

    #[test]
    fn pool_respects_its_capacity() {
        let pool = BufPool::new(2, 128);
        let bufs: Vec<_> = (0..8).map(|_| Pooled::new(pool.clone())).collect();
        drop(bufs);
        assert_eq!(pool.free.lock().len(), 2);
    }

    #[test]
    fn oversized_buffers_are_dropped_rather_than_retained() {
        let pool = BufPool::new(4, 100);
        {
            let mut b = Pooled::new(pool.clone());
            b.resize(10_000, 0);
        }
        assert_eq!(pool.free.lock().len(), 0);
    }

    #[test]
    fn concurrent_take_and_put_is_sound() {
        let pool = BufPool::new(8, 4096);
        let threads: Vec<_> = (0..8)
            .map(|_| {
                let pool = pool.clone();
                std::thread::spawn(move || {
                    for _ in 0..1000 {
                        let mut b = Pooled::new(pool.clone());
                        b.extend_from_slice(&[7u8; 100]);
                        assert_eq!(b.len(), 100);
                    }
                })
            })
            .collect();
        for t in threads {
            t.join().unwrap();
        }
        assert!(pool.free.lock().len() <= 8);
    }
}