1use crate::{BLOCK_SIZE, Extents};
6use anyhow::{Error, anyhow};
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 (mut splittable, handle) = SplittableBuffer::new(buffer);
224
225 for extent in extents.iter_extents(offset) {
226 if extent.logical_range().start >= actual_end {
227 break;
228 }
229 let slice_start = max(extent.logical_range().start, offset);
230 let slice_end = min(extent.logical_range().end, actual_end);
231 let len_in_buf = (slice_end - slice_start) as usize;
232 let mut child_buf = splittable.take_prefix(len_in_buf);
233
234 if let Some(dev_offset) = extent.device_offset() {
235 let extent_dev_offset = dev_offset + (slice_start - extent.logical_range().start);
236 let handle_clone = handle.clone();
237 let context_clone = context.clone();
238 if let Err(e) = service.read_blocks(
239 extent_dev_offset,
240 child_buf,
241 Box::new(move |res| match res {
242 Ok(buf) => {
243 drop(buf);
244 if let Some(merged) = handle_clone.into_buffer() {
245 ReadContext::on_block_read_completed(
246 &context_clone,
247 chunk_index,
248 Ok(merged),
249 );
250 }
251 }
252 Err(e) => {
253 ReadContext::on_block_read_completed(
254 &context_clone,
255 chunk_index,
256 Err(e),
257 );
258 }
259 }),
260 ) {
261 ReadContext::on_block_read_completed(&context, chunk_index, Err(e));
262 return;
263 }
264 } else {
265 child_buf.fill(0);
266 }
267 }
268
269 let remaining_len = splittable.remaining_range().len();
270 if remaining_len > 0 {
271 ReadContext::on_block_read_completed(
272 &context,
273 chunk_index,
274 Err(anyhow!(
275 "read_aligned_range: requested range extends {remaining_len} bytes beyond the \
276 end of extents mappings"
277 )),
278 );
279 return;
280 }
281
282 drop(splittable);
283 if let Some(merged) = handle.into_buffer() {
284 ReadContext::on_block_read_completed(&context, chunk_index, Ok(merged));
285 }
286
287 offset = actual_end;
288 }
289
290 context.lock().total_chunks = Some(next_chunk_index);
291}
292
293#[cfg(test)]
294pub(crate) mod tests {
295 use super::*;
296 use crate::Extent;
297 use std::cmp::min;
298 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
299 use std::thread;
300 use std::time::Duration;
301 use storage_device::buffer_allocator::{BufferAllocator, BufferSource};
302
303 pub(crate) struct FakeBlockService {
304 allocator: Arc<BufferAllocator>,
305 device_data: Mutex<Vec<u8>>,
306 cap: Option<usize>,
307 }
308
309 impl FakeBlockService {
310 pub(crate) fn new(device_data: Vec<u8>) -> Self {
311 Self::new_with_cap(device_data, None)
312 }
313
314 pub(crate) fn new_with_cap(device_data: Vec<u8>, cap: Option<usize>) -> Self {
315 Self::new_with_pool_size(device_data, 1024 * 1024, cap)
316 }
317
318 pub(crate) fn new_with_pool_size(
319 device_data: Vec<u8>,
320 pool_size: usize,
321 cap: Option<usize>,
322 ) -> Self {
323 let source = BufferSource::new(pool_size);
324 let allocator = Arc::new(BufferAllocator::new(4096, source));
325 Self { allocator, device_data: Mutex::new(device_data), cap }
326 }
327 }
328
329 impl BlockService for FakeBlockService {
330 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
331 assert!(max_len > 0, "FakeBlockService::allocate_buffer: max_len must be > 0");
332 let len = match self.cap {
333 Some(cap) => min(max_len, cap),
334 None => max_len,
335 };
336 self.allocator.allocate_buffer_sync_owned(len)
337 }
338
339 fn read_blocks(
340 &self,
341 device_offset: u64,
342 mut dest_buffer: OwnedBuffer,
343 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
344 ) -> Result<(), Error> {
345 assert_eq!(
346 device_offset % BLOCK_SIZE,
347 0,
348 "FakeBlockService::read_blocks: device_offset ({device_offset}) must be block \
349 aligned"
350 );
351 let len = dest_buffer.len();
352 assert_eq!(
353 len % BLOCK_SIZE as usize,
354 0,
355 "FakeBlockService::read_blocks: dest_buffer.len() ({len}) must be block aligned"
356 );
357 let start = device_offset as usize;
358 let end = start + dest_buffer.len();
359 let data = self.device_data.lock();
360 dest_buffer.copy_from_slice(&data[start..end]);
361 on_complete(Ok(dest_buffer));
362 Ok(())
363 }
364 }
365
366 #[test]
367 #[should_panic(expected = "must be block aligned")]
368 fn test_read_aligned_range_unaligned_panics() {
369 let service = Arc::new(FakeBlockService::new(vec![0u8; 8192]));
370 let extents = vec![Extent::new(0..8192, Some(0))];
371 let encoded = Extents::encode_extents(&extents);
372 let mappings = Extents::from_encoded(&encoded).unwrap();
373
374 read_aligned_range(&mappings, 1000..5000, &*service, |_| ControlFlow::Continue(()));
375 }
376
377 #[test]
378 #[allow(clippy::reversed_empty_ranges)]
379 fn test_read_aligned_range_empty_range_errors() {
380 let service = FakeBlockService::new(vec![0u8; 8192]);
381 let extents = vec![Extent::new(0..8192, Some(0))];
382 let encoded = Extents::encode_extents(&extents);
383 let mappings = Extents::from_encoded(&encoded).unwrap();
384
385 let err_count = Arc::new(AtomicUsize::new(0));
386 let err_count_clone = err_count.clone();
387 read_aligned_range(&mappings, 4096..4096, &service, move |res| {
388 assert!(res.is_err());
389 err_count_clone.fetch_add(1, Ordering::Relaxed);
390 ControlFlow::Continue(())
391 });
392 assert_eq!(err_count.load(Ordering::Relaxed), 1);
393
394 let err_count_clone2 = err_count.clone();
395 read_aligned_range(&mappings, 8192..4096, &service, move |res| {
396 assert!(res.is_err());
397 err_count_clone2.fetch_add(1, Ordering::Relaxed);
398 ControlFlow::Continue(())
399 });
400 assert_eq!(err_count.load(Ordering::Relaxed), 2);
401 }
402
403 #[test]
404 fn test_read_aligned_range_beyond_extents_errors() {
405 let service = FakeBlockService::new(vec![0u8; 16384]);
406 let extents = vec![Extent::new(0..8192, Some(0))];
408 let encoded = Extents::encode_extents(&extents);
409 let mappings = Extents::from_encoded(&encoded).unwrap();
410
411 let err_occurred = Arc::new(AtomicBool::new(false));
412 let err_occurred_clone = err_occurred.clone();
413 read_aligned_range(&mappings, 0..16384, &service, move |res| {
414 if res.is_err() {
415 err_occurred_clone.store(true, Ordering::Relaxed);
416 }
417 ControlFlow::Continue(())
418 });
419 assert!(err_occurred.load(Ordering::Relaxed));
420 }
421
422 #[test]
423 fn test_read_aligned_range_single_extent() {
424 let mut device_data = vec![0u8; 16384];
425 device_data[4096..8192].copy_from_slice(&[42u8; 4096]);
426 let service = Arc::new(FakeBlockService::new(device_data));
427
428 let extents = vec![Extent::new(0..4096, Some(4096))];
429 let encoded = Extents::encode_extents(&extents);
430 let mappings = Extents::from_encoded(&encoded).unwrap();
431
432 let completed = Arc::new(AtomicBool::new(false));
433 let completed_clone = completed.clone();
434
435 read_aligned_range(&mappings, 0..4096, &*service, move |res| {
436 let buffer = res.expect("read_aligned_range should succeed");
437 assert_eq!(buffer.len(), 4096);
438 assert!(buffer.as_ptr_slice().iter_as::<u8>().all(|b| b.read() == 42));
439 completed_clone.store(true, Ordering::Relaxed);
440 ControlFlow::Continue(())
441 });
442
443 assert!(completed.load(Ordering::Relaxed));
444 }
445
446 #[test]
447 fn test_read_aligned_range_multi_extent() {
448 let mut device_data = vec![0u8; 16384];
449 device_data[0..4096].fill(1);
450 device_data[8192..12288].fill(2);
451 let service = Arc::new(FakeBlockService::new(device_data));
452
453 let extents = vec![Extent::new(0..4096, Some(0)), Extent::new(4096..8192, Some(8192))];
454 let encoded = Extents::encode_extents(&extents);
455 let mappings = Extents::from_encoded(&encoded).unwrap();
456
457 let completed = Arc::new(AtomicBool::new(false));
458 let completed_clone = completed.clone();
459 let mut received_bytes = Vec::new();
460
461 read_aligned_range(&mappings, 0..8192, &*service, move |res| {
462 let buffer = res.expect("read_aligned_range should succeed");
463 buffer.as_ptr_slice().append_to(&mut received_bytes);
464 if received_bytes.len() == 8192 {
465 assert!(received_bytes[0..4096].iter().all(|&b| b == 1));
466 assert!(received_bytes[4096..8192].iter().all(|&b| b == 2));
467 completed_clone.store(true, Ordering::Relaxed);
468 }
469 ControlFlow::Continue(())
470 });
471
472 assert!(completed.load(Ordering::Relaxed));
473 }
474
475 #[test]
476 fn test_read_aligned_range_sparse_extent() {
477 let service = Arc::new(FakeBlockService::new(vec![0u8; 4096]));
478
479 let extents = vec![Extent::new(0..4096, None)];
480 let encoded = Extents::encode_extents(&extents);
481 let mappings = Extents::from_encoded(&encoded).unwrap();
482
483 let completed = Arc::new(AtomicBool::new(false));
484 let completed_clone = completed.clone();
485
486 read_aligned_range(&mappings, 0..4096, &*service, move |res| {
487 let buffer = res.expect("read_aligned_range should succeed");
488 assert_eq!(buffer.len(), 4096);
489 assert!(buffer.as_ptr_slice().iter_as::<u64>().all(|b| b.read() == 0));
490 completed_clone.store(true, Ordering::Relaxed);
491 ControlFlow::Continue(())
492 });
493
494 assert!(completed.load(Ordering::Relaxed));
495 }
496
497 #[test]
498 fn test_read_aligned_range_partial_buffer_chaining() {
499 let mut device_data = vec![0u8; 16384];
500 device_data[0..4096].fill(1);
501 device_data[4096..8192].fill(2);
502 let service = Arc::new(FakeBlockService::new_with_cap(device_data, Some(4096)));
504
505 let extents = vec![Extent::new(0..8192, Some(0))];
506 let encoded = Extents::encode_extents(&extents);
507 let mappings = Extents::from_encoded(&encoded).unwrap();
508
509 let completed = Arc::new(AtomicBool::new(false));
510 let completed_clone = completed.clone();
511 let mut chunk_count = 0;
512 let mut received_bytes = Vec::new();
513
514 read_aligned_range(&mappings, 0..8192, &*service, move |res| {
515 let buffer = res.expect("read_aligned_range should succeed");
516 chunk_count += 1;
517 buffer.as_ptr_slice().append_to(&mut received_bytes);
518 if received_bytes.len() == 8192 {
519 assert_eq!(chunk_count, 2, "Should have streamed in 2 rounds of 4096 bytes");
520 assert!(received_bytes[0..4096].iter().all(|&b| b == 1));
521 assert!(received_bytes[4096..8192].iter().all(|&b| b == 2));
522 completed_clone.store(true, Ordering::Relaxed);
523 }
524 ControlFlow::Continue(())
525 });
526
527 assert!(completed.load(Ordering::Relaxed));
528 }
529
530 struct OutOfOrderBlockService {
531 inner: FakeBlockService,
532 delayed: Mutex<Vec<Box<dyn FnOnce() + Send>>>,
533 }
534
535 impl BlockService for OutOfOrderBlockService {
536 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
537 self.inner.allocate_buffer(max_len)
538 }
539
540 fn read_blocks(
541 &self,
542 device_offset: u64,
543 dest_buffer: OwnedBuffer,
544 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
545 ) -> Result<(), Error> {
546 if device_offset == 0 {
547 let inner = &self.inner;
549 let data = inner.device_data.lock();
550 let start = device_offset as usize;
551 let end = start + dest_buffer.len();
552 let mut buf = dest_buffer;
553 buf.copy_from_slice(&data[start..end]);
554 self.delayed.lock().push(Box::new(move || on_complete(Ok(buf))));
555 Ok(())
556 } else if device_offset == 4096 {
557 self.inner.read_blocks(device_offset, dest_buffer, on_complete)?;
559 Ok(())
560 } else {
561 self.inner.read_blocks(device_offset, dest_buffer, on_complete)?;
565 let mut delayed = self.delayed.lock();
566 for cb in delayed.drain(..) {
567 cb();
568 }
569 Ok(())
570 }
571 }
572 }
573
574 #[test]
575 fn test_read_aligned_range_out_of_order_reordering() {
576 let mut device_data = vec![0u8; 16384];
577 device_data[0..4096].fill(10);
578 device_data[4096..8192].fill(20);
579 device_data[8192..12288].fill(30);
580 let inner = FakeBlockService::new_with_cap(device_data, Some(4096));
581 let service = Arc::new(OutOfOrderBlockService { inner, delayed: Mutex::new(Vec::new()) });
582
583 let extents = vec![
584 Extent::new(0..4096, Some(0)),
585 Extent::new(4096..8192, Some(4096)),
586 Extent::new(8192..12288, Some(8192)),
587 ];
588 let encoded = Extents::encode_extents(&extents);
589 let mappings = Extents::from_encoded(&encoded).unwrap();
590
591 let completed = Arc::new(AtomicBool::new(false));
592 let completed_clone = completed.clone();
593 let mut received_chunks = Vec::new();
594
595 read_aligned_range(&mappings, 0..12288, &*service, move |res| {
596 let buffer = res.expect("read_aligned_range should succeed");
597 let first_byte = buffer.as_ptr_slice().to_vec()[0];
598 received_chunks.push(first_byte);
599 if received_chunks.len() == 3 {
600 assert_eq!(received_chunks, vec![10, 20, 30]);
602 completed_clone.store(true, Ordering::Relaxed);
603 }
604 ControlFlow::Continue(())
605 });
606
607 assert!(completed.load(Ordering::Relaxed));
608 }
609
610 struct ThreadedBlockService {
611 inner: FakeBlockService,
612 }
613
614 impl BlockService for ThreadedBlockService {
615 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
616 self.inner.allocate_buffer(max_len)
617 }
618
619 fn read_blocks(
620 &self,
621 device_offset: u64,
622 mut dest_buffer: OwnedBuffer,
623 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
624 ) -> Result<(), Error> {
625 let start = device_offset as usize;
626 let end = start + dest_buffer.len();
627 let data = self.inner.device_data.lock();
628 dest_buffer.copy_from_slice(&data[start..end]);
629 drop(data);
630
631 thread::spawn(move || {
634 thread::sleep(Duration::from_millis(15));
635 on_complete(Ok(dest_buffer));
636 });
637 Ok(())
638 }
639 }
640
641 #[test]
642 fn test_read_aligned_range_limited_memory_blocking_allocator() {
643 let mut device_data = vec![0u8; 16384];
644 device_data[0..4096].fill(1);
645 device_data[4096..8192].fill(2);
646 device_data[8192..12288].fill(3);
647 device_data[12288..16384].fill(4);
648
649 let inner = FakeBlockService::new_with_pool_size(device_data, 8192, Some(4096));
655 let service = Arc::new(ThreadedBlockService { inner });
656
657 let extents = vec![Extent::new(0..16384, Some(0))];
658 let encoded = Extents::encode_extents(&extents);
659 let mappings = Extents::from_encoded(&encoded).unwrap();
660
661 let completed = Arc::new(AtomicBool::new(false));
662 let completed_clone = completed.clone();
663 let mut received_bytes = Vec::new();
664
665 read_aligned_range(&mappings, 0..16384, &*service, move |res| {
666 let buffer = res.expect("read_aligned_range should succeed");
667 buffer.as_ptr_slice().append_to(&mut received_bytes);
668 if received_bytes.len() == 16384 {
669 assert!(received_bytes[0..4096].iter().all(|&b| b == 1));
670 assert!(received_bytes[4096..8192].iter().all(|&b| b == 2));
671 assert!(received_bytes[8192..12288].iter().all(|&b| b == 3));
672 assert!(received_bytes[12288..16384].iter().all(|&b| b == 4));
673 completed_clone.store(true, Ordering::Relaxed);
674 }
675 ControlFlow::Continue(())
676 });
677
678 for _ in 0..100 {
680 if completed.load(Ordering::Relaxed) {
681 break;
682 }
683 thread::sleep(Duration::from_millis(10));
684 }
685 assert!(completed.load(Ordering::Relaxed));
686 }
687
688 #[test]
689 fn test_read_aligned_range_max_read_buffer_size_capping() {
690 struct CapturingBlockService {
691 inner: FakeBlockService,
692 requested_lens: Mutex<Vec<usize>>,
693 }
694
695 impl BlockService for CapturingBlockService {
696 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
697 self.requested_lens.lock().push(max_len);
698 self.inner.allocate_buffer(max_len)
699 }
700
701 fn read_blocks(
702 &self,
703 device_offset: u64,
704 dest_buffer: OwnedBuffer,
705 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
706 ) -> Result<(), Error> {
707 self.inner.read_blocks(device_offset, dest_buffer, on_complete)
708 }
709 }
710
711 let inner = FakeBlockService::new(vec![0u8; MAX_READ_BUFFER_SIZE + 8192]);
712 let service = CapturingBlockService { inner, requested_lens: Mutex::new(Vec::new()) };
713 let extents = vec![Extent::new(0..(MAX_READ_BUFFER_SIZE + 8192) as u64, Some(0))];
714 let mappings = Extents::from_encoded(&Extents::encode_extents(&extents)).unwrap();
715
716 read_aligned_range(&mappings, 0..(MAX_READ_BUFFER_SIZE + 8192) as u64, &service, |_| {
717 ControlFlow::Continue(())
718 });
719
720 assert_eq!(*service.requested_lens.lock(), vec![MAX_READ_BUFFER_SIZE, 8192]);
721 }
722
723 #[test]
724 fn test_read_aligned_range_sync_read_blocks_error() {
725 struct SyncErrorBlockService(FakeBlockService);
726 impl BlockService for SyncErrorBlockService {
727 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
728 self.0.allocate_buffer(max_len)
729 }
730 fn read_blocks(
731 &self,
732 _offset: u64,
733 _dest: OwnedBuffer,
734 _on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
735 ) -> Result<(), Error> {
736 Err(anyhow!("synchronous read_blocks failure"))
737 }
738 }
739 let service = SyncErrorBlockService(FakeBlockService::new(vec![0u8; 4096]));
740 let extents = vec![Extent::new(0..4096, Some(0))];
741 let mappings = Extents::from_encoded(&Extents::encode_extents(&extents)).unwrap();
742
743 let err_received = Arc::new(AtomicBool::new(false));
744 let err_clone = err_received.clone();
745 read_aligned_range(&mappings, 0..4096, &service, move |res| {
746 assert!(res.is_err());
747 err_clone.store(true, Ordering::Relaxed);
748 ControlFlow::Continue(())
749 });
750 assert!(err_received.load(Ordering::Relaxed));
751 }
752
753 #[test]
754 fn test_read_aligned_range_async_callback_errors() {
755 struct AsyncErrorBlockService(FakeBlockService, u64);
756 impl BlockService for AsyncErrorBlockService {
757 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
758 self.0.allocate_buffer(max_len)
759 }
760 fn read_blocks(
761 &self,
762 device_offset: u64,
763 dest_buffer: OwnedBuffer,
764 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
765 ) -> Result<(), Error> {
766 if device_offset == self.1 {
767 on_complete(Err(anyhow!("async block read error on offset {device_offset}")));
768 Ok(())
769 } else {
770 self.0.read_blocks(device_offset, dest_buffer, on_complete)
771 }
772 }
773 }
774
775 let service = AsyncErrorBlockService(FakeBlockService::new(vec![0u8; 8192]), 0);
777 let extents = vec![Extent::new(0..8192, Some(0))];
778 let mappings = Extents::from_encoded(&Extents::encode_extents(&extents)).unwrap();
779 let err_count = Arc::new(AtomicUsize::new(0));
780 let err_clone = err_count.clone();
781 read_aligned_range(&mappings, 0..8192, &service, move |res| {
782 if res.is_err() {
783 err_clone.fetch_add(1, Ordering::Relaxed);
784 }
785 ControlFlow::Continue(())
786 });
787 assert_eq!(err_count.load(Ordering::Relaxed), 1);
788
789 let service = Arc::new(OutOfOrderBlockService {
791 inner: FakeBlockService::new(vec![0u8; 8192]),
792 delayed: Mutex::new(Vec::new()),
793 });
794 struct FailingDelayedService(Arc<OutOfOrderBlockService>);
796 impl BlockService for FailingDelayedService {
797 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
798 self.0.allocate_buffer(max_len)
799 }
800 fn read_blocks(
801 &self,
802 device_offset: u64,
803 dest_buffer: OwnedBuffer,
804 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
805 ) -> Result<(), Error> {
806 if device_offset == 0 {
807 self.0.delayed.lock().push(Box::new(move || {
808 on_complete(Err(anyhow!("delayed chunk 0 failure")))
809 }));
810 Ok(())
811 } else {
812 let res = self.0.read_blocks(device_offset, dest_buffer, on_complete);
813 let mut delayed = self.0.delayed.lock();
814 for cb in delayed.drain(..) {
815 cb();
816 }
817 res
818 }
819 }
820 }
821 let failing_service = FailingDelayedService(service.clone());
822 let err_count = Arc::new(AtomicUsize::new(0));
823 let err_clone = err_count.clone();
824 read_aligned_range(&mappings, 0..8192, &failing_service, move |res| {
825 if res.is_err() {
826 err_clone.fetch_add(1, Ordering::Relaxed);
827 }
828 ControlFlow::Continue(())
829 });
830 let mut delayed = service.delayed.lock();
831 for cb in delayed.drain(..) {
832 cb();
833 }
834 assert_eq!(err_count.load(Ordering::Relaxed), 1);
835 }
836
837 #[test]
838 fn test_read_context_dropped_before_completion() {
839 struct DroppingBlockService(Mutex<Vec<Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>>>);
840 impl BlockService for DroppingBlockService {
841 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
842 FakeBlockService::new(vec![0u8; max_len]).allocate_buffer(max_len)
843 }
844 fn read_blocks(
845 &self,
846 _offset: u64,
847 _dest: OwnedBuffer,
848 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
849 ) -> Result<(), Error> {
850 self.0.lock().push(on_complete);
851 Ok(())
852 }
853 }
854 let service = DroppingBlockService(Mutex::new(Vec::new()));
855 let extents = vec![Extent::new(0..4096, Some(0))];
856 let mappings = Extents::from_encoded(&Extents::encode_extents(&extents)).unwrap();
857
858 let err_msg = Arc::new(Mutex::new(String::new()));
859 let err_msg_clone = err_msg.clone();
860 read_aligned_range(&mappings, 0..4096, &service, move |res| {
861 if let Err(e) = res {
862 *err_msg_clone.lock() = e.to_string();
863 }
864 ControlFlow::Continue(())
865 });
866 service.0.lock().clear();
867 assert_eq!(*err_msg.lock(), "ReadContext dropped before completion");
868 }
869
870 #[test]
871 fn test_read_aligned_range_multi_iteration_error_break() {
872 struct FirstIterSyncErrorService(FakeBlockService);
873 impl BlockService for FirstIterSyncErrorService {
874 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
875 self.0.allocate_buffer(max_len)
876 }
877 fn read_blocks(
878 &self,
879 device_offset: u64,
880 dest_buffer: OwnedBuffer,
881 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
882 ) -> Result<(), Error> {
883 if device_offset == 0 {
884 Err(anyhow!("sync failure on first iteration"))
885 } else {
886 self.0.read_blocks(device_offset, dest_buffer, on_complete)
887 }
888 }
889 }
890 let inner = FakeBlockService::new_with_cap(vec![0u8; 8192], Some(4096));
891 let service = FirstIterSyncErrorService(inner);
892 let extents = vec![Extent::new(0..8192, Some(0))];
893 let mappings = Extents::from_encoded(&Extents::encode_extents(&extents)).unwrap();
894
895 let err_received = Arc::new(AtomicBool::new(false));
896 let err_clone = err_received.clone();
897 read_aligned_range(&mappings, 0..8192, &service, move |res| {
898 if res.is_err() {
899 err_clone.store(true, Ordering::Relaxed);
900 }
901 ControlFlow::Continue(())
902 });
903 assert!(err_received.load(Ordering::Relaxed));
904 }
905
906 #[test]
907 fn test_read_aligned_range_extent_beyond_actual_end_break() {
908 let inner = FakeBlockService::new_with_cap(vec![0u8; 8192], Some(4096));
909 let extents = vec![Extent::new(0..4096, Some(0)), Extent::new(4096..8192, Some(4096))];
910 let mappings = Extents::from_encoded(&Extents::encode_extents(&extents)).unwrap();
911
912 let count = Arc::new(AtomicUsize::new(0));
913 let count_clone = count.clone();
914 read_aligned_range(&mappings, 0..8192, &inner, move |res| {
915 if res.is_ok() {
916 count_clone.fetch_add(1, Ordering::Relaxed);
917 }
918 ControlFlow::Continue(())
919 });
920 assert_eq!(count.load(Ordering::Relaxed), 2);
921 }
922
923 #[test]
924 fn test_read_aligned_range_merged_extents() {
925 let mut device_data = vec![0u8; 32768];
926 device_data[4096..8192].fill(0xAA);
927 device_data[8192..20480].fill(0xBB);
928 let service = Arc::new(FakeBlockService::new(device_data));
929
930 let extents = vec![Extent::new(0..4096, Some(4096)), Extent::new(4096..16384, Some(8192))];
931 let mappings = Extents::from_encoded(&Extents::encode_extents(&extents)).unwrap();
932
933 let count = Arc::new(AtomicUsize::new(0));
934 let count_clone = count.clone();
935 let completed = Arc::new(AtomicBool::new(false));
936 let completed_clone = completed.clone();
937
938 read_aligned_range(&mappings, 0..16384, &*service, move |res| {
939 count_clone.fetch_add(1, Ordering::Relaxed);
940 let buffer = res.expect("merged read should succeed");
941 assert_eq!(buffer.len(), 16384);
942 assert!(
943 buffer.as_ptr_slice().subslice(0..4096).iter_as::<u8>().all(|b| b.read() == 0xAA)
944 );
945 assert!(
946 buffer
947 .as_ptr_slice()
948 .subslice(4096..16384)
949 .iter_as::<u8>()
950 .all(|b| b.read() == 0xBB)
951 );
952 completed_clone.store(true, Ordering::Relaxed);
953 ControlFlow::Continue(())
954 });
955
956 assert_eq!(count.load(Ordering::Relaxed), 1);
957 assert!(completed.load(Ordering::Relaxed));
958 }
959
960 #[test]
961 fn test_read_aligned_range_callback_error_propagation() {
962 let service = Arc::new(FakeBlockService::new(vec![0u8; 8192]));
963 let extents = vec![Extent::new(0..8192, Some(0))];
964 let encoded = Extents::encode_extents(&extents);
965 let mappings = Extents::from_encoded(&encoded).unwrap();
966
967 let call_count = Arc::new(AtomicUsize::new(0));
968 let call_count_clone = call_count.clone();
969 read_aligned_range(&mappings, 0..8192, &*service, move |_res| {
970 call_count_clone.fetch_add(1, Ordering::Relaxed);
971 ControlFlow::Break(())
972 });
973
974 assert_eq!(call_count.load(Ordering::Relaxed), 1);
975 }
976}