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 let slice = buffer.as_slice();
438 assert_eq!(slice.len(), 4096);
439 assert!(slice.iter().all(|&b| b == 42));
440 completed_clone.store(true, Ordering::Relaxed);
441 ControlFlow::Continue(())
442 });
443
444 assert!(completed.load(Ordering::Relaxed));
445 }
446
447 #[test]
448 fn test_read_aligned_range_multi_extent() {
449 let mut device_data = vec![0u8; 16384];
450 device_data[0..4096].fill(1);
451 device_data[8192..12288].fill(2);
452 let service = Arc::new(FakeBlockService::new(device_data));
453
454 let extents = vec![Extent::new(0..4096, Some(0)), Extent::new(4096..8192, Some(8192))];
455 let encoded = Extents::encode_extents(&extents);
456 let mappings = Extents::from_encoded(&encoded).unwrap();
457
458 let completed = Arc::new(AtomicBool::new(false));
459 let completed_clone = completed.clone();
460 let mut received_bytes = Vec::new();
461
462 read_aligned_range(&mappings, 0..8192, &*service, move |res| {
463 let buffer = res.expect("read_aligned_range should succeed");
464 received_bytes.extend_from_slice(buffer.as_slice());
465 if received_bytes.len() == 8192 {
466 assert!(received_bytes[0..4096].iter().all(|&b| b == 1));
467 assert!(received_bytes[4096..8192].iter().all(|&b| b == 2));
468 completed_clone.store(true, Ordering::Relaxed);
469 }
470 ControlFlow::Continue(())
471 });
472
473 assert!(completed.load(Ordering::Relaxed));
474 }
475
476 #[test]
477 fn test_read_aligned_range_sparse_extent() {
478 let service = Arc::new(FakeBlockService::new(vec![0u8; 4096]));
479
480 let extents = vec![Extent::new(0..4096, None)];
481 let encoded = Extents::encode_extents(&extents);
482 let mappings = Extents::from_encoded(&encoded).unwrap();
483
484 let completed = Arc::new(AtomicBool::new(false));
485 let completed_clone = completed.clone();
486
487 read_aligned_range(&mappings, 0..4096, &*service, move |res| {
488 let buffer = res.expect("read_aligned_range should succeed");
489 let slice = buffer.as_slice();
490 assert_eq!(slice.len(), 4096);
491 assert!(slice.iter().all(|&b| b == 0));
492 completed_clone.store(true, Ordering::Relaxed);
493 ControlFlow::Continue(())
494 });
495
496 assert!(completed.load(Ordering::Relaxed));
497 }
498
499 #[test]
500 fn test_read_aligned_range_partial_buffer_chaining() {
501 let mut device_data = vec![0u8; 16384];
502 device_data[0..4096].fill(1);
503 device_data[4096..8192].fill(2);
504 let service = Arc::new(FakeBlockService::new_with_cap(device_data, Some(4096)));
506
507 let extents = vec![Extent::new(0..8192, Some(0))];
508 let encoded = Extents::encode_extents(&extents);
509 let mappings = Extents::from_encoded(&encoded).unwrap();
510
511 let completed = Arc::new(AtomicBool::new(false));
512 let completed_clone = completed.clone();
513 let mut chunk_count = 0;
514 let mut received_bytes = Vec::new();
515
516 read_aligned_range(&mappings, 0..8192, &*service, move |res| {
517 let buffer = res.expect("read_aligned_range should succeed");
518 chunk_count += 1;
519 received_bytes.extend_from_slice(buffer.as_slice());
520 if received_bytes.len() == 8192 {
521 assert_eq!(chunk_count, 2, "Should have streamed in 2 rounds of 4096 bytes");
522 assert!(received_bytes[0..4096].iter().all(|&b| b == 1));
523 assert!(received_bytes[4096..8192].iter().all(|&b| b == 2));
524 completed_clone.store(true, Ordering::Relaxed);
525 }
526 ControlFlow::Continue(())
527 });
528
529 assert!(completed.load(Ordering::Relaxed));
530 }
531
532 struct OutOfOrderBlockService {
533 inner: FakeBlockService,
534 delayed: Mutex<Vec<Box<dyn FnOnce() + Send>>>,
535 }
536
537 impl BlockService for OutOfOrderBlockService {
538 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
539 self.inner.allocate_buffer(max_len)
540 }
541
542 fn read_blocks(
543 &self,
544 device_offset: u64,
545 dest_buffer: OwnedBuffer,
546 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
547 ) -> Result<(), Error> {
548 if device_offset == 0 {
549 let inner = &self.inner;
551 let data = inner.device_data.lock();
552 let start = device_offset as usize;
553 let end = start + dest_buffer.len();
554 let mut buf = dest_buffer;
555 buf.copy_from_slice(&data[start..end]);
556 self.delayed.lock().push(Box::new(move || on_complete(Ok(buf))));
557 Ok(())
558 } else if device_offset == 4096 {
559 self.inner.read_blocks(device_offset, dest_buffer, on_complete)?;
561 Ok(())
562 } else {
563 self.inner.read_blocks(device_offset, dest_buffer, on_complete)?;
567 let mut delayed = self.delayed.lock();
568 for cb in delayed.drain(..) {
569 cb();
570 }
571 Ok(())
572 }
573 }
574 }
575
576 #[test]
577 fn test_read_aligned_range_out_of_order_reordering() {
578 let mut device_data = vec![0u8; 16384];
579 device_data[0..4096].fill(10);
580 device_data[4096..8192].fill(20);
581 device_data[8192..12288].fill(30);
582 let inner = FakeBlockService::new_with_cap(device_data, Some(4096));
583 let service = Arc::new(OutOfOrderBlockService { inner, delayed: Mutex::new(Vec::new()) });
584
585 let extents = vec![
586 Extent::new(0..4096, Some(0)),
587 Extent::new(4096..8192, Some(4096)),
588 Extent::new(8192..12288, Some(8192)),
589 ];
590 let encoded = Extents::encode_extents(&extents);
591 let mappings = Extents::from_encoded(&encoded).unwrap();
592
593 let completed = Arc::new(AtomicBool::new(false));
594 let completed_clone = completed.clone();
595 let mut received_chunks = Vec::new();
596
597 read_aligned_range(&mappings, 0..12288, &*service, move |res| {
598 let buffer = res.expect("read_aligned_range should succeed");
599 let first_byte = buffer.as_slice()[0];
600 received_chunks.push(first_byte);
601 if received_chunks.len() == 3 {
602 assert_eq!(received_chunks, vec![10, 20, 30]);
604 completed_clone.store(true, Ordering::Relaxed);
605 }
606 ControlFlow::Continue(())
607 });
608
609 assert!(completed.load(Ordering::Relaxed));
610 }
611
612 struct ThreadedBlockService {
613 inner: FakeBlockService,
614 }
615
616 impl BlockService for ThreadedBlockService {
617 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
618 self.inner.allocate_buffer(max_len)
619 }
620
621 fn read_blocks(
622 &self,
623 device_offset: u64,
624 mut dest_buffer: OwnedBuffer,
625 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
626 ) -> Result<(), Error> {
627 let start = device_offset as usize;
628 let end = start + dest_buffer.len();
629 let data = self.inner.device_data.lock();
630 dest_buffer.copy_from_slice(&data[start..end]);
631 drop(data);
632
633 thread::spawn(move || {
636 thread::sleep(Duration::from_millis(15));
637 on_complete(Ok(dest_buffer));
638 });
639 Ok(())
640 }
641 }
642
643 #[test]
644 fn test_read_aligned_range_limited_memory_blocking_allocator() {
645 let mut device_data = vec![0u8; 16384];
646 device_data[0..4096].fill(1);
647 device_data[4096..8192].fill(2);
648 device_data[8192..12288].fill(3);
649 device_data[12288..16384].fill(4);
650
651 let inner = FakeBlockService::new_with_pool_size(device_data, 8192, Some(4096));
657 let service = Arc::new(ThreadedBlockService { inner });
658
659 let extents = vec![Extent::new(0..16384, Some(0))];
660 let encoded = Extents::encode_extents(&extents);
661 let mappings = Extents::from_encoded(&encoded).unwrap();
662
663 let completed = Arc::new(AtomicBool::new(false));
664 let completed_clone = completed.clone();
665 let mut received_bytes = Vec::new();
666
667 read_aligned_range(&mappings, 0..16384, &*service, move |res| {
668 let buffer = res.expect("read_aligned_range should succeed");
669 received_bytes.extend_from_slice(buffer.as_slice());
670 if received_bytes.len() == 16384 {
671 assert!(received_bytes[0..4096].iter().all(|&b| b == 1));
672 assert!(received_bytes[4096..8192].iter().all(|&b| b == 2));
673 assert!(received_bytes[8192..12288].iter().all(|&b| b == 3));
674 assert!(received_bytes[12288..16384].iter().all(|&b| b == 4));
675 completed_clone.store(true, Ordering::Relaxed);
676 }
677 ControlFlow::Continue(())
678 });
679
680 for _ in 0..100 {
682 if completed.load(Ordering::Relaxed) {
683 break;
684 }
685 thread::sleep(Duration::from_millis(10));
686 }
687 assert!(completed.load(Ordering::Relaxed));
688 }
689
690 #[test]
691 fn test_read_aligned_range_max_read_buffer_size_capping() {
692 struct CapturingBlockService {
693 inner: FakeBlockService,
694 requested_lens: Mutex<Vec<usize>>,
695 }
696
697 impl BlockService for CapturingBlockService {
698 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
699 self.requested_lens.lock().push(max_len);
700 self.inner.allocate_buffer(max_len)
701 }
702
703 fn read_blocks(
704 &self,
705 device_offset: u64,
706 dest_buffer: OwnedBuffer,
707 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
708 ) -> Result<(), Error> {
709 self.inner.read_blocks(device_offset, dest_buffer, on_complete)
710 }
711 }
712
713 let inner = FakeBlockService::new(vec![0u8; MAX_READ_BUFFER_SIZE + 8192]);
714 let service = CapturingBlockService { inner, requested_lens: Mutex::new(Vec::new()) };
715 let extents = vec![Extent::new(0..(MAX_READ_BUFFER_SIZE + 8192) as u64, Some(0))];
716 let mappings = Extents::from_encoded(&Extents::encode_extents(&extents)).unwrap();
717
718 read_aligned_range(&mappings, 0..(MAX_READ_BUFFER_SIZE + 8192) as u64, &service, |_| {
719 ControlFlow::Continue(())
720 });
721
722 assert_eq!(*service.requested_lens.lock(), vec![MAX_READ_BUFFER_SIZE, 8192]);
723 }
724
725 #[test]
726 fn test_read_aligned_range_sync_read_blocks_error() {
727 struct SyncErrorBlockService(FakeBlockService);
728 impl BlockService for SyncErrorBlockService {
729 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
730 self.0.allocate_buffer(max_len)
731 }
732 fn read_blocks(
733 &self,
734 _offset: u64,
735 _dest: OwnedBuffer,
736 _on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
737 ) -> Result<(), Error> {
738 Err(anyhow!("synchronous read_blocks failure"))
739 }
740 }
741 let service = SyncErrorBlockService(FakeBlockService::new(vec![0u8; 4096]));
742 let extents = vec![Extent::new(0..4096, Some(0))];
743 let mappings = Extents::from_encoded(&Extents::encode_extents(&extents)).unwrap();
744
745 let err_received = Arc::new(AtomicBool::new(false));
746 let err_clone = err_received.clone();
747 read_aligned_range(&mappings, 0..4096, &service, move |res| {
748 assert!(res.is_err());
749 err_clone.store(true, Ordering::Relaxed);
750 ControlFlow::Continue(())
751 });
752 assert!(err_received.load(Ordering::Relaxed));
753 }
754
755 #[test]
756 fn test_read_aligned_range_async_callback_errors() {
757 struct AsyncErrorBlockService(FakeBlockService, u64);
758 impl BlockService for AsyncErrorBlockService {
759 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
760 self.0.allocate_buffer(max_len)
761 }
762 fn read_blocks(
763 &self,
764 device_offset: u64,
765 dest_buffer: OwnedBuffer,
766 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
767 ) -> Result<(), Error> {
768 if device_offset == self.1 {
769 on_complete(Err(anyhow!("async block read error on offset {device_offset}")));
770 Ok(())
771 } else {
772 self.0.read_blocks(device_offset, dest_buffer, on_complete)
773 }
774 }
775 }
776
777 let service = AsyncErrorBlockService(FakeBlockService::new(vec![0u8; 8192]), 0);
779 let extents = vec![Extent::new(0..8192, Some(0))];
780 let mappings = Extents::from_encoded(&Extents::encode_extents(&extents)).unwrap();
781 let err_count = Arc::new(AtomicUsize::new(0));
782 let err_clone = err_count.clone();
783 read_aligned_range(&mappings, 0..8192, &service, move |res| {
784 if res.is_err() {
785 err_clone.fetch_add(1, Ordering::Relaxed);
786 }
787 ControlFlow::Continue(())
788 });
789 assert_eq!(err_count.load(Ordering::Relaxed), 1);
790
791 let service = Arc::new(OutOfOrderBlockService {
793 inner: FakeBlockService::new(vec![0u8; 8192]),
794 delayed: Mutex::new(Vec::new()),
795 });
796 struct FailingDelayedService(Arc<OutOfOrderBlockService>);
798 impl BlockService for FailingDelayedService {
799 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
800 self.0.allocate_buffer(max_len)
801 }
802 fn read_blocks(
803 &self,
804 device_offset: u64,
805 dest_buffer: OwnedBuffer,
806 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
807 ) -> Result<(), Error> {
808 if device_offset == 0 {
809 self.0.delayed.lock().push(Box::new(move || {
810 on_complete(Err(anyhow!("delayed chunk 0 failure")))
811 }));
812 Ok(())
813 } else {
814 let res = self.0.read_blocks(device_offset, dest_buffer, on_complete);
815 let mut delayed = self.0.delayed.lock();
816 for cb in delayed.drain(..) {
817 cb();
818 }
819 res
820 }
821 }
822 }
823 let failing_service = FailingDelayedService(service.clone());
824 let err_count = Arc::new(AtomicUsize::new(0));
825 let err_clone = err_count.clone();
826 read_aligned_range(&mappings, 0..8192, &failing_service, move |res| {
827 if res.is_err() {
828 err_clone.fetch_add(1, Ordering::Relaxed);
829 }
830 ControlFlow::Continue(())
831 });
832 let mut delayed = service.delayed.lock();
833 for cb in delayed.drain(..) {
834 cb();
835 }
836 assert_eq!(err_count.load(Ordering::Relaxed), 1);
837 }
838
839 #[test]
840 fn test_read_context_dropped_before_completion() {
841 struct DroppingBlockService(Mutex<Vec<Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>>>);
842 impl BlockService for DroppingBlockService {
843 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
844 FakeBlockService::new(vec![0u8; max_len]).allocate_buffer(max_len)
845 }
846 fn read_blocks(
847 &self,
848 _offset: u64,
849 _dest: OwnedBuffer,
850 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
851 ) -> Result<(), Error> {
852 self.0.lock().push(on_complete);
853 Ok(())
854 }
855 }
856 let service = DroppingBlockService(Mutex::new(Vec::new()));
857 let extents = vec![Extent::new(0..4096, Some(0))];
858 let mappings = Extents::from_encoded(&Extents::encode_extents(&extents)).unwrap();
859
860 let err_msg = Arc::new(Mutex::new(String::new()));
861 let err_msg_clone = err_msg.clone();
862 read_aligned_range(&mappings, 0..4096, &service, move |res| {
863 if let Err(e) = res {
864 *err_msg_clone.lock() = e.to_string();
865 }
866 ControlFlow::Continue(())
867 });
868 service.0.lock().clear();
869 assert_eq!(*err_msg.lock(), "ReadContext dropped before completion");
870 }
871
872 #[test]
873 fn test_read_aligned_range_multi_iteration_error_break() {
874 struct FirstIterSyncErrorService(FakeBlockService);
875 impl BlockService for FirstIterSyncErrorService {
876 fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
877 self.0.allocate_buffer(max_len)
878 }
879 fn read_blocks(
880 &self,
881 device_offset: u64,
882 dest_buffer: OwnedBuffer,
883 on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
884 ) -> Result<(), Error> {
885 if device_offset == 0 {
886 Err(anyhow!("sync failure on first iteration"))
887 } else {
888 self.0.read_blocks(device_offset, dest_buffer, on_complete)
889 }
890 }
891 }
892 let inner = FakeBlockService::new_with_cap(vec![0u8; 8192], Some(4096));
893 let service = FirstIterSyncErrorService(inner);
894 let extents = vec![Extent::new(0..8192, Some(0))];
895 let mappings = Extents::from_encoded(&Extents::encode_extents(&extents)).unwrap();
896
897 let err_received = Arc::new(AtomicBool::new(false));
898 let err_clone = err_received.clone();
899 read_aligned_range(&mappings, 0..8192, &service, move |res| {
900 if res.is_err() {
901 err_clone.store(true, Ordering::Relaxed);
902 }
903 ControlFlow::Continue(())
904 });
905 assert!(err_received.load(Ordering::Relaxed));
906 }
907
908 #[test]
909 fn test_read_aligned_range_extent_beyond_actual_end_break() {
910 let inner = FakeBlockService::new_with_cap(vec![0u8; 8192], Some(4096));
911 let extents = vec![Extent::new(0..4096, Some(0)), Extent::new(4096..8192, Some(4096))];
912 let mappings = Extents::from_encoded(&Extents::encode_extents(&extents)).unwrap();
913
914 let count = Arc::new(AtomicUsize::new(0));
915 let count_clone = count.clone();
916 read_aligned_range(&mappings, 0..8192, &inner, move |res| {
917 if res.is_ok() {
918 count_clone.fetch_add(1, Ordering::Relaxed);
919 }
920 ControlFlow::Continue(())
921 });
922 assert_eq!(count.load(Ordering::Relaxed), 2);
923 }
924
925 #[test]
926 fn test_read_aligned_range_merged_extents() {
927 let mut device_data = vec![0u8; 32768];
928 device_data[4096..8192].fill(0xAA);
929 device_data[8192..20480].fill(0xBB);
930 let service = Arc::new(FakeBlockService::new(device_data));
931
932 let extents = vec![Extent::new(0..4096, Some(4096)), Extent::new(4096..16384, Some(8192))];
933 let mappings = Extents::from_encoded(&Extents::encode_extents(&extents)).unwrap();
934
935 let count = Arc::new(AtomicUsize::new(0));
936 let count_clone = count.clone();
937 let completed = Arc::new(AtomicBool::new(false));
938 let completed_clone = completed.clone();
939
940 read_aligned_range(&mappings, 0..16384, &*service, move |res| {
941 count_clone.fetch_add(1, Ordering::Relaxed);
942 let buffer = res.expect("merged read should succeed");
943 assert_eq!(buffer.len(), 16384);
944 assert!(buffer.as_slice()[0..4096].iter().all(|&b| b == 0xAA));
945 assert!(buffer.as_slice()[4096..16384].iter().all(|&b| b == 0xBB));
946 completed_clone.store(true, Ordering::Relaxed);
947 ControlFlow::Continue(())
948 });
949
950 assert_eq!(count.load(Ordering::Relaxed), 1);
951 assert!(completed.load(Ordering::Relaxed));
952 }
953
954 #[test]
955 fn test_read_aligned_range_callback_error_propagation() {
956 let service = Arc::new(FakeBlockService::new(vec![0u8; 8192]));
957 let extents = vec![Extent::new(0..8192, Some(0))];
958 let encoded = Extents::encode_extents(&extents);
959 let mappings = Extents::from_encoded(&encoded).unwrap();
960
961 let call_count = Arc::new(AtomicUsize::new(0));
962 let call_count_clone = call_count.clone();
963 read_aligned_range(&mappings, 0..8192, &*service, move |_res| {
964 call_count_clone.fetch_add(1, Ordering::Relaxed);
965 ControlFlow::Break(())
966 });
967
968 assert_eq!(call_count.load(Ordering::Relaxed), 1);
969 }
970}