1use std::sync::atomic::{AtomicU64, Ordering};
10
11pub trait RangeReader {
13 fn len(&self) -> u64;
15
16 fn is_empty(&self) -> bool {
18 self.len() == 0
19 }
20
21 fn read_at(&self, offset: u64, len: u64) -> std::io::Result<Vec<u8>>;
24
25 fn read_many(&self, ranges: &[(u64, u64)]) -> std::io::Result<Vec<Vec<u8>>> {
31 ranges
32 .iter()
33 .map(|&(offset, len)| self.read_at(offset, len))
34 .collect()
35 }
36
37 fn concurrency(&self) -> usize {
43 1
44 }
45}
46
47impl<R: RangeReader + ?Sized> RangeReader for std::sync::Arc<R> {
50 fn len(&self) -> u64 {
51 (**self).len()
52 }
53
54 fn read_at(&self, offset: u64, len: u64) -> std::io::Result<Vec<u8>> {
55 (**self).read_at(offset, len)
56 }
57
58 fn read_many(&self, ranges: &[(u64, u64)]) -> std::io::Result<Vec<Vec<u8>>> {
59 (**self).read_many(ranges)
60 }
61
62 fn concurrency(&self) -> usize {
63 (**self).concurrency()
64 }
65}
66
67pub struct SliceReader<'a> {
69 data: &'a [u8],
70}
71
72impl<'a> SliceReader<'a> {
73 pub fn new(data: &'a [u8]) -> Self {
74 Self { data }
75 }
76}
77
78impl RangeReader for SliceReader<'_> {
79 fn len(&self) -> u64 {
80 self.data.len() as u64
81 }
82
83 fn read_at(&self, offset: u64, len: u64) -> std::io::Result<Vec<u8>> {
84 let start = offset as usize;
85 let end = start
86 .checked_add(len as usize)
87 .filter(|&e| e <= self.data.len())
88 .ok_or_else(|| {
89 std::io::Error::new(std::io::ErrorKind::UnexpectedEof, "range out of bounds")
90 })?;
91 Ok(self.data[start..end].to_vec())
92 }
93}
94
95pub struct CountingReader<R> {
100 inner: R,
101 requests: AtomicU64,
102 bytes: AtomicU64,
103}
104
105impl<R: RangeReader> CountingReader<R> {
106 pub fn new(inner: R) -> Self {
107 Self {
108 inner,
109 requests: AtomicU64::new(0),
110 bytes: AtomicU64::new(0),
111 }
112 }
113
114 pub fn requests(&self) -> u64 {
116 self.requests.load(Ordering::Relaxed)
117 }
118
119 pub fn bytes_read(&self) -> u64 {
121 self.bytes.load(Ordering::Relaxed)
122 }
123}
124
125impl<R: RangeReader> RangeReader for CountingReader<R> {
126 fn len(&self) -> u64 {
127 self.inner.len()
128 }
129
130 fn read_at(&self, offset: u64, len: u64) -> std::io::Result<Vec<u8>> {
131 let out = self.inner.read_at(offset, len)?;
132 self.requests.fetch_add(1, Ordering::Relaxed);
133 self.bytes.fetch_add(out.len() as u64, Ordering::Relaxed);
134 Ok(out)
135 }
136
137 fn read_many(&self, ranges: &[(u64, u64)]) -> std::io::Result<Vec<Vec<u8>>> {
141 let out = self.inner.read_many(ranges)?;
142 self.requests.fetch_add(out.len() as u64, Ordering::Relaxed);
143 self.bytes
144 .fetch_add(out.iter().map(|b| b.len() as u64).sum(), Ordering::Relaxed);
145 Ok(out)
146 }
147
148 fn concurrency(&self) -> usize {
149 self.inner.concurrency()
150 }
151}
152
153#[cfg(test)]
154mod tests {
155 use super::*;
156
157 #[test]
158 fn slice_reader_serves_ranges_and_bounds_check() {
159 let data = (0u8..=255).collect::<Vec<_>>();
160 let r = SliceReader::new(&data);
161 assert_eq!(r.len(), 256);
162 assert_eq!(r.read_at(10, 4).unwrap(), vec![10, 11, 12, 13]);
163 assert!(r.read_at(254, 10).is_err()); }
165
166 #[test]
167 fn counting_reader_tallies() {
168 let data = vec![0u8; 100];
169 let r = CountingReader::new(SliceReader::new(&data));
170 r.read_at(0, 10).unwrap();
171 r.read_at(50, 20).unwrap();
172 assert_eq!(r.requests(), 2);
173 assert_eq!(r.bytes_read(), 30);
174 }
175}