1use 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
14pub 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 pub fn extents(&self) -> &Extents {
32 &self.extents
33 }
34
35 pub fn uncompressed_size(&self) -> u64 {
37 self.uncompressed_size
38 }
39
40 pub fn compression_info(&self) -> Option<&CompressionInfo> {
42 self.compression_info.as_deref()
43 }
44
45 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 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 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 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 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 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}