1use crate::block::BlockRead;
26use crate::error::Error;
27use std::io::{self, Read, Seek, SeekFrom};
28
29pub struct BlockReadStreamer<T: BlockRead> {
36 inner: T,
37 pos: u64,
38}
39
40impl<T: BlockRead> BlockReadStreamer<T> {
41 pub fn new(inner: T) -> Self {
43 Self { inner, pos: 0 }
44 }
45
46 pub fn with_position(inner: T, pos: u64) -> Self {
50 Self { inner, pos }
51 }
52
53 pub fn position(&self) -> u64 {
55 self.pos
56 }
57
58 pub fn get_ref(&self) -> &T {
60 &self.inner
61 }
62
63 pub fn into_inner(self) -> T {
65 self.inner
66 }
67}
68
69impl<T: BlockRead> Read for BlockReadStreamer<T> {
70 fn read(&mut self, buf: &mut [u8]) -> io::Result<usize> {
71 let size = self.inner.size_bytes();
72 if self.pos >= size {
73 return Ok(0);
74 }
75 let remaining = size - self.pos;
76 let n = std::cmp::min(buf.len() as u64, remaining) as usize;
77 if n == 0 {
78 return Ok(0);
79 }
80 self.inner
81 .read_at(self.pos, &mut buf[..n])
82 .map_err(fs_core_error_to_io)?;
83 self.pos += n as u64;
84 Ok(n)
85 }
86}
87
88impl<T: BlockRead> Seek for BlockReadStreamer<T> {
89 fn seek(&mut self, pos: SeekFrom) -> io::Result<u64> {
90 let new_pos = match pos {
95 SeekFrom::Start(n) => n,
96 SeekFrom::End(n) => offset_from(self.inner.size_bytes(), n)?,
97 SeekFrom::Current(n) => offset_from(self.pos, n)?,
98 };
99 self.pos = new_pos;
100 Ok(new_pos)
101 }
102}
103
104fn offset_from(base: u64, delta: i64) -> io::Result<u64> {
105 if delta >= 0 {
106 base.checked_add(delta as u64)
107 .ok_or_else(|| io::Error::new(io::ErrorKind::InvalidInput, "seek offset overflows u64"))
108 } else {
109 let abs = delta.unsigned_abs();
110 base.checked_sub(abs).ok_or_else(|| {
111 io::Error::new(
112 io::ErrorKind::InvalidInput,
113 "seek would place cursor before byte 0",
114 )
115 })
116 }
117}
118
119fn fs_core_error_to_io(e: Error) -> io::Error {
120 match e {
121 Error::Io(io) => io,
122 Error::ShortRead { offset, want, got } => io::Error::new(
123 io::ErrorKind::UnexpectedEof,
124 format!("short read at {offset}: wanted {want} got {got}"),
125 ),
126 Error::OutOfBounds { offset, len, size } => io::Error::new(
127 io::ErrorKind::UnexpectedEof,
128 format!("{offset}+{len} past device size {size}"),
129 ),
130 Error::ReadOnly => io::Error::new(io::ErrorKind::PermissionDenied, "device is read-only"),
131 Error::Custom(s) => io::Error::other(s),
132 }
133}
134
135#[cfg(test)]
136mod tests {
137 use super::*;
138 use crate::error::Result as FsResult;
139 use crate::test_device::Bytes;
140 use std::sync::{Arc, Mutex};
141
142 struct AlwaysFails;
145 impl BlockRead for AlwaysFails {
146 fn read_at(&self, _offset: u64, _buf: &mut [u8]) -> FsResult<()> {
147 Err(Error::Custom("simulated failure".into()))
148 }
149 fn size_bytes(&self) -> u64 {
150 1024
151 }
152 }
153
154 fn fixture() -> Bytes {
155 let mut v = vec![0u8; 32];
156 for (i, b) in v.iter_mut().enumerate() {
157 *b = i as u8;
158 }
159 Bytes(Mutex::new(v))
160 }
161
162 #[test]
163 fn read_to_end_returns_full_contents() {
164 let mut s = BlockReadStreamer::new(fixture());
165 let mut out = Vec::new();
166 let n = s.read_to_end(&mut out).unwrap();
167 assert_eq!(n, 32);
168 assert_eq!(out.len(), 32);
169 assert_eq!(out[0], 0);
170 assert_eq!(out[31], 31);
171 }
172
173 #[test]
174 fn partial_end_read_is_clamped_not_errored() {
175 let mut s = BlockReadStreamer::with_position(fixture(), 30);
176 let mut buf = [0u8; 16];
177 let n = s.read(&mut buf).unwrap();
178 assert_eq!(n, 2);
179 assert_eq!(&buf[..2], &[30, 31]);
180 assert_eq!(s.position(), 32);
181 }
182
183 #[test]
184 fn read_at_eof_returns_zero() {
185 let mut s = BlockReadStreamer::with_position(fixture(), 32);
186 let mut buf = [0u8; 8];
187 assert_eq!(s.read(&mut buf).unwrap(), 0);
188 assert_eq!(s.read(&mut buf).unwrap(), 0);
190 }
191
192 #[test]
193 fn read_past_eof_position_returns_zero() {
194 let mut s = BlockReadStreamer::with_position(fixture(), 9_999);
195 let mut buf = [0u8; 8];
196 assert_eq!(s.read(&mut buf).unwrap(), 0);
197 }
198
199 #[test]
200 fn zero_length_buf_returns_zero() {
201 let mut s = BlockReadStreamer::new(fixture());
202 let mut buf: [u8; 0] = [];
203 assert_eq!(s.read(&mut buf).unwrap(), 0);
204 assert_eq!(s.position(), 0);
205 }
206
207 #[test]
208 fn position_advances_after_read() {
209 let mut s = BlockReadStreamer::new(fixture());
210 let mut buf = [0u8; 4];
211 s.read_exact(&mut buf).unwrap();
212 assert_eq!(s.position(), 4);
213 assert_eq!(buf, [0, 1, 2, 3]);
214
215 s.read_exact(&mut buf).unwrap();
216 assert_eq!(s.position(), 8);
217 assert_eq!(buf, [4, 5, 6, 7]);
218 }
219
220 #[test]
221 fn seek_start_jumps_absolute() {
222 let mut s = BlockReadStreamer::new(fixture());
223 let p = s.seek(SeekFrom::Start(10)).unwrap();
224 assert_eq!(p, 10);
225 assert_eq!(s.position(), 10);
226 let mut buf = [0u8; 2];
227 s.read_exact(&mut buf).unwrap();
228 assert_eq!(buf, [10, 11]);
229 }
230
231 #[test]
232 fn seek_end_jumps_relative_to_size() {
233 let mut s = BlockReadStreamer::new(fixture());
234 let p = s.seek(SeekFrom::End(-4)).unwrap();
235 assert_eq!(p, 28);
236 let mut buf = [0u8; 4];
237 s.read_exact(&mut buf).unwrap();
238 assert_eq!(buf, [28, 29, 30, 31]);
239 }
240
241 #[test]
242 fn seek_current_jumps_relative_to_cursor() {
243 let mut s = BlockReadStreamer::with_position(fixture(), 10);
244 let p = s.seek(SeekFrom::Current(5)).unwrap();
245 assert_eq!(p, 15);
246 let p = s.seek(SeekFrom::Current(-3)).unwrap();
247 assert_eq!(p, 12);
248 }
249
250 #[test]
251 fn seek_past_end_is_allowed_then_read_returns_zero() {
252 let mut s = BlockReadStreamer::new(fixture());
253 assert_eq!(s.seek(SeekFrom::Start(1_000_000)).unwrap(), 1_000_000);
254 let mut buf = [0u8; 4];
255 assert_eq!(s.read(&mut buf).unwrap(), 0);
256 }
257
258 #[test]
259 fn seek_before_zero_is_invalid_input() {
260 let mut s = BlockReadStreamer::new(fixture());
261 let err = s.seek(SeekFrom::Current(-1)).unwrap_err();
262 assert_eq!(err.kind(), io::ErrorKind::InvalidInput);
263
264 let err = s.seek(SeekFrom::End(-99_999)).unwrap_err();
265 assert_eq!(err.kind(), io::ErrorKind::InvalidInput);
266 }
267
268 #[test]
269 fn works_through_arc_dyn_blockread() {
270 let dev: Arc<dyn BlockRead> = Arc::new(fixture());
271 let mut s = BlockReadStreamer::new(dev);
272 let mut out = Vec::new();
273 s.read_to_end(&mut out).unwrap();
274 assert_eq!(out.len(), 32);
275 }
276
277 #[test]
278 fn works_through_borrowed_reference() {
279 let dev = fixture();
280 {
281 let mut s = BlockReadStreamer::new(&dev as &dyn BlockRead);
282 let mut buf = [0u8; 8];
283 s.read_exact(&mut buf).unwrap();
284 assert_eq!(buf, [0, 1, 2, 3, 4, 5, 6, 7]);
285 }
286 assert_eq!(dev.size_bytes(), 32);
288 }
289
290 #[test]
291 fn into_inner_returns_wrapped_device() {
292 let s = BlockReadStreamer::new(fixture());
293 let inner = s.into_inner();
294 assert_eq!(inner.size_bytes(), 32);
295 }
296
297 #[test]
298 fn get_ref_exposes_inner_without_consuming() {
299 let s = BlockReadStreamer::new(fixture());
300 assert_eq!(s.get_ref().size_bytes(), 32);
301 assert_eq!(s.position(), 0);
303 }
304
305 #[test]
306 fn error_from_inner_propagates_as_io_error() {
307 let mut s = BlockReadStreamer::new(AlwaysFails);
308 let mut buf = [0u8; 8];
309 let err = s.read(&mut buf).unwrap_err();
310 assert_eq!(err.kind(), io::ErrorKind::Other);
312 assert!(err.to_string().contains("simulated failure"));
313 }
314
315 #[test]
316 fn offset_from_handles_add_sub_and_overflow() {
317 assert_eq!(offset_from(10, 5).unwrap(), 15);
318 assert_eq!(offset_from(10, -4).unwrap(), 6);
319 assert_eq!(offset_from(10, 0).unwrap(), 10);
320
321 let err = offset_from(u64::MAX, 1).unwrap_err();
323 assert_eq!(err.kind(), io::ErrorKind::InvalidInput);
324
325 let err = offset_from(3, -4).unwrap_err();
327 assert_eq!(err.kind(), io::ErrorKind::InvalidInput);
328 }
329
330 #[test]
331 fn error_mapping_covers_every_variant() {
332 let io_err = fs_core_error_to_io(Error::Io(io::Error::new(
334 io::ErrorKind::NotFound,
335 "missing",
336 )));
337 assert_eq!(io_err.kind(), io::ErrorKind::NotFound);
338
339 let sr = fs_core_error_to_io(Error::ShortRead {
340 offset: 4,
341 want: 8,
342 got: 2,
343 });
344 assert_eq!(sr.kind(), io::ErrorKind::UnexpectedEof);
345 assert!(sr.to_string().contains("short read"));
346
347 let oob = fs_core_error_to_io(Error::OutOfBounds {
348 offset: 16,
349 len: 4,
350 size: 8,
351 });
352 assert_eq!(oob.kind(), io::ErrorKind::UnexpectedEof);
353
354 let ro = fs_core_error_to_io(Error::ReadOnly);
355 assert_eq!(ro.kind(), io::ErrorKind::PermissionDenied);
356
357 let custom = fs_core_error_to_io(Error::Custom("boom".into()));
358 assert_eq!(custom.kind(), io::ErrorKind::Other);
359 assert!(custom.to_string().contains("boom"));
360 }
361}