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
153pub struct OffsetReader<R> {
159 inner: R,
160 base: u64,
161}
162
163impl<R> OffsetReader<R> {
164 pub fn new(inner: R, base: u64) -> Self {
165 Self { inner, base }
166 }
167 pub fn base(&self) -> u64 {
168 self.base
169 }
170 pub fn into_inner(self) -> R {
171 self.inner
172 }
173}
174
175impl<R: RangeReader> RangeReader for OffsetReader<R> {
176 fn len(&self) -> u64 {
177 self.inner.len().saturating_sub(self.base)
178 }
179
180 fn read_at(&self, offset: u64, len: u64) -> std::io::Result<Vec<u8>> {
181 self.inner.read_at(self.base + offset, len)
182 }
183
184 fn read_many(&self, ranges: &[(u64, u64)]) -> std::io::Result<Vec<Vec<u8>>> {
185 let shifted: Vec<(u64, u64)> = ranges.iter().map(|&(o, l)| (self.base + o, l)).collect();
186 self.inner.read_many(&shifted)
187 }
188
189 fn concurrency(&self) -> usize {
190 self.inner.concurrency()
191 }
192}
193
194pub const POLYGLOT_MARKER: &[u8] = b"RETE-BASE:";
200pub const POLYGLOT_DIGITS: usize = 16;
202
203pub fn detect_polyglot_base(head: &[u8]) -> Option<u64> {
206 let pos = head
207 .windows(POLYGLOT_MARKER.len())
208 .position(|w| w == POLYGLOT_MARKER)?;
209 let start = pos + POLYGLOT_MARKER.len();
210 let digits = head.get(start..start + POLYGLOT_DIGITS)?;
211 std::str::from_utf8(digits).ok()?.parse::<u64>().ok()
212}
213
214#[cfg(test)]
215mod tests {
216 use super::*;
217
218 #[test]
219 fn slice_reader_serves_ranges_and_bounds_check() {
220 let data = (0u8..=255).collect::<Vec<_>>();
221 let r = SliceReader::new(&data);
222 assert_eq!(r.len(), 256);
223 assert_eq!(r.read_at(10, 4).unwrap(), vec![10, 11, 12, 13]);
224 assert!(r.read_at(254, 10).is_err()); }
226
227 #[test]
228 fn counting_reader_tallies() {
229 let data = vec![0u8; 100];
230 let r = CountingReader::new(SliceReader::new(&data));
231 r.read_at(0, 10).unwrap();
232 r.read_at(50, 20).unwrap();
233 assert_eq!(r.requests(), 2);
234 assert_eq!(r.bytes_read(), 30);
235 }
236}