runsync-transfer 2026.1.0

High-throughput P2P file transfer engine: adaptive compression, end-to-end AEAD, parallel chunked pipeline over QUIC or any async transport.
Documentation
//! Source-side reads.

use crate::error::Result;
use std::fs::File;
use std::io;
use std::path::Path;
use std::sync::Arc;

/// A file opened for concurrent positional reads.
///
/// Cloning is cheap and shares one descriptor; `pread` carries its own offset,
/// so concurrent readers never contend on a file position.
#[derive(Clone)]
pub struct ChunkReader {
    file: Arc<File>,
    len: u64,
}

impl ChunkReader {
    pub fn open(path: &Path) -> Result<Self> {
        let file = File::open(path)
            .map_err(|e| io::Error::new(e.kind(), format!("{}: {e}", path.display())))?;
        let len = file.metadata()?.len();
        let r = Self {
            file: Arc::new(file),
            len,
        };
        r.advise_sequential();
        Ok(r)
    }

    pub fn len(&self) -> u64 {
        self.len
    }

    pub fn is_empty(&self) -> bool {
        self.len == 0
    }

    /// Fill `buf` from `offset`, returning the number of bytes read.
    ///
    /// A short read is only accepted at end of file; anywhere else it means the
    /// file changed underneath us, which is an error rather than silent
    /// truncation of the transfer.
    pub fn read_at(&self, offset: u64, buf: &mut [u8]) -> io::Result<usize> {
        let mut total = 0usize;
        while total < buf.len() {
            let n = pread(&self.file, &mut buf[total..], offset + total as u64)?;
            if n == 0 {
                break;
            }
            total += n;
        }
        Ok(total)
    }

    /// Hint that access is sequential. Lets the kernel read ahead aggressively,
    /// which is a large win on rotational and network-backed storage and free
    /// elsewhere. Advisory: failure is not an error.
    fn advise_sequential(&self) {
        #[cfg(target_os = "linux")]
        {
            use std::os::unix::io::AsRawFd;
            unsafe {
                libc::posix_fadvise(self.file.as_raw_fd(), 0, 0, libc::POSIX_FADV_SEQUENTIAL);
            }
        }
        #[cfg(target_os = "macos")]
        {
            use std::os::unix::io::AsRawFd;
            unsafe {
                libc::fcntl(self.file.as_raw_fd(), libc::F_RDAHEAD, 1);
            }
        }
    }
}

#[cfg(unix)]
fn pread(file: &File, buf: &mut [u8], offset: u64) -> io::Result<usize> {
    use std::os::unix::fs::FileExt;
    file.read_at(buf, offset)
}

#[cfg(windows)]
fn pread(file: &File, buf: &mut [u8], offset: u64) -> io::Result<usize> {
    use std::os::windows::fs::FileExt;
    file.seek_read(buf, offset)
}

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

    #[test]
    fn reads_at_arbitrary_offsets() {
        let mut f = tempfile::NamedTempFile::new().unwrap();
        let data: Vec<u8> = (0..=255u8).cycle().take(100_000).collect();
        f.write_all(&data).unwrap();
        f.flush().unwrap();

        let r = ChunkReader::open(f.path()).unwrap();
        assert_eq!(r.len(), 100_000);

        let mut buf = vec![0u8; 4096];
        let n = r.read_at(50_000, &mut buf).unwrap();
        assert_eq!(n, 4096);
        assert_eq!(&buf[..], &data[50_000..54_096]);

        // Tail: short read is correct here and only here.
        let n = r.read_at(99_000, &mut buf).unwrap();
        assert_eq!(n, 1000);
        assert_eq!(&buf[..n], &data[99_000..]);

        // Past the end.
        assert_eq!(r.read_at(200_000, &mut buf).unwrap(), 0);
    }

    #[test]
    fn concurrent_reads_do_not_interfere() {
        let mut f = tempfile::NamedTempFile::new().unwrap();
        let data: Vec<u8> = (0..1_000_000u32).map(|i| (i % 251) as u8).collect();
        f.write_all(&data).unwrap();
        f.flush().unwrap();

        let r = ChunkReader::open(f.path()).unwrap();
        let handles: Vec<_> = (0..8)
            .map(|i| {
                let r = r.clone();
                let expect = data.clone();
                std::thread::spawn(move || {
                    for round in 0..50 {
                        let off = ((i * 50 + round) * 2048) as u64 % 900_000;
                        let mut buf = vec![0u8; 2048];
                        let n = r.read_at(off, &mut buf).unwrap();
                        assert_eq!(&buf[..n], &expect[off as usize..off as usize + n]);
                    }
                })
            })
            .collect();
        for h in handles {
            h.join().unwrap();
        }
    }

    #[test]
    fn empty_file() {
        let f = tempfile::NamedTempFile::new().unwrap();
        let r = ChunkReader::open(f.path()).unwrap();
        assert!(r.is_empty());
        let mut buf = [0u8; 16];
        assert_eq!(r.read_at(0, &mut buf).unwrap(), 0);
    }
}