1use crate::{BLOCK_SIZE, Extents, File};
6use anyhow::{Error, anyhow, bail};
7use fuchsia_sync::Mutex;
8use std::cmp::{Ordering, Reverse, max, min};
9use std::collections::BinaryHeap;
10use std::ops::{ControlFlow, Range};
11use std::sync::Arc;
12
13pub const MAX_READ_BUFFER_SIZE: usize = 1024 * 1024;
15
16pub use storage_device::SplittableBuffer;
17pub use storage_device::buffer::OwnedBuffer;
18
19pub trait BlockService: Send + Sync {
21 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer;
27
28 fn read_blocks(
31 &self,
32 device_offset: u64,
33 dest_buffer: OwnedBuffer,
34 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
35 ) -> Result<(), Error>;
36}
37
38struct BufferedChunk {
39 index: usize,
40 buffer: OwnedBuffer,
41}
42
43impl PartialEq for BufferedChunk {
44 fn eq(&self, other: &Self) -> bool {
45 self.index == other.index
46 }
47}
48
49impl Eq for BufferedChunk {}
50
51impl PartialOrd for BufferedChunk {
52 fn partial_cmp(&self, other: &Self) -> Option<Ordering> {
53 Some(self.cmp(other))
54 }
55}
56
57impl Ord for BufferedChunk {
58 fn cmp(&self, other: &Self) -> Ordering {
59 self.index.cmp(&other.index)
60 }
61}
62
63struct ReadContext<F>
68where
69 F: FnMut(Result<OwnedBuffer, Error>) -> ControlFlow<()> + Send + 'static,
70{
71 callback: Option<F>,
72
73 next_expected_chunk: usize,
75
76 buffered_chunks: BinaryHeap<Reverse<BufferedChunk>>,
78
79 total_chunks: Option<usize>,
81}
82
83impl<F> ReadContext<F>
84where
85 F: FnMut(Result<OwnedBuffer, Error>) -> ControlFlow<()> + Send + 'static,
86{
87 fn on_block_read_completed(
92 context: &Mutex<ReadContext<F>>,
93 chunk_index: usize,
94 res: Result<OwnedBuffer, Error>,
95 ) {
96 let mut guard = context.lock();
97 if guard.callback.is_none() {
98 return;
99 }
100
101 let mut current_buf = match res {
102 Ok(buf) => buf,
103 Err(e) => {
104 guard.buffered_chunks.clear();
105 let mut cb = guard.callback.take().unwrap();
106 drop(guard);
107 let _ = cb(Err(e));
108 return;
109 }
110 };
111
112 if chunk_index != guard.next_expected_chunk {
113 guard
114 .buffered_chunks
115 .push(Reverse(BufferedChunk { index: chunk_index, buffer: current_buf }));
116 return;
117 }
118
119 loop {
120 guard.next_expected_chunk += 1;
121 let cb = guard.callback.as_mut().unwrap();
122 if cb(Ok(current_buf)).is_break() {
123 guard.buffered_chunks.clear();
124 guard.callback = None;
125 return;
126 }
127
128 if let Some(top) = guard.buffered_chunks.peek()
129 && top.0.index == guard.next_expected_chunk
130 {
131 current_buf = guard.buffered_chunks.pop().unwrap().0.buffer;
132 } else {
133 break;
134 }
135 }
136 }
137}
138
139impl<F> Drop for ReadContext<F>
140where
141 F: FnMut(Result<OwnedBuffer, Error>) -> ControlFlow<()> + Send + 'static,
142{
143 fn drop(&mut self) {
144 if let Some(mut callback) = self.callback.take() {
145 let incomplete = match self.total_chunks {
146 Some(total) => {
147 self.next_expected_chunk != total || !self.buffered_chunks.is_empty()
148 }
149 None => true,
150 };
151 if incomplete {
152 let _ = callback(Err(anyhow!("ReadContext dropped before completion")));
153 }
154 }
155 }
156}
157
158pub fn read_aligned_range<F>(
172 extents: &Extents,
173 range: Range<u64>,
174 service: &(impl BlockService + ?Sized),
175 mut callback: F,
176) where
177 F: FnMut(Result<OwnedBuffer, Error>) -> ControlFlow<()> + Send + 'static,
178{
179 let mut offset = range.start;
180 let end = range.end;
181 assert_eq!(
182 offset % BLOCK_SIZE,
183 0,
184 "read_aligned_range: range.start ({offset}) must be block aligned ({BLOCK_SIZE})"
185 );
186 assert_eq!(
187 end % BLOCK_SIZE,
188 0,
189 "read_aligned_range: range.end ({end}) must be block aligned ({BLOCK_SIZE})"
190 );
191
192 if range.is_empty() {
193 let _ = callback(Err(anyhow!(
194 "read_aligned_range: range {range:?} must have a non-zero length"
195 )));
196 return;
197 }
198
199 let context = Arc::new(Mutex::new(ReadContext {
200 callback: Some(callback),
201 next_expected_chunk: 0,
202 buffered_chunks: BinaryHeap::new(),
203 total_chunks: None,
204 }));
205
206 let mut next_chunk_index = 0usize;
207
208 while offset < end {
209 if context.lock().callback.is_none() {
210 break;
211 }
212
213 let total_len = min(end - offset, MAX_READ_BUFFER_SIZE as u64) as usize;
214 let buffer = service.allocate_buffer(total_len);
215 let actual_len = buffer.len();
216 assert!(actual_len > 0 && actual_len <= total_len && actual_len % BLOCK_SIZE as usize == 0);
217
218 let actual_end = offset + actual_len as u64;
219
220 let chunk_index = next_chunk_index;
221 next_chunk_index += 1;
222
223 let context_clone = context.clone();
224 if let Err(e) = buffer.split(
225 |splittable| {
226 for extent in extents.iter_extents(offset) {
227 if extent.logical_range().start >= actual_end {
228 break;
229 }
230 let slice_start = max(extent.logical_range().start, offset);
231 let slice_end = min(extent.logical_range().end, actual_end);
232 let len_in_buf = (slice_end - slice_start) as usize;
233
234 if let Some(dev_offset) = extent.device_offset() {
235 let extent_dev_offset =
236 dev_offset + (slice_start - extent.logical_range().start);
237 let (child_buf, sub_handle) = splittable.take_prefix(len_in_buf);
238 service.read_blocks(
239 extent_dev_offset,
240 child_buf,
241 Box::new(move |res| {
242 sub_handle.merge(|| {
243 let _buf = res?;
244 Ok(())
245 });
246 }),
247 )?;
248 } else {
249 splittable.fill_zeros(len_in_buf);
250 }
251 }
252
253 let remaining_len = splittable.remaining_range().len();
254 if remaining_len > 0 {
255 bail!(
256 "read_aligned_range: requested range extends {remaining_len} bytes \
257 beyond the end of extents mappings"
258 );
259 }
260
261 Ok(())
262 },
263 move |res| ReadContext::on_block_read_completed(&context_clone, chunk_index, res),
264 ) {
265 ReadContext::on_block_read_completed(&context, chunk_index, Err(e));
266 return;
267 }
268
269 offset = actual_end;
270 }
271
272 context.lock().total_chunks = Some(next_chunk_index);
273}
274
275pub fn read_buffer_from_extents(
281 extents: &Extents,
282 logical_offset: u64,
283 dest_buffer: OwnedBuffer,
284 service: &(impl BlockService + ?Sized),
285 on_complete: impl FnOnce(Result<OwnedBuffer, Error>) + Send + 'static,
286) -> Result<(), Error> {
287 let total_len = dest_buffer.len() as u64;
288 if total_len == 0 {
289 on_complete(Ok(dest_buffer));
290 return Ok(());
291 }
292
293 dest_buffer.split(
294 |splittable| {
295 let end_offset = logical_offset + total_len;
296
297 if let Some(extent) = extents.iter_extents(logical_offset).next() {
299 let logical = extent.logical_range();
300 if logical.start <= logical_offset && logical.end >= end_offset {
301 if let Some(dev_offset) = extent.device_offset() {
302 let extent_dev_offset = dev_offset + (logical_offset - logical.start);
303 let (child_buf, sub_handle) = splittable.take_prefix(total_len as usize);
304 return service.read_blocks(
305 extent_dev_offset,
306 child_buf,
307 Box::new(move |res| {
308 sub_handle.merge(|| {
309 let _buf = res?;
310 Ok(())
311 });
312 }),
313 );
314 } else {
315 splittable.fill_zeros(total_len as usize);
317 return Ok(());
318 }
319 }
320 }
321
322 let mut current_offset = logical_offset;
324 for extent in extents.iter_extents(logical_offset) {
325 if extent.logical_range().start >= end_offset {
326 break;
327 }
328 let slice_start = max(extent.logical_range().start, current_offset);
329 let slice_end = min(extent.logical_range().end, end_offset);
330 if slice_start >= slice_end {
331 continue;
332 }
333
334 if slice_start > current_offset {
335 let hole_len = (slice_start - current_offset) as usize;
336 splittable.fill_zeros(hole_len);
337 }
338
339 let len_in_buf = (slice_end - slice_start) as usize;
340 if let Some(dev_offset) = extent.device_offset() {
341 let extent_dev_offset =
342 dev_offset + (slice_start - extent.logical_range().start);
343 let (child_buf, sub_handle) = splittable.take_prefix(len_in_buf);
344 service.read_blocks(
345 extent_dev_offset,
346 child_buf,
347 Box::new(move |res| {
348 sub_handle.merge(|| {
349 let _buf = res?;
350 Ok(())
351 });
352 }),
353 )?;
354 } else {
355 splittable.fill_zeros(len_in_buf);
356 }
357 current_offset = slice_end;
358 }
359
360 if current_offset < end_offset {
361 let hole_len = (end_offset - current_offset) as usize;
362 splittable.fill_zeros(hole_len);
363 }
364
365 Ok(())
366 },
367 on_complete,
368 )
369}
370
371pub struct ChildBlockService<S: ?Sized> {
374 parent_service: Arc<S>,
375 file: Arc<File>,
376}
377
378impl<S: BlockService + ?Sized> ChildBlockService<S> {
379 pub fn new(parent_service: Arc<S>, file: Arc<File>) -> Self {
380 Self { parent_service, file }
381 }
382
383 pub fn file(&self) -> &Arc<File> {
384 &self.file
385 }
386}
387
388impl<S: BlockService + ?Sized> BlockService for ChildBlockService<S> {
389 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
390 self.parent_service.allocate_buffer(max_len)
391 }
392
393 fn read_blocks(
394 &self,
395 device_offset: u64,
396 dest_buffer: OwnedBuffer,
397 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
398 ) -> Result<(), Error> {
399 read_buffer_from_extents(
400 self.file.extents(),
401 device_offset,
402 dest_buffer,
403 self.parent_service.as_ref(),
404 on_complete,
405 )
406 }
407}
408
409#[cfg(test)]
410pub(crate) mod tests {
411 use super::*;
412 use crate::{Extent, Transform};
413 use std::cmp::min;
414 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
415 use std::thread;
416 use std::time::Duration;
417 use storage_device::buffer_allocator::{BufferAllocator, BufferSource};
418
419 pub(crate) struct FakeBlockService {
420 allocator: Arc<BufferAllocator>,
421 device_data: Mutex<Vec<u8>>,
422 cap: Option<usize>,
423 }
424
425 impl FakeBlockService {
426 pub(crate) fn new(device_data: Vec<u8>) -> Self {
427 Self::new_with_cap(device_data, None)
428 }
429
430 pub(crate) fn new_with_cap(device_data: Vec<u8>, cap: Option<usize>) -> Self {
431 Self::new_with_pool_size(device_data, 1024 * 1024, cap)
432 }
433
434 pub(crate) fn new_with_pool_size(
435 device_data: Vec<u8>,
436 pool_size: usize,
437 cap: Option<usize>,
438 ) -> Self {
439 let source = BufferSource::new(pool_size);
440 let allocator = Arc::new(BufferAllocator::new(4096, source));
441 Self { allocator, device_data: Mutex::new(device_data), cap }
442 }
443 }
444
445 impl BlockService for FakeBlockService {
446 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
447 assert!(max_len > 0, "FakeBlockService::allocate_buffer: max_len must be > 0");
448 let len = match self.cap {
449 Some(cap) => min(max_len, cap),
450 None => max_len,
451 };
452 self.allocator.allocate_buffer_sync_owned(len)
453 }
454
455 fn read_blocks(
456 &self,
457 device_offset: u64,
458 mut dest_buffer: OwnedBuffer,
459 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
460 ) -> Result<(), Error> {
461 assert_eq!(
462 device_offset % 512,
463 0,
464 "FakeBlockService::read_blocks: device_offset ({device_offset}) must be 512-byte \
465 aligned"
466 );
467 let len = dest_buffer.len();
468 assert_eq!(
469 len % BLOCK_SIZE as usize,
470 0,
471 "FakeBlockService::read_blocks: dest_buffer.len() ({len}) must be block aligned"
472 );
473 let start = device_offset as usize;
474 let end = start + dest_buffer.len();
475 let data = self.device_data.lock();
476 dest_buffer.copy_from_slice(&data[start..end]);
477 on_complete(Ok(dest_buffer));
478 Ok(())
479 }
480 }
481
482 #[test]
483 #[should_panic(expected = "must be block aligned")]
484 fn test_read_aligned_range_unaligned_panics() {
485 let service = Arc::new(FakeBlockService::new(vec![0u8; 8192]));
486 let mappings = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
487
488 read_aligned_range(&mappings, 1000..5000, &*service, |_| ControlFlow::Continue(()));
489 }
490
491 #[test]
492 #[allow(clippy::reversed_empty_ranges)]
493 fn test_read_aligned_range_empty_range_errors() {
494 let service = FakeBlockService::new(vec![0u8; 8192]);
495 let mappings = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
496
497 let err_count = Arc::new(AtomicUsize::new(0));
498 let err_count_clone = err_count.clone();
499 read_aligned_range(&mappings, 4096..4096, &service, move |res| {
500 assert!(res.is_err());
501 err_count_clone.fetch_add(1, Ordering::Relaxed);
502 ControlFlow::Continue(())
503 });
504 assert_eq!(err_count.load(Ordering::Relaxed), 1);
505
506 let err_count_clone2 = err_count.clone();
507 read_aligned_range(&mappings, 8192..4096, &service, move |res| {
508 assert!(res.is_err());
509 err_count_clone2.fetch_add(1, Ordering::Relaxed);
510 ControlFlow::Continue(())
511 });
512 assert_eq!(err_count.load(Ordering::Relaxed), 2);
513 }
514
515 #[test]
516 fn test_read_aligned_range_beyond_extents_errors() {
517 let service = FakeBlockService::new(vec![0u8; 16384]);
518 let mappings = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
520
521 let err_occurred = Arc::new(AtomicBool::new(false));
522 let err_occurred_clone = err_occurred.clone();
523 read_aligned_range(&mappings, 0..16384, &service, move |res| {
524 if res.is_err() {
525 err_occurred_clone.store(true, Ordering::Relaxed);
526 }
527 ControlFlow::Continue(())
528 });
529 assert!(err_occurred.load(Ordering::Relaxed));
530 }
531
532 #[test]
533 fn test_read_aligned_range_single_extent() {
534 let mut device_data = vec![0u8; 16384];
535 device_data[4096..8192].copy_from_slice(&[42u8; 4096]);
536 let service = Arc::new(FakeBlockService::new(device_data));
537
538 let mappings = Extents::try_new([Extent::new(0..4096, Some(4096))], 0).unwrap();
539
540 let completed = Arc::new(AtomicBool::new(false));
541 let completed_clone = completed.clone();
542
543 read_aligned_range(&mappings, 0..4096, &*service, move |res| {
544 let buffer = res.expect("read_aligned_range should succeed");
545 assert_eq!(buffer.len(), 4096);
546 assert!(buffer.as_ptr_slice().iter_as::<u8>().all(|b| b == 42));
547 completed_clone.store(true, Ordering::Relaxed);
548 ControlFlow::Continue(())
549 });
550
551 assert!(completed.load(Ordering::Relaxed));
552 }
553
554 #[test]
555 fn test_read_aligned_range_multi_extent() {
556 let mut device_data = vec![0u8; 16384];
557 device_data[0..4096].fill(1);
558 device_data[8192..12288].fill(2);
559 let service = Arc::new(FakeBlockService::new(device_data));
560
561 let mappings = Extents::try_new(
562 [Extent::new(0..4096, Some(0)), Extent::new(4096..8192, Some(8192))],
563 0,
564 )
565 .unwrap();
566
567 let completed = Arc::new(AtomicBool::new(false));
568 let completed_clone = completed.clone();
569 let mut received_bytes = Vec::new();
570
571 read_aligned_range(&mappings, 0..8192, &*service, move |res| {
572 let buffer = res.expect("read_aligned_range should succeed");
573 buffer.as_ptr_slice().append_to(&mut received_bytes);
574 if received_bytes.len() == 8192 {
575 assert!(received_bytes[0..4096].iter().all(|&b| b == 1));
576 assert!(received_bytes[4096..8192].iter().all(|&b| b == 2));
577 completed_clone.store(true, Ordering::Relaxed);
578 }
579 ControlFlow::Continue(())
580 });
581
582 assert!(completed.load(Ordering::Relaxed));
583 }
584
585 #[test]
586 fn test_read_aligned_range_sparse_extent() {
587 let service = Arc::new(FakeBlockService::new(vec![0u8; 4096]));
588
589 let mappings = Extents::try_new([Extent::new(0..4096, None)], 0).unwrap();
590
591 let completed = Arc::new(AtomicBool::new(false));
592 let completed_clone = completed.clone();
593
594 read_aligned_range(&mappings, 0..4096, &*service, move |res| {
595 let buffer = res.expect("read_aligned_range should succeed");
596 assert_eq!(buffer.len(), 4096);
597 assert!(buffer.as_ptr_slice().iter_as::<u64>().all(|b| b == 0));
598 completed_clone.store(true, Ordering::Relaxed);
599 ControlFlow::Continue(())
600 });
601
602 assert!(completed.load(Ordering::Relaxed));
603 }
604
605 #[test]
606 fn test_read_aligned_range_partial_buffer_chaining() {
607 let mut device_data = vec![0u8; 16384];
608 device_data[0..4096].fill(1);
609 device_data[4096..8192].fill(2);
610 let service = Arc::new(FakeBlockService::new_with_cap(device_data, Some(4096)));
612
613 let mappings = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
614
615 let completed = Arc::new(AtomicBool::new(false));
616 let completed_clone = completed.clone();
617 let mut chunk_count = 0;
618 let mut received_bytes = Vec::new();
619
620 read_aligned_range(&mappings, 0..8192, &*service, move |res| {
621 let buffer = res.expect("read_aligned_range should succeed");
622 chunk_count += 1;
623 buffer.as_ptr_slice().append_to(&mut received_bytes);
624 if received_bytes.len() == 8192 {
625 assert_eq!(chunk_count, 2, "Should have streamed in 2 rounds of 4096 bytes");
626 assert!(received_bytes[0..4096].iter().all(|&b| b == 1));
627 assert!(received_bytes[4096..8192].iter().all(|&b| b == 2));
628 completed_clone.store(true, Ordering::Relaxed);
629 }
630 ControlFlow::Continue(())
631 });
632
633 assert!(completed.load(Ordering::Relaxed));
634 }
635
636 struct OutOfOrderBlockService {
637 inner: FakeBlockService,
638 delayed: Mutex<Vec<Box<dyn FnOnce() + Send>>>,
639 }
640
641 impl BlockService for OutOfOrderBlockService {
642 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
643 self.inner.allocate_buffer(max_len)
644 }
645
646 fn read_blocks(
647 &self,
648 device_offset: u64,
649 dest_buffer: OwnedBuffer,
650 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
651 ) -> Result<(), Error> {
652 if device_offset == 0 {
653 let inner = &self.inner;
655 let data = inner.device_data.lock();
656 let start = device_offset as usize;
657 let end = start + dest_buffer.len();
658 let mut buf = dest_buffer;
659 buf.copy_from_slice(&data[start..end]);
660 self.delayed.lock().push(Box::new(move || on_complete(Ok(buf))));
661 Ok(())
662 } else if device_offset == 4096 {
663 self.inner.read_blocks(device_offset, dest_buffer, on_complete)?;
665 Ok(())
666 } else {
667 self.inner.read_blocks(device_offset, dest_buffer, on_complete)?;
671 let mut delayed = self.delayed.lock();
672 for cb in delayed.drain(..) {
673 cb();
674 }
675 Ok(())
676 }
677 }
678 }
679
680 #[test]
681 fn test_read_aligned_range_out_of_order_reordering() {
682 let mut device_data = vec![0u8; 16384];
683 device_data[0..4096].fill(10);
684 device_data[4096..8192].fill(20);
685 device_data[8192..12288].fill(30);
686 let inner = FakeBlockService::new_with_cap(device_data, Some(4096));
687 let service = Arc::new(OutOfOrderBlockService { inner, delayed: Mutex::new(Vec::new()) });
688
689 let mappings = Extents::try_new(
690 [
691 Extent::new(0..4096, Some(0)),
692 Extent::new(4096..8192, Some(4096)),
693 Extent::new(8192..12288, Some(8192)),
694 ],
695 0,
696 )
697 .unwrap();
698
699 let completed = Arc::new(AtomicBool::new(false));
700 let completed_clone = completed.clone();
701 let mut received_chunks = Vec::new();
702
703 read_aligned_range(&mappings, 0..12288, &*service, move |res| {
704 let buffer = res.expect("read_aligned_range should succeed");
705 let first_byte = buffer.as_ptr_slice().to_vec()[0];
706 received_chunks.push(first_byte);
707 if received_chunks.len() == 3 {
708 assert_eq!(received_chunks, vec![10, 20, 30]);
710 completed_clone.store(true, Ordering::Relaxed);
711 }
712 ControlFlow::Continue(())
713 });
714
715 assert!(completed.load(Ordering::Relaxed));
716 }
717
718 struct ThreadedBlockService {
719 inner: FakeBlockService,
720 }
721
722 impl BlockService for ThreadedBlockService {
723 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
724 self.inner.allocate_buffer(max_len)
725 }
726
727 fn read_blocks(
728 &self,
729 device_offset: u64,
730 mut dest_buffer: OwnedBuffer,
731 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
732 ) -> Result<(), Error> {
733 let start = device_offset as usize;
734 let end = start + dest_buffer.len();
735 let data = self.inner.device_data.lock();
736 dest_buffer.copy_from_slice(&data[start..end]);
737 drop(data);
738
739 thread::spawn(move || {
742 thread::sleep(Duration::from_millis(15));
743 on_complete(Ok(dest_buffer));
744 });
745 Ok(())
746 }
747 }
748
749 #[test]
750 fn test_read_aligned_range_limited_memory_blocking_allocator() {
751 let mut device_data = vec![0u8; 16384];
752 device_data[0..4096].fill(1);
753 device_data[4096..8192].fill(2);
754 device_data[8192..12288].fill(3);
755 device_data[12288..16384].fill(4);
756
757 let inner = FakeBlockService::new_with_pool_size(device_data, 8192, Some(4096));
763 let service = Arc::new(ThreadedBlockService { inner });
764
765 let mappings = Extents::try_new([Extent::new(0..16384, Some(0))], 0).unwrap();
766
767 let completed = Arc::new(AtomicBool::new(false));
768 let completed_clone = completed.clone();
769 let mut received_bytes = Vec::new();
770
771 read_aligned_range(&mappings, 0..16384, &*service, move |res| {
772 let buffer = res.expect("read_aligned_range should succeed");
773 buffer.as_ptr_slice().append_to(&mut received_bytes);
774 if received_bytes.len() == 16384 {
775 assert!(received_bytes[0..4096].iter().all(|&b| b == 1));
776 assert!(received_bytes[4096..8192].iter().all(|&b| b == 2));
777 assert!(received_bytes[8192..12288].iter().all(|&b| b == 3));
778 assert!(received_bytes[12288..16384].iter().all(|&b| b == 4));
779 completed_clone.store(true, Ordering::Relaxed);
780 }
781 ControlFlow::Continue(())
782 });
783
784 for _ in 0..100 {
786 if completed.load(Ordering::Relaxed) {
787 break;
788 }
789 thread::sleep(Duration::from_millis(10));
790 }
791 assert!(completed.load(Ordering::Relaxed));
792 }
793
794 #[test]
795 fn test_read_aligned_range_max_read_buffer_size_capping() {
796 struct CapturingBlockService {
797 inner: FakeBlockService,
798 requested_lens: Mutex<Vec<usize>>,
799 }
800
801 impl BlockService for CapturingBlockService {
802 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
803 self.requested_lens.lock().push(max_len);
804 self.inner.allocate_buffer(max_len)
805 }
806
807 fn read_blocks(
808 &self,
809 device_offset: u64,
810 dest_buffer: OwnedBuffer,
811 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
812 ) -> Result<(), Error> {
813 self.inner.read_blocks(device_offset, dest_buffer, on_complete)
814 }
815 }
816
817 let inner = FakeBlockService::new(vec![0u8; MAX_READ_BUFFER_SIZE + 8192]);
818 let service = CapturingBlockService { inner, requested_lens: Mutex::new(Vec::new()) };
819 let mappings =
820 Extents::try_new([Extent::new(0..(MAX_READ_BUFFER_SIZE + 8192) as u64, Some(0))], 0)
821 .unwrap();
822
823 read_aligned_range(&mappings, 0..(MAX_READ_BUFFER_SIZE + 8192) as u64, &service, |_| {
824 ControlFlow::Continue(())
825 });
826
827 assert_eq!(*service.requested_lens.lock(), vec![MAX_READ_BUFFER_SIZE, 8192]);
828 }
829
830 #[test]
831 fn test_read_aligned_range_sync_read_blocks_error() {
832 struct SyncErrorBlockService(FakeBlockService);
833 impl BlockService for SyncErrorBlockService {
834 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
835 self.0.allocate_buffer(max_len)
836 }
837 fn read_blocks(
838 &self,
839 _offset: u64,
840 _dest: OwnedBuffer,
841 _on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
842 ) -> Result<(), Error> {
843 Err(anyhow!("synchronous read_blocks failure"))
844 }
845 }
846 let service = SyncErrorBlockService(FakeBlockService::new(vec![0u8; 4096]));
847 let mappings = Extents::try_new([Extent::new(0..4096, Some(0))], 0).unwrap();
848
849 let err_received = Arc::new(AtomicBool::new(false));
850 let err_clone = err_received.clone();
851 read_aligned_range(&mappings, 0..4096, &service, move |res| {
852 assert!(res.is_err());
853 err_clone.store(true, Ordering::Relaxed);
854 ControlFlow::Continue(())
855 });
856 assert!(err_received.load(Ordering::Relaxed));
857 }
858
859 #[test]
860 fn test_read_aligned_range_async_callback_errors() {
861 struct AsyncErrorBlockService(FakeBlockService, u64);
862 impl BlockService for AsyncErrorBlockService {
863 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
864 self.0.allocate_buffer(max_len)
865 }
866 fn read_blocks(
867 &self,
868 device_offset: u64,
869 dest_buffer: OwnedBuffer,
870 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
871 ) -> Result<(), Error> {
872 if device_offset == self.1 {
873 on_complete(Err(anyhow!("async block read error on offset {device_offset}")));
874 Ok(())
875 } else {
876 self.0.read_blocks(device_offset, dest_buffer, on_complete)
877 }
878 }
879 }
880
881 let service = AsyncErrorBlockService(FakeBlockService::new(vec![0u8; 8192]), 0);
883 let mappings = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
884 let err_count = Arc::new(AtomicUsize::new(0));
885 let err_clone = err_count.clone();
886 read_aligned_range(&mappings, 0..8192, &service, move |res| {
887 if res.is_err() {
888 err_clone.fetch_add(1, Ordering::Relaxed);
889 }
890 ControlFlow::Continue(())
891 });
892 assert_eq!(err_count.load(Ordering::Relaxed), 1);
893
894 let service = Arc::new(OutOfOrderBlockService {
896 inner: FakeBlockService::new(vec![0u8; 8192]),
897 delayed: Mutex::new(Vec::new()),
898 });
899 struct FailingDelayedService(Arc<OutOfOrderBlockService>);
901 impl BlockService for FailingDelayedService {
902 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
903 self.0.allocate_buffer(max_len)
904 }
905 fn read_blocks(
906 &self,
907 device_offset: u64,
908 dest_buffer: OwnedBuffer,
909 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
910 ) -> Result<(), Error> {
911 if device_offset == 0 {
912 self.0.delayed.lock().push(Box::new(move || {
913 on_complete(Err(anyhow!("delayed chunk 0 failure")))
914 }));
915 Ok(())
916 } else {
917 let res = self.0.read_blocks(device_offset, dest_buffer, on_complete);
918 let mut delayed = self.0.delayed.lock();
919 for cb in delayed.drain(..) {
920 cb();
921 }
922 res
923 }
924 }
925 }
926 let failing_service = FailingDelayedService(service.clone());
927 let err_count = Arc::new(AtomicUsize::new(0));
928 let err_clone = err_count.clone();
929 read_aligned_range(&mappings, 0..8192, &failing_service, move |res| {
930 if res.is_err() {
931 err_clone.fetch_add(1, Ordering::Relaxed);
932 }
933 ControlFlow::Continue(())
934 });
935 let mut delayed = service.delayed.lock();
936 for cb in delayed.drain(..) {
937 cb();
938 }
939 assert_eq!(err_count.load(Ordering::Relaxed), 1);
940 }
941
942 #[test]
943 fn test_read_context_dropped_before_completion() {
944 struct DroppingBlockService(Mutex<Vec<Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>>>);
945 impl BlockService for DroppingBlockService {
946 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
947 FakeBlockService::new(vec![0u8; max_len]).allocate_buffer(max_len)
948 }
949 fn read_blocks(
950 &self,
951 _offset: u64,
952 _dest: OwnedBuffer,
953 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
954 ) -> Result<(), Error> {
955 self.0.lock().push(on_complete);
956 Ok(())
957 }
958 }
959 let service = DroppingBlockService(Mutex::new(Vec::new()));
960 let mappings = Extents::try_new([Extent::new(0..4096, Some(0))], 0).unwrap();
961
962 let err_msg = Arc::new(Mutex::new(String::new()));
963 let err_msg_clone = err_msg.clone();
964 read_aligned_range(&mappings, 0..4096, &service, move |res| {
965 if let Err(e) = res {
966 *err_msg_clone.lock() = e.to_string();
967 }
968 ControlFlow::Continue(())
969 });
970 service.0.lock().clear();
971 assert_eq!(*err_msg.lock(), "Read sub-request dropped before completion");
972 }
973
974 #[test]
975 fn test_read_aligned_range_multi_iteration_error_break() {
976 struct FirstIterSyncErrorService(FakeBlockService);
977 impl BlockService for FirstIterSyncErrorService {
978 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
979 self.0.allocate_buffer(max_len)
980 }
981 fn read_blocks(
982 &self,
983 device_offset: u64,
984 dest_buffer: OwnedBuffer,
985 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
986 ) -> Result<(), Error> {
987 if device_offset == 0 {
988 Err(anyhow!("sync failure on first iteration"))
989 } else {
990 self.0.read_blocks(device_offset, dest_buffer, on_complete)
991 }
992 }
993 }
994 let inner = FakeBlockService::new_with_cap(vec![0u8; 8192], Some(4096));
995 let service = FirstIterSyncErrorService(inner);
996 let mappings = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
997
998 let err_received = Arc::new(AtomicBool::new(false));
999 let err_clone = err_received.clone();
1000 read_aligned_range(&mappings, 0..8192, &service, move |res| {
1001 if res.is_err() {
1002 err_clone.store(true, Ordering::Relaxed);
1003 }
1004 ControlFlow::Continue(())
1005 });
1006 assert!(err_received.load(Ordering::Relaxed));
1007 }
1008
1009 #[test]
1010 fn test_read_aligned_range_extent_beyond_actual_end_break() {
1011 let inner = FakeBlockService::new_with_cap(vec![0u8; 8192], Some(4096));
1012 let mappings = Extents::try_new(
1013 [Extent::new(0..4096, Some(0)), Extent::new(4096..8192, Some(4096))],
1014 0,
1015 )
1016 .unwrap();
1017
1018 let count = Arc::new(AtomicUsize::new(0));
1019 let count_clone = count.clone();
1020 read_aligned_range(&mappings, 0..8192, &inner, move |res| {
1021 if res.is_ok() {
1022 count_clone.fetch_add(1, Ordering::Relaxed);
1023 }
1024 ControlFlow::Continue(())
1025 });
1026 assert_eq!(count.load(Ordering::Relaxed), 2);
1027 }
1028
1029 #[test]
1030 fn test_read_aligned_range_merged_extents() {
1031 let mut device_data = vec![0u8; 32768];
1032 device_data[4096..8192].fill(0xAA);
1033 device_data[8192..20480].fill(0xBB);
1034 let service = Arc::new(FakeBlockService::new(device_data));
1035
1036 let mappings = Extents::try_new(
1037 [Extent::new(0..4096, Some(4096)), Extent::new(4096..16384, Some(8192))],
1038 0,
1039 )
1040 .unwrap();
1041
1042 let count = Arc::new(AtomicUsize::new(0));
1043 let count_clone = count.clone();
1044 let completed = Arc::new(AtomicBool::new(false));
1045 let completed_clone = completed.clone();
1046
1047 read_aligned_range(&mappings, 0..16384, &*service, move |res| {
1048 count_clone.fetch_add(1, Ordering::Relaxed);
1049 let buffer = res.expect("merged read should succeed");
1050 assert_eq!(buffer.len(), 16384);
1051 assert!(buffer.as_ptr_slice().subslice(0..4096).iter_as::<u8>().all(|b| b == 0xAA));
1052 assert!(buffer.as_ptr_slice().subslice(4096..16384).iter_as::<u8>().all(|b| b == 0xBB));
1053 completed_clone.store(true, Ordering::Relaxed);
1054 ControlFlow::Continue(())
1055 });
1056
1057 assert_eq!(count.load(Ordering::Relaxed), 1);
1058 assert!(completed.load(Ordering::Relaxed));
1059 }
1060
1061 #[test]
1062 fn test_read_aligned_range_callback_error_propagation() {
1063 let service = Arc::new(FakeBlockService::new(vec![0u8; 8192]));
1064 let mappings = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
1065
1066 let call_count = Arc::new(AtomicUsize::new(0));
1067 let call_count_clone = call_count.clone();
1068 read_aligned_range(&mappings, 0..8192, &*service, move |_res| {
1069 call_count_clone.fetch_add(1, Ordering::Relaxed);
1070 ControlFlow::Break(())
1071 });
1072
1073 assert_eq!(call_count.load(Ordering::Relaxed), 1);
1074 }
1075
1076 #[test]
1077 fn test_read_buffer_from_extents_single_extent() {
1078 let mut device_data = vec![0u8; 16384];
1079 device_data[4096..8192].fill(0xAA);
1080 let service = Arc::new(FakeBlockService::new(device_data));
1081
1082 let mappings = Extents::try_new([Extent::new(0..4096, Some(4096))], 0).unwrap();
1084
1085 let dest = service.allocate_buffer(4096);
1086 let (tx, rx) = std::sync::mpsc::channel();
1087 read_buffer_from_extents(&mappings, 0, dest, service.as_ref(), move |res| {
1088 tx.send(res).unwrap();
1089 })
1090 .unwrap();
1091
1092 let buffer = rx.recv().unwrap().unwrap();
1093 assert_eq!(buffer.len(), 4096);
1094 assert!(buffer.as_ptr_slice().iter_as::<u8>().all(|b| b == 0xAA));
1095 }
1096
1097 #[test]
1098 fn test_read_buffer_from_extents_sparse_hole_fast_path() {
1099 let service = Arc::new(FakeBlockService::new(vec![0u8; 16384]));
1100
1101 let mappings = Extents::try_new([Extent::new(0..4096, None)], 0).unwrap();
1103
1104 let dest = service.allocate_buffer(4096);
1105 let (tx, rx) = std::sync::mpsc::channel();
1106 read_buffer_from_extents(&mappings, 0, dest, service.as_ref(), move |res| {
1107 tx.send(res).unwrap();
1108 })
1109 .unwrap();
1110
1111 let buffer = rx.recv().unwrap().unwrap();
1112 assert_eq!(buffer.len(), 4096);
1113 assert!(buffer.as_ptr_slice().iter_as::<u8>().all(|b| b == 0x00));
1114 }
1115
1116 #[test]
1117 fn test_read_buffer_from_extents_multi_extent_and_holes() {
1118 let mut device_data = vec![0u8; 32768];
1119 device_data[0..4096].fill(0x11);
1120 device_data[8192..12288].fill(0x22);
1121 let service = Arc::new(FakeBlockService::new(device_data));
1122
1123 let mappings = Extents::try_new(
1128 [
1129 Extent::new(0..4096, Some(0)),
1130 Extent::new(4096..8192, None),
1131 Extent::new(8192..12288, Some(8192)),
1132 ],
1133 0,
1134 )
1135 .unwrap();
1136
1137 let dest = service.allocate_buffer(12288);
1138 let (tx, rx) = std::sync::mpsc::channel();
1139 read_buffer_from_extents(&mappings, 0, dest, service.as_ref(), move |res| {
1140 tx.send(res).unwrap();
1141 })
1142 .unwrap();
1143
1144 let buffer = rx.recv().unwrap().unwrap();
1145 assert_eq!(buffer.len(), 12288);
1146 assert!(buffer.as_ptr_slice().subslice(0..4096).iter_as::<u8>().all(|b| b == 0x11));
1147 assert!(buffer.as_ptr_slice().subslice(4096..8192).iter_as::<u8>().all(|b| b == 0x00));
1148 assert!(buffer.as_ptr_slice().subslice(8192..12288).iter_as::<u8>().all(|b| b == 0x22));
1149 }
1150
1151 #[test]
1152 fn test_read_buffer_from_extents_sync_error_disarms_callback() {
1153 struct SecondChunkSyncErrorService {
1154 inner: FakeBlockService,
1155 }
1156 impl BlockService for SecondChunkSyncErrorService {
1157 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
1158 self.inner.allocate_buffer(max_len)
1159 }
1160 fn read_blocks(
1161 &self,
1162 device_offset: u64,
1163 dest_buffer: OwnedBuffer,
1164 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
1165 ) -> Result<(), Error> {
1166 if device_offset == 8192 {
1167 Err(anyhow!("sync failure on second chunk"))
1168 } else {
1169 self.inner.read_blocks(device_offset, dest_buffer, on_complete)
1170 }
1171 }
1172 }
1173
1174 let inner = FakeBlockService::new(vec![0u8; 32768]);
1175 let service = Arc::new(SecondChunkSyncErrorService { inner });
1176
1177 let mappings = Extents::try_new(
1179 [Extent::new(0..4096, Some(0)), Extent::new(4096..8192, Some(8192))],
1180 0,
1181 )
1182 .unwrap();
1183
1184 let dest = service.allocate_buffer(8192);
1185 let completed = Arc::new(AtomicBool::new(false));
1186 let completed_clone = completed.clone();
1187
1188 let res = read_buffer_from_extents(&mappings, 0, dest, service.as_ref(), move |_res| {
1189 completed_clone.store(true, Ordering::Relaxed);
1190 });
1191
1192 assert!(res.is_err());
1194 assert!(!completed.load(Ordering::Relaxed));
1197 }
1198
1199 #[test]
1200 fn test_read_buffer_from_extents_in_flight_chunk_succeeds_after_sync_error() {
1201 struct InFlightFirstChunkService {
1202 inner: FakeBlockService,
1203 pending_first_chunk:
1204 Mutex<Option<(OwnedBuffer, Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>)>>,
1205 }
1206 impl BlockService for InFlightFirstChunkService {
1207 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
1208 self.inner.allocate_buffer(max_len)
1209 }
1210 fn read_blocks(
1211 &self,
1212 device_offset: u64,
1213 dest_buffer: OwnedBuffer,
1214 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
1215 ) -> Result<(), Error> {
1216 if device_offset == 0 {
1217 *self.pending_first_chunk.lock() = Some((dest_buffer, on_complete));
1219 Ok(())
1220 } else if device_offset == 8192 {
1221 Err(anyhow!("sync failure on second chunk"))
1223 } else {
1224 self.inner.read_blocks(device_offset, dest_buffer, on_complete)
1225 }
1226 }
1227 }
1228
1229 let inner = FakeBlockService::new(vec![0u8; 32768]);
1230 let service =
1231 Arc::new(InFlightFirstChunkService { inner, pending_first_chunk: Mutex::new(None) });
1232
1233 let mappings = Extents::try_new(
1235 [Extent::new(0..4096, Some(0)), Extent::new(4096..8192, Some(8192))],
1236 0,
1237 )
1238 .unwrap();
1239
1240 let dest = service.allocate_buffer(8192);
1241 let completed = Arc::new(AtomicBool::new(false));
1242 let completed_clone = completed.clone();
1243
1244 let res = read_buffer_from_extents(&mappings, 0, dest, service.as_ref(), move |_res| {
1245 completed_clone.store(true, Ordering::Relaxed);
1246 });
1247
1248 assert!(res.is_err());
1250 assert!(!completed.load(Ordering::Relaxed));
1251
1252 let (buf, cb) = service.pending_first_chunk.lock().take().expect("chunk 1 pending");
1254 cb(Ok(buf));
1255
1256 assert!(!completed.load(Ordering::Relaxed));
1258 }
1259
1260 #[test]
1261 fn test_read_buffer_from_extents_async_error_disarms_subsequent_chunks() {
1262 struct AsyncErrorService {
1263 inner: FakeBlockService,
1264 pending_chunks:
1265 Mutex<Vec<(OwnedBuffer, Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>, u64)>>,
1266 }
1267 impl BlockService for AsyncErrorService {
1268 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
1269 self.inner.allocate_buffer(max_len)
1270 }
1271 fn read_blocks(
1272 &self,
1273 device_offset: u64,
1274 dest_buffer: OwnedBuffer,
1275 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
1276 ) -> Result<(), Error> {
1277 self.pending_chunks.lock().push((dest_buffer, on_complete, device_offset));
1278 Ok(())
1279 }
1280 }
1281
1282 let inner = FakeBlockService::new(vec![0u8; 32768]);
1283 let service = Arc::new(AsyncErrorService { inner, pending_chunks: Mutex::new(Vec::new()) });
1284
1285 let mappings = Extents::try_new(
1286 [Extent::new(0..4096, Some(0)), Extent::new(4096..8192, Some(8192))],
1287 0,
1288 )
1289 .unwrap();
1290
1291 let dest = service.allocate_buffer(8192);
1292 let err_count = Arc::new(AtomicUsize::new(0));
1293 let ok_count = Arc::new(AtomicUsize::new(0));
1294 let err_count_clone = err_count.clone();
1295 let ok_count_clone = ok_count.clone();
1296
1297 let res = read_buffer_from_extents(&mappings, 0, dest, service.as_ref(), move |res| {
1298 if res.is_ok() {
1299 ok_count_clone.fetch_add(1, Ordering::Relaxed);
1300 } else {
1301 err_count_clone.fetch_add(1, Ordering::Relaxed);
1302 }
1303 });
1304 assert!(res.is_ok());
1305
1306 let mut chunks = service.pending_chunks.lock();
1307 assert_eq!(chunks.len(), 2);
1308 let (buf1, cb1, _) = chunks.remove(0);
1309 let (buf2, cb2, _) = chunks.remove(0);
1310 drop(chunks);
1311
1312 drop(buf1);
1314 cb1(Err(anyhow!("async failure on chunk 1")));
1315 assert_eq!(err_count.load(Ordering::Relaxed), 1);
1316 assert_eq!(ok_count.load(Ordering::Relaxed), 0);
1317
1318 cb2(Ok(buf2));
1320 assert_eq!(err_count.load(Ordering::Relaxed), 1);
1322 assert_eq!(ok_count.load(Ordering::Relaxed), 0);
1323 }
1324
1325 #[test]
1326 fn test_child_block_service() {
1327 let mut device_data = vec![0u8; 16384];
1328 device_data[4096..8192].fill(0x55);
1329 let parent_service = Arc::new(FakeBlockService::new(device_data));
1330
1331 let mappings = Extents::try_new([Extent::new(0..4096, Some(4096))], 0).unwrap();
1332 let file = Arc::new(File::new(mappings, 4096, Transform::None));
1333
1334 let child_service = ChildBlockService::new(parent_service.clone(), file.clone());
1335 assert_eq!(child_service.file().extents().iter_extents(0).count(), 1);
1336
1337 let dest = child_service.allocate_buffer(4096);
1338 let (tx, rx) = std::sync::mpsc::channel();
1339 child_service
1340 .read_blocks(
1341 0,
1342 dest,
1343 Box::new(move |res| {
1344 tx.send(res).unwrap();
1345 }),
1346 )
1347 .unwrap();
1348
1349 let buffer = rx.recv().unwrap().unwrap();
1350 assert_eq!(buffer.len(), 4096);
1351 assert!(buffer.as_ptr_slice().iter_as::<u8>().all(|b| b == 0x55));
1352 }
1353
1354 #[test]
1355 fn test_read_buffer_from_extents_non_aligned_start() {
1356 let base_device_offset = 17408u64; let mut device_data = vec![0u8; 32768];
1358 device_data[base_device_offset as usize..(base_device_offset as usize + 4096)].fill(0x77);
1359 let service = Arc::new(FakeBlockService::new(device_data));
1360
1361 let mappings =
1362 Extents::try_new([Extent::new(0..4096, Some(base_device_offset))], base_device_offset)
1363 .unwrap();
1364
1365 let dest = service.allocate_buffer(4096);
1366 let (tx, rx) = std::sync::mpsc::channel();
1367 read_buffer_from_extents(&mappings, 0, dest, service.as_ref(), move |res| {
1368 tx.send(res).unwrap();
1369 })
1370 .unwrap();
1371
1372 let buffer = rx.recv().unwrap().unwrap();
1373 assert_eq!(buffer.len(), 4096);
1374 assert!(buffer.as_ptr_slice().iter_as::<u8>().all(|b| b == 0x77));
1375 }
1376}