Skip to main content

mapping/
reader.rs

1// Copyright 2026 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5use 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
13/// Maximum number of bytes to read per single storage buffer request (1 MiB).
14pub const MAX_READ_BUFFER_SIZE: usize = 1024 * 1024;
15
16pub use storage_device::SplittableBuffer;
17pub use storage_device::buffer::OwnedBuffer;
18
19/// Interface for block storage services providing read buffers and block read operations.
20pub trait BlockService: Send + Sync {
21    /// Allocates a read buffer of at most `max_len` bytes (`OwnedBuffer`).
22    ///
23    /// `max_len` must be greater than zero. Implementations must guarantee that the returned
24    /// buffer satisfies `0 < buffer.len() <= max_len` and that `buffer.len()` is a multiple of
25    /// `BLOCK_SIZE`.
26    fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer;
27
28    /// Submits a read for the specified range from `device_offset` into `dest_buffer`.
29    /// When the read completes, `on_complete` is invoked with the result containing the buffer.
30    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
63/// Tracks the state of a streaming read operation.
64///
65/// Ensures out-of-order completions from parallel block read requests across extents and buffer
66/// allocations are delivered to `callback` in strict ascending logical order.
67struct ReadContext<F>
68where
69    F: FnMut(Result<OwnedBuffer, Error>) -> ControlFlow<()> + Send + 'static,
70{
71    callback: Option<F>,
72
73    /// Sequence index (`k`) of the next chunk due to be delivered to `callback`.
74    next_expected_chunk: usize,
75
76    /// Out-of-order successful completions waiting for earlier chunks to finish before delivery.
77    buffered_chunks: BinaryHeap<Reverse<BufferedChunk>>,
78
79    /// Total number of buffers allocated across all iterations of `read_aligned_range`.
80    total_chunks: Option<usize>,
81}
82
83impl<F> ReadContext<F>
84where
85    F: FnMut(Result<OwnedBuffer, Error>) -> ControlFlow<()> + Send + 'static,
86{
87    /// Invoked when the asynchronous read for chunk sequence index `chunk_index` completes.
88    ///
89    /// Buffers out-of-order completions and delivers consecutive chunks (`0, 1, 2, ...`) in order
90    /// to `callback`. Any error immediately aborts future deliveries and clears buffered chunks.
91    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
158/// Reads the specified block-aligned logical `range` (`start..end`) into one or more `OwnedBuffer`
159/// chunks with zero copies, streaming each chunk sequentially to `callback`.
160///
161/// Both `range.start` and `range.end` must be multiples of `BLOCK_SIZE`. Callers requiring
162/// unaligned sub-ranges must align `range` to block boundaries and handle leading or trailing
163/// byte trimming/zeroing in `callback`.
164///
165/// Allocates buffers and submits block read requests immediately inside a loop until `range`
166/// is fully covered. If `service.allocate_buffer` returns a partial buffer or blocks due to memory
167/// pool exhaustion, the function makes progress by allocating and submitting what is available,
168/// then blocking inside `allocate_buffer` when needed until background completions free memory back
169/// to the pool. Out-of-order completions across requests are buffered so `callback` always receives
170/// chunks strictly in ascending logical offset sequence.
171pub 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        // extents only cover 0..8192, but we request 0..16384.
407        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        // Request up to 8192 bytes, but cap allocator to 4096 bytes per round.
505        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                // Delay chunk 0 until after both chunk 1 and chunk 2 complete.
550                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                // Immediately complete chunk 1, adding index 1 to `buffered_chunks`.
560                self.inner.read_blocks(device_offset, dest_buffer, on_complete)?;
561                Ok(())
562            } else {
563                // Immediately complete chunk 2, adding index 2 to `buffered_chunks` and exercising
564                // `BufferedChunk::cmp` when sorting multiple out-of-order elements in the heap.
565                // Then execute delayed chunk 0 to trigger sequential delivery of 0, 1, and 2.
566                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                // Verify strict delivery order: chunk 0 (10), then chunk 1 (20), then chunk 2 (30).
603                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            // Spawn on a background thread so the upfront read_range loop can run right ahead,
634            // and block on allocate_buffer when memory is full until this thread drops its buffer.
635            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        // Pool size is only 8192 bytes, capped at 4096 bytes per buffer.
652        // Requesting 16384 bytes will allocate buffers 0 and 1 (using all 8192 bytes of the pool),
653        // submit their reads right away inside read_aligned_range's while loop, and then
654        // block when trying to allocate buffer 2 until background threads complete and free
655        // buffer 0.
656        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        // Wait briefly for background threads to finish completing all 4 chunks.
681        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        // Test in-order error on chunk 0.
778        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        // Test out-of-order error where chunk 1 (4096) succeeds first, then chunk 0 (0) fails.
792        let service = Arc::new(OutOfOrderBlockService {
793            inner: FakeBlockService::new(vec![0u8; 8192]),
794            delayed: Mutex::new(Vec::new()),
795        });
796        // We wrap OutOfOrderBlockService to make chunk 0 fail when drained.
797        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}