use crate::error::Result;
use std::fs::File;
use std::io;
use std::path::Path;
use std::sync::Arc;
#[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
}
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)
}
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]);
let n = r.read_at(99_000, &mut buf).unwrap();
assert_eq!(n, 1000);
assert_eq!(&buf[..n], &data[99_000..]);
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);
}
}