Skip to main content

runsync_transfer/io/
reader.rs

1//! Source-side reads.
2
3use crate::error::Result;
4use std::fs::File;
5use std::io;
6use std::path::Path;
7use std::sync::Arc;
8
9/// A file opened for concurrent positional reads.
10///
11/// Cloning is cheap and shares one descriptor; `pread` carries its own offset,
12/// so concurrent readers never contend on a file position.
13#[derive(Clone)]
14pub struct ChunkReader {
15    file: Arc<File>,
16    len: u64,
17}
18
19impl ChunkReader {
20    pub fn open(path: &Path) -> Result<Self> {
21        let file = File::open(path)
22            .map_err(|e| io::Error::new(e.kind(), format!("{}: {e}", path.display())))?;
23        let len = file.metadata()?.len();
24        let r = Self {
25            file: Arc::new(file),
26            len,
27        };
28        r.advise_sequential();
29        Ok(r)
30    }
31
32    pub fn len(&self) -> u64 {
33        self.len
34    }
35
36    pub fn is_empty(&self) -> bool {
37        self.len == 0
38    }
39
40    /// Fill `buf` from `offset`, returning the number of bytes read.
41    ///
42    /// A short read is only accepted at end of file; anywhere else it means the
43    /// file changed underneath us, which is an error rather than silent
44    /// truncation of the transfer.
45    pub fn read_at(&self, offset: u64, buf: &mut [u8]) -> io::Result<usize> {
46        let mut total = 0usize;
47        while total < buf.len() {
48            let n = pread(&self.file, &mut buf[total..], offset + total as u64)?;
49            if n == 0 {
50                break;
51            }
52            total += n;
53        }
54        Ok(total)
55    }
56
57    /// Hint that access is sequential. Lets the kernel read ahead aggressively,
58    /// which is a large win on rotational and network-backed storage and free
59    /// elsewhere. Advisory: failure is not an error.
60    fn advise_sequential(&self) {
61        #[cfg(target_os = "linux")]
62        {
63            use std::os::unix::io::AsRawFd;
64            unsafe {
65                libc::posix_fadvise(self.file.as_raw_fd(), 0, 0, libc::POSIX_FADV_SEQUENTIAL);
66            }
67        }
68        #[cfg(target_os = "macos")]
69        {
70            use std::os::unix::io::AsRawFd;
71            unsafe {
72                libc::fcntl(self.file.as_raw_fd(), libc::F_RDAHEAD, 1);
73            }
74        }
75    }
76}
77
78#[cfg(unix)]
79fn pread(file: &File, buf: &mut [u8], offset: u64) -> io::Result<usize> {
80    use std::os::unix::fs::FileExt;
81    file.read_at(buf, offset)
82}
83
84#[cfg(windows)]
85fn pread(file: &File, buf: &mut [u8], offset: u64) -> io::Result<usize> {
86    use std::os::windows::fs::FileExt;
87    file.seek_read(buf, offset)
88}
89
90#[cfg(test)]
91mod tests {
92    use super::*;
93    use std::io::Write;
94
95    #[test]
96    fn reads_at_arbitrary_offsets() {
97        let mut f = tempfile::NamedTempFile::new().unwrap();
98        let data: Vec<u8> = (0..=255u8).cycle().take(100_000).collect();
99        f.write_all(&data).unwrap();
100        f.flush().unwrap();
101
102        let r = ChunkReader::open(f.path()).unwrap();
103        assert_eq!(r.len(), 100_000);
104
105        let mut buf = vec![0u8; 4096];
106        let n = r.read_at(50_000, &mut buf).unwrap();
107        assert_eq!(n, 4096);
108        assert_eq!(&buf[..], &data[50_000..54_096]);
109
110        // Tail: short read is correct here and only here.
111        let n = r.read_at(99_000, &mut buf).unwrap();
112        assert_eq!(n, 1000);
113        assert_eq!(&buf[..n], &data[99_000..]);
114
115        // Past the end.
116        assert_eq!(r.read_at(200_000, &mut buf).unwrap(), 0);
117    }
118
119    #[test]
120    fn concurrent_reads_do_not_interfere() {
121        let mut f = tempfile::NamedTempFile::new().unwrap();
122        let data: Vec<u8> = (0..1_000_000u32).map(|i| (i % 251) as u8).collect();
123        f.write_all(&data).unwrap();
124        f.flush().unwrap();
125
126        let r = ChunkReader::open(f.path()).unwrap();
127        let handles: Vec<_> = (0..8)
128            .map(|i| {
129                let r = r.clone();
130                let expect = data.clone();
131                std::thread::spawn(move || {
132                    for round in 0..50 {
133                        let off = ((i * 50 + round) * 2048) as u64 % 900_000;
134                        let mut buf = vec![0u8; 2048];
135                        let n = r.read_at(off, &mut buf).unwrap();
136                        assert_eq!(&buf[..n], &expect[off as usize..off as usize + n]);
137                    }
138                })
139            })
140            .collect();
141        for h in handles {
142            h.join().unwrap();
143        }
144    }
145
146    #[test]
147    fn empty_file() {
148        let f = tempfile::NamedTempFile::new().unwrap();
149        let r = ChunkReader::open(f.path()).unwrap();
150        assert!(r.is_empty());
151        let mut buf = [0u8; 16];
152        assert_eq!(r.read_at(0, &mut buf).unwrap(), 0);
153    }
154}