Skip to main content

mapping/
blob.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::Extents;
6use crate::reader::{BlockService, read_aligned_range};
7use delivery_blob::compression::{CompressionInfo, StreamingDecompressor};
8use std::cmp::min;
9use std::ops::{ControlFlow, Range};
10use std::sync::Arc;
11
12use delivery_blob::DataBuffer;
13
14/// A mapped blob containing extents and decompression metadata.
15pub struct Blob {
16    extents: Extents,
17    uncompressed_size: u64,
18    compression_info: Option<Arc<CompressionInfo>>,
19}
20
21impl Blob {
22    pub fn new(
23        extents: Extents,
24        uncompressed_size: u64,
25        compression_info: Option<CompressionInfo>,
26    ) -> Self {
27        Self { extents, uncompressed_size, compression_info: compression_info.map(Arc::new) }
28    }
29
30    /// Returns the extents mapping logical offsets to device offsets.
31    pub fn extents(&self) -> &Extents {
32        &self.extents
33    }
34
35    /// Returns the uncompressed size of the blob in bytes.
36    pub fn uncompressed_size(&self) -> u64 {
37        self.uncompressed_size
38    }
39
40    /// Returns decompression metadata if the blob is compressed.
41    pub fn compression_info(&self) -> Option<&CompressionInfo> {
42        self.compression_info.as_deref()
43    }
44
45    /// Streams and decodes the specified uncompressed `range` into the provided `dest_buf`.
46    ///
47    /// For uncompressed blobs, both `range.start` and `range.end` must be multiples of
48    /// `BLOCK_SIZE`. For compressed blobs, `range.start` must be a multiple of the compression
49    /// chunk size, and `range.end` must either be a multiple of the chunk size or equal to
50    /// `uncompressed_size`.
51    pub fn read_range(
52        &self,
53        range: Range<u64>,
54        service: &(impl BlockService + ?Sized),
55        mut dest_buf: impl DataBuffer,
56    ) {
57        if range.is_empty() {
58            return;
59        }
60
61        match &self.compression_info {
62            None => {
63                let mut current_offset = range.start;
64                let uncompressed_size = self.uncompressed_size;
65
66                read_aligned_range(&self.extents, range, service, move |res| {
67                    let buffer = match res {
68                        Ok(buf) => buf,
69                        Err(_) => {
70                            return ControlFlow::Break(());
71                        }
72                    };
73                    let valid_len =
74                        min(buffer.len() as u64, uncompressed_size.saturating_sub(current_offset))
75                            as usize;
76                    if valid_len > 0 {
77                        dest_buf
78                            .mut_ptr_slice()
79                            .subslice_mut(0..valid_len)
80                            .copy_from_slice(&buffer.as_slice()[..valid_len]);
81                        if dest_buf.commit(valid_len).is_err() {
82                            return ControlFlow::Break(());
83                        }
84                    }
85                    current_offset += buffer.len() as u64;
86                    ControlFlow::Continue(())
87                });
88            }
89            Some(info) => {
90                let info = Arc::clone(info);
91                let Ok((mut decompressor, aligned_range)) =
92                    StreamingDecompressor::new(info, range, self.uncompressed_size, dest_buf)
93                else {
94                    // The range must be out of range. This should be handled when `dest_buf`
95                    // is dropped.
96                    return;
97                };
98
99                read_aligned_range(&self.extents, aligned_range, service, move |res| {
100                    let buffer = match res {
101                        Ok(buf) => buf,
102                        Err(_) => {
103                            return ControlFlow::Break(());
104                        }
105                    };
106                    if decompressor.push(buffer.as_slice()).is_err() {
107                        return ControlFlow::Break(());
108                    }
109                    ControlFlow::Continue(())
110                });
111            }
112        }
113    }
114}
115
116#[cfg(test)]
117mod tests {
118    use super::*;
119    use crate::reader::tests::FakeBlockService;
120    use crate::{BLOCK_SIZE, Extent};
121    use anyhow::Error;
122    use delivery_blob::compression::{
123        ChunkedArchiveError, ChunkedArchiveOptions, CompressionAlgorithm,
124    };
125    use fuchsia_sync::Mutex;
126    use std::sync::Arc;
127    use storage_ptr_slice::MutPtrByteSlice;
128
129    #[derive(Default)]
130    struct TestVecBufferInner {
131        commits: Vec<(u64, usize)>,
132        output: Vec<u8>,
133    }
134
135    #[derive(Clone)]
136    struct TestVecBufferReceiver(Arc<Mutex<TestVecBufferInner>>);
137
138    impl TestVecBufferReceiver {
139        fn commits(&self) -> Vec<(u64, usize)> {
140            self.0.lock().commits.clone()
141        }
142
143        fn output(&self) -> Vec<u8> {
144            self.0.lock().output.clone()
145        }
146    }
147
148    struct TestVecBuffer {
149        data: Vec<u8>,
150        committed_len: usize,
151        offset: u64,
152        receiver: TestVecBufferReceiver,
153    }
154
155    impl TestVecBuffer {
156        fn new(size: usize) -> (Self, TestVecBufferReceiver) {
157            Self::new_with_offset(size, 0)
158        }
159
160        fn new_with_offset(size: usize, offset: u64) -> (Self, TestVecBufferReceiver) {
161            let receiver =
162                TestVecBufferReceiver(Arc::new(Mutex::new(TestVecBufferInner::default())));
163            let buf = Self {
164                data: vec![0u8; size],
165                committed_len: 0,
166                offset,
167                receiver: receiver.clone(),
168            };
169            (buf, receiver)
170        }
171    }
172
173    impl Drop for TestVecBuffer {
174        fn drop(&mut self) {
175            self.receiver.0.lock().output = std::mem::take(&mut self.data);
176        }
177    }
178
179    impl DataBuffer for TestVecBuffer {
180        fn mut_ptr_slice(&mut self) -> MutPtrByteSlice<'_> {
181            let remaining = &mut self.data[self.committed_len..];
182            MutPtrByteSlice::from(remaining)
183        }
184
185        fn commit(&mut self, size: usize) -> Result<(), ChunkedArchiveError> {
186            self.receiver.0.lock().commits.push((self.offset, size));
187            self.offset += size as u64;
188            self.committed_len += size;
189            Ok(())
190        }
191    }
192
193    #[test]
194    fn test_read_range_uncompressed() {
195        let block_count = 8;
196        let mut expected_data = vec![0u8; (block_count as u64 * BLOCK_SIZE) as usize];
197        for (i, byte) in expected_data.iter_mut().enumerate() {
198            *byte = (i % 255) as u8;
199        }
200        let service = FakeBlockService::new(expected_data.clone());
201
202        let extents = Extents::encode_extents(&[Extent::new(0..(8 * BLOCK_SIZE), Some(0))]);
203        let extents = Extents::from_encoded(&extents).unwrap();
204        let blob = Arc::new(Blob::new(extents, 8 * BLOCK_SIZE, None));
205
206        let (dest_buf, rx) = TestVecBuffer::new(expected_data.len());
207        blob.read_range(0..(8 * BLOCK_SIZE), &service, dest_buf);
208
209        assert_eq!(rx.commits(), vec![(0, (8 * BLOCK_SIZE) as usize)]);
210        assert_eq!(rx.output(), expected_data);
211    }
212
213    #[test]
214    fn test_read_range_compressed_zstd() {
215        let uncompressed_size = 32768 * 2 + 1024;
216        let mut uncompressed_data = vec![0u8; uncompressed_size];
217        for (i, byte) in uncompressed_data.iter_mut().enumerate() {
218            *byte = ((i * 7) % 255) as u8;
219        }
220
221        let options =
222            ChunkedArchiveOptions::V3 { compression_algorithm: CompressionAlgorithm::Zstd };
223        let archive =
224            delivery_blob::compression::ChunkedArchive::new(&uncompressed_data, options).unwrap();
225
226        let mut compressed_offsets = vec![0];
227        let mut compressed_data = vec![];
228        for chunk in archive.chunks() {
229            compressed_data.extend_from_slice(&chunk.compressed_data);
230            compressed_offsets.push(compressed_data.len() as u64);
231        }
232        compressed_offsets.pop();
233
234        let chunk_size = archive.chunk_size();
235        let stored_size = compressed_data.len() as u64;
236        let stored_blocks = stored_size.div_ceil(BLOCK_SIZE);
237        let mut device_data = vec![0u8; (stored_blocks * BLOCK_SIZE) as usize];
238        device_data[..compressed_data.len()].copy_from_slice(&compressed_data);
239        let service = FakeBlockService::new(device_data);
240
241        let extents =
242            Extents::encode_extents(&[Extent::new(0..(stored_blocks * BLOCK_SIZE), Some(0))]);
243        let extents = Extents::from_encoded(&extents).unwrap();
244        let compression_info = CompressionInfo::new(
245            chunk_size as u64,
246            stored_size,
247            &compressed_offsets,
248            CompressionAlgorithm::Zstd,
249        )
250        .unwrap();
251        let blob = Arc::new(Blob::new(extents, uncompressed_size as u64, Some(compression_info)));
252
253        let dest_alloc_size = uncompressed_size.next_multiple_of(chunk_size);
254        let (dest_buf, rx) = TestVecBuffer::new(dest_alloc_size);
255        blob.read_range(0..(uncompressed_size as u64), &service, dest_buf);
256
257        assert_eq!(
258            rx.commits(),
259            vec![(0, chunk_size), (chunk_size as u64, chunk_size), (chunk_size as u64 * 2, 1024)]
260        );
261        assert_eq!(&rx.output()[..uncompressed_size], &uncompressed_data[..]);
262    }
263
264    #[test]
265    fn test_read_range_compressed_lz4_split_across_buffers() {
266        let uncompressed_size = 32768 * 2;
267        let mut uncompressed_data = vec![0u8; uncompressed_size];
268        for (i, byte) in uncompressed_data.iter_mut().enumerate() {
269            *byte = ((i * 13) % 255) as u8;
270        }
271
272        let options =
273            ChunkedArchiveOptions::V3 { compression_algorithm: CompressionAlgorithm::Lz4 };
274        let archive =
275            delivery_blob::compression::ChunkedArchive::new(&uncompressed_data, options).unwrap();
276
277        let mut compressed_offsets = vec![0];
278        let mut compressed_data = vec![];
279        for chunk in archive.chunks() {
280            compressed_data.extend_from_slice(&chunk.compressed_data);
281            compressed_offsets.push(compressed_data.len() as u64);
282        }
283        compressed_offsets.pop();
284
285        let chunk_size = archive.chunk_size();
286        let stored_size = compressed_data.len() as u64;
287        let stored_blocks = stored_size.div_ceil(BLOCK_SIZE);
288        let mut device_data = vec![0u8; (stored_blocks * BLOCK_SIZE) as usize];
289        device_data[..compressed_data.len()].copy_from_slice(&compressed_data);
290
291        // Force a small block allocation limit (e.g. 4096 bytes) so that read_aligned_range
292        // splits the compressed chunks across multiple consecutive OwnedBuffers!
293        let service = FakeBlockService::new_with_cap(device_data, Some(4096));
294
295        let extents =
296            Extents::encode_extents(&[Extent::new(0..(stored_blocks * BLOCK_SIZE), Some(0))]);
297        let extents = Extents::from_encoded(&extents).unwrap();
298        let compression_info = CompressionInfo::new(
299            chunk_size as u64,
300            stored_size,
301            &compressed_offsets,
302            CompressionAlgorithm::Lz4,
303        )
304        .unwrap();
305        let blob = Arc::new(Blob::new(extents, uncompressed_size as u64, Some(compression_info)));
306
307        let (dest_buf, rx) = TestVecBuffer::new(uncompressed_size);
308        blob.read_range(0..(uncompressed_size as u64), &service, dest_buf);
309
310        assert_eq!(rx.commits(), vec![(0, chunk_size), (chunk_size as u64, chunk_size)]);
311        assert_eq!(rx.output(), uncompressed_data);
312    }
313
314    #[test]
315    fn test_read_range_invalid_range_noop() {
316        let service = FakeBlockService::new(vec![0u8; 8192]);
317        let extents = Extents::encode_extents(&[Extent::new(0..8192, Some(0))]);
318        let extents = Extents::from_encoded(&extents).unwrap();
319        let blob = Arc::new(Blob::new(extents, 8192, None));
320
321        let (dest_buf, rx) = TestVecBuffer::new_with_offset(0, 4096);
322        // start >= end should be a no-op returning Ok(())
323        blob.read_range(4096..4096, &service, dest_buf);
324        assert_eq!(rx.commits().len(), 0);
325    }
326
327    #[test]
328    fn test_blob_getters() {
329        let extents_raw = Extents::encode_extents(&[Extent::new(0..8192, Some(0))]);
330        let extents = Extents::from_encoded(&extents_raw).unwrap();
331        let uncompressed_size = 8192u64;
332
333        let blob_uncompressed = Blob::new(extents, uncompressed_size, None);
334        assert_eq!(blob_uncompressed.uncompressed_size(), 8192);
335        assert!(blob_uncompressed.compression_info().is_none());
336
337        let compression_info =
338            CompressionInfo::new(32768, 4096, &[0], CompressionAlgorithm::Zstd).unwrap();
339        let blob_compressed = Blob::new(
340            Extents::from_encoded(&extents_raw).unwrap(),
341            uncompressed_size,
342            Some(compression_info),
343        );
344        assert!(blob_compressed.compression_info().is_some());
345    }
346
347    #[test]
348    fn test_read_range_block_service_error_returns_err() {
349        struct FailingBlockService;
350        impl BlockService for FailingBlockService {
351            fn allocate_buffer(&self, max_len: usize) -> storage_device::buffer::OwnedBuffer {
352                FakeBlockService::new(vec![0u8; max_len]).allocate_buffer(max_len)
353            }
354            fn read_blocks(
355                &self,
356                _device_offset: u64,
357                _dest_buffer: storage_device::buffer::OwnedBuffer,
358                _on_complete: Box<
359                    dyn FnOnce(Result<storage_device::buffer::OwnedBuffer, Error>) + Send,
360                >,
361            ) -> Result<(), Error> {
362                Err(anyhow::anyhow!("block read failure"))
363            }
364        }
365
366        let extents = Extents::encode_extents(&[Extent::new(0..8192, Some(0))]);
367        let extents = Extents::from_encoded(&extents).unwrap();
368        let blob = Blob::new(extents, 8192, None);
369
370        let (dest_buf, rx) = TestVecBuffer::new(8192);
371
372        blob.read_range(0..8192, &FailingBlockService, dest_buf);
373        assert_eq!(rx.commits().len(), 0);
374    }
375
376    #[test]
377    fn test_read_range_uncompressed_multi_chunk() {
378        let block_count = 4;
379        let mut expected_data = vec![0u8; (block_count as u64 * BLOCK_SIZE) as usize];
380        for (i, byte) in expected_data.iter_mut().enumerate() {
381            *byte = ((i * 11) % 255) as u8;
382        }
383        // Force capping to 4096 bytes per buffer allocation so read_range processes
384        // 4 separate chunks.
385        let service = FakeBlockService::new_with_cap(expected_data.clone(), Some(4096));
386
387        let extents =
388            Extents::encode_extents(&[Extent::new(0..(block_count * BLOCK_SIZE), Some(0))]);
389        let extents = Extents::from_encoded(&extents).unwrap();
390        let blob = Blob::new(extents, block_count * BLOCK_SIZE, None);
391
392        let (dest_buf, rx) = TestVecBuffer::new(expected_data.len());
393        blob.read_range(0..(block_count * BLOCK_SIZE), &service, dest_buf);
394
395        assert_eq!(rx.commits().len(), 4);
396        assert_eq!(rx.output(), expected_data);
397    }
398
399    #[test]
400    fn test_read_range_uncompressed_unaligned_uncompressed_size() {
401        let uncompressed_size = 5000u64;
402        let mut expected_data = vec![0u8; 8192];
403        for (i, byte) in expected_data.iter_mut().enumerate() {
404            *byte = (i % 251) as u8;
405        }
406        let service = FakeBlockService::new(expected_data.clone());
407
408        let extents = Extents::encode_extents(&[Extent::new(0..8192, Some(0))]);
409        let extents = Extents::from_encoded(&extents).unwrap();
410        let blob = Blob::new(extents, uncompressed_size, None);
411
412        let (dest_buf, rx) = TestVecBuffer::new(8192);
413        blob.read_range(0..8192, &service, dest_buf);
414
415        assert_eq!(rx.commits(), vec![(0, 5000)]);
416        assert_eq!(&rx.output()[..5000], &expected_data[..5000]);
417    }
418
419    #[test]
420    fn test_read_range_compressed_tail_chunk_only() {
421        let uncompressed_size = 32768 * 2 + 1024;
422        let mut uncompressed_data = vec![0u8; uncompressed_size];
423        for (i, byte) in uncompressed_data.iter_mut().enumerate() {
424            *byte = ((i * 7) % 251) as u8;
425        }
426
427        let options =
428            ChunkedArchiveOptions::V3 { compression_algorithm: CompressionAlgorithm::Zstd };
429        let archive =
430            delivery_blob::compression::ChunkedArchive::new(&uncompressed_data, options).unwrap();
431
432        let mut compressed_offsets = vec![0];
433        let mut compressed_data = vec![];
434        for chunk in archive.chunks() {
435            compressed_data.extend_from_slice(&chunk.compressed_data);
436            compressed_offsets.push(compressed_data.len() as u64);
437        }
438        compressed_offsets.pop();
439
440        let chunk_size = archive.chunk_size();
441        let stored_size = compressed_data.len() as u64;
442        let stored_blocks = stored_size.div_ceil(BLOCK_SIZE);
443        let mut device_data = vec![0u8; (stored_blocks * BLOCK_SIZE) as usize];
444        device_data[..compressed_data.len()].copy_from_slice(&compressed_data);
445        let service = FakeBlockService::new(device_data);
446
447        let extents =
448            Extents::encode_extents(&[Extent::new(0..(stored_blocks * BLOCK_SIZE), Some(0))]);
449        let extents = Extents::from_encoded(&extents).unwrap();
450        let compression_info = CompressionInfo::new(
451            chunk_size as u64,
452            stored_size,
453            &compressed_offsets,
454            CompressionAlgorithm::Zstd,
455        )
456        .unwrap();
457        let blob = Blob::new(extents, uncompressed_size as u64, Some(compression_info));
458
459        let tail_start = chunk_size as u64 * 2;
460        let (dest_buf, rx) = TestVecBuffer::new_with_offset(32768, tail_start);
461        blob.read_range(tail_start..(uncompressed_size as u64), &service, dest_buf);
462
463        assert_eq!(rx.commits(), vec![(tail_start, 1024)]);
464        assert_eq!(&rx.output()[..1024], &uncompressed_data[65536..]);
465    }
466
467    #[test]
468    fn test_read_range_compressed_partial_final_chunk_zero_tail() {
469        let uncompressed_size = 32768 + 1024;
470        let mut uncompressed_data = vec![0u8; uncompressed_size];
471        for (i, byte) in uncompressed_data.iter_mut().enumerate() {
472            *byte = ((i * 13) % 251) as u8;
473        }
474
475        let options =
476            ChunkedArchiveOptions::V3 { compression_algorithm: CompressionAlgorithm::Zstd };
477        let archive =
478            delivery_blob::compression::ChunkedArchive::new(&uncompressed_data, options).unwrap();
479
480        let mut compressed_offsets = vec![0];
481        let mut compressed_data = vec![];
482        for chunk in archive.chunks() {
483            compressed_data.extend_from_slice(&chunk.compressed_data);
484            compressed_offsets.push(compressed_data.len() as u64);
485        }
486        compressed_offsets.pop();
487
488        let chunk_size = archive.chunk_size();
489        let stored_size = compressed_data.len() as u64;
490        let stored_blocks = stored_size.div_ceil(BLOCK_SIZE);
491        let mut device_data = vec![0u8; (stored_blocks * BLOCK_SIZE) as usize];
492        device_data[..compressed_data.len()].copy_from_slice(&compressed_data);
493        let service = FakeBlockService::new(device_data);
494
495        let extents =
496            Extents::encode_extents(&[Extent::new(0..(stored_blocks * BLOCK_SIZE), Some(0))]);
497        let extents = Extents::from_encoded(&extents).unwrap();
498        let compression_info = CompressionInfo::new(
499            chunk_size as u64,
500            stored_size,
501            &compressed_offsets,
502            CompressionAlgorithm::Zstd,
503        )
504        .unwrap();
505        let blob = Blob::new(extents, uncompressed_size as u64, Some(compression_info));
506
507        // Pre-fill destination buffer with 0xFF bytes to verify tail zeroing
508        let (mut dest_buf, rx) = TestVecBuffer::new(65536);
509        dest_buf.data.fill(0xFF);
510        blob.read_range(0..(uncompressed_size as u64), &service, dest_buf);
511
512        assert_eq!(rx.commits(), vec![(0, chunk_size), (chunk_size as u64, 1024)]);
513        assert_eq!(&rx.output()[..uncompressed_size], &uncompressed_data[..]);
514        assert_eq!(&rx.output()[uncompressed_size..65536], &[0u8; 31744]);
515    }
516}