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            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        // Request up to 8192 bytes, but cap allocator to 4096 bytes per round.
503        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                // Delay chunk 0 until after both chunk 1 and chunk 2 complete.
548                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                // Immediately complete chunk 1, adding index 1 to `buffered_chunks`.
558                self.inner.read_blocks(device_offset, dest_buffer, on_complete)?;
559                Ok(())
560            } else {
561                // Immediately complete chunk 2, adding index 2 to `buffered_chunks` and exercising
562                // `BufferedChunk::cmp` when sorting multiple out-of-order elements in the heap.
563                // Then execute delayed chunk 0 to trigger sequential delivery of 0, 1, and 2.
564                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                // Verify strict delivery order: chunk 0 (10), then chunk 1 (20), then chunk 2 (30).
601                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            // Spawn on a background thread so the upfront read_range loop can run right ahead,
632            // and block on allocate_buffer when memory is full until this thread drops its buffer.
633            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        // Pool size is only 8192 bytes, capped at 4096 bytes per buffer.
650        // Requesting 16384 bytes will allocate buffers 0 and 1 (using all 8192 bytes of the pool),
651        // submit their reads right away inside read_aligned_range's while loop, and then
652        // block when trying to allocate buffer 2 until background threads complete and free
653        // buffer 0.
654        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        // Wait briefly for background threads to finish completing all 4 chunks.
679        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        // Test in-order error on chunk 0.
776        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        // Test out-of-order error where chunk 1 (4096) succeeds first, then chunk 0 (0) fails.
790        let service = Arc::new(OutOfOrderBlockService {
791            inner: FakeBlockService::new(vec![0u8; 8192]),
792            delayed: Mutex::new(Vec::new()),
793        });
794        // We wrap OutOfOrderBlockService to make chunk 0 fail when drained.
795        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}