runsync_transfer/io/
reader.rs1use crate::error::Result;
4use std::fs::File;
5use std::io;
6use std::path::Path;
7use std::sync::Arc;
8
9#[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 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 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 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 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}