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, File};
6use anyhow::{Error, anyhow, bail};
7use fuchsia_sync::Mutex;
8use std::cmp::{Ordering, Reverse, max, min};
9use std::collections::BinaryHeap;
10use std::ops::{ControlFlow, Range};
11use std::sync::Arc;
12
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 context_clone = context.clone();
224        if let Err(e) = buffer.split(
225            |splittable| {
226                for extent in extents.iter_extents(offset) {
227                    if extent.logical_range().start >= actual_end {
228                        break;
229                    }
230                    let slice_start = max(extent.logical_range().start, offset);
231                    let slice_end = min(extent.logical_range().end, actual_end);
232                    let len_in_buf = (slice_end - slice_start) as usize;
233
234                    if let Some(dev_offset) = extent.device_offset() {
235                        let extent_dev_offset =
236                            dev_offset + (slice_start - extent.logical_range().start);
237                        let (child_buf, sub_handle) = splittable.take_prefix(len_in_buf);
238                        service.read_blocks(
239                            extent_dev_offset,
240                            child_buf,
241                            Box::new(move |res| {
242                                sub_handle.merge(|| {
243                                    let _buf = res?;
244                                    Ok(())
245                                });
246                            }),
247                        )?;
248                    } else {
249                        splittable.fill_zeros(len_in_buf);
250                    }
251                }
252
253                let remaining_len = splittable.remaining_range().len();
254                if remaining_len > 0 {
255                    bail!(
256                        "read_aligned_range: requested range extends {remaining_len} bytes \
257                         beyond the end of extents mappings"
258                    );
259                }
260
261                Ok(())
262            },
263            move |res| ReadContext::on_block_read_completed(&context_clone, chunk_index, res),
264        ) {
265            ReadContext::on_block_read_completed(&context, chunk_index, Err(e));
266            return;
267        }
268
269        offset = actual_end;
270    }
271
272    context.lock().total_chunks = Some(next_chunk_index);
273}
274
275/// Reads data for `dest_buffer` starting at `logical_offset` according to `extents` from `service`.
276///
277/// If the range is contiguous within a single extent, issues a single read. If the range spans
278/// multiple extents or holes, splits `dest_buffer` across the extents and reconstructs the
279/// original buffer when all reads complete.
280pub fn read_buffer_from_extents(
281    extents: &Extents,
282    logical_offset: u64,
283    dest_buffer: OwnedBuffer,
284    service: &(impl BlockService + ?Sized),
285    on_complete: impl FnOnce(Result<OwnedBuffer, Error>) + Send + 'static,
286) -> Result<(), Error> {
287    let total_len = dest_buffer.len() as u64;
288    if total_len == 0 {
289        on_complete(Ok(dest_buffer));
290        return Ok(());
291    }
292
293    dest_buffer.split(
294        |splittable| {
295            let end_offset = logical_offset + total_len;
296
297            // Fast path: Check if the entire range is covered by a single contiguous extent.
298            if let Some(extent) = extents.iter_extents(logical_offset).next() {
299                let logical = extent.logical_range();
300                if logical.start <= logical_offset && logical.end >= end_offset {
301                    if let Some(dev_offset) = extent.device_offset() {
302                        let extent_dev_offset = dev_offset + (logical_offset - logical.start);
303                        let (child_buf, sub_handle) = splittable.take_prefix(total_len as usize);
304                        return service.read_blocks(
305                            extent_dev_offset,
306                            child_buf,
307                            Box::new(move |res| {
308                                sub_handle.merge(|| {
309                                    let _buf = res?;
310                                    Ok(())
311                                });
312                            }),
313                        );
314                    } else {
315                        // Sparse hole covering the full range.
316                        splittable.fill_zeros(total_len as usize);
317                        return Ok(());
318                    }
319                }
320            }
321
322            // Multi-extent path using SplittableBuffer.
323            let mut current_offset = logical_offset;
324            for extent in extents.iter_extents(logical_offset) {
325                if extent.logical_range().start >= end_offset {
326                    break;
327                }
328                let slice_start = max(extent.logical_range().start, current_offset);
329                let slice_end = min(extent.logical_range().end, end_offset);
330                if slice_start >= slice_end {
331                    continue;
332                }
333
334                if slice_start > current_offset {
335                    let hole_len = (slice_start - current_offset) as usize;
336                    splittable.fill_zeros(hole_len);
337                }
338
339                let len_in_buf = (slice_end - slice_start) as usize;
340                if let Some(dev_offset) = extent.device_offset() {
341                    let extent_dev_offset =
342                        dev_offset + (slice_start - extent.logical_range().start);
343                    let (child_buf, sub_handle) = splittable.take_prefix(len_in_buf);
344                    service.read_blocks(
345                        extent_dev_offset,
346                        child_buf,
347                        Box::new(move |res| {
348                            sub_handle.merge(|| {
349                                let _buf = res?;
350                                Ok(())
351                            });
352                        }),
353                    )?;
354                } else {
355                    splittable.fill_zeros(len_in_buf);
356                }
357                current_offset = slice_end;
358            }
359
360            if current_offset < end_offset {
361                let hole_len = (end_offset - current_offset) as usize;
362                splittable.fill_zeros(hole_len);
363            }
364
365            Ok(())
366        },
367        on_complete,
368    )
369}
370
371/// A [`BlockService`] that maps logical block offsets through the extents of a mapped [`File`]
372/// in a parent session.
373pub struct ChildBlockService<S: ?Sized> {
374    parent_service: Arc<S>,
375    file: Arc<File>,
376}
377
378impl<S: BlockService + ?Sized> ChildBlockService<S> {
379    pub fn new(parent_service: Arc<S>, file: Arc<File>) -> Self {
380        Self { parent_service, file }
381    }
382
383    pub fn file(&self) -> &Arc<File> {
384        &self.file
385    }
386}
387
388impl<S: BlockService + ?Sized> BlockService for ChildBlockService<S> {
389    fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
390        self.parent_service.allocate_buffer(max_len)
391    }
392
393    fn read_blocks(
394        &self,
395        device_offset: u64,
396        dest_buffer: OwnedBuffer,
397        on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
398    ) -> Result<(), Error> {
399        read_buffer_from_extents(
400            self.file.extents(),
401            device_offset,
402            dest_buffer,
403            self.parent_service.as_ref(),
404            on_complete,
405        )
406    }
407}
408
409#[cfg(test)]
410pub(crate) mod tests {
411    use super::*;
412    use crate::{Extent, Transform};
413    use std::cmp::min;
414    use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
415    use std::thread;
416    use std::time::Duration;
417    use storage_device::buffer_allocator::{BufferAllocator, BufferSource};
418
419    pub(crate) struct FakeBlockService {
420        allocator: Arc<BufferAllocator>,
421        device_data: Mutex<Vec<u8>>,
422        cap: Option<usize>,
423    }
424
425    impl FakeBlockService {
426        pub(crate) fn new(device_data: Vec<u8>) -> Self {
427            Self::new_with_cap(device_data, None)
428        }
429
430        pub(crate) fn new_with_cap(device_data: Vec<u8>, cap: Option<usize>) -> Self {
431            Self::new_with_pool_size(device_data, 1024 * 1024, cap)
432        }
433
434        pub(crate) fn new_with_pool_size(
435            device_data: Vec<u8>,
436            pool_size: usize,
437            cap: Option<usize>,
438        ) -> Self {
439            let source = BufferSource::new(pool_size);
440            let allocator = Arc::new(BufferAllocator::new(4096, source));
441            Self { allocator, device_data: Mutex::new(device_data), cap }
442        }
443    }
444
445    impl BlockService for FakeBlockService {
446        fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
447            assert!(max_len > 0, "FakeBlockService::allocate_buffer: max_len must be > 0");
448            let len = match self.cap {
449                Some(cap) => min(max_len, cap),
450                None => max_len,
451            };
452            self.allocator.allocate_buffer_sync_owned(len)
453        }
454
455        fn read_blocks(
456            &self,
457            device_offset: u64,
458            mut dest_buffer: OwnedBuffer,
459            on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
460        ) -> Result<(), Error> {
461            assert_eq!(
462                device_offset % 512,
463                0,
464                "FakeBlockService::read_blocks: device_offset ({device_offset}) must be 512-byte \
465                 aligned"
466            );
467            let len = dest_buffer.len();
468            assert_eq!(
469                len % BLOCK_SIZE as usize,
470                0,
471                "FakeBlockService::read_blocks: dest_buffer.len() ({len}) must be block aligned"
472            );
473            let start = device_offset as usize;
474            let end = start + dest_buffer.len();
475            let data = self.device_data.lock();
476            dest_buffer.copy_from_slice(&data[start..end]);
477            on_complete(Ok(dest_buffer));
478            Ok(())
479        }
480    }
481
482    #[test]
483    #[should_panic(expected = "must be block aligned")]
484    fn test_read_aligned_range_unaligned_panics() {
485        let service = Arc::new(FakeBlockService::new(vec![0u8; 8192]));
486        let mappings = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
487
488        read_aligned_range(&mappings, 1000..5000, &*service, |_| ControlFlow::Continue(()));
489    }
490
491    #[test]
492    #[allow(clippy::reversed_empty_ranges)]
493    fn test_read_aligned_range_empty_range_errors() {
494        let service = FakeBlockService::new(vec![0u8; 8192]);
495        let mappings = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
496
497        let err_count = Arc::new(AtomicUsize::new(0));
498        let err_count_clone = err_count.clone();
499        read_aligned_range(&mappings, 4096..4096, &service, move |res| {
500            assert!(res.is_err());
501            err_count_clone.fetch_add(1, Ordering::Relaxed);
502            ControlFlow::Continue(())
503        });
504        assert_eq!(err_count.load(Ordering::Relaxed), 1);
505
506        let err_count_clone2 = err_count.clone();
507        read_aligned_range(&mappings, 8192..4096, &service, move |res| {
508            assert!(res.is_err());
509            err_count_clone2.fetch_add(1, Ordering::Relaxed);
510            ControlFlow::Continue(())
511        });
512        assert_eq!(err_count.load(Ordering::Relaxed), 2);
513    }
514
515    #[test]
516    fn test_read_aligned_range_beyond_extents_errors() {
517        let service = FakeBlockService::new(vec![0u8; 16384]);
518        // extents only cover 0..8192, but we request 0..16384.
519        let mappings = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
520
521        let err_occurred = Arc::new(AtomicBool::new(false));
522        let err_occurred_clone = err_occurred.clone();
523        read_aligned_range(&mappings, 0..16384, &service, move |res| {
524            if res.is_err() {
525                err_occurred_clone.store(true, Ordering::Relaxed);
526            }
527            ControlFlow::Continue(())
528        });
529        assert!(err_occurred.load(Ordering::Relaxed));
530    }
531
532    #[test]
533    fn test_read_aligned_range_single_extent() {
534        let mut device_data = vec![0u8; 16384];
535        device_data[4096..8192].copy_from_slice(&[42u8; 4096]);
536        let service = Arc::new(FakeBlockService::new(device_data));
537
538        let mappings = Extents::try_new([Extent::new(0..4096, Some(4096))], 0).unwrap();
539
540        let completed = Arc::new(AtomicBool::new(false));
541        let completed_clone = completed.clone();
542
543        read_aligned_range(&mappings, 0..4096, &*service, move |res| {
544            let buffer = res.expect("read_aligned_range should succeed");
545            assert_eq!(buffer.len(), 4096);
546            assert!(buffer.as_ptr_slice().iter_as::<u8>().all(|b| b == 42));
547            completed_clone.store(true, Ordering::Relaxed);
548            ControlFlow::Continue(())
549        });
550
551        assert!(completed.load(Ordering::Relaxed));
552    }
553
554    #[test]
555    fn test_read_aligned_range_multi_extent() {
556        let mut device_data = vec![0u8; 16384];
557        device_data[0..4096].fill(1);
558        device_data[8192..12288].fill(2);
559        let service = Arc::new(FakeBlockService::new(device_data));
560
561        let mappings = Extents::try_new(
562            [Extent::new(0..4096, Some(0)), Extent::new(4096..8192, Some(8192))],
563            0,
564        )
565        .unwrap();
566
567        let completed = Arc::new(AtomicBool::new(false));
568        let completed_clone = completed.clone();
569        let mut received_bytes = Vec::new();
570
571        read_aligned_range(&mappings, 0..8192, &*service, move |res| {
572            let buffer = res.expect("read_aligned_range should succeed");
573            buffer.as_ptr_slice().append_to(&mut received_bytes);
574            if received_bytes.len() == 8192 {
575                assert!(received_bytes[0..4096].iter().all(|&b| b == 1));
576                assert!(received_bytes[4096..8192].iter().all(|&b| b == 2));
577                completed_clone.store(true, Ordering::Relaxed);
578            }
579            ControlFlow::Continue(())
580        });
581
582        assert!(completed.load(Ordering::Relaxed));
583    }
584
585    #[test]
586    fn test_read_aligned_range_sparse_extent() {
587        let service = Arc::new(FakeBlockService::new(vec![0u8; 4096]));
588
589        let mappings = Extents::try_new([Extent::new(0..4096, None)], 0).unwrap();
590
591        let completed = Arc::new(AtomicBool::new(false));
592        let completed_clone = completed.clone();
593
594        read_aligned_range(&mappings, 0..4096, &*service, move |res| {
595            let buffer = res.expect("read_aligned_range should succeed");
596            assert_eq!(buffer.len(), 4096);
597            assert!(buffer.as_ptr_slice().iter_as::<u64>().all(|b| b == 0));
598            completed_clone.store(true, Ordering::Relaxed);
599            ControlFlow::Continue(())
600        });
601
602        assert!(completed.load(Ordering::Relaxed));
603    }
604
605    #[test]
606    fn test_read_aligned_range_partial_buffer_chaining() {
607        let mut device_data = vec![0u8; 16384];
608        device_data[0..4096].fill(1);
609        device_data[4096..8192].fill(2);
610        // Request up to 8192 bytes, but cap allocator to 4096 bytes per round.
611        let service = Arc::new(FakeBlockService::new_with_cap(device_data, Some(4096)));
612
613        let mappings = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
614
615        let completed = Arc::new(AtomicBool::new(false));
616        let completed_clone = completed.clone();
617        let mut chunk_count = 0;
618        let mut received_bytes = Vec::new();
619
620        read_aligned_range(&mappings, 0..8192, &*service, move |res| {
621            let buffer = res.expect("read_aligned_range should succeed");
622            chunk_count += 1;
623            buffer.as_ptr_slice().append_to(&mut received_bytes);
624            if received_bytes.len() == 8192 {
625                assert_eq!(chunk_count, 2, "Should have streamed in 2 rounds of 4096 bytes");
626                assert!(received_bytes[0..4096].iter().all(|&b| b == 1));
627                assert!(received_bytes[4096..8192].iter().all(|&b| b == 2));
628                completed_clone.store(true, Ordering::Relaxed);
629            }
630            ControlFlow::Continue(())
631        });
632
633        assert!(completed.load(Ordering::Relaxed));
634    }
635
636    struct OutOfOrderBlockService {
637        inner: FakeBlockService,
638        delayed: Mutex<Vec<Box<dyn FnOnce() + Send>>>,
639    }
640
641    impl BlockService for OutOfOrderBlockService {
642        fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
643            self.inner.allocate_buffer(max_len)
644        }
645
646        fn read_blocks(
647            &self,
648            device_offset: u64,
649            dest_buffer: OwnedBuffer,
650            on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
651        ) -> Result<(), Error> {
652            if device_offset == 0 {
653                // Delay chunk 0 until after both chunk 1 and chunk 2 complete.
654                let inner = &self.inner;
655                let data = inner.device_data.lock();
656                let start = device_offset as usize;
657                let end = start + dest_buffer.len();
658                let mut buf = dest_buffer;
659                buf.copy_from_slice(&data[start..end]);
660                self.delayed.lock().push(Box::new(move || on_complete(Ok(buf))));
661                Ok(())
662            } else if device_offset == 4096 {
663                // Immediately complete chunk 1, adding index 1 to `buffered_chunks`.
664                self.inner.read_blocks(device_offset, dest_buffer, on_complete)?;
665                Ok(())
666            } else {
667                // Immediately complete chunk 2, adding index 2 to `buffered_chunks` and exercising
668                // `BufferedChunk::cmp` when sorting multiple out-of-order elements in the heap.
669                // Then execute delayed chunk 0 to trigger sequential delivery of 0, 1, and 2.
670                self.inner.read_blocks(device_offset, dest_buffer, on_complete)?;
671                let mut delayed = self.delayed.lock();
672                for cb in delayed.drain(..) {
673                    cb();
674                }
675                Ok(())
676            }
677        }
678    }
679
680    #[test]
681    fn test_read_aligned_range_out_of_order_reordering() {
682        let mut device_data = vec![0u8; 16384];
683        device_data[0..4096].fill(10);
684        device_data[4096..8192].fill(20);
685        device_data[8192..12288].fill(30);
686        let inner = FakeBlockService::new_with_cap(device_data, Some(4096));
687        let service = Arc::new(OutOfOrderBlockService { inner, delayed: Mutex::new(Vec::new()) });
688
689        let mappings = Extents::try_new(
690            [
691                Extent::new(0..4096, Some(0)),
692                Extent::new(4096..8192, Some(4096)),
693                Extent::new(8192..12288, Some(8192)),
694            ],
695            0,
696        )
697        .unwrap();
698
699        let completed = Arc::new(AtomicBool::new(false));
700        let completed_clone = completed.clone();
701        let mut received_chunks = Vec::new();
702
703        read_aligned_range(&mappings, 0..12288, &*service, move |res| {
704            let buffer = res.expect("read_aligned_range should succeed");
705            let first_byte = buffer.as_ptr_slice().to_vec()[0];
706            received_chunks.push(first_byte);
707            if received_chunks.len() == 3 {
708                // Verify strict delivery order: chunk 0 (10), then chunk 1 (20), then chunk 2 (30).
709                assert_eq!(received_chunks, vec![10, 20, 30]);
710                completed_clone.store(true, Ordering::Relaxed);
711            }
712            ControlFlow::Continue(())
713        });
714
715        assert!(completed.load(Ordering::Relaxed));
716    }
717
718    struct ThreadedBlockService {
719        inner: FakeBlockService,
720    }
721
722    impl BlockService for ThreadedBlockService {
723        fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
724            self.inner.allocate_buffer(max_len)
725        }
726
727        fn read_blocks(
728            &self,
729            device_offset: u64,
730            mut dest_buffer: OwnedBuffer,
731            on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
732        ) -> Result<(), Error> {
733            let start = device_offset as usize;
734            let end = start + dest_buffer.len();
735            let data = self.inner.device_data.lock();
736            dest_buffer.copy_from_slice(&data[start..end]);
737            drop(data);
738
739            // Spawn on a background thread so the upfront read_range loop can run right ahead,
740            // and block on allocate_buffer when memory is full until this thread drops its buffer.
741            thread::spawn(move || {
742                thread::sleep(Duration::from_millis(15));
743                on_complete(Ok(dest_buffer));
744            });
745            Ok(())
746        }
747    }
748
749    #[test]
750    fn test_read_aligned_range_limited_memory_blocking_allocator() {
751        let mut device_data = vec![0u8; 16384];
752        device_data[0..4096].fill(1);
753        device_data[4096..8192].fill(2);
754        device_data[8192..12288].fill(3);
755        device_data[12288..16384].fill(4);
756
757        // Pool size is only 8192 bytes, capped at 4096 bytes per buffer.
758        // Requesting 16384 bytes will allocate buffers 0 and 1 (using all 8192 bytes of the pool),
759        // submit their reads right away inside read_aligned_range's while loop, and then
760        // block when trying to allocate buffer 2 until background threads complete and free
761        // buffer 0.
762        let inner = FakeBlockService::new_with_pool_size(device_data, 8192, Some(4096));
763        let service = Arc::new(ThreadedBlockService { inner });
764
765        let mappings = Extents::try_new([Extent::new(0..16384, Some(0))], 0).unwrap();
766
767        let completed = Arc::new(AtomicBool::new(false));
768        let completed_clone = completed.clone();
769        let mut received_bytes = Vec::new();
770
771        read_aligned_range(&mappings, 0..16384, &*service, move |res| {
772            let buffer = res.expect("read_aligned_range should succeed");
773            buffer.as_ptr_slice().append_to(&mut received_bytes);
774            if received_bytes.len() == 16384 {
775                assert!(received_bytes[0..4096].iter().all(|&b| b == 1));
776                assert!(received_bytes[4096..8192].iter().all(|&b| b == 2));
777                assert!(received_bytes[8192..12288].iter().all(|&b| b == 3));
778                assert!(received_bytes[12288..16384].iter().all(|&b| b == 4));
779                completed_clone.store(true, Ordering::Relaxed);
780            }
781            ControlFlow::Continue(())
782        });
783
784        // Wait briefly for background threads to finish completing all 4 chunks.
785        for _ in 0..100 {
786            if completed.load(Ordering::Relaxed) {
787                break;
788            }
789            thread::sleep(Duration::from_millis(10));
790        }
791        assert!(completed.load(Ordering::Relaxed));
792    }
793
794    #[test]
795    fn test_read_aligned_range_max_read_buffer_size_capping() {
796        struct CapturingBlockService {
797            inner: FakeBlockService,
798            requested_lens: Mutex<Vec<usize>>,
799        }
800
801        impl BlockService for CapturingBlockService {
802            fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
803                self.requested_lens.lock().push(max_len);
804                self.inner.allocate_buffer(max_len)
805            }
806
807            fn read_blocks(
808                &self,
809                device_offset: u64,
810                dest_buffer: OwnedBuffer,
811                on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
812            ) -> Result<(), Error> {
813                self.inner.read_blocks(device_offset, dest_buffer, on_complete)
814            }
815        }
816
817        let inner = FakeBlockService::new(vec![0u8; MAX_READ_BUFFER_SIZE + 8192]);
818        let service = CapturingBlockService { inner, requested_lens: Mutex::new(Vec::new()) };
819        let mappings =
820            Extents::try_new([Extent::new(0..(MAX_READ_BUFFER_SIZE + 8192) as u64, Some(0))], 0)
821                .unwrap();
822
823        read_aligned_range(&mappings, 0..(MAX_READ_BUFFER_SIZE + 8192) as u64, &service, |_| {
824            ControlFlow::Continue(())
825        });
826
827        assert_eq!(*service.requested_lens.lock(), vec![MAX_READ_BUFFER_SIZE, 8192]);
828    }
829
830    #[test]
831    fn test_read_aligned_range_sync_read_blocks_error() {
832        struct SyncErrorBlockService(FakeBlockService);
833        impl BlockService for SyncErrorBlockService {
834            fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
835                self.0.allocate_buffer(max_len)
836            }
837            fn read_blocks(
838                &self,
839                _offset: u64,
840                _dest: OwnedBuffer,
841                _on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
842            ) -> Result<(), Error> {
843                Err(anyhow!("synchronous read_blocks failure"))
844            }
845        }
846        let service = SyncErrorBlockService(FakeBlockService::new(vec![0u8; 4096]));
847        let mappings = Extents::try_new([Extent::new(0..4096, Some(0))], 0).unwrap();
848
849        let err_received = Arc::new(AtomicBool::new(false));
850        let err_clone = err_received.clone();
851        read_aligned_range(&mappings, 0..4096, &service, move |res| {
852            assert!(res.is_err());
853            err_clone.store(true, Ordering::Relaxed);
854            ControlFlow::Continue(())
855        });
856        assert!(err_received.load(Ordering::Relaxed));
857    }
858
859    #[test]
860    fn test_read_aligned_range_async_callback_errors() {
861        struct AsyncErrorBlockService(FakeBlockService, u64);
862        impl BlockService for AsyncErrorBlockService {
863            fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
864                self.0.allocate_buffer(max_len)
865            }
866            fn read_blocks(
867                &self,
868                device_offset: u64,
869                dest_buffer: OwnedBuffer,
870                on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
871            ) -> Result<(), Error> {
872                if device_offset == self.1 {
873                    on_complete(Err(anyhow!("async block read error on offset {device_offset}")));
874                    Ok(())
875                } else {
876                    self.0.read_blocks(device_offset, dest_buffer, on_complete)
877                }
878            }
879        }
880
881        // Test in-order error on chunk 0.
882        let service = AsyncErrorBlockService(FakeBlockService::new(vec![0u8; 8192]), 0);
883        let mappings = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
884        let err_count = Arc::new(AtomicUsize::new(0));
885        let err_clone = err_count.clone();
886        read_aligned_range(&mappings, 0..8192, &service, move |res| {
887            if res.is_err() {
888                err_clone.fetch_add(1, Ordering::Relaxed);
889            }
890            ControlFlow::Continue(())
891        });
892        assert_eq!(err_count.load(Ordering::Relaxed), 1);
893
894        // Test out-of-order error where chunk 1 (4096) succeeds first, then chunk 0 (0) fails.
895        let service = Arc::new(OutOfOrderBlockService {
896            inner: FakeBlockService::new(vec![0u8; 8192]),
897            delayed: Mutex::new(Vec::new()),
898        });
899        // We wrap OutOfOrderBlockService to make chunk 0 fail when drained.
900        struct FailingDelayedService(Arc<OutOfOrderBlockService>);
901        impl BlockService for FailingDelayedService {
902            fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
903                self.0.allocate_buffer(max_len)
904            }
905            fn read_blocks(
906                &self,
907                device_offset: u64,
908                dest_buffer: OwnedBuffer,
909                on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
910            ) -> Result<(), Error> {
911                if device_offset == 0 {
912                    self.0.delayed.lock().push(Box::new(move || {
913                        on_complete(Err(anyhow!("delayed chunk 0 failure")))
914                    }));
915                    Ok(())
916                } else {
917                    let res = self.0.read_blocks(device_offset, dest_buffer, on_complete);
918                    let mut delayed = self.0.delayed.lock();
919                    for cb in delayed.drain(..) {
920                        cb();
921                    }
922                    res
923                }
924            }
925        }
926        let failing_service = FailingDelayedService(service.clone());
927        let err_count = Arc::new(AtomicUsize::new(0));
928        let err_clone = err_count.clone();
929        read_aligned_range(&mappings, 0..8192, &failing_service, move |res| {
930            if res.is_err() {
931                err_clone.fetch_add(1, Ordering::Relaxed);
932            }
933            ControlFlow::Continue(())
934        });
935        let mut delayed = service.delayed.lock();
936        for cb in delayed.drain(..) {
937            cb();
938        }
939        assert_eq!(err_count.load(Ordering::Relaxed), 1);
940    }
941
942    #[test]
943    fn test_read_context_dropped_before_completion() {
944        struct DroppingBlockService(Mutex<Vec<Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>>>);
945        impl BlockService for DroppingBlockService {
946            fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
947                FakeBlockService::new(vec![0u8; max_len]).allocate_buffer(max_len)
948            }
949            fn read_blocks(
950                &self,
951                _offset: u64,
952                _dest: OwnedBuffer,
953                on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
954            ) -> Result<(), Error> {
955                self.0.lock().push(on_complete);
956                Ok(())
957            }
958        }
959        let service = DroppingBlockService(Mutex::new(Vec::new()));
960        let mappings = Extents::try_new([Extent::new(0..4096, Some(0))], 0).unwrap();
961
962        let err_msg = Arc::new(Mutex::new(String::new()));
963        let err_msg_clone = err_msg.clone();
964        read_aligned_range(&mappings, 0..4096, &service, move |res| {
965            if let Err(e) = res {
966                *err_msg_clone.lock() = e.to_string();
967            }
968            ControlFlow::Continue(())
969        });
970        service.0.lock().clear();
971        assert_eq!(*err_msg.lock(), "Read sub-request dropped before completion");
972    }
973
974    #[test]
975    fn test_read_aligned_range_multi_iteration_error_break() {
976        struct FirstIterSyncErrorService(FakeBlockService);
977        impl BlockService for FirstIterSyncErrorService {
978            fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
979                self.0.allocate_buffer(max_len)
980            }
981            fn read_blocks(
982                &self,
983                device_offset: u64,
984                dest_buffer: OwnedBuffer,
985                on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
986            ) -> Result<(), Error> {
987                if device_offset == 0 {
988                    Err(anyhow!("sync failure on first iteration"))
989                } else {
990                    self.0.read_blocks(device_offset, dest_buffer, on_complete)
991                }
992            }
993        }
994        let inner = FakeBlockService::new_with_cap(vec![0u8; 8192], Some(4096));
995        let service = FirstIterSyncErrorService(inner);
996        let mappings = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
997
998        let err_received = Arc::new(AtomicBool::new(false));
999        let err_clone = err_received.clone();
1000        read_aligned_range(&mappings, 0..8192, &service, move |res| {
1001            if res.is_err() {
1002                err_clone.store(true, Ordering::Relaxed);
1003            }
1004            ControlFlow::Continue(())
1005        });
1006        assert!(err_received.load(Ordering::Relaxed));
1007    }
1008
1009    #[test]
1010    fn test_read_aligned_range_extent_beyond_actual_end_break() {
1011        let inner = FakeBlockService::new_with_cap(vec![0u8; 8192], Some(4096));
1012        let mappings = Extents::try_new(
1013            [Extent::new(0..4096, Some(0)), Extent::new(4096..8192, Some(4096))],
1014            0,
1015        )
1016        .unwrap();
1017
1018        let count = Arc::new(AtomicUsize::new(0));
1019        let count_clone = count.clone();
1020        read_aligned_range(&mappings, 0..8192, &inner, move |res| {
1021            if res.is_ok() {
1022                count_clone.fetch_add(1, Ordering::Relaxed);
1023            }
1024            ControlFlow::Continue(())
1025        });
1026        assert_eq!(count.load(Ordering::Relaxed), 2);
1027    }
1028
1029    #[test]
1030    fn test_read_aligned_range_merged_extents() {
1031        let mut device_data = vec![0u8; 32768];
1032        device_data[4096..8192].fill(0xAA);
1033        device_data[8192..20480].fill(0xBB);
1034        let service = Arc::new(FakeBlockService::new(device_data));
1035
1036        let mappings = Extents::try_new(
1037            [Extent::new(0..4096, Some(4096)), Extent::new(4096..16384, Some(8192))],
1038            0,
1039        )
1040        .unwrap();
1041
1042        let count = Arc::new(AtomicUsize::new(0));
1043        let count_clone = count.clone();
1044        let completed = Arc::new(AtomicBool::new(false));
1045        let completed_clone = completed.clone();
1046
1047        read_aligned_range(&mappings, 0..16384, &*service, move |res| {
1048            count_clone.fetch_add(1, Ordering::Relaxed);
1049            let buffer = res.expect("merged read should succeed");
1050            assert_eq!(buffer.len(), 16384);
1051            assert!(buffer.as_ptr_slice().subslice(0..4096).iter_as::<u8>().all(|b| b == 0xAA));
1052            assert!(buffer.as_ptr_slice().subslice(4096..16384).iter_as::<u8>().all(|b| b == 0xBB));
1053            completed_clone.store(true, Ordering::Relaxed);
1054            ControlFlow::Continue(())
1055        });
1056
1057        assert_eq!(count.load(Ordering::Relaxed), 1);
1058        assert!(completed.load(Ordering::Relaxed));
1059    }
1060
1061    #[test]
1062    fn test_read_aligned_range_callback_error_propagation() {
1063        let service = Arc::new(FakeBlockService::new(vec![0u8; 8192]));
1064        let mappings = Extents::try_new([Extent::new(0..8192, Some(0))], 0).unwrap();
1065
1066        let call_count = Arc::new(AtomicUsize::new(0));
1067        let call_count_clone = call_count.clone();
1068        read_aligned_range(&mappings, 0..8192, &*service, move |_res| {
1069            call_count_clone.fetch_add(1, Ordering::Relaxed);
1070            ControlFlow::Break(())
1071        });
1072
1073        assert_eq!(call_count.load(Ordering::Relaxed), 1);
1074    }
1075
1076    #[test]
1077    fn test_read_buffer_from_extents_single_extent() {
1078        let mut device_data = vec![0u8; 16384];
1079        device_data[4096..8192].fill(0xAA);
1080        let service = Arc::new(FakeBlockService::new(device_data));
1081
1082        // Logical 0..4096 maps to device 4096..8192.
1083        let mappings = Extents::try_new([Extent::new(0..4096, Some(4096))], 0).unwrap();
1084
1085        let dest = service.allocate_buffer(4096);
1086        let (tx, rx) = std::sync::mpsc::channel();
1087        read_buffer_from_extents(&mappings, 0, dest, service.as_ref(), move |res| {
1088            tx.send(res).unwrap();
1089        })
1090        .unwrap();
1091
1092        let buffer = rx.recv().unwrap().unwrap();
1093        assert_eq!(buffer.len(), 4096);
1094        assert!(buffer.as_ptr_slice().iter_as::<u8>().all(|b| b == 0xAA));
1095    }
1096
1097    #[test]
1098    fn test_read_buffer_from_extents_sparse_hole_fast_path() {
1099        let service = Arc::new(FakeBlockService::new(vec![0u8; 16384]));
1100
1101        // Logical 0..4096 is a sparse hole (None).
1102        let mappings = Extents::try_new([Extent::new(0..4096, None)], 0).unwrap();
1103
1104        let dest = service.allocate_buffer(4096);
1105        let (tx, rx) = std::sync::mpsc::channel();
1106        read_buffer_from_extents(&mappings, 0, dest, service.as_ref(), move |res| {
1107            tx.send(res).unwrap();
1108        })
1109        .unwrap();
1110
1111        let buffer = rx.recv().unwrap().unwrap();
1112        assert_eq!(buffer.len(), 4096);
1113        assert!(buffer.as_ptr_slice().iter_as::<u8>().all(|b| b == 0x00));
1114    }
1115
1116    #[test]
1117    fn test_read_buffer_from_extents_multi_extent_and_holes() {
1118        let mut device_data = vec![0u8; 32768];
1119        device_data[0..4096].fill(0x11);
1120        device_data[8192..12288].fill(0x22);
1121        let service = Arc::new(FakeBlockService::new(device_data));
1122
1123        // Child maps:
1124        // 0..4096 -> device 0..4096 (0x11)
1125        // 4096..8192 -> hole (0x00)
1126        // 8192..12288 -> device 8192..12288 (0x22)
1127        let mappings = Extents::try_new(
1128            [
1129                Extent::new(0..4096, Some(0)),
1130                Extent::new(4096..8192, None),
1131                Extent::new(8192..12288, Some(8192)),
1132            ],
1133            0,
1134        )
1135        .unwrap();
1136
1137        let dest = service.allocate_buffer(12288);
1138        let (tx, rx) = std::sync::mpsc::channel();
1139        read_buffer_from_extents(&mappings, 0, dest, service.as_ref(), move |res| {
1140            tx.send(res).unwrap();
1141        })
1142        .unwrap();
1143
1144        let buffer = rx.recv().unwrap().unwrap();
1145        assert_eq!(buffer.len(), 12288);
1146        assert!(buffer.as_ptr_slice().subslice(0..4096).iter_as::<u8>().all(|b| b == 0x11));
1147        assert!(buffer.as_ptr_slice().subslice(4096..8192).iter_as::<u8>().all(|b| b == 0x00));
1148        assert!(buffer.as_ptr_slice().subslice(8192..12288).iter_as::<u8>().all(|b| b == 0x22));
1149    }
1150
1151    #[test]
1152    fn test_read_buffer_from_extents_sync_error_disarms_callback() {
1153        struct SecondChunkSyncErrorService {
1154            inner: FakeBlockService,
1155        }
1156        impl BlockService for SecondChunkSyncErrorService {
1157            fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
1158                self.inner.allocate_buffer(max_len)
1159            }
1160            fn read_blocks(
1161                &self,
1162                device_offset: u64,
1163                dest_buffer: OwnedBuffer,
1164                on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
1165            ) -> Result<(), Error> {
1166                if device_offset == 8192 {
1167                    Err(anyhow!("sync failure on second chunk"))
1168                } else {
1169                    self.inner.read_blocks(device_offset, dest_buffer, on_complete)
1170                }
1171            }
1172        }
1173
1174        let inner = FakeBlockService::new(vec![0u8; 32768]);
1175        let service = Arc::new(SecondChunkSyncErrorService { inner });
1176
1177        // Two extents: 0..4096 (offset 0) and 4096..8192 (offset 8192).
1178        let mappings = Extents::try_new(
1179            [Extent::new(0..4096, Some(0)), Extent::new(4096..8192, Some(8192))],
1180            0,
1181        )
1182        .unwrap();
1183
1184        let dest = service.allocate_buffer(8192);
1185        let completed = Arc::new(AtomicBool::new(false));
1186        let completed_clone = completed.clone();
1187
1188        let res = read_buffer_from_extents(&mappings, 0, dest, service.as_ref(), move |_res| {
1189            completed_clone.store(true, Ordering::Relaxed);
1190        });
1191
1192        // The dispatch must fail synchronously with Err:
1193        assert!(res.is_err());
1194        // The callback must NOT have been invoked (and must not be invoked by chunk 1 in-flight
1195        // completion):
1196        assert!(!completed.load(Ordering::Relaxed));
1197    }
1198
1199    #[test]
1200    fn test_read_buffer_from_extents_in_flight_chunk_succeeds_after_sync_error() {
1201        struct InFlightFirstChunkService {
1202            inner: FakeBlockService,
1203            pending_first_chunk:
1204                Mutex<Option<(OwnedBuffer, Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>)>>,
1205        }
1206        impl BlockService for InFlightFirstChunkService {
1207            fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
1208                self.inner.allocate_buffer(max_len)
1209            }
1210            fn read_blocks(
1211                &self,
1212                device_offset: u64,
1213                dest_buffer: OwnedBuffer,
1214                on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
1215            ) -> Result<(), Error> {
1216                if device_offset == 0 {
1217                    // Keep chunk 1 in flight in the background without invoking callback yet.
1218                    *self.pending_first_chunk.lock() = Some((dest_buffer, on_complete));
1219                    Ok(())
1220                } else if device_offset == 8192 {
1221                    // Chunk 2 fails synchronously.
1222                    Err(anyhow!("sync failure on second chunk"))
1223                } else {
1224                    self.inner.read_blocks(device_offset, dest_buffer, on_complete)
1225                }
1226            }
1227        }
1228
1229        let inner = FakeBlockService::new(vec![0u8; 32768]);
1230        let service =
1231            Arc::new(InFlightFirstChunkService { inner, pending_first_chunk: Mutex::new(None) });
1232
1233        // Two extents: 0..4096 (offset 0) and 4096..8192 (offset 8192).
1234        let mappings = Extents::try_new(
1235            [Extent::new(0..4096, Some(0)), Extent::new(4096..8192, Some(8192))],
1236            0,
1237        )
1238        .unwrap();
1239
1240        let dest = service.allocate_buffer(8192);
1241        let completed = Arc::new(AtomicBool::new(false));
1242        let completed_clone = completed.clone();
1243
1244        let res = read_buffer_from_extents(&mappings, 0, dest, service.as_ref(), move |_res| {
1245            completed_clone.store(true, Ordering::Relaxed);
1246        });
1247
1248        // The dispatch must fail synchronously with Err:
1249        assert!(res.is_err());
1250        assert!(!completed.load(Ordering::Relaxed));
1251
1252        // Now complete the in-flight chunk 1 with Ok after dispatch has already failed.
1253        let (buf, cb) = service.pending_first_chunk.lock().take().expect("chunk 1 pending");
1254        cb(Ok(buf));
1255
1256        // The user callback must NOT have been called.
1257        assert!(!completed.load(Ordering::Relaxed));
1258    }
1259
1260    #[test]
1261    fn test_read_buffer_from_extents_async_error_disarms_subsequent_chunks() {
1262        struct AsyncErrorService {
1263            inner: FakeBlockService,
1264            pending_chunks:
1265                Mutex<Vec<(OwnedBuffer, Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>, u64)>>,
1266        }
1267        impl BlockService for AsyncErrorService {
1268            fn allocate_buffer(&self, max_len: usize) -> OwnedBuffer {
1269                self.inner.allocate_buffer(max_len)
1270            }
1271            fn read_blocks(
1272                &self,
1273                device_offset: u64,
1274                dest_buffer: OwnedBuffer,
1275                on_complete: Box<dyn FnOnce(Result<OwnedBuffer, Error>) + Send>,
1276            ) -> Result<(), Error> {
1277                self.pending_chunks.lock().push((dest_buffer, on_complete, device_offset));
1278                Ok(())
1279            }
1280        }
1281
1282        let inner = FakeBlockService::new(vec![0u8; 32768]);
1283        let service = Arc::new(AsyncErrorService { inner, pending_chunks: Mutex::new(Vec::new()) });
1284
1285        let mappings = Extents::try_new(
1286            [Extent::new(0..4096, Some(0)), Extent::new(4096..8192, Some(8192))],
1287            0,
1288        )
1289        .unwrap();
1290
1291        let dest = service.allocate_buffer(8192);
1292        let err_count = Arc::new(AtomicUsize::new(0));
1293        let ok_count = Arc::new(AtomicUsize::new(0));
1294        let err_count_clone = err_count.clone();
1295        let ok_count_clone = ok_count.clone();
1296
1297        let res = read_buffer_from_extents(&mappings, 0, dest, service.as_ref(), move |res| {
1298            if res.is_ok() {
1299                ok_count_clone.fetch_add(1, Ordering::Relaxed);
1300            } else {
1301                err_count_clone.fetch_add(1, Ordering::Relaxed);
1302            }
1303        });
1304        assert!(res.is_ok());
1305
1306        let mut chunks = service.pending_chunks.lock();
1307        assert_eq!(chunks.len(), 2);
1308        let (buf1, cb1, _) = chunks.remove(0);
1309        let (buf2, cb2, _) = chunks.remove(0);
1310        drop(chunks);
1311
1312        // Chunk 1 fails asynchronously:
1313        drop(buf1);
1314        cb1(Err(anyhow!("async failure on chunk 1")));
1315        assert_eq!(err_count.load(Ordering::Relaxed), 1);
1316        assert_eq!(ok_count.load(Ordering::Relaxed), 0);
1317
1318        // Chunk 2 completes successfully afterwards:
1319        cb2(Ok(buf2));
1320        // Callback must NOT have been called again (ok_count remains 0, err_count remains 1):
1321        assert_eq!(err_count.load(Ordering::Relaxed), 1);
1322        assert_eq!(ok_count.load(Ordering::Relaxed), 0);
1323    }
1324
1325    #[test]
1326    fn test_child_block_service() {
1327        let mut device_data = vec![0u8; 16384];
1328        device_data[4096..8192].fill(0x55);
1329        let parent_service = Arc::new(FakeBlockService::new(device_data));
1330
1331        let mappings = Extents::try_new([Extent::new(0..4096, Some(4096))], 0).unwrap();
1332        let file = Arc::new(File::new(mappings, 4096, Transform::None));
1333
1334        let child_service = ChildBlockService::new(parent_service.clone(), file.clone());
1335        assert_eq!(child_service.file().extents().iter_extents(0).count(), 1);
1336
1337        let dest = child_service.allocate_buffer(4096);
1338        let (tx, rx) = std::sync::mpsc::channel();
1339        child_service
1340            .read_blocks(
1341                0,
1342                dest,
1343                Box::new(move |res| {
1344                    tx.send(res).unwrap();
1345                }),
1346            )
1347            .unwrap();
1348
1349        let buffer = rx.recv().unwrap().unwrap();
1350        assert_eq!(buffer.len(), 4096);
1351        assert!(buffer.as_ptr_slice().iter_as::<u8>().all(|b| b == 0x55));
1352    }
1353
1354    #[test]
1355    fn test_read_buffer_from_extents_non_aligned_start() {
1356        let base_device_offset = 17408u64; // e.g. LBA 34 on 512-byte sector disk
1357        let mut device_data = vec![0u8; 32768];
1358        device_data[base_device_offset as usize..(base_device_offset as usize + 4096)].fill(0x77);
1359        let service = Arc::new(FakeBlockService::new(device_data));
1360
1361        let mappings =
1362            Extents::try_new([Extent::new(0..4096, Some(base_device_offset))], base_device_offset)
1363                .unwrap();
1364
1365        let dest = service.allocate_buffer(4096);
1366        let (tx, rx) = std::sync::mpsc::channel();
1367        read_buffer_from_extents(&mappings, 0, dest, service.as_ref(), move |res| {
1368            tx.send(res).unwrap();
1369        })
1370        .unwrap();
1371
1372        let buffer = rx.recv().unwrap().unwrap();
1373        assert_eq!(buffer.len(), 4096);
1374        assert!(buffer.as_ptr_slice().iter_as::<u8>().all(|b| b == 0x77));
1375    }
1376}